diff --git a/apps/api/src/controllers/v0/admin/metrics.ts b/apps/api/src/controllers/v0/admin/metrics.ts new file mode 100644 index 000000000..43b6288d4 --- /dev/null +++ b/apps/api/src/controllers/v0/admin/metrics.ts @@ -0,0 +1,30 @@ +import type { Request, Response } from "express"; +import { redisEvictConnection } from "../../../services/redis"; + +export async function metricsController(_: Request, res: Response) { + let cursor: string = "0"; + const metrics: Record = {}; + do { + const res = await redisEvictConnection.sscan("concurrency-limit-queues", cursor); + cursor = res[0]; + + const keys = res[1]; + + for (const key of keys) { + const jobCount = await redisEvictConnection.zcard(key); + + if (jobCount === 0) { + await redisEvictConnection.srem("concurrency-limit-queues", key); + } else { + const teamId = key.split(":")[1]; + metrics[teamId] = jobCount; + } + } + } while (cursor !== "0"); + + res.contentType("text/plain").send(`\ +# HELP concurrency_limit_queue_job_count The number of jobs in the concurrency limit queue per team +# TYPE concurrency_limit_queue_job_count gauge +${Object.entries(metrics).map(([key, value]) => `concurrency_limit_queue_job_count{team_id="${key}"} ${value}`).join("\n")} +`); +} \ No newline at end of file diff --git a/apps/api/src/controllers/v0/crawl-cancel.ts b/apps/api/src/controllers/v0/crawl-cancel.ts index ff8b88040..3d5249b9c 100644 --- a/apps/api/src/controllers/v0/crawl-cancel.ts +++ b/apps/api/src/controllers/v0/crawl-cancel.ts @@ -1,7 +1,6 @@ import { Request, Response } from "express"; import { authenticateUser } from "../auth"; import { RateLimiterMode } from "../../../src/types"; -import { supabase_service } from "../../../src/services/supabase"; import { logger } from "../../../src/lib/logger"; import { getCrawl, saveCrawl } from "../../../src/lib/crawl-redis"; import * as Sentry from "@sentry/node"; @@ -11,8 +10,6 @@ configDotenv(); export async function crawlCancelController(req: Request, res: Response) { try { - const useDbAuthentication = process.env.USE_DB_AUTHENTICATION === "true"; - const auth = await authenticateUser(req, res, RateLimiterMode.CrawlStatus); if (!auth.success) { return res.status(auth.status).json({ error: auth.error }); @@ -36,19 +33,8 @@ export async function crawlCancelController(req: Request, res: Response) { } // check if the job belongs to the team - if (useDbAuthentication) { - const { data, error: supaError } = await supabase_service - .from("bulljobs_teams") - .select("*") - .eq("job_id", req.params.jobId) - .eq("team_id", team_id); - if (supaError) { - return res.status(500).json({ error: supaError.message }); - } - - if (data.length === 0) { - return res.status(403).json({ error: "Unauthorized" }); - } + if (sc.team_id !== team_id) { + return res.status(403).json({ error: "Unauthorized" }); } try { diff --git a/apps/api/src/controllers/v0/crawl.ts b/apps/api/src/controllers/v0/crawl.ts index 6b0eb9796..cf0421330 100644 --- a/apps/api/src/controllers/v0/crawl.ts +++ b/apps/api/src/controllers/v0/crawl.ts @@ -4,7 +4,6 @@ import { authenticateUser } from "../auth"; import { RateLimiterMode } from "../../../src/types"; import { addScrapeJob } from "../../../src/services/queue-jobs"; import { isUrlBlocked } from "../../../src/scraper/WebScraper/utils/blocklist"; -import { logCrawl } from "../../../src/services/logging/crawl_log"; import { validateIdempotencyKey } from "../../../src/services/idempotency/validate"; import { createIdempotencyKey } from "../../../src/services/idempotency/create"; import { @@ -160,8 +159,6 @@ export async function crawlController(req: Request, res: Response) { // } // } - await logCrawl(id, team_id); - const { scrapeOptions, internalOptions } = fromV0ScrapeOptions( pageOptions, undefined, diff --git a/apps/api/src/controllers/v1/batch-scrape.ts b/apps/api/src/controllers/v1/batch-scrape.ts index 2f8a3e5ed..fef7033be 100644 --- a/apps/api/src/controllers/v1/batch-scrape.ts +++ b/apps/api/src/controllers/v1/batch-scrape.ts @@ -17,7 +17,6 @@ import { saveCrawl, StoredCrawl, } from "../../lib/crawl-redis"; -import { logCrawl } from "../../services/logging/crawl_log"; import { getJobPriority } from "../../lib/job-priority"; import { addScrapeJobs } from "../../services/queue-jobs"; import { callWebhook } from "../../services/webhook"; @@ -97,10 +96,6 @@ export async function batchScrapeController( account: req.account, }); - if (!req.body.appendToId) { - await logCrawl(id, req.auth.team_id); - } - const { scrapeOptions, internalOptions } = fromV1ScrapeOptions(req.body, req.body.timeout, req.auth.team_id); const sc: StoredCrawl = req.body.appendToId diff --git a/apps/api/src/controllers/v1/crawl-cancel.ts b/apps/api/src/controllers/v1/crawl-cancel.ts index 00af8b31f..8ddd9fdba 100644 --- a/apps/api/src/controllers/v1/crawl-cancel.ts +++ b/apps/api/src/controllers/v1/crawl-cancel.ts @@ -1,5 +1,4 @@ import { Response } from "express"; -import { supabase_service } from "../../services/supabase"; import { logger } from "../../lib/logger"; import { getCrawl, saveCrawl } from "../../lib/crawl-redis"; import * as Sentry from "@sentry/node"; @@ -12,27 +11,13 @@ export async function crawlCancelController( res: Response, ) { try { - const useDbAuthentication = process.env.USE_DB_AUTHENTICATION === "true"; - const sc = await getCrawl(req.params.jobId); if (!sc) { return res.status(404).json({ error: "Job not found" }); } - // check if the job belongs to the team - if (useDbAuthentication) { - const { data, error: supaError } = await supabase_service - .from("bulljobs_teams") - .select("*") - .eq("job_id", req.params.jobId) - .eq("team_id", req.auth.team_id); - if (supaError) { - return res.status(500).json({ error: supaError.message }); - } - - if (data.length === 0) { - return res.status(403).json({ error: "Unauthorized" }); - } + if (sc.team_id !== req.auth.team_id) { + return res.status(403).json({ error: "Unauthorized" }); } try { diff --git a/apps/api/src/controllers/v1/crawl.ts b/apps/api/src/controllers/v1/crawl.ts index c7a87c373..62df33703 100644 --- a/apps/api/src/controllers/v1/crawl.ts +++ b/apps/api/src/controllers/v1/crawl.ts @@ -8,7 +8,6 @@ import { toLegacyCrawlerOptions, } from "./types"; import { crawlToCrawler, saveCrawl, StoredCrawl } from "../../lib/crawl-redis"; -import { logCrawl } from "../../services/logging/crawl_log"; import { _addScrapeJobToBullMQ } from "../../services/queue-jobs"; import { logger as _logger } from "../../lib/logger"; import { fromV1ScrapeOptions } from "../v2/types"; @@ -44,8 +43,6 @@ export async function crawlController( account: req.account, }); - await logCrawl(id, req.auth.team_id); - let { remainingCredits } = req.account!; const useDbAuthentication = process.env.USE_DB_AUTHENTICATION === "true"; if (!useDbAuthentication) { diff --git a/apps/api/src/controllers/v1/map.ts b/apps/api/src/controllers/v1/map.ts index ee48fed68..61e2c4fa3 100644 --- a/apps/api/src/controllers/v1/map.ts +++ b/apps/api/src/controllers/v1/map.ts @@ -258,7 +258,7 @@ export async function getMapResults({ links = links .map((x) => { try { - return checkAndUpdateURLForMap(x).url.trim(); + return checkAndUpdateURLForMap(x, crawlerOptions.ignoreQueryParameters ?? true).url.trim(); } catch (_) { return null; } diff --git a/apps/api/src/controllers/v1/types.ts b/apps/api/src/controllers/v1/types.ts index 05c01156d..5e61566ce 100644 --- a/apps/api/src/controllers/v1/types.ts +++ b/apps/api/src/controllers/v1/types.ts @@ -795,12 +795,14 @@ export type CrawlRequestInput = z.input; // Note: Map types have been transitioned to v2/types.ts while maintaining backwards compatibility export const mapRequestSchema = crawlerOptions + .omit({ ignoreQueryParameters: true }) .extend({ url, origin: z.string().optional().default("api"), integration: z.nativeEnum(IntegrationEnum).optional().transform(val => val || null), includeSubdomains: z.boolean().default(true), search: z.string().optional(), + ignoreQueryParameters: z.boolean().default(true), ignoreSitemap: z.boolean().default(false), sitemapOnly: z.boolean().default(false), limit: z.number().min(1).max(30000).default(5000), diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index bee6925b5..6571b67dd 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -39,6 +39,10 @@ const { ExpressAdapter } = require("@bull-board/express"); const numCPUs = process.env.ENV === "local" ? 2 : os.cpus().length; logger.info(`Number of CPUs: ${numCPUs} available`); +logger.info("Network info dump", { + networkInterfaces: os.networkInterfaces(), +}); + // Install cacheable lookup for all other requests cacheableLookup.install(http.globalAgent); cacheableLookup.install(https.globalAgent); diff --git a/apps/api/src/lib/concurrency-limit.ts b/apps/api/src/lib/concurrency-limit.ts index 022f3523b..7fbb7837b 100644 --- a/apps/api/src/lib/concurrency-limit.ts +++ b/apps/api/src/lib/concurrency-limit.ts @@ -75,11 +75,9 @@ export async function pushConcurrencyLimitedJob( timeout: number, now: number = Date.now(), ) { - await redisEvictConnection.zadd( - constructQueueKey(team_id), - now + timeout, - JSON.stringify(job), - ); + const queueKey = constructQueueKey(team_id); + await redisEvictConnection.zadd(queueKey, now + timeout, JSON.stringify(job)); + await redisEvictConnection.sadd("concurrency-limit-queues", queueKey); } export async function getConcurrencyLimitedJobs( diff --git a/apps/api/src/lib/validateUrl.ts b/apps/api/src/lib/validateUrl.ts index 6216e0ec5..e8fe866b5 100644 --- a/apps/api/src/lib/validateUrl.ts +++ b/apps/api/src/lib/validateUrl.ts @@ -126,7 +126,7 @@ export function isSameSubdomain(url: string, baseUrl: string) { return domain1 === domain2 && subdomain1 === subdomain2; } -export const checkAndUpdateURLForMap = (url: string) => { +export const checkAndUpdateURLForMap = (url: string, ignoreQueryParameters: boolean = false) => { if (!protocolIncluded(url)) { url = `http://${url}`; } @@ -147,7 +147,10 @@ export const checkAndUpdateURLForMap = (url: string) => { } // remove any query params - // url = url.split("?")[0].trim(); + if (ignoreQueryParameters) { + url = url.split("?")[0].trim(); + typedUrlObj.search = ""; + } return { urlObj: typedUrlObj, url: url }; }; diff --git a/apps/api/src/routes/admin.ts b/apps/api/src/routes/admin.ts index a8de40fd4..25f0880ad 100644 --- a/apps/api/src/routes/admin.ts +++ b/apps/api/src/routes/admin.ts @@ -13,6 +13,7 @@ import { cclogController } from "../controllers/v0/admin/cclog"; import { indexQueuePrometheus } from "../controllers/v0/admin/index-queue-prometheus"; import { zdrcleanerController } from "../controllers/v0/admin/zdrcleaner"; import { triggerPrecrawl } from "../controllers/v0/admin/precrawl"; +import { metricsController } from "../controllers/v0/admin/metrics"; export const adminRouter = express.Router(); @@ -68,3 +69,8 @@ adminRouter.get( `/admin/${process.env.BULL_AUTH_KEY}/precrawl`, wrap(triggerPrecrawl), ); + +adminRouter.get( + `/admin/${process.env.BULL_AUTH_KEY}/metrics`, + wrap(metricsController), +); \ No newline at end of file diff --git a/apps/api/src/scraper/scrapeURL/engines/pdf/index.ts b/apps/api/src/scraper/scrapeURL/engines/pdf/index.ts index b86d1cb2e..fcd1ce0af 100644 --- a/apps/api/src/scraper/scrapeURL/engines/pdf/index.ts +++ b/apps/api/src/scraper/scrapeURL/engines/pdf/index.ts @@ -52,6 +52,7 @@ async function scrapePDFWithRunPodMU( meta.abort.throwIfAborted(); + const podStart = await robustFetch({ url: "https://api.runpod.ai/v2/" + process.env.RUNPOD_MU_POD_ID + "/runsync", diff --git a/apps/api/src/services/logging/crawl_log.ts b/apps/api/src/services/logging/crawl_log.ts deleted file mode 100644 index 86f885293..000000000 --- a/apps/api/src/services/logging/crawl_log.ts +++ /dev/null @@ -1,22 +0,0 @@ -import { supabase_service } from "../supabase"; -import { logger } from "../../../src/lib/logger"; -import { configDotenv } from "dotenv"; -configDotenv(); - -export async function logCrawl(job_id: string, team_id: string) { - const useDbAuthentication = process.env.USE_DB_AUTHENTICATION === "true"; - if (useDbAuthentication) { - try { - const { data, error } = await supabase_service - .from("bulljobs_teams") - .insert([ - { - job_id: job_id, - team_id: team_id, - }, - ]); - } catch (error) { - logger.error(`Error logging crawl job to supabase:\n${error}`); - } - } -}