From c3761309746b2b9832826c90f4dae4e8e752dd85 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Sat, 18 Oct 2025 19:57:11 +0200 Subject: [PATCH] feat(nuq): per-owner-per-group concurrency limiting (#2302) --- .github/CODEOWNERS | 7 +- apps/api/src/controllers/v0/admin/cclog.ts | 67 -- .../v0/admin/concurrency-queue-migration.ts | 88 ++ .../src/controllers/v0/admin/crawl-check.ts | 18 +- apps/api/src/controllers/v0/crawl.ts | 37 +- apps/api/src/controllers/v1/batch-scrape.ts | 7 + .../src/controllers/v1/concurrency-check.ts | 50 +- .../api/src/controllers/v1/crawl-status-ws.ts | 15 +- apps/api/src/controllers/v1/crawl.ts | 12 +- apps/api/src/controllers/v1/queue-status.ts | 20 +- apps/api/src/controllers/v2/batch-scrape.ts | 7 + .../src/controllers/v2/concurrency-check.ts | 55 +- .../api/src/controllers/v2/crawl-status-ws.ts | 15 +- apps/api/src/controllers/v2/crawl.ts | 11 +- apps/api/src/controllers/v2/queue-status.ts | 20 +- apps/api/src/lib/concurrency-limit.ts | 358 ------- apps/api/src/lib/crawl-redis.ts | 10 +- apps/api/src/routes/admin.ts | 12 +- .../api/src/services/indexing/index-worker.ts | 4 +- apps/api/src/services/queue-jobs.ts | 383 +------- apps/api/src/services/worker/crawl-logic.ts | 2 +- apps/api/src/services/worker/nuq.ts | 923 ++++++++++++++++-- apps/api/src/services/worker/scrape-worker.ts | 132 +-- apps/nuq-postgres/nuq.sql | 68 +- 24 files changed, 1218 insertions(+), 1103 deletions(-) delete mode 100644 apps/api/src/controllers/v0/admin/cclog.ts create mode 100644 apps/api/src/controllers/v0/admin/concurrency-queue-migration.ts delete mode 100644 apps/api/src/lib/concurrency-limit.ts diff --git a/.github/CODEOWNERS b/.github/CODEOWNERS index 959141572..ecba85ef9 100644 --- a/.github/CODEOWNERS +++ b/.github/CODEOWNERS @@ -12,7 +12,7 @@ /apps/api/src/controllers/v2/crawl* @mogery ### /batch/scrape -/apps/api/src/controllers/v2/batch-scrape* @mogery +/apps/api/src/controllers/v2/batch-scrape* @mogery ### /extract /apps/api/src/controllers/v2/extract* @nickscamara @@ -36,7 +36,7 @@ /apps/api/src/controllers/v1/crawl* @mogery ### /batch/scrape -/apps/api/src/controllers/v1/batch-scrape* @mogery +/apps/api/src/controllers/v1/batch-scrape* @mogery ### /extract /apps/api/src/controllers/v1/extract* @nickscamara @@ -77,9 +77,6 @@ ### remnants of WebScraper/WebCrawler /apps/api/src/scraper/WebScraper/* @mogery @nickscamara -## concurrency limits -/apps/api/src/lib/concurrency-limit.ts @mogery @nickscamara - ## BullMQ-related code /apps/api/src/services/queue-worker.ts @mogery @nickscamara /apps/api/src/services/worker/scrape-worker.ts @mogery diff --git a/apps/api/src/controllers/v0/admin/cclog.ts b/apps/api/src/controllers/v0/admin/cclog.ts deleted file mode 100644 index 76d2c415d..000000000 --- a/apps/api/src/controllers/v0/admin/cclog.ts +++ /dev/null @@ -1,67 +0,0 @@ -import { redisEvictConnection } from "../../../services/redis"; -import { supabase_service } from "../../../services/supabase"; -import { logger as _logger } from "../../../lib/logger"; -import { Request, Response } from "express"; - -async function cclog() { - const logger = _logger.child({ - module: "cclog", - }); - - let cursor = 0; - do { - const result = await redisEvictConnection.scan( - cursor, - "MATCH", - "concurrency-limiter:*", - "COUNT", - 100000, - ); - cursor = parseInt(result[0], 10); - const usable = result[1].filter(x => !x.includes("preview_")); - - logger.info("Stepped", { cursor, usable: usable.length }); - - if (usable.length > 0) { - const entries: { - team_id: string; - concurrency: number; - created_at: Date; - }[] = []; - - for (const x of usable) { - const at = new Date(); - const concurrency = await redisEvictConnection.zrangebyscore( - x, - Date.now(), - Infinity, - ); - if (concurrency) { - entries.push({ - team_id: x.split(":")[1], - concurrency: concurrency.length, - created_at: at, - }); - } - } - - try { - await supabase_service.from("concurrency_log").insert(entries); - } catch (e) { - logger.error("Error inserting", { error: e }); - } - } - } while (cursor != 0); -} - -export async function cclogController(req: Request, res: Response) { - try { - await cclog(); - res.status(200).json({ ok: true }); - } catch (e) { - _logger.error("Error", { module: "cclog", error: e }); - res.status(500).json({ - message: "Error", - }); - } -} diff --git a/apps/api/src/controllers/v0/admin/concurrency-queue-migration.ts b/apps/api/src/controllers/v0/admin/concurrency-queue-migration.ts new file mode 100644 index 000000000..65827d4c1 --- /dev/null +++ b/apps/api/src/controllers/v0/admin/concurrency-queue-migration.ts @@ -0,0 +1,88 @@ +import type { Request, Response } from "express"; +import { logger as _logger } from "../../../lib/logger"; +import { redisEvictConnection } from "../../../services/redis"; +import { crawlGroup, scrapeQueue } from "../../../services/worker/nuq"; +import { getCrawl } from "../../../lib/crawl-redis"; + +export async function migrateConcurrencyQueue(_: Request, res: Response) { + const logger = _logger.child({ + module: "admin", + method: "migrateConcurrencyQueue" + }); + + let crawlsCursor = "0"; + + do { + const crawlsScan = await redisEvictConnection.sscan( + "ongoing_crawls", + crawlsCursor + ); + crawlsCursor = crawlsScan[0]; + + for (const crawlId of crawlsScan[1]) { + const crawlData = (await getCrawl(crawlId)) ?? { maxConcurrency: undefined } + + logger.info("Migrating crawl", { crawlId }); + + await crawlGroup.addGroup(crawlId, [ + { + queue: scrapeQueue, + maxConcurrency: crawlData.maxConcurrency ?? undefined, + }, + ]); + } + + } while (crawlsCursor !== "0"); + + let queuesCursor = "0"; + + do { + const queuesScan = await redisEvictConnection.sscan( + "concurrency-limit-queues", + queuesCursor, + ); + queuesCursor = queuesScan[0]; + + for (const queueKey of queuesScan[1]) { + if (queueKey.startsWith("concurrency-limit-queue:preview_")) { + logger.warn("Skipping preview queue", { queueKey }); + continue; + } + + let queueCursor = "0"; + + logger.info("Migrating a queue", { queueKey }); + + do { + const queueScan = await redisEvictConnection.zscan( + queueKey, + queueCursor, + ); + queueCursor = queueScan[0]; + + for (let i = 0; i < queueScan[1].length; i += 2) { + const jobData = JSON.parse(queueScan[1][i]); + const jobScore = queueScan[1][i + 1]; + + logger.info("Migrating job", { queueKey, teamId: jobData?.data?.teamId, zeroDataRetention: jobData?.data?.zeroDataRetention, scrapeId: jobData?.id }); + + const success = await scrapeQueue.tryAddJob(jobData.id, jobData.data, { + priority: jobData.priority, + listenable: jobData.listenable, + ownerId: jobData.data.team_id ?? undefined, + groupId: jobData.data.crawl_id ?? undefined, + timesOutAt: jobScore === "inf" ? undefined : new Date(parseInt(jobScore, 10)), + }); + + if (success === null) { + logger.warn("Failed to migrate job due to conflict", { queueKey, teamId: jobData?.data?.teamId, zeroDataRetention: jobData?.data?.zeroDataRetention, scrapeId: jobData?.id }); + } + } + } while (queueCursor !== "0"); + } + } while (queuesCursor !== "0"); + + logger.info("Migration complete! 🎉"); + + res.json({ ok: true }); +} diff --git a/apps/api/src/controllers/v0/admin/crawl-check.ts b/apps/api/src/controllers/v0/admin/crawl-check.ts index 75b42c5e4..b462a1e7d 100644 --- a/apps/api/src/controllers/v0/admin/crawl-check.ts +++ b/apps/api/src/controllers/v0/admin/crawl-check.ts @@ -6,7 +6,6 @@ import { StoredCrawl, } from "../../../lib/crawl-redis"; import { supabase_service } from "../../../services/supabase"; -import { getConcurrencyLimitedJobs } from "../../../lib/concurrency-limit"; import { scrapeQueue } from "../../../services/worker/nuq"; type AnalyzedCrawlPass1 = { @@ -217,23 +216,8 @@ export async function crawlCheckController(req: Request, res: Response) { await Promise.all(activeCrawls.map(analyzeCrawlPass1)) ).filter(result => result.success); - const teamIds = [...new Set(firstPass.map(result => result.teamId))]; - const teamConcurrencyQueuedJobs = Object.fromEntries( - await Promise.all( - teamIds.map(async teamId => [ - teamId, - await getConcurrencyLimitedJobs(teamId), - ]), - ), - ); - const results = await Promise.all( - firstPass.map(result => - analyzeCrawlPass2( - { ...result }, - teamConcurrencyQueuedJobs[result.teamId], - ), - ), + firstPass.map(result => analyzeCrawlPass2({ ...result }, new Set())), ); // Clean results which should be "solved" diff --git a/apps/api/src/controllers/v0/crawl.ts b/apps/api/src/controllers/v0/crawl.ts index 59c38eb89..980492b65 100644 --- a/apps/api/src/controllers/v0/crawl.ts +++ b/apps/api/src/controllers/v0/crawl.ts @@ -33,6 +33,7 @@ import { ZodError } from "zod"; import { BLOCKLISTED_URL_MESSAGE } from "../../lib/strings"; import { fromV0ScrapeOptions } from "../v2/types"; import { isSelfHosted } from "../../lib/deployment"; +import { crawlGroup, scrapeQueue } from "../../services/worker/nuq"; export async function crawlController(req: Request, res: Response) { try { @@ -146,35 +147,6 @@ export async function crawlController(req: Request, res: Response) { }); } - // if (mode === "single_urls" && !url.includes(",")) { // NOTE: do we need this? - // try { - // const a = new WebScraperDataProvider(); - // await a.setOptions({ - // jobId: uuidv4(), - // mode: "single_urls", - // urls: [url], - // crawlerOptions: { ...crawlerOptions, returnOnlyUrls: true }, - // pageOptions: pageOptions, - // }); - - // const docs = await a.getDocuments(false, (progress) => { - // job.updateProgress({ - // current: progress.current, - // total: progress.total, - // current_step: "SCRAPING", - // current_url: progress.currentDocumentUrl, - // }); - // }); - // return res.json({ - // success: true, - // documents: docs, - // }); - // } catch (error) { - // logger.error(error); - // return res.status(500).json({ error: error.message }); - // } - // } - const { scrapeOptions, internalOptions } = fromV0ScrapeOptions( pageOptions, undefined, @@ -203,6 +175,13 @@ export async function crawlController(req: Request, res: Response) { sc.robots = await crawler.getRobotsTxt(); } catch (_) {} + await crawlGroup.addGroup(id, [ + { + queue: scrapeQueue, + maxConcurrency: undefined, + }, + ]); + await saveCrawl(id, sc); await markCrawlActive(id); diff --git a/apps/api/src/controllers/v1/batch-scrape.ts b/apps/api/src/controllers/v1/batch-scrape.ts index c12ff7a9f..cea629d38 100644 --- a/apps/api/src/controllers/v1/batch-scrape.ts +++ b/apps/api/src/controllers/v1/batch-scrape.ts @@ -25,6 +25,7 @@ import { BLOCKLISTED_URL_MESSAGE } from "../../lib/strings"; import { isUrlBlocked } from "../../scraper/WebScraper/utils/blocklist"; import { fromV1ScrapeOptions } from "../v2/types"; import { checkPermissions } from "../../lib/permissions"; +import { crawlGroup, scrapeQueue } from "../../services/worker/nuq"; export async function batchScrapeController( req: RequestWithAuth<{}, BatchScrapeResponse, BatchScrapeRequest>, @@ -136,6 +137,12 @@ export async function batchScrapeController( }; if (!req.body.appendToId) { + await crawlGroup.addGroup(id, [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ]); await saveCrawl(id, sc); await markCrawlActive(id); } diff --git a/apps/api/src/controllers/v1/concurrency-check.ts b/apps/api/src/controllers/v1/concurrency-check.ts index 517aed515..26b11f74e 100644 --- a/apps/api/src/controllers/v1/concurrency-check.ts +++ b/apps/api/src/controllers/v1/concurrency-check.ts @@ -1,27 +1,59 @@ import { + AuthCreditUsageChunkFromTeam, ConcurrencyCheckParams, ConcurrencyCheckResponse, RequestWithAuth, } from "./types"; import { Response } from "express"; -import { redisEvictConnection } from "../../../src/services/redis"; +import { scrapeQueue } from "../../services/worker/nuq"; +import { getACUCTeam } from "../auth"; +import { RateLimiterMode } from "../../types"; // Basically just middleware and error wrapping export async function concurrencyCheckController( req: RequestWithAuth, res: Response, ) { - const concurrencyLimiterKey = "concurrency-limiter:" + req.auth.team_id; - const now = Date.now(); - const activeJobsOfTeam = await redisEvictConnection.zrangebyscore( - concurrencyLimiterKey, - now, - Infinity, + if (!req.acuc) { + return res.status(401).json({ + success: false, + error: "Unauthorized", + }); + } + + const ownerConcurrency = await scrapeQueue.getOwnerConcurrency( + req.auth.team_id, ); + let maxConcurrency: number | null = ownerConcurrency?.maxConcurrency ?? null; + + if (maxConcurrency === null) { + let otherACUC: AuthCreditUsageChunkFromTeam | null = null; + if (!req.acuc.is_extract) { + otherACUC = await getACUCTeam( + req.auth.team_id, + false, + true, + RateLimiterMode.Extract, + ); + } else { + otherACUC = await getACUCTeam( + req.auth.team_id, + false, + true, + RateLimiterMode.Crawl, + ); + } + + maxConcurrency = Math.max( + req.acuc.concurrency, + otherACUC?.concurrency ?? 0, + ); + } + return res.status(200).json({ success: true, - concurrency: activeJobsOfTeam.length, - maxConcurrency: req.acuc?.concurrency ?? 0, + concurrency: ownerConcurrency?.currentConcurrency ?? 0, + maxConcurrency, }); } diff --git a/apps/api/src/controllers/v1/crawl-status-ws.ts b/apps/api/src/controllers/v1/crawl-status-ws.ts index 6c4c6155b..d96d5c1b9 100644 --- a/apps/api/src/controllers/v1/crawl-status-ws.ts +++ b/apps/api/src/controllers/v1/crawl-status-ws.ts @@ -18,7 +18,6 @@ import { } from "../../lib/crawl-redis"; import { getJobs, PseudoJob } from "./crawl-status"; import * as Sentry from "@sentry/node"; -import { getConcurrencyLimitedJobs } from "../../lib/concurrency-limit"; import { scrapeQueue, NuQJobStatus } from "../../services/worker/nuq"; import { getErrorContactMessage } from "../../lib/deployment"; @@ -111,10 +110,9 @@ async function crawlStatusWS( setTimeout(loop, 1000); - let [_doneJobIDs, jobIDs, throttledJobsSet] = await Promise.all([ + let [_doneJobIDs, jobIDs] = await Promise.all([ getDoneJobsOrdered(req.params.jobId), getCrawlJobs(req.params.jobId), - getConcurrencyLimitedJobs(req.auth.team_id), ]); doneJobIDs = _doneJobIDs; @@ -124,15 +122,10 @@ async function crawlStatusWS( const validJobIDs: string[] = []; for (const id of jobIDs) { - if (throttledJobsSet.has(id)) { - validJobStatuses.push([id, "queued"]); + const job = jobs.get(id); + if (job && job.status !== "failed") { + validJobStatuses.push([id, job.status]); validJobIDs.push(id); - } else { - const job = jobs.get(id); - if (job && job.status !== "failed") { - validJobStatuses.push([id, job.status]); - validJobIDs.push(id); - } } } diff --git a/apps/api/src/controllers/v1/crawl.ts b/apps/api/src/controllers/v1/crawl.ts index e83489826..7abc54872 100644 --- a/apps/api/src/controllers/v1/crawl.ts +++ b/apps/api/src/controllers/v1/crawl.ts @@ -13,10 +13,11 @@ import { StoredCrawl, markCrawlActive, } from "../../lib/crawl-redis"; -import { _addScrapeJobToBullMQ } from "../../services/queue-jobs"; +import { addScrapeJob } from "../../services/queue-jobs"; import { logger as _logger } from "../../lib/logger"; import { fromV1ScrapeOptions } from "../v2/types"; import { checkPermissions } from "../../lib/permissions"; +import { crawlGroup, scrapeQueue } from "../../services/worker/nuq"; export async function crawlController( req: RequestWithAuth<{}, CrawlResponse, CrawlRequest>, @@ -135,11 +136,18 @@ export async function crawlController( }); } + await crawlGroup.addGroup(id, [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ]); + await saveCrawl(id, sc); await markCrawlActive(id); - await _addScrapeJobToBullMQ( + await addScrapeJob( { mode: "kickoff" as const, url: req.body.url, diff --git a/apps/api/src/controllers/v1/queue-status.ts b/apps/api/src/controllers/v1/queue-status.ts index 63e0157fc..46c46107d 100644 --- a/apps/api/src/controllers/v1/queue-status.ts +++ b/apps/api/src/controllers/v1/queue-status.ts @@ -3,12 +3,7 @@ import { getACUCTeam } from "../auth"; import { AuthCreditUsageChunkFromTeam, RequestWithAuth } from "./types"; import { Response } from "express"; import { redisEvictConnection } from "../../services/redis"; -import { - cleanOldConcurrencyLimitedJobs, - cleanOldConcurrencyLimitEntries, - getConcurrencyLimitActiveJobsCount, - getConcurrencyQueueJobsCount, -} from "../../lib/concurrency-limit"; +import { scrapeQueue } from "../../services/worker/nuq"; type QueueStatusResponse = { success: boolean; @@ -40,12 +35,7 @@ export async function queueStatusController( ); } - await cleanOldConcurrencyLimitEntries(req.auth.team_id); - const activeJobsOfTeam = await getConcurrencyLimitActiveJobsCount( - req.auth.team_id, - ); - await cleanOldConcurrencyLimitedJobs(req.auth.team_id); - const queuedJobsOfTeam = await getConcurrencyQueueJobsCount(req.auth.team_id); + const jobCounts = await scrapeQueue.getOwnerJobCounts(req.auth.team_id); const mostRecentSuccess = await redisEvictConnection.get( "most-recent-success:" + req.auth.team_id, @@ -54,9 +44,9 @@ export async function queueStatusController( return res.status(200).json({ success: true, - jobsInQueue: activeJobsOfTeam + queuedJobsOfTeam, - activeJobsInQueue: activeJobsOfTeam, - waitingJobsInQueue: queuedJobsOfTeam, + jobsInQueue: jobCounts.active + jobCounts.queued, + activeJobsInQueue: jobCounts.active, + waitingJobsInQueue: jobCounts.queued, maxConcurrency: Math.max( req.acuc?.concurrency ?? 1, otherACUC?.concurrency ?? 1, diff --git a/apps/api/src/controllers/v2/batch-scrape.ts b/apps/api/src/controllers/v2/batch-scrape.ts index 04a478bb1..cf57573e6 100644 --- a/apps/api/src/controllers/v2/batch-scrape.ts +++ b/apps/api/src/controllers/v2/batch-scrape.ts @@ -25,6 +25,7 @@ import { logger as _logger } from "../../lib/logger"; import { BLOCKLISTED_URL_MESSAGE } from "../../lib/strings"; import { isUrlBlocked } from "../../scraper/WebScraper/utils/blocklist"; import { checkPermissions } from "../../lib/permissions"; +import { crawlGroup, scrapeQueue } from "../../services/worker/nuq"; export async function batchScrapeController( req: RequestWithAuth<{}, BatchScrapeResponse, BatchScrapeRequest>, @@ -129,6 +130,12 @@ export async function batchScrapeController( }; if (!req.body.appendToId) { + await crawlGroup.addGroup(id, [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ]); await saveCrawl(id, sc); await markCrawlActive(id); } diff --git a/apps/api/src/controllers/v2/concurrency-check.ts b/apps/api/src/controllers/v2/concurrency-check.ts index 17edca1c4..04b6980ec 100644 --- a/apps/api/src/controllers/v2/concurrency-check.ts +++ b/apps/api/src/controllers/v2/concurrency-check.ts @@ -5,9 +5,9 @@ import { } from "./types"; import { AuthCreditUsageChunkFromTeam } from "../v1/types"; import { Response } from "express"; -import { redisEvictConnection } from "../../../src/services/redis"; import { getACUCTeam } from "../auth"; import { RateLimiterMode } from "../../types"; +import { scrapeQueue } from "../../services/worker/nuq"; // Basically just middleware and error wrapping export async function concurrencyCheckController( @@ -21,34 +21,39 @@ export async function concurrencyCheckController( }); } - let otherACUC: AuthCreditUsageChunkFromTeam | null = null; - if (!req.acuc.is_extract) { - otherACUC = await getACUCTeam( - req.auth.team_id, - false, - true, - RateLimiterMode.Extract, - ); - } else { - otherACUC = await getACUCTeam( - req.auth.team_id, - false, - true, - RateLimiterMode.Crawl, + const ownerConcurrency = await scrapeQueue.getOwnerConcurrency( + req.auth.team_id, + ); + + let maxConcurrency: number | null = ownerConcurrency?.maxConcurrency ?? null; + + if (maxConcurrency === null) { + let otherACUC: AuthCreditUsageChunkFromTeam | null = null; + if (!req.acuc.is_extract) { + otherACUC = await getACUCTeam( + req.auth.team_id, + false, + true, + RateLimiterMode.Extract, + ); + } else { + otherACUC = await getACUCTeam( + req.auth.team_id, + false, + true, + RateLimiterMode.Crawl, + ); + } + + maxConcurrency = Math.max( + req.acuc.concurrency, + otherACUC?.concurrency ?? 0, ); } - const concurrencyLimiterKey = "concurrency-limiter:" + req.auth.team_id; - const now = Date.now(); - const activeJobsOfTeam = await redisEvictConnection.zrangebyscore( - concurrencyLimiterKey, - now, - Infinity, - ); - return res.status(200).json({ success: true, - concurrency: activeJobsOfTeam.length, - maxConcurrency: Math.max(req.acuc.concurrency, otherACUC?.concurrency ?? 0), + concurrency: ownerConcurrency?.currentConcurrency ?? 0, + maxConcurrency, }); } diff --git a/apps/api/src/controllers/v2/crawl-status-ws.ts b/apps/api/src/controllers/v2/crawl-status-ws.ts index 91c1dd809..df5bfe5df 100644 --- a/apps/api/src/controllers/v2/crawl-status-ws.ts +++ b/apps/api/src/controllers/v2/crawl-status-ws.ts @@ -18,7 +18,6 @@ import { } from "../../lib/crawl-redis"; import { getJobs, PseudoJob } from "./crawl-status"; import * as Sentry from "@sentry/node"; -import { getConcurrencyLimitedJobs } from "../../lib/concurrency-limit"; import { scrapeQueue, NuQJobStatus } from "../../services/worker/nuq"; import { getErrorContactMessage } from "../../lib/deployment"; @@ -114,10 +113,9 @@ async function crawlStatusWS( setTimeout(loop, 1000); - let [_doneJobIDs, jobIDs, throttledJobsSet] = await Promise.all([ + let [_doneJobIDs, jobIDs] = await Promise.all([ getDoneJobsOrdered(req.params.jobId), getCrawlJobs(req.params.jobId), - getConcurrencyLimitedJobs(req.auth.team_id), ]); doneJobIDs = _doneJobIDs; @@ -127,15 +125,10 @@ async function crawlStatusWS( const validJobIDs: string[] = []; for (const id of jobIDs) { - if (throttledJobsSet.has(id)) { - validJobStatuses.push([id, "queued"]); + const job = jobs.get(id); + if (job && job.status !== "failed") { + validJobStatuses.push([id, job.status]); validJobIDs.push(id); - } else { - const job = jobs.get(id); - if (job && job.status !== "failed") { - validJobStatuses.push([id, job.status]); - validJobIDs.push(id); - } } } diff --git a/apps/api/src/controllers/v2/crawl.ts b/apps/api/src/controllers/v2/crawl.ts index 2a74b00fa..37b172c6e 100644 --- a/apps/api/src/controllers/v2/crawl.ts +++ b/apps/api/src/controllers/v2/crawl.ts @@ -13,13 +13,14 @@ import { StoredCrawl, markCrawlActive, } from "../../lib/crawl-redis"; -import { _addScrapeJobToBullMQ } from "../../services/queue-jobs"; +import { addScrapeJob } from "../../services/queue-jobs"; import { logger as _logger } from "../../lib/logger"; import { generateCrawlerOptionsFromPrompt } from "../../scraper/scrapeURL/transformers/llmExtract"; import { CostTracking } from "../../lib/cost-tracking"; import { checkPermissions } from "../../lib/permissions"; import { buildPromptWithWebsiteStructure } from "../../lib/map-utils"; import { modifyCrawlUrl } from "../../utils/url-utils"; +import { crawlGroup, scrapeQueue } from "../../services/worker/nuq"; export async function crawlController( req: RequestWithAuth<{}, CrawlResponse, CrawlRequest>, @@ -202,11 +203,17 @@ export async function crawlController( }); } + await crawlGroup.addGroup(id, [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ]); await saveCrawl(id, sc); await markCrawlActive(id); - await _addScrapeJobToBullMQ( + await addScrapeJob( { url: req.body.url, mode: "kickoff" as const, diff --git a/apps/api/src/controllers/v2/queue-status.ts b/apps/api/src/controllers/v2/queue-status.ts index e52542dce..1c11851f4 100644 --- a/apps/api/src/controllers/v2/queue-status.ts +++ b/apps/api/src/controllers/v2/queue-status.ts @@ -4,12 +4,7 @@ import { RequestWithAuth } from "./types"; import { AuthCreditUsageChunkFromTeam } from "../v1/types"; import { Response } from "express"; import { redisEvictConnection } from "../../services/redis"; -import { - cleanOldConcurrencyLimitedJobs, - cleanOldConcurrencyLimitEntries, - getConcurrencyLimitActiveJobsCount, - getConcurrencyQueueJobsCount, -} from "../../lib/concurrency-limit"; +import { scrapeQueue } from "../../services/worker/nuq"; type QueueStatusResponse = { success: boolean; @@ -41,12 +36,7 @@ export async function queueStatusController( ); } - await cleanOldConcurrencyLimitEntries(req.auth.team_id); - const activeJobsOfTeam = await getConcurrencyLimitActiveJobsCount( - req.auth.team_id, - ); - await cleanOldConcurrencyLimitedJobs(req.auth.team_id); - const queuedJobsOfTeam = await getConcurrencyQueueJobsCount(req.auth.team_id); + const jobCounts = await scrapeQueue.getOwnerJobCounts(req.auth.team_id); const mostRecentSuccess = await redisEvictConnection.get( "most-recent-success:" + req.auth.team_id, @@ -55,9 +45,9 @@ export async function queueStatusController( return res.status(200).json({ success: true, - jobsInQueue: activeJobsOfTeam + queuedJobsOfTeam, - activeJobsInQueue: activeJobsOfTeam, - waitingJobsInQueue: queuedJobsOfTeam, + jobsInQueue: jobCounts.active + jobCounts.queued, + activeJobsInQueue: jobCounts.active, + waitingJobsInQueue: jobCounts.queued, maxConcurrency: Math.max( req.acuc?.concurrency ?? 1, otherACUC?.concurrency ?? 1, diff --git a/apps/api/src/lib/concurrency-limit.ts b/apps/api/src/lib/concurrency-limit.ts deleted file mode 100644 index c70e30416..000000000 --- a/apps/api/src/lib/concurrency-limit.ts +++ /dev/null @@ -1,358 +0,0 @@ -import { RateLimiterMode } from "../types"; -import { redisEvictConnection } from "../services/redis"; -import { getACUCTeam } from "../controllers/auth"; -import { getCrawl, StoredCrawl } from "./crawl-redis"; -import { logger } from "./logger"; -import { abTestJob } from "../services/ab-test"; -import { scrapeQueue, type NuQJob } from "../services/worker/nuq"; - -const constructKey = (team_id: string) => "concurrency-limiter:" + team_id; -const constructQueueKey = (team_id: string) => - "concurrency-limit-queue:" + team_id; - -const constructCrawlKey = (crawl_id: string) => - "crawl-concurrency-limiter:" + crawl_id; - -export async function cleanOldConcurrencyLimitEntries( - team_id: string, - now: number = Date.now(), -) { - await redisEvictConnection.zremrangebyscore( - constructKey(team_id), - -Infinity, - now, - ); -} - -export async function getConcurrencyLimitActiveJobsCount( - team_id: string, -): Promise { - return await redisEvictConnection.zcount( - constructKey(team_id), - Date.now(), - Infinity, - ); -} - -export async function getConcurrencyLimitActiveJobs( - team_id: string, - now: number = Date.now(), -): Promise { - return await redisEvictConnection.zrangebyscore( - constructKey(team_id), - now, - Infinity, - ); -} - -export async function pushConcurrencyLimitActiveJob( - team_id: string, - id: string, - timeout: number, - now: number = Date.now(), -) { - await redisEvictConnection.zadd(constructKey(team_id), now + timeout, id); -} - -async function removeConcurrencyLimitActiveJob(team_id: string, id: string) { - await redisEvictConnection.zrem(constructKey(team_id), id); -} - -type ConcurrencyLimitedJob = { - id: string; - data: any; - priority: number; - listenable: boolean; -}; - -export async function cleanOldConcurrencyLimitedJobs( - team_id: string, - now: number = Date.now(), -) { - await redisEvictConnection.zremrangebyscore( - constructQueueKey(team_id), - -Infinity, - now, - ); -} - -export async function pushConcurrencyLimitedJob( - team_id: string, - job: ConcurrencyLimitedJob, - timeout: number, - now: number = Date.now(), -) { - 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(team_id: string) { - return new Set( - (await redisEvictConnection.zrange(constructQueueKey(team_id), 0, -1)).map( - x => JSON.parse(x).id, - ), - ); -} - -export async function getConcurrencyQueueJobsCount( - team_id: string, -): Promise { - return await redisEvictConnection.zcount( - constructQueueKey(team_id), - Date.now(), - Infinity, - ); -} - -async function cleanOldCrawlConcurrencyLimitEntries( - crawl_id: string, - now: number = Date.now(), -) { - await redisEvictConnection.zremrangebyscore( - constructCrawlKey(crawl_id), - -Infinity, - now, - ); -} - -export async function getCrawlConcurrencyLimitActiveJobs( - crawl_id: string, - now: number = Date.now(), -): Promise { - return await redisEvictConnection.zrangebyscore( - constructCrawlKey(crawl_id), - now, - Infinity, - ); -} - -export async function pushCrawlConcurrencyLimitActiveJob( - crawl_id: string, - id: string, - timeout: number, - now: number = Date.now(), -) { - await redisEvictConnection.zadd( - constructCrawlKey(crawl_id), - now + timeout, - id, - ); -} - -async function removeCrawlConcurrencyLimitActiveJob( - crawl_id: string, - id: string, -) { - await redisEvictConnection.zrem(constructCrawlKey(crawl_id), id); -} - -/** - * Grabs the next job from the team's concurrency limit queue. Handles crawl concurrency limits. - * - * This function may only be called once the outer code has verified that the team has not reached its concurrency limit. - * - * @param teamId - * @returns A job that can be run, or null if there are no more jobs to run. - */ -async function getNextConcurrentJob( - teamId: string, - i = 0, -): Promise<{ - job: ConcurrencyLimitedJob; - timeout: number; -} | null> { - let finalJobs: { - job: ConcurrencyLimitedJob; - _member: string; - timeout: number; - }[] = []; - - const crawlCache = new Map(); - let cursor: string = "0"; - - do { - const scanResult = await redisEvictConnection.zscan( - constructQueueKey(teamId), - cursor, - "COUNT", - 20, - ); - cursor = scanResult[0]; - const results = scanResult[1]; - - for (let i = 0; i < results.length; i += 2) { - const res = { - job: JSON.parse(results[i]), - _member: results[i], - timeout: - results[i + 1] === "inf" ? Infinity : parseFloat(results[i + 1]), - }; - - // If the job is associated with a crawl ID, we need to check if the crawl has a max concurrency limit - if (res.job.data.crawl_id) { - const sc = - crawlCache.get(res.job.data.crawl_id) ?? - (await getCrawl(res.job.data.crawl_id)); - if (sc !== null) { - crawlCache.set(res.job.data.crawl_id, sc); - } - - const maxCrawlConcurrency = - sc === null - ? null - : typeof sc.crawlerOptions?.delay === "number" && - sc.crawlerOptions.delay > 0 - ? 1 - : (sc.maxConcurrency ?? null); - - if (maxCrawlConcurrency !== null) { - // If the crawl has a max concurrency limit, we need to check if the crawl has reached the limit - const currentActiveConcurrency = ( - await getCrawlConcurrencyLimitActiveJobs(res.job.data.crawl_id) - ).length; - if (currentActiveConcurrency < maxCrawlConcurrency) { - // If we're under the max concurrency limit, we can run the job - finalJobs.push(res); - } - } else { - // If the crawl has no max concurrency limit, we can run the job - finalJobs.push(res); - } - } else { - // If the job is not associated with a crawl ID, we can run the job - finalJobs.push(res); - } - } - } while (finalJobs.length === 0 && cursor !== "0"); - - let finalJob: (typeof finalJobs)[number] | null = null; - if (finalJobs.length > 0) { - for (const job of finalJobs) { - const res = await redisEvictConnection.zrem( - constructQueueKey(teamId), - job._member, - ); - if (res !== 0) { - finalJob = job; - break; - } - } - - if (finalJob === null) { - // It's normal for this to happen, but if it happens too many times, we should log a warning - if (i > 100) { - logger.error( - "Failed to remove job from concurrency limit queue, hard bailing", - { - teamId, - jobIds: finalJobs.map(x => x.job.id), - zeroDataRetention: finalJobs.some( - x => x.job.data?.zeroDataRetention, - ), - i, - }, - ); - return null; - } else if (i > 15) { - logger.warn("Failed to remove job from concurrency limit queue", { - teamId, - jobIds: finalJobs.map(x => x.job.id), - zeroDataRetention: finalJobs.some(x => x.job.data?.zeroDataRetention), - i, - }); - } - - return await new Promise((resolve, reject) => - setTimeout( - () => { - getNextConcurrentJob(teamId, i + 1) - .then(resolve) - .catch(reject); - }, - Math.floor(Math.random() * 300), - ), - ); // Stagger the workers off to break up the clump that causes the race condition - } else { - logger.debug("Removed job from concurrency limit queue", { - teamId, - jobId: finalJob.job.id, - zeroDataRetention: finalJob.job.data?.zeroDataRetention, - i, - }); - } - } - - return finalJob; -} - -/** - * Called when a job associated with a concurrency queue is done. - * - * @param job The BullMQ job that is done. - */ -export async function concurrentJobDone(job: NuQJob) { - if (job.id && job.data && job.data.team_id) { - await removeConcurrencyLimitActiveJob(job.data.team_id, job.id); - await cleanOldConcurrencyLimitEntries(job.data.team_id); - - if (job.data.crawl_id) { - await removeCrawlConcurrencyLimitActiveJob(job.data.crawl_id, job.id); - await cleanOldCrawlConcurrencyLimitEntries(job.data.crawl_id); - } - - const maxTeamConcurrency = - ( - await getACUCTeam( - job.data.team_id, - false, - true, - job.data.is_extract ? RateLimiterMode.Extract : RateLimiterMode.Crawl, - ) - )?.concurrency ?? 2; - const currentActiveConcurrency = ( - await getConcurrencyLimitActiveJobs(job.data.team_id) - ).length; - - if (currentActiveConcurrency < maxTeamConcurrency) { - const nextJob = await getNextConcurrentJob(job.data.team_id); - if (nextJob !== null) { - await pushConcurrencyLimitActiveJob( - job.data.team_id, - nextJob.job.id, - 60 * 1000, - ); - - if (nextJob.job.data.crawl_id) { - await pushCrawlConcurrencyLimitActiveJob( - nextJob.job.data.crawl_id, - nextJob.job.id, - 60 * 1000, - ); - - const sc = await getCrawl(nextJob.job.data.crawl_id); - if (sc !== null && typeof sc.crawlerOptions?.delay === "number") { - await new Promise(resolve => - setTimeout(resolve, sc.crawlerOptions.delay * 1000), - ); - } - } - - abTestJob(nextJob.job.data); - - await scrapeQueue.addJob( - nextJob.job.id, - { - ...nextJob.job.data, - concurrencyLimitHit: true, - }, - { - priority: nextJob.job.priority, - listenable: nextJob.job.listenable, - ownerId: nextJob.job.data.team_id, - }, - ); - } - } - } -} diff --git a/apps/api/src/lib/crawl-redis.ts b/apps/api/src/lib/crawl-redis.ts index ff4ee837f..1b256ea67 100644 --- a/apps/api/src/lib/crawl-redis.ts +++ b/apps/api/src/lib/crawl-redis.ts @@ -49,14 +49,8 @@ export async function saveCrawl(id: string, crawl: StoredCrawl) { }); } -export async function recordRobotsBlocked( - crawlId: string, - url: string, -) { - await redisEvictConnection.sadd( - "crawl:" + crawlId + ":robots_blocked", - url, - ); +export async function recordRobotsBlocked(crawlId: string, url: string) { + await redisEvictConnection.sadd("crawl:" + crawlId + ":robots_blocked", url); await redisEvictConnection.expire( "crawl:" + crawlId + ":robots_blocked", 24 * 60 * 60, diff --git a/apps/api/src/routes/admin.ts b/apps/api/src/routes/admin.ts index 370e98291..3f2ae85e5 100644 --- a/apps/api/src/routes/admin.ts +++ b/apps/api/src/routes/admin.ts @@ -3,7 +3,6 @@ import { redisHealthController } from "../controllers/v0/admin/redis-health"; import { authMiddleware, blocklistMiddleware, wrap } from "./shared"; import { acucCacheClearController } from "../controllers/v0/admin/acuc-cache-clear"; import { checkFireEngine } from "../controllers/v0/admin/check-fire-engine"; -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"; @@ -13,6 +12,7 @@ import { } from "../controllers/v0/admin/metrics"; import { crawlCheckController } from "../controllers/v0/admin/crawl-check"; import { realtimeSearchController } from "../controllers/v2/f-search"; +import { migrateConcurrencyQueue } from "../controllers/v0/admin/concurrency-queue-migration"; export const adminRouter = express.Router(); @@ -31,11 +31,6 @@ adminRouter.get( wrap(checkFireEngine), ); -adminRouter.get( - `/admin/${process.env.BULL_AUTH_KEY}/cclog`, - wrap(cclogController), -); - adminRouter.get( `/admin/${process.env.BULL_AUTH_KEY}/zdrcleaner`, wrap(zdrcleanerController), @@ -70,3 +65,8 @@ adminRouter.post( `/admin/${process.env.BULL_AUTH_KEY}/fsearch`, wrap(realtimeSearchController), ); + +adminRouter.get( + `/admin/${process.env.BULL_AUTH_KEY}/migrate-concurrency-queue`, + wrap(migrateConcurrencyQueue), +); diff --git a/apps/api/src/services/indexing/index-worker.ts b/apps/api/src/services/indexing/index-worker.ts index e662e999b..74b5ad358 100644 --- a/apps/api/src/services/indexing/index-worker.ts +++ b/apps/api/src/services/indexing/index-worker.ts @@ -34,7 +34,7 @@ import { toV0CrawlerOptions, } from "../../controllers/v2/types"; import { StoredCrawl, crawlToCrawler, saveCrawl } from "../../lib/crawl-redis"; -import { _addScrapeJobToBullMQ } from "../queue-jobs"; +import { addScrapeJob } from "../queue-jobs"; import { BullMQOtel } from "bullmq-otel"; import { withSpan, setSpanAttributes } from "../../lib/otel-tracer"; @@ -492,7 +492,7 @@ const processPrecrawlJob = async (token: string, job: Job) => { await saveCrawl(crawlId, sc); - await _addScrapeJobToBullMQ( + await addScrapeJob( { url: url, mode: "kickoff" as const, diff --git a/apps/api/src/services/queue-jobs.ts b/apps/api/src/services/queue-jobs.ts index 99bd7d0de..6648c3f89 100644 --- a/apps/api/src/services/queue-jobs.ts +++ b/apps/api/src/services/queue-jobs.ts @@ -1,21 +1,8 @@ import { v4 as uuidv4 } from "uuid"; -import { NotificationType, RateLimiterMode, ScrapeJobData } from "../types"; -import { - cleanOldConcurrencyLimitEntries, - getConcurrencyLimitActiveJobs, - getConcurrencyQueueJobsCount, - getCrawlConcurrencyLimitActiveJobs, - pushConcurrencyLimitActiveJob, - pushConcurrencyLimitedJob, - pushCrawlConcurrencyLimitActiveJob, -} from "../lib/concurrency-limit"; +import { ScrapeJobData } from "../types"; import { logger as _logger } from "../lib/logger"; -import { sendNotificationWithCustomDays } from "./notification/email_notification"; -import { shouldSendConcurrencyLimitNotification } from "./notification/notification-check"; -import { getACUCTeam } from "../controllers/auth"; import { getJobFromGCS, removeJobFromGCS } from "../lib/gcs-jobs"; import { Document } from "../controllers/v1/types"; -import { getCrawl } from "../lib/crawl-redis"; import { Logger } from "winston"; import { ScrapeJobTimeoutError, TransportableError } from "../lib/error"; import { deserializeTransportableError } from "../lib/error-serde"; @@ -23,175 +10,6 @@ import { abTestJob } from "./ab-test"; import { NuQJob, scrapeQueue } from "./worker/nuq"; import { serializeTraceContext } from "../lib/otel-tracer"; -/** - * Checks if a job is a crawl or batch scrape based on its options - * @param options The job options containing crawlerOptions and crawl_id - * @returns true if the job is either a crawl or batch scrape - */ -function isCrawlOrBatchScrape(options: { - crawlerOptions?: any; - crawl_id?: string; -}): boolean { - // If crawlerOptions exists, it's a crawl - // If crawl_id exists but no crawlerOptions, it's a batch scrape - return !!options.crawlerOptions || !!options.crawl_id; -} - -async function _addScrapeJobToConcurrencyQueue( - webScraperOptions: any, - jobId: string, - priority: number = 0, - listenable: boolean = false, -) { - await pushConcurrencyLimitedJob( - webScraperOptions.team_id, - { - id: jobId, - data: webScraperOptions, - priority, - listenable, - }, - webScraperOptions.crawl_id - ? Infinity - : (webScraperOptions.scrapeOptions?.timeout ?? 60 * 1000), - ); -} - -export async function _addScrapeJobToBullMQ( - webScraperOptions: ScrapeJobData, - jobId: string, - priority: number = 0, - listenable: boolean = false, -): Promise> { - if (webScraperOptions.mode === "single_urls") { - abTestJob(webScraperOptions); - } - - if (webScraperOptions && webScraperOptions.team_id) { - await pushConcurrencyLimitActiveJob( - webScraperOptions.team_id, - jobId, - 60 * 1000, - ); // 60s default timeout - - if (webScraperOptions.crawl_id) { - const sc = await getCrawl(webScraperOptions.crawl_id); - if (sc?.crawlerOptions?.delay || sc?.maxConcurrency) { - await pushCrawlConcurrencyLimitActiveJob( - webScraperOptions.crawl_id, - jobId, - 60 * 1000, - ); - } - } - } - - return await scrapeQueue.addJob(jobId, webScraperOptions, { - priority, - listenable, - ownerId: webScraperOptions.team_id, - }); -} - -async function addScrapeJobRaw( - webScraperOptions: ScrapeJobData, - jobId: string, - priority: number = 0, - directToBullMQ: boolean = false, - listenable: boolean = false, -): Promise | null> { - let concurrencyLimited: "yes" | "yes-crawl" | "no" | null = null; - let currentActiveConcurrency = 0; - let maxConcurrency = 0; - - if (directToBullMQ) { - concurrencyLimited = "no"; - } else { - if (webScraperOptions.crawl_id) { - const crawl = await getCrawl(webScraperOptions.crawl_id); - const concurrencyLimit = !crawl - ? null - : crawl.crawlerOptions?.delay === undefined && - crawl.maxConcurrency === undefined - ? null - : (crawl.maxConcurrency ?? 1); - - if (concurrencyLimit !== null) { - const crawlConcurrency = ( - await getCrawlConcurrencyLimitActiveJobs(webScraperOptions.crawl_id) - ).length; - const freeSlots = Math.max(concurrencyLimit - crawlConcurrency, 0); - if (freeSlots === 0) { - concurrencyLimited = "yes-crawl"; - } - } - } - - if (concurrencyLimited === null) { - const now = Date.now(); - const maxConcurrency = - ( - await getACUCTeam( - webScraperOptions.team_id, - false, - true, - webScraperOptions.mode === "single_urls" && - webScraperOptions.from_extract - ? RateLimiterMode.Extract - : RateLimiterMode.Crawl, - ) - )?.concurrency ?? 2; - await cleanOldConcurrencyLimitEntries(webScraperOptions.team_id, now); - const currentActiveConcurrency = ( - await getConcurrencyLimitActiveJobs(webScraperOptions.team_id, now) - ).length; - concurrencyLimited = - currentActiveConcurrency >= maxConcurrency ? "yes" : "no"; - } - } - - if (concurrencyLimited === "yes" || concurrencyLimited === "yes-crawl") { - if (concurrencyLimited === "yes") { - // Detect if they hit their concurrent limit - // If above by 2x, send them an email - // No need to 2x as if there are more than the max concurrency in the concurrency queue, it is already 2x - const concurrencyQueueJobs = await getConcurrencyQueueJobsCount( - webScraperOptions.team_id, - ); - if (concurrencyQueueJobs > maxConcurrency) { - // logger.info("Concurrency limited 2x (single) - ", "Concurrency queue jobs: ", concurrencyQueueJobs, "Max concurrency: ", maxConcurrency, "Team ID: ", webScraperOptions.team_id); - - // Only send notification if it's not a crawl or batch scrape - const shouldSendNotification = - await shouldSendConcurrencyLimitNotification( - webScraperOptions.team_id, - ); - if (shouldSendNotification) { - sendNotificationWithCustomDays( - webScraperOptions.team_id, - NotificationType.CONCURRENCY_LIMIT_REACHED, - 15, - false, - true, - ).catch(error => { - _logger.error( - "Error sending notification (concurrency limit reached)", - { error }, - ); - }); - } - } - } - - webScraperOptions.concurrencyLimited = true; - - await _addScrapeJobToConcurrencyQueue(webScraperOptions, jobId, priority, listenable); - return null; - } else { - return await _addScrapeJobToBullMQ(webScraperOptions, jobId, priority, listenable); - } -} - export async function addScrapeJob( webScraperOptions: ScrapeJobData, jobId: string = uuidv4(), @@ -206,13 +24,22 @@ export async function addScrapeJob( traceContext, }; - return await addScrapeJobRaw( - optionsWithTrace, - jobId, + if (webScraperOptions.mode === "single_urls") { + abTestJob(webScraperOptions); + } + + return await scrapeQueue.addJob(jobId, webScraperOptions, { priority, - directToBullMQ, listenable, - ); + ownerId: webScraperOptions.team_id, + groupId: webScraperOptions.crawl_id ?? undefined, + timesOutAt: webScraperOptions.crawl_id + ? undefined + : new Date( + Date.now() + + ((webScraperOptions as any).scrapeOptions?.timeout ?? 300) * 1000, + ), + }); } export async function addScrapeJobs( @@ -228,173 +55,27 @@ export async function addScrapeJobs( // Capture trace context for all jobs const traceContext = serializeTraceContext(); - const jobsByTeam = new Map< - string, - { - jobId: string; - data: ScrapeJobData; - priority: number; - listenable?: boolean; - }[] - >(); - for (const job of jobs) { - if (!jobsByTeam.has(job.data.team_id)) { - jobsByTeam.set(job.data.team_id, []); + if (job.data.mode === "single_urls") { + abTestJob(job.data); } - jobsByTeam.get(job.data.team_id)!.push(job); } - for (const [teamId, teamJobs] of jobsByTeam) { - // == Buckets for jobs == - let jobsForcedToCQ: { - data: ScrapeJobData; - jobId: string; - priority: number; - listenable?: boolean; - }[] = []; - - let jobsPotentiallyInCQ: { - data: ScrapeJobData; - jobId: string; - priority: number; - listenable?: boolean; - }[] = []; - - // == Select jobs by crawl ID == - const jobsByCrawlID = new Map< - string, - { - data: ScrapeJobData; - jobId: string; - priority: number; - listenable?: boolean; - }[] - >(); - - const jobsWithoutCrawlID: { - data: ScrapeJobData; - jobId: string; - priority: number; - listenable?: boolean; - }[] = []; - - for (const job of teamJobs) { - if (job.data.crawl_id) { - if (!jobsByCrawlID.has(job.data.crawl_id)) { - jobsByCrawlID.set(job.data.crawl_id, []); - } - jobsByCrawlID.get(job.data.crawl_id)!.push(job); - } else { - jobsWithoutCrawlID.push(job); - } - } - - // == Select jobs by crawl ID == - for (const [crawlID, crawlJobs] of jobsByCrawlID) { - const crawl = await getCrawl(crawlID); - const concurrencyLimit = !crawl - ? null - : crawl.crawlerOptions?.delay === undefined && - crawl.maxConcurrency === undefined - ? null - : (crawl.maxConcurrency ?? 1); - - if (concurrencyLimit === null) { - // All jobs may be in the CQ depending on the global team concurrency limit - jobsPotentiallyInCQ.push(...crawlJobs); - } else { - const crawlConcurrency = ( - await getCrawlConcurrencyLimitActiveJobs(crawlID) - ).length; - const freeSlots = Math.max(concurrencyLimit - crawlConcurrency, 0); - - // The first n jobs may be in the CQ depending on the global team concurrency limit - jobsPotentiallyInCQ.push(...crawlJobs.slice(0, freeSlots)); - - // Every job after that must be in the CQ, as the crawl concurrency limit has been reached - jobsForcedToCQ.push(...crawlJobs.slice(freeSlots)); - } - } - - // All jobs without a crawl ID may be in the CQ depending on the global team concurrency limit - jobsPotentiallyInCQ.push(...jobsWithoutCrawlID); - - const now = Date.now(); - const maxConcurrency = - ( - await getACUCTeam( - teamId, - false, - true, - jobs[0].data.mode === "single_urls" && jobs[0].data.from_extract - ? RateLimiterMode.Extract - : RateLimiterMode.Crawl, - ) - )?.concurrency ?? 2; - await cleanOldConcurrencyLimitEntries(teamId, now); - - const currentActiveConcurrency = ( - await getConcurrencyLimitActiveJobs(teamId, now) - ).length; - - const countCanBeDirectlyAdded = Math.max( - maxConcurrency - currentActiveConcurrency, - 0, - ); - - const addToBull = jobsPotentiallyInCQ.slice(0, countCanBeDirectlyAdded); - const addToCQ = jobsPotentiallyInCQ - .slice(countCanBeDirectlyAdded) - .concat(jobsForcedToCQ); - - // equals 2x the max concurrency - if (jobsPotentiallyInCQ.length - countCanBeDirectlyAdded > maxConcurrency) { - // logger.info(`Concurrency limited 2x (multiple) - Concurrency queue jobs: ${addToCQ.length} Max concurrency: ${maxConcurrency} Team ID: ${jobs[0].data.team_id}`); - // Only send notification if it's not a crawl or batch scrape - if (!isCrawlOrBatchScrape(jobs[0].data)) { - const shouldSendNotification = - await shouldSendConcurrencyLimitNotification(jobs[0].data.team_id); - if (shouldSendNotification) { - sendNotificationWithCustomDays( - jobs[0].data.team_id, - NotificationType.CONCURRENCY_LIMIT_REACHED, - 15, - false, - true, - ).catch(error => { - _logger.error( - "Error sending notification (concurrency limit reached)", - { error }, - ); - }); - } - } - } - - await Promise.all( - addToCQ.map(async job => { - const size = JSON.stringify(job.data).length; - await _addScrapeJobToConcurrencyQueue( - { ...job.data, traceContext }, - job.jobId, - job.priority, - job.listenable, - ); - }), - ); - - await Promise.all( - addToBull.map(async job => { - await _addScrapeJobToBullMQ( - { ...job.data, traceContext }, - job.jobId, - job.priority, - job.listenable, - ); - }), - ); - } + await scrapeQueue.addJobs( + jobs.map(job => ({ + data: { + ...job.data, + traceContext, + }, + id: job.jobId, + options: { + priority: job.priority, + listenable: job.listenable, + ownerId: job.data.team_id, + groupId: job.data.crawl_id ?? undefined, + }, + })), + ); } export async function waitForJob( diff --git a/apps/api/src/services/worker/crawl-logic.ts b/apps/api/src/services/worker/crawl-logic.ts index af3489d44..4131f4388 100644 --- a/apps/api/src/services/worker/crawl-logic.ts +++ b/apps/api/src/services/worker/crawl-logic.ts @@ -18,7 +18,7 @@ import { v4 as uuidv4 } from "uuid"; import { addScrapeJobs } from "../queue-jobs"; import { getJobs } from "../../controllers/v1/crawl-status"; import { logJob } from "../logging/log_job"; -import { createWebhookSender, WebhookEvent } from "../webhook"; +import { createWebhookSender, WebhookEvent } from "../webhook/index"; import { hasFormatOfType } from "../../lib/format-utils"; import type { NuQJob } from "./nuq"; import { ScrapeJobData } from "../../types"; diff --git a/apps/api/src/services/worker/nuq.ts b/apps/api/src/services/worker/nuq.ts index 2cf95c4b5..b4aedab85 100644 --- a/apps/api/src/services/worker/nuq.ts +++ b/apps/api/src/services/worker/nuq.ts @@ -30,6 +30,8 @@ export type NuQJob = { failedReason?: string; lock?: string; ownerId?: string; + groupId?: string; + timesOutAt?: Date; }; const listenChannelId = process.env.NUQ_POD_NAME ?? "main"; @@ -37,19 +39,21 @@ const listenChannelId = process.env.NUQ_POD_NAME ?? "main"; // === Queue type NuQOptions = { - concurrencyLimit?: false | "per-owner"; + concurrencyLimit?: false | "per-owner" | "per-owner-per-group"; }; type NuQJobOptions = { listenable?: boolean; priority?: number; ownerId?: string; + groupId?: string; + timesOutAt?: Date; }; class NuQ { constructor( public readonly queueName: string, - public readonly options: NuQOptions = {}, + public readonly options: NuQOptions = { concurrencyLimit: false }, ) {} // === Listener @@ -281,7 +285,7 @@ class NuQ { await this.startSender(); if (this.sender) { - await this.sender.channel.sendToQueue( + this.sender.channel.sendToQueue( this.queueName + ".listen." + listenChannelId, Buffer.from(status, "utf8"), { @@ -301,7 +305,7 @@ class NuQ { await this.startSender(); if (this.sender) { - await this.sender.channel.sendToQueue( + this.sender.channel.sendToQueue( this.queueName + ".prefetch", Buffer.from(JSON.stringify(job), "utf8"), { @@ -330,6 +334,8 @@ class NuQ { "failedreason", "lock", "owner_id", + "group_id", + "times_out_at", ]; private rowToJob(row: any): NuQJob | null { @@ -346,6 +352,8 @@ class NuQ { failedReason: row.failedreason ?? undefined, lock: row.lock ?? undefined, ownerId: row.owner_id ?? undefined, + groupId: row.group_id ?? undefined, + timesOutAt: row.times_out_at ? new Date(row.times_out_at) : undefined, }; } @@ -517,6 +525,67 @@ class NuQ { } // === Producer + public async tryAddJob( + id: string, + data: JobData, + options: NuQJobOptions = {}, + ): Promise | null> { + return withSpan("nuq.tryAddJob", async span => { + const bareOwnerId = options.ownerId ?? undefined; + const normalizedOwnerId = bareOwnerId + ? uuidValidate(bareOwnerId) + ? bareOwnerId + : uuidv5(bareOwnerId, "b208cbac-8bdf-4599-bf17-da78426e3f7c") // preview namespace + : null; + + setSpanAttributes(span, { + "nuq.queue_name": this.queueName, + "nuq.job_id": id, + "nuq.priority": options.priority ?? 0, + "nuq.zero_data_retention": (data as any)?.zeroDataRetention ?? false, + "nuq.listenable": options.listenable ?? false, + }); + + const start = Date.now(); + try { + const result = this.rowToJob( + ( + await nuqPool.query( + `INSERT INTO ${this.queueName} (id, data, priority, listen_channel_id, owner_id, group_id, times_out_at) VALUES ($1, $2, $3, $4, $5, $6, $7) ON CONFLICT DO NOTHING RETURNING ${this.jobReturning.join(", ")};`, + [ + id, + data, + options.priority ?? 0, + options.listenable ? listenChannelId : null, + normalizedOwnerId ?? null, + options.groupId ?? null, + options.timesOutAt ? options.timesOutAt.toISOString() : null, + ], + ) + ).rows[0], + )!; + + setSpanAttributes(span, { + "nuq.job_created": result !== null, + }); + + return result; + } finally { + const duration = Date.now() - start; + setSpanAttributes(span, { + "nuq.duration_ms": duration, + }); + logger.info("nuqAddJob metrics", { + module: "nuq/metrics", + method: "nuqAddJob", + duration, + scrapeId: id, + zeroDataRetention: (data as any)?.zeroDataRetention ?? false, + }); + } + }); + } + public async addJob( id: string, data: JobData, @@ -543,13 +612,15 @@ class NuQ { const result = this.rowToJob( ( await nuqPool.query( - `INSERT INTO ${this.queueName} (id, data, priority, listen_channel_id, owner_id) VALUES ($1, $2, $3, $4, $5) RETURNING ${this.jobReturning.join(", ")};`, + `INSERT INTO ${this.queueName} (id, data, priority, listen_channel_id, owner_id, group_id, times_out_at) VALUES ($1, $2, $3, $4, $5, $6, $7) RETURNING ${this.jobReturning.join(", ")};`, [ id, data, options.priority ?? 0, options.listenable ? listenChannelId : null, normalizedOwnerId ?? null, + options.groupId ?? null, + options.timesOutAt ? options.timesOutAt.toISOString() : null, ], ) ).rows[0], @@ -576,6 +647,94 @@ class NuQ { }); } + public async addJobs( + jobs: Array<{ + id: string; + data: JobData; + options?: NuQJobOptions; + }>, + _logger: Logger = logger, + ): Promise[]> { + if (jobs.length === 0) return []; + + return withSpan("nuq.addJobs", async span => { + setSpanAttributes(span, { + "nuq.queue_name": this.queueName, + "nuq.job_count": jobs.length, + }); + + const start = Date.now(); + try { + // Prepare arrays for bulk insert + const ids: string[] = []; + const dataArray: JobData[] = []; + const priorities: number[] = []; + const listenChannelIds: (string | null)[] = []; + const ownerIds: (string | null)[] = []; + const groupIds: (string | null)[] = []; + const timesOutAts: (string | null)[] = []; + + for (const job of jobs) { + const bareOwnerId = job.options?.ownerId ?? undefined; + const normalizedOwnerId = bareOwnerId + ? uuidValidate(bareOwnerId) + ? bareOwnerId + : uuidv5(bareOwnerId, "b208cbac-8bdf-4599-bf17-da78426e3f7c") // preview namespace + : null; + + ids.push(job.id); + dataArray.push(job.data); + priorities.push(job.options?.priority ?? 0); + listenChannelIds.push( + job.options?.listenable ? listenChannelId : null, + ); + ownerIds.push(normalizedOwnerId); + groupIds.push(job.options?.groupId ?? null); + timesOutAts.push( + job.options?.timesOutAt + ? job.options.timesOutAt.toISOString() + : null, + ); + } + + // Bulk insert using UNNEST + const result = await nuqPool.query( + `INSERT INTO ${this.queueName} (id, data, priority, listen_channel_id, owner_id, group_id, times_out_at) + SELECT * FROM UNNEST($1::uuid[], $2::jsonb[], $3::int[], $4::text[], $5::uuid[], $6::uuid[], $7::timestamptz[]) + RETURNING ${this.jobReturning.join(", ")};`, + [ + ids, + dataArray, + priorities, + listenChannelIds, + ownerIds, + groupIds, + timesOutAts, + ], + ); + + const createdJobs = result.rows.map(row => this.rowToJob(row)!); + + setSpanAttributes(span, { + "nuq.jobs_created": createdJobs.length, + }); + + return createdJobs; + } finally { + const duration = Date.now() - start; + setSpanAttributes(span, { + "nuq.duration_ms": duration, + }); + _logger.info("nuqAddJobs metrics", { + module: "nuq/metrics", + method: "nuqAddJobs", + duration, + jobCount: jobs.length, + }); + } + }); + } + private readonly nuqWaitMode = process.env.NUQ_WAIT_MODE === "listen" || process.env.NUQ_RABBITMQ_URL ? ("listen" as const) @@ -694,49 +853,62 @@ class NuQ { public async prefetchJobs(_logger: Logger = logger): Promise { const start = Date.now(); try { - // With per-owner concurrency limiting: smaller batch size reduces lock contention. - // Use blocking advisory locks to ensure jobs aren't skipped. - const queryGetNext = ` - SELECT ${this.jobReturning.map(x => `${this.queueName}.${x}`).join(", ")} FROM ${this.queueName} - WHERE ${this.queueName}.status = 'queued'::nuq.job_status - ORDER BY ${this.queueName}.priority ASC, ${this.queueName}.created_at ASC - FOR UPDATE SKIP LOCKED LIMIT 100 - `; - - const queryUpdateStatus = ` - UPDATE ${this.queueName} q - SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() - FROM next - WHERE q.id = next.id - RETURNING ${this.jobReturning.map(x => `q.${x}`).join(", ")} - `; - let updateQuery: string; if (this.options.concurrencyLimit === "per-owner") { updateQuery = ` - WITH next_initial AS ( - ${queryGetNext} + WITH queued_owners AS ( + SELECT DISTINCT owner_id + FROM ${this.queueName} + WHERE status = 'queued'::nuq.job_status + ), + available_capacity AS ( + SELECT + qo.owner_id, + CASE + WHEN qo.owner_id IS NULL THEN 999999 + WHEN oc.max_concurrency IS NOT NULL THEN GREATEST(0, oc.max_concurrency - oc.current_concurrency) + ELSE ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(qo.owner_id) + END as slots + FROM queued_owners qo + LEFT JOIN ${this.queueName}_owner_concurrency oc ON qo.owner_id = oc.id ), distinct_owners AS ( - SELECT DISTINCT owner_id - FROM next_initial - WHERE owner_id IS NOT NULL - ORDER BY owner_id -- Deterministic lock order prevents deadlocks + SELECT owner_id + FROM available_capacity + WHERE owner_id IS NOT NULL AND slots > 0 + ORDER BY owner_id ), - acquired_locks AS ( + acquired_owner_locks AS ( SELECT owner_id, pg_advisory_xact_lock(hashtext(owner_id::text)) as dummy FROM distinct_owners ), - next AS ( - SELECT n.* - FROM next_initial n - WHERE n.owner_id IS NULL - OR EXISTS (SELECT 1 FROM acquired_locks WHERE owner_id = n.owner_id) + available_capacity_locked AS ( + SELECT * + FROM available_capacity ac + WHERE ac.slots > 0 + AND (ac.owner_id IS NULL OR EXISTS (SELECT 1 FROM acquired_owner_locks WHERE acquired_owner_locks.owner_id = ac.owner_id)) + ), + selected_jobs AS ( + SELECT j.id + FROM available_capacity_locked ac + CROSS JOIN LATERAL ( + SELECT j.id + FROM ${this.queueName} j + WHERE j.status = 'queued'::nuq.job_status + AND j.owner_id IS NOT DISTINCT FROM ac.owner_id + ORDER BY j.priority ASC, j.created_at ASC + FOR UPDATE SKIP LOCKED + LIMIT ac.slots + ) j + LIMIT 100 ), updated AS ( - ${queryUpdateStatus} + UPDATE ${this.queueName} q + SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() + WHERE q.id IN (SELECT id FROM selected_jobs) + RETURNING ${this.jobReturning.map(x => `q.${x}`).join(", ")} ), owner_counts AS ( SELECT owner_id, COUNT(*)::int8 as job_count @@ -748,11 +920,9 @@ class NuQ { SELECT oc.owner_id, oc.job_count, - COALESCE( - (SELECT max_concurrency FROM ${this.queueName}_owner_concurrency WHERE id = oc.owner_id), - ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(oc.owner_id) - ) as max_concurrency + COALESCE(ocon.max_concurrency, ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(oc.owner_id)) as max_concurrency FROM owner_counts oc + LEFT JOIN ${this.queueName}_owner_concurrency ocon ON oc.owner_id = ocon.id ), owner_upsert AS ( INSERT INTO ${this.queueName}_owner_concurrency (id, current_concurrency, max_concurrency) @@ -763,14 +933,137 @@ class NuQ { current_concurrency = ${this.queueName}_owner_concurrency.current_concurrency + EXCLUDED.current_concurrency, max_concurrency = EXCLUDED.max_concurrency ) - SELECT * FROM updated; + SELECT ${this.jobReturning.map(x => `updated.${x}`).join(", ")} FROM updated; + `; + } else if (this.options.concurrencyLimit === "per-owner-per-group") { + updateQuery = ` + WITH queued_combinations AS ( + SELECT DISTINCT owner_id, group_id + FROM ${this.queueName} + WHERE status = 'queued'::nuq.job_status + ), + available_capacity AS ( + SELECT + qc.owner_id, + qc.group_id, + CASE + WHEN qc.owner_id IS NULL AND qc.group_id IS NULL THEN 999999 + WHEN qc.owner_id IS NULL AND gc.max_concurrency IS NOT NULL THEN GREATEST(0, gc.max_concurrency - gc.current_concurrency) + WHEN qc.owner_id IS NULL THEN 999999 + WHEN qc.group_id IS NULL AND oc.max_concurrency IS NOT NULL THEN GREATEST(0, oc.max_concurrency - oc.current_concurrency) + WHEN qc.group_id IS NULL THEN ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(qc.owner_id) + ELSE LEAST( + CASE WHEN oc.max_concurrency IS NOT NULL THEN GREATEST(0, oc.max_concurrency - oc.current_concurrency) ELSE ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(qc.owner_id) END, + CASE WHEN gc.max_concurrency IS NOT NULL THEN GREATEST(0, gc.max_concurrency - gc.current_concurrency) ELSE 999999 END + ) + END as slots + FROM queued_combinations qc + LEFT JOIN ${this.queueName}_owner_concurrency oc ON qc.owner_id = oc.id + LEFT JOIN ${this.queueName}_group_concurrency gc ON qc.group_id = gc.id + ), + distinct_owners AS ( + SELECT DISTINCT owner_id + FROM available_capacity + WHERE owner_id IS NOT NULL AND slots > 0 + ORDER BY owner_id + ), + acquired_owner_locks AS ( + SELECT + owner_id, + pg_advisory_xact_lock(hashtext(owner_id::text)) as dummy + FROM distinct_owners + ), + distinct_groups AS ( + SELECT DISTINCT group_id + FROM available_capacity + WHERE group_id IS NOT NULL AND slots > 0 + ORDER BY group_id + ), + acquired_group_locks AS ( + SELECT + group_id, + pg_advisory_xact_lock(hashtext(group_id::text)) as dummy + FROM distinct_groups + ), + available_capacity_locked AS ( + SELECT * + FROM available_capacity ac + WHERE ac.slots > 0 + AND (ac.owner_id IS NULL OR EXISTS (SELECT 1 FROM acquired_owner_locks WHERE acquired_owner_locks.owner_id = ac.owner_id)) + AND (ac.group_id IS NULL OR EXISTS (SELECT 1 FROM acquired_group_locks WHERE acquired_group_locks.group_id = ac.group_id)) + ), + selected_jobs AS ( + SELECT j.id + FROM available_capacity_locked ac + CROSS JOIN LATERAL ( + SELECT j.id + FROM ${this.queueName} j + WHERE j.status = 'queued'::nuq.job_status + AND j.owner_id IS NOT DISTINCT FROM ac.owner_id + AND j.group_id IS NOT DISTINCT FROM ac.group_id + ORDER BY j.priority ASC, j.created_at ASC + FOR UPDATE SKIP LOCKED + LIMIT ac.slots + ) j + LIMIT 100 + ), + updated AS ( + UPDATE ${this.queueName} q + SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() + WHERE q.id IN (SELECT id FROM selected_jobs) + RETURNING ${this.jobReturning.map(x => `q.${x}`).join(", ")} + ), + owner_counts AS ( + SELECT owner_id, COUNT(*)::int8 as job_count + FROM updated + WHERE owner_id IS NOT NULL + GROUP BY owner_id + ), + owner_counts_with_max AS ( + SELECT + oc.owner_id, + oc.job_count, + COALESCE(ocon.max_concurrency, ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(oc.owner_id)) as max_concurrency + FROM owner_counts oc + LEFT JOIN ${this.queueName}_owner_concurrency ocon ON oc.owner_id = ocon.id + ), + owner_upsert AS ( + INSERT INTO ${this.queueName}_owner_concurrency (id, current_concurrency, max_concurrency) + SELECT owner_id, job_count, max_concurrency + FROM owner_counts_with_max + ON CONFLICT (id) + DO UPDATE SET + current_concurrency = ${this.queueName}_owner_concurrency.current_concurrency + EXCLUDED.current_concurrency, + max_concurrency = EXCLUDED.max_concurrency + ), + group_counts AS ( + SELECT group_id, COUNT(*)::int8 as job_count + FROM updated + WHERE group_id IS NOT NULL + GROUP BY group_id + ), + group_update AS ( + UPDATE ${this.queueName}_group_concurrency + SET current_concurrency = ${this.queueName}_group_concurrency.current_concurrency + group_counts.job_count + FROM group_counts + WHERE ${this.queueName}_group_concurrency.id = group_counts.group_id + ) + SELECT ${this.jobReturning.map(x => `updated.${x}`).join(", ")} FROM updated; `; } else { updateQuery = ` - WITH next AS ( - ${queryGetNext} + WITH selected_jobs AS ( + SELECT j.id + FROM ${this.queueName} j + WHERE j.status = 'queued'::nuq.job_status + ORDER BY j.priority ASC, j.created_at ASC + FOR UPDATE SKIP LOCKED + LIMIT 100 ) - ${queryUpdateStatus} + UPDATE ${this.queueName} q + SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() + WHERE q.id IN (SELECT id FROM selected_jobs) + RETURNING ${this.jobReturning.join(", ")}; `; } @@ -828,56 +1121,69 @@ class NuQ { } } - const queryGetNext = ` - SELECT ${this.jobReturning.join(", ")} FROM ${this.queueName} - WHERE ${this.queueName}.status = 'queued'::nuq.job_status - ORDER BY ${this.queueName}.priority ASC, ${this.queueName}.created_at ASC - FOR UPDATE SKIP LOCKED LIMIT 1 - `; - - const queryUpdateStatus = ` - UPDATE ${this.queueName} q - SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() - FROM next - WHERE q.id = next.id - RETURNING ${this.jobReturning.map(x => `q.${x}`).join(", ")} - `; - let updateQuery: string; if (this.options.concurrencyLimit === "per-owner") { updateQuery = ` - WITH next_initial AS ( - ${queryGetNext} + WITH queued_owners AS ( + SELECT DISTINCT owner_id + FROM ${this.queueName} + WHERE status = 'queued'::nuq.job_status + ), + available_capacity AS ( + SELECT + qo.owner_id, + CASE + WHEN qo.owner_id IS NULL THEN 999999 + WHEN oc.max_concurrency IS NOT NULL THEN GREATEST(0, oc.max_concurrency - oc.current_concurrency) + ELSE ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(qo.owner_id) + END as slots + FROM queued_owners qo + LEFT JOIN ${this.queueName}_owner_concurrency oc ON qo.owner_id = oc.id ), distinct_owners AS ( - SELECT DISTINCT owner_id - FROM next_initial - WHERE owner_id IS NOT NULL - ORDER BY owner_id -- Deterministic lock order prevents deadlocks + SELECT owner_id + FROM available_capacity + WHERE owner_id IS NOT NULL AND slots > 0 + ORDER BY owner_id ), - acquired_locks AS ( + acquired_owner_locks AS ( SELECT owner_id, pg_advisory_xact_lock(hashtext(owner_id::text)) as dummy FROM distinct_owners ), - next AS ( - SELECT n.* - FROM next_initial n - WHERE n.owner_id IS NULL - OR EXISTS (SELECT 1 FROM acquired_locks WHERE owner_id = n.owner_id) + available_capacity_locked AS ( + SELECT * + FROM available_capacity ac + WHERE ac.slots > 0 + AND (ac.owner_id IS NULL OR EXISTS (SELECT 1 FROM acquired_owner_locks WHERE acquired_owner_locks.owner_id = ac.owner_id)) + ), + selected_jobs AS ( + SELECT j.id + FROM available_capacity_locked ac + CROSS JOIN LATERAL ( + SELECT j.id + FROM ${this.queueName} j + WHERE j.status = 'queued'::nuq.job_status + AND j.owner_id IS NOT DISTINCT FROM ac.owner_id + ORDER BY j.priority ASC, j.created_at ASC + FOR UPDATE SKIP LOCKED + LIMIT ac.slots + ) j + LIMIT 1 ), updated AS ( - ${queryUpdateStatus} + UPDATE ${this.queueName} q + SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() + WHERE q.id IN (SELECT id FROM selected_jobs) + RETURNING ${this.jobReturning.map(x => `q.${x}`).join(", ")} ), updated_with_max AS ( SELECT u.*, - COALESCE( - (SELECT max_concurrency FROM ${this.queueName}_owner_concurrency WHERE id = u.owner_id), - ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(u.owner_id) - ) as max_concurrency + COALESCE(ocon.max_concurrency, ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(u.owner_id)) as max_concurrency FROM updated u + LEFT JOIN ${this.queueName}_owner_concurrency ocon ON u.owner_id = ocon.id WHERE u.owner_id IS NOT NULL ), owner_increment AS ( @@ -889,14 +1195,137 @@ class NuQ { current_concurrency = ${this.queueName}_owner_concurrency.current_concurrency + 1, max_concurrency = EXCLUDED.max_concurrency ) - SELECT * FROM updated; + SELECT ${this.jobReturning.map(x => `updated.${x}`).join(", ")} FROM updated; + `; + } else if (this.options.concurrencyLimit === "per-owner-per-group") { + updateQuery = ` + WITH queued_combinations AS ( + SELECT DISTINCT owner_id, group_id + FROM ${this.queueName} + WHERE status = 'queued'::nuq.job_status + ), + available_capacity AS ( + SELECT + qc.owner_id, + qc.group_id, + CASE + WHEN qc.owner_id IS NULL AND qc.group_id IS NULL THEN 999999 + WHEN qc.owner_id IS NULL AND gc.max_concurrency IS NOT NULL THEN GREATEST(0, gc.max_concurrency - gc.current_concurrency) + WHEN qc.owner_id IS NULL THEN 999999 + WHEN qc.group_id IS NULL AND oc.max_concurrency IS NOT NULL THEN GREATEST(0, oc.max_concurrency - oc.current_concurrency) + WHEN qc.group_id IS NULL THEN ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(qc.owner_id) + ELSE LEAST( + CASE WHEN oc.max_concurrency IS NOT NULL THEN GREATEST(0, oc.max_concurrency - oc.current_concurrency) ELSE ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(qc.owner_id) END, + CASE WHEN gc.max_concurrency IS NOT NULL THEN GREATEST(0, gc.max_concurrency - gc.current_concurrency) ELSE 999999 END + ) + END as slots + FROM queued_combinations qc + LEFT JOIN ${this.queueName}_owner_concurrency oc ON qc.owner_id = oc.id + LEFT JOIN ${this.queueName}_group_concurrency gc ON qc.group_id = gc.id + ), + distinct_owners AS ( + SELECT DISTINCT owner_id + FROM available_capacity + WHERE owner_id IS NOT NULL AND slots > 0 + ORDER BY owner_id + ), + acquired_owner_locks AS ( + SELECT + owner_id, + pg_advisory_xact_lock(hashtext(owner_id::text)) as dummy + FROM distinct_owners + ), + distinct_groups AS ( + SELECT DISTINCT group_id + FROM available_capacity + WHERE group_id IS NOT NULL AND slots > 0 + ORDER BY group_id + ), + acquired_group_locks AS ( + SELECT + group_id, + pg_advisory_xact_lock(hashtext(group_id::text)) as dummy + FROM distinct_groups + ), + available_capacity_locked AS ( + SELECT * + FROM available_capacity ac + WHERE ac.slots > 0 + AND (ac.owner_id IS NULL OR EXISTS (SELECT 1 FROM acquired_owner_locks WHERE owner_id = ac.owner_id)) + AND (ac.group_id IS NULL OR EXISTS (SELECT 1 FROM acquired_group_locks WHERE group_id = ac.group_id)) + ), + selected_jobs AS ( + SELECT j.id + FROM available_capacity_locked ac + CROSS JOIN LATERAL ( + SELECT j.id + FROM ${this.queueName} j + WHERE j.status = 'queued'::nuq.job_status + AND j.owner_id IS NOT DISTINCT FROM ac.owner_id + AND j.group_id IS NOT DISTINCT FROM ac.group_id + ORDER BY j.priority ASC, j.created_at ASC + FOR UPDATE SKIP LOCKED + LIMIT ac.slots + ) j + LIMIT 1 + ), + updated AS ( + UPDATE ${this.queueName} q + SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() + WHERE q.id IN (SELECT id FROM selected_jobs) + RETURNING ${this.jobReturning.map(x => `q.${x}`).join(", ")} + ), + owner_counts AS ( + SELECT owner_id, COUNT(*)::int8 as job_count + FROM updated + WHERE owner_id IS NOT NULL + GROUP BY owner_id + ), + owner_counts_with_max AS ( + SELECT + oc.owner_id, + oc.job_count, + COALESCE(ocon.max_concurrency, ${this.queueName.replaceAll(".", "_")}_owner_resolve_max_concurrency(oc.owner_id)) as max_concurrency + FROM owner_counts oc + LEFT JOIN ${this.queueName}_owner_concurrency ocon ON oc.owner_id = ocon.id + ), + owner_upsert AS ( + INSERT INTO ${this.queueName}_owner_concurrency (id, current_concurrency, max_concurrency) + SELECT owner_id, job_count, max_concurrency + FROM owner_counts_with_max + ON CONFLICT (id) + DO UPDATE SET + current_concurrency = ${this.queueName}_owner_concurrency.current_concurrency + EXCLUDED.current_concurrency, + max_concurrency = EXCLUDED.max_concurrency + ), + group_counts AS ( + SELECT group_id, COUNT(*)::int8 as job_count + FROM updated + WHERE group_id IS NOT NULL + GROUP BY group_id + ), + group_update AS ( + UPDATE ${this.queueName}_group_concurrency + SET current_concurrency = ${this.queueName}_group_concurrency.current_concurrency + group_counts.job_count + FROM group_counts + WHERE ${this.queueName}_group_concurrency.id = group_counts.group_id + ) + SELECT ${this.jobReturning.map(x => `updated.${x}`).join(", ")} FROM updated; `; } else { updateQuery = ` - WITH next AS ( - ${queryGetNext} + WITH selected_jobs AS ( + SELECT j.id + FROM ${this.queueName} j + WHERE j.status = 'queued'::nuq.job_status + ORDER BY j.priority ASC, j.created_at ASC + FOR UPDATE SKIP LOCKED + LIMIT 1 ) - ${queryUpdateStatus} + UPDATE ${this.queueName} q + SET status = 'active'::nuq.job_status, lock = gen_random_uuid(), locked_at = now() + WHERE q.id IN (SELECT id FROM selected_jobs) + RETURNING ${this.jobReturning.join(", ")}; `; } @@ -967,6 +1396,28 @@ class NuQ { ) SELECT * FROM updated; `; + } else if (this.options.concurrencyLimit === "per-owner-per-group") { + updateQuery = ` + WITH updated AS ( + UPDATE ${this.queueName} + SET status = 'completed'::nuq.job_status, lock = null, locked_at = null, finished_at = now(), returnvalue = $3 + WHERE id = $1 AND lock = $2 + RETURNING id, listen_channel_id, owner_id, group_id + ), + owner_decrement AS ( + UPDATE ${this.queueName}_owner_concurrency + SET current_concurrency = GREATEST(0, current_concurrency - 1) + FROM updated + WHERE updated.owner_id IS NOT NULL AND ${this.queueName}_owner_concurrency.id = updated.owner_id + ), + group_decrement AS ( + UPDATE ${this.queueName}_group_concurrency + SET current_concurrency = GREATEST(0, current_concurrency - 1) + FROM updated + WHERE updated.group_id IS NOT NULL AND ${this.queueName}_group_concurrency.id = updated.group_id + ) + SELECT * FROM updated; + `; } else { updateQuery = `UPDATE ${this.queueName} SET status = 'completed'::nuq.job_status, lock = null, locked_at = null, finished_at = now(), returnvalue = $3 WHERE id = $1 AND lock = $2 RETURNING id, listen_channel_id;`; } @@ -1039,7 +1490,6 @@ class NuQ { try { let updateQuery: string; if (this.options.concurrencyLimit === "per-owner") { - // Update job and decrement owner concurrency in one query updateQuery = ` WITH updated AS ( UPDATE ${this.queueName} @@ -1055,6 +1505,28 @@ class NuQ { ) SELECT * FROM updated; `; + } else if (this.options.concurrencyLimit === "per-owner-per-group") { + updateQuery = ` + WITH updated AS ( + UPDATE ${this.queueName} + SET status = 'failed'::nuq.job_status, lock = null, locked_at = null, finished_at = now(), failedreason = $3 + WHERE id = $1 AND lock = $2 + RETURNING id, listen_channel_id, owner_id, group_id + ), + owner_decrement AS ( + UPDATE ${this.queueName}_owner_concurrency + SET current_concurrency = GREATEST(0, current_concurrency - 1) + FROM updated + WHERE updated.owner_id IS NOT NULL AND ${this.queueName}_owner_concurrency.id = updated.owner_id + ), + group_decrement AS ( + UPDATE ${this.queueName}_group_concurrency + SET current_concurrency = GREATEST(0, current_concurrency - 1) + FROM updated + WHERE updated.group_id IS NOT NULL AND ${this.queueName}_group_concurrency.id = updated.group_id + ) + SELECT * FROM updated; + `; } else { updateQuery = `UPDATE ${this.queueName} SET status = 'failed'::nuq.job_status, lock = null, locked_at = null, finished_at = now(), failedreason = $3 WHERE id = $1 AND lock = $2 RETURNING id, listen_channel_id;`; } @@ -1110,12 +1582,192 @@ class NuQ { }); } + public async getOwnerJobCounts( + ownerId: string, + _logger: Logger = logger, + ): Promise<{ active: number; queued: number }> { + const start = Date.now(); + try { + const result = await nuqPool.query( + `SELECT + COALESCE(SUM(CASE WHEN status = 'active'::nuq.job_status THEN 1 ELSE 0 END), 0)::int as active, + COALESCE(SUM(CASE WHEN status = 'queued'::nuq.job_status THEN 1 ELSE 0 END), 0)::int as queued + FROM ${this.queueName} + WHERE owner_id = $1;`, + [ownerId], + ); + + return { + active: result.rows[0]?.active ?? 0, + queued: result.rows[0]?.queued ?? 0, + }; + } finally { + _logger.info("nuqGetOwnerJobCounts metrics", { + module: "nuq/metrics", + method: "nuqGetOwnerJobCounts", + duration: Date.now() - start, + ownerId, + }); + } + } + + public async getOwnerConcurrency( + ownerId: string, + _logger: Logger = logger, + ): Promise<{ + currentConcurrency: number; + maxConcurrency: number; + } | null> { + const start = Date.now(); + try { + const result = await nuqPool.query( + `SELECT current_concurrency, max_concurrency + FROM ${this.queueName}_owner_concurrency + WHERE id = $1;`, + [ownerId], + ); + + if (result.rows.length === 0) { + return null; + } + + return { + currentConcurrency: result.rows[0].current_concurrency, + maxConcurrency: result.rows[0].max_concurrency, + }; + } finally { + _logger.info("nuqGetOwnerConcurrency metrics", { + module: "nuq/metrics", + method: "nuqGetOwnerConcurrency", + duration: Date.now() - start, + ownerId, + }); + } + } + // === Metrics public async getMetrics(): Promise { const start = Date.now(); - const result = await nuqPool.query( - `SELECT status, COUNT(id) as count FROM ${this.queueName} GROUP BY status ORDER BY count DESC;`, - ); + + let query: string; + if (this.options.concurrencyLimit === "per-owner") { + query = ` + WITH owner_status AS ( + SELECT + id, + GREATEST(0, max_concurrency - current_concurrency) as available_slots + FROM ${this.queueName}_owner_concurrency + ), + queued_per_owner AS ( + SELECT + owner_id, + COUNT(*) as queued_count + FROM ${this.queueName} + WHERE status = 'queued'::nuq.job_status AND owner_id IS NOT NULL + GROUP BY owner_id + ), + owner_breakdown AS ( + SELECT + LEAST(qpo.queued_count, COALESCE(os.available_slots, 999999)) as can_run, + GREATEST(0, qpo.queued_count - COALESCE(os.available_slots, 0)) as blocked + FROM queued_per_owner qpo + LEFT JOIN owner_status os ON qpo.owner_id = os.id + ), + non_queued AS ( + SELECT status::text, COUNT(*) as count + FROM ${this.queueName} + WHERE status != 'queued'::nuq.job_status + GROUP BY status + ), + totals AS ( + SELECT + COALESCE(SUM(can_run), 0) as total_queued, + COALESCE(SUM(blocked), 0) as total_limited + FROM owner_breakdown + ) + SELECT status, count FROM non_queued + UNION ALL + SELECT 'queued' as status, total_queued::bigint as count + FROM totals + WHERE total_queued > 0 + UNION ALL + SELECT 'concurrency-limited' as status, total_limited::bigint as count + FROM totals + WHERE total_limited > 0 + ORDER BY count DESC; + `; + } else if (this.options.concurrencyLimit === "per-owner-per-group") { + query = ` + WITH owner_status AS ( + SELECT + id, + GREATEST(0, max_concurrency - current_concurrency) as available_slots + FROM ${this.queueName}_owner_concurrency + ), + group_status AS ( + SELECT + id, + GREATEST(0, max_concurrency - current_concurrency) as available_slots + FROM ${this.queueName}_group_concurrency + ), + queued_per_owner_group AS ( + SELECT + owner_id, + group_id, + COUNT(*) as queued_count + FROM ${this.queueName} + WHERE status = 'queued'::nuq.job_status + GROUP BY owner_id, group_id + ), + owner_group_breakdown AS ( + SELECT + LEAST( + qpog.queued_count, + LEAST( + COALESCE(os.available_slots, 999999), + COALESCE(gs.available_slots, 999999) + ) + ) as can_run, + GREATEST( + 0, + qpog.queued_count - LEAST( + COALESCE(os.available_slots, 999999), + COALESCE(gs.available_slots, 999999) + ) + ) as blocked + FROM queued_per_owner_group qpog + LEFT JOIN owner_status os ON qpog.owner_id = os.id + LEFT JOIN group_status gs ON qpog.group_id = gs.id + ), + non_queued AS ( + SELECT status::text, COUNT(*) as count + FROM ${this.queueName} + WHERE status != 'queued'::nuq.job_status + GROUP BY status + ), + totals AS ( + SELECT + COALESCE(SUM(can_run), 0) as total_queued, + COALESCE(SUM(blocked), 0) as total_limited + FROM owner_group_breakdown + ) + SELECT status, count FROM non_queued + UNION ALL + SELECT 'queued' as status, total_queued::bigint as count + FROM totals + WHERE total_queued > 0 + UNION ALL + SELECT 'concurrency-limited' as status, total_limited::bigint as count + FROM totals + WHERE total_limited > 0 + ORDER BY count DESC; + `; + } else { + query = `SELECT status, COUNT(id) as count FROM ${this.queueName} GROUP BY status ORDER BY count DESC;`; + } + + const result = await nuqPool.query(query); + logger.info("nuqGetMetrics metrics", { module: "nuq/metrics", method: "nuqGetMetrics", @@ -1150,6 +1802,102 @@ class NuQ { } } +// === Group + +type NuQGroupOptions = { + memberQueues: NuQ[]; + finishQueue?: NuQ; + groupTTL: number; +}; + +type NuQGroupConcurrencySettings = { + queue: NuQ; + maxConcurrency?: number; +}; + +type NuQGroupInstance = { + id: string; + status: "active" | "completed"; + createdAt: Date; + finishedAt?: Date; + expiresAt?: Date; +}; + +class NuQGroup { + constructor( + public readonly groupName: string, + public readonly options: NuQGroupOptions, + ) {} + + private groupReturning = [ + "id", + "status", + "created_at", + "finished_at", + "expires_at", + ]; + + private rowToGroup(row: any): NuQGroupInstance | null { + if (!row) return null; + return { + id: row.id, + status: row.status, + createdAt: new Date(row.created_at), + finishedAt: row.finished_at ? new Date(row.finished_at) : undefined, + expiresAt: row.expires_at ? new Date(row.expires_at) : undefined, + }; + } + + public async addGroup( + id: string, + maxConcurrency: NuQGroupConcurrencySettings[], + ): Promise { + const client = await nuqPool.connect(); + + await client.query("BEGIN"); + + try { + const insert = await client.query( + `INSERT INTO ${this.groupName} (id) VALUES ($1) RETURNING ${this.groupReturning.join(", ")};`, + [id], + ); + + if (maxConcurrency.length > 0) { + for (const entry of maxConcurrency) { + if (entry.queue.options.concurrencyLimit === "per-owner-per-group") { + await client.query( + `INSERT INTO ${entry.queue.queueName}_group_concurrency (id, current_concurrency, max_concurrency) VALUES ($1, 0, $2);`, + [id, entry.maxConcurrency ?? null], + ); + } + } + } + + await client.query("COMMIT"); + client.release(); + + return this.rowToGroup(insert.rows[0])!; + } catch (e) { + await client.query("ROLLBACK"); + client.release(e); + throw e; + } + } + + public async getGroup(id: string): Promise { + return this.rowToGroup( + ( + await nuqPool.query( + `SELECT ${this.groupReturning.join(", ")} FROM ${this.groupName} WHERE id = $1 LIMIT 1;`, + [id], + ) + ).rows[0], + ); + } +} + +// === Utilities + export function nuqGetLocalMetrics(): string { return `# HELP nuq_pool_waiting_count Number of requests waiting in the pool\n# TYPE nuq_pool_waiting_count gauge\nnuq_pool_waiting_count ${nuqPool.waitingCount}\n # HELP nuq_pool_idle_count Number of connections idle in the pool\n# TYPE nuq_pool_idle_count gauge\nnuq_pool_idle_count ${nuqPool.idleCount}\n @@ -1172,7 +1920,14 @@ export async function nuqHealthCheck(): Promise { // === Instances export const scrapeQueue = new NuQ("nuq.queue_scrape", { - concurrencyLimit: "per-owner", + concurrencyLimit: "per-owner-per-group", +}); +// export const crawlFinishQueue = new NuQ("nuq.queue_crawl_finish"); + +export const crawlGroup = new NuQGroup("nuq.group_crawl", { + memberQueues: [scrapeQueue], + // finishQueue: crawlFinishQueue, + groupTTL: 24 * 60 * 60, }); // === Cleanup diff --git a/apps/api/src/services/worker/scrape-worker.ts b/apps/api/src/services/worker/scrape-worker.ts index 14b888a40..f86dfef41 100644 --- a/apps/api/src/services/worker/scrape-worker.ts +++ b/apps/api/src/services/worker/scrape-worker.ts @@ -4,10 +4,6 @@ import http from "http"; import https from "https"; import { logger as _logger } from "../../lib/logger"; -import { - concurrentJobDone, - pushConcurrencyLimitActiveJob, -} from "../../lib/concurrency-limit"; import { addJobPriority, deleteJobPriority } from "../../lib/job-priority"; import { cacheableLookup } from "../../scraper/scrapeURL/lib/cacheableLookup"; import { v4 as uuidv4 } from "uuid"; @@ -27,17 +23,13 @@ import { StoredCrawl, } from "../../lib/crawl-redis"; import { redisEvictConnection } from "../redis"; -import { - _addScrapeJobToBullMQ, - addScrapeJob, - addScrapeJobs, -} from "../queue-jobs"; +import { addScrapeJob, addScrapeJobs } from "../queue-jobs"; import psl from "psl"; import { getJobPriority } from "../../lib/job-priority"; import { Document, scrapeOptions, TeamFlags } from "../../controllers/v2/types"; import { hasFormatOfType } from "../../lib/format-utils"; import { getACUCTeam } from "../../controllers/auth"; -import { createWebhookSender, WebhookEvent } from "../webhook"; +import { createWebhookSender, WebhookEvent } from "../webhook/index"; import { CustomError } from "../../lib/custom-error"; import { startWebScraperPipeline } from "../../main/runWebScraper"; import { CostTracking } from "../../lib/cost-tracking"; @@ -338,10 +330,7 @@ async function processJob(job: NuQJob) { // Store robots blocked URLs in Redis set for (const [url, reason] of links.denialReasons) { if (reason === "URL blocked by robots.txt") { - await recordRobotsBlocked( - job.data.crawl_id, - url - ); + await recordRobotsBlocked(job.data.crawl_id, url); } } @@ -552,15 +541,12 @@ async function processJob(job: NuQJob) { error instanceof Error && error.message === "URL blocked by robots.txt" ) { - await recordRobotsBlocked( - job.data.crawl_id, - job.data.url, - ); + await recordRobotsBlocked(job.data.crawl_id, job.data.url); } } catch (e) { logger.debug("Failed to record top-level robots block", { e }); } - + if (job.data.crawl_id) { const sc = (await getCrawl(job.data.crawl_id)) as StoredCrawl; @@ -764,7 +750,7 @@ async function addKickoffSitemapJob( } const jobId = uuidv4(); - await _addScrapeJobToBullMQ( + await addScrapeJob( { mode: "kickoff_sitemap" as const, team_id: sourceJob.data.team_id, @@ -1111,75 +1097,55 @@ export const processJobInternal = async (job: NuQJob) => { async function processJobWithTracing(job: NuQJob, logger: any) { try { + await addJobPriority(job.data.team_id, job.id); try { - let extendLockInterval: NodeJS.Timeout | null = null; - if (job.data?.mode !== "kickoff" && job.data?.team_id) { - extendLockInterval = setInterval(async () => { - await pushConcurrencyLimitActiveJob( - job.data.team_id, - job.id, - 60 * 1000, - ); // 60s lock renew, just like in the queue - }, jobLockExtendInterval); - } - - await addJobPriority(job.data.team_id, job.id); - try { - if (job.data.mode === "kickoff") { - const result = await processKickoffJob( - job as NuQJob, - ); - if (result.success) { - return null; - } else { - throw (result as any).error; - } - } else if (job.data.mode === "kickoff_sitemap") { - const result = await processKickoffSitemapJob( - job as NuQJob, - ); - if (result.success) { - return null; - } else { - throw (result as any).error; - } + if (job.data.mode === "kickoff") { + const result = await processKickoffJob(job as NuQJob); + if (result.success) { + return null; } else { - const result = await processJob(job as NuQJob); - if (result.success) { - try { - if (job.data.team_id) { - await redisEvictConnection.set( - "most-recent-success:" + job.data.team_id, - new Date().toISOString(), - "EX", - 60 * 60 * 24, - ); - } - } catch (e) { - logger.warn("Failed to set most recent success", { error: e }); - } - - try { - if (process.env.GCS_BUCKET_NAME) { - logger.debug("Job succeeded -- putting null in Redis"); - return null; - } else { - logger.debug("Job succeeded -- putting result in Redis"); - return result.document; - } - } catch (e) {} - } else { - throw (result as any).error; - } + throw (result as any).error; } - } finally { - await deleteJobPriority(job.data.team_id, job.id); - if (extendLockInterval) { - clearInterval(extendLockInterval); + } else if (job.data.mode === "kickoff_sitemap") { + const result = await processKickoffSitemapJob( + job as NuQJob, + ); + if (result.success) { + return null; + } else { + throw (result as any).error; + } + } else { + const result = await processJob(job as NuQJob); + if (result.success) { + try { + if (job.data.team_id) { + await redisEvictConnection.set( + "most-recent-success:" + job.data.team_id, + new Date().toISOString(), + "EX", + 60 * 60 * 24, + ); + } + } catch (e) { + logger.warn("Failed to set most recent success", { error: e }); + } + + try { + if (process.env.GCS_BUCKET_NAME) { + logger.debug("Job succeeded -- putting null in Redis"); + return null; + } else { + logger.debug("Job succeeded -- putting result in Redis"); + return result.document; + } + } catch (e) {} + } else { + throw (result as any).error; } } } finally { - await concurrentJobDone(job); + await deleteJobPriority(job.data.team_id, job.id); } } catch (error) { logger.debug("Job failed", { error }); diff --git a/apps/nuq-postgres/nuq.sql b/apps/nuq-postgres/nuq.sql index a0098f150..834390d0d 100644 --- a/apps/nuq-postgres/nuq.sql +++ b/apps/nuq-postgres/nuq.sql @@ -9,6 +9,12 @@ EXCEPTION WHEN duplicate_object THEN null; END $$; +DO $$ BEGIN + CREATE TYPE nuq.group_status AS ENUM ('active', 'completed', 'cancelled'); +EXCEPTION + WHEN duplicate_object THEN null; +END $$; + CREATE TABLE IF NOT EXISTS nuq.queue_scrape ( id uuid NOT NULL DEFAULT gen_random_uuid(), status nuq.job_status NOT NULL DEFAULT 'queued'::nuq.job_status, @@ -23,6 +29,8 @@ CREATE TABLE IF NOT EXISTS nuq.queue_scrape ( returnvalue jsonb, -- only for selfhost failedreason text, -- only for selfhost owner_id uuid, + group_id uuid, + times_out_at timestamp with time zone, CONSTRAINT queue_scrape_pkey PRIMARY KEY (id) ); @@ -36,6 +44,8 @@ CREATE INDEX IF NOT EXISTS queue_scrape_active_locked_at_idx ON nuq.queue_scrape CREATE INDEX IF NOT EXISTS nuq_queue_scrape_queued_optimal_2_idx ON nuq.queue_scrape (priority ASC, created_at ASC, id) WHERE (status = 'queued'::nuq.job_status); CREATE INDEX IF NOT EXISTS nuq_queue_scrape_failed_created_at_idx ON nuq.queue_scrape USING btree (created_at) WHERE (status = 'failed'::nuq.job_status); CREATE INDEX IF NOT EXISTS nuq_queue_scrape_completed_created_at_idx ON nuq.queue_scrape USING btree (created_at) WHERE (status = 'completed'::nuq.job_status); +CREATE INDEX IF NOT EXISTS nuq_queue_scrape_queued_owner_idx ON nuq.queue_scrape (owner_id, priority ASC, created_at ASC) WHERE (status = 'queued'::nuq.job_status); +CREATE INDEX IF NOT EXISTS nuq_queue_scrape_queued_owner_group_idx ON nuq.queue_scrape (owner_id, group_id, priority ASC, created_at ASC) WHERE (status = 'queued'::nuq.job_status); CREATE TABLE IF NOT EXISTS nuq.queue_scrape_owner_concurrency ( id uuid NOT NULL, @@ -59,12 +69,19 @@ AS $$ SELECT COALESCE((SELECT max_concurrency FROM nuq.queue_scrape_owner_concurrency_source WHERE id = owner_id LIMIT 1), 100)::int8; $$; +CREATE TABLE IF NOT EXISTS nuq.queue_scrape_group_concurrency ( + id uuid NOT NULL, + current_concurrency int8 NOT NULL, + max_concurrency int8, + CONSTRAINT queue_scrape_group_concurrency_pkey PRIMARY KEY (id) +); + SELECT cron.schedule('nuq_queue_scrape_clean_completed', '*/5 * * * *', $$ - DELETE FROM nuq.queue_scrape WHERE nuq.queue_scrape.status = 'completed'::nuq.job_status AND nuq.queue_scrape.created_at < now() - interval '1 hour'; + DELETE FROM nuq.queue_scrape WHERE nuq.queue_scrape.status = 'completed'::nuq.job_status AND nuq.queue_scrape.group_id IS NULL AND nuq.queue_scrape.created_at < now() - interval '1 hour'; $$); SELECT cron.schedule('nuq_queue_scrape_clean_failed', '*/5 * * * *', $$ - DELETE FROM nuq.queue_scrape WHERE nuq.queue_scrape.status = 'failed'::nuq.job_status AND nuq.queue_scrape.created_at < now() - interval '6 hours'; + DELETE FROM nuq.queue_scrape WHERE nuq.queue_scrape.status = 'failed'::nuq.job_status AND nuq.queue_scrape.group_id IS NULL AND nuq.queue_scrape.created_at < now() - interval '6 hours'; $$); SELECT cron.schedule('nuq_queue_scrape_lock_reaper', '15 seconds', $$ @@ -141,3 +158,50 @@ SELECT cron.schedule('nuq_queue_scrape_concurrency_sync', '*/5 * * * *', $$ UPDATE nuq.queue_scrape_owner_concurrency SET max_concurrency = (SELECT nuq_queue_scrape_owner_resolve_max_concurrency(nuq.queue_scrape_owner_concurrency.id)); $$); + +SELECT cron.schedule('nuq_queue_scrape_timeout', '* * * * *', $$ + UPDATE nuq.queue_scrape + SET status = 'failed', finished_at = now(), failedreason = 'SCRAPE_TIMEOUT|{"stack":"Error: Scrape timed out\n in DB","message":"Scrape timed out"}' + WHERE status = 'queued' AND times_out_at < now(); +$$); + +CREATE TABLE IF NOT EXISTS nuq.group_crawl ( + id uuid NOT NULL DEFAULT gen_random_uuid(), + status nuq.group_status NOT NULL DEFAULT 'active'::nuq.group_status, + created_at timestamp with time zone NOT NULL DEFAULT now(), + finished_at timestamp with time zone, + expires_at timestamp with time zone, + CONSTRAINT group_crawl_pkey PRIMARY KEY (id) +); + +CREATE OR REPLACE FUNCTION nuq_queue_scrape_check_group_completion() +RETURNS TRIGGER +LANGUAGE plpgsql +AS $$ +BEGIN + IF NEW.group_id IS NOT NULL THEN + UPDATE nuq.group_crawl + SET status = 'completed'::nuq.group_status, + finished_at = now(), + expires_at = now() + interval '24 hours' + WHERE id = NEW.group_id + AND status != 'completed'::nuq.group_status + AND NOT EXISTS ( + SELECT 1 + FROM nuq.queue_scrape + WHERE group_id = NEW.group_id + AND status NOT IN ('completed'::nuq.job_status, 'failed'::nuq.job_status) + ); + END IF; + + RETURN NEW; +END; +$$; + +-- Trigger to automatically mark groups as completed +CREATE OR REPLACE TRIGGER nuq_queue_scrape_group_completion_trigger +AFTER UPDATE OF status ON nuq.queue_scrape +FOR EACH ROW +WHEN (NEW.status IN ('completed'::nuq.job_status, 'failed'::nuq.job_status) + AND OLD.status NOT IN ('completed'::nuq.job_status, 'failed'::nuq.job_status)) +EXECUTE FUNCTION nuq_queue_scrape_check_group_completion();