From 864b03c1bb07faacb2c2498c95ff39adf086178c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Sun, 19 Oct 2025 17:37:48 +0200 Subject: [PATCH] feat(nuq): better group support --- apps/api/src/controllers/v0/crawl.ts | 15 +- apps/api/src/controllers/v1/batch-scrape.ts | 15 +- apps/api/src/controllers/v2/batch-scrape.ts | 15 +- apps/api/src/controllers/v2/crawl.ts | 15 +- .../api/src/services/indexing/index-worker.ts | 15 +- apps/api/src/services/worker/nuq.ts | 156 +++++++++++++----- apps/nuq-postgres/nuq.sql | 1 + 7 files changed, 157 insertions(+), 75 deletions(-) diff --git a/apps/api/src/controllers/v0/crawl.ts b/apps/api/src/controllers/v0/crawl.ts index 980492b65..e9e1e80ec 100644 --- a/apps/api/src/controllers/v0/crawl.ts +++ b/apps/api/src/controllers/v0/crawl.ts @@ -175,12 +175,15 @@ export async function crawlController(req: Request, res: Response) { sc.robots = await crawler.getRobotsTxt(); } catch (_) {} - await crawlGroup.addGroup(id, [ - { - queue: scrapeQueue, - maxConcurrency: undefined, - }, - ]); + await crawlGroup.addGroup(id, { + concurrency: [ + { + queue: scrapeQueue, + maxConcurrency: undefined, + }, + ], + ownerId: sc.team_id, + }); await saveCrawl(id, sc); diff --git a/apps/api/src/controllers/v1/batch-scrape.ts b/apps/api/src/controllers/v1/batch-scrape.ts index cea629d38..ed56d5397 100644 --- a/apps/api/src/controllers/v1/batch-scrape.ts +++ b/apps/api/src/controllers/v1/batch-scrape.ts @@ -137,12 +137,15 @@ export async function batchScrapeController( }; if (!req.body.appendToId) { - await crawlGroup.addGroup(id, [ - { - queue: scrapeQueue, - maxConcurrency: sc.maxConcurrency, - }, - ]); + await crawlGroup.addGroup(id, { + concurrency: [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ], + ownerId: sc.team_id, + }); await saveCrawl(id, sc); await markCrawlActive(id); } diff --git a/apps/api/src/controllers/v2/batch-scrape.ts b/apps/api/src/controllers/v2/batch-scrape.ts index cf57573e6..d0b362c21 100644 --- a/apps/api/src/controllers/v2/batch-scrape.ts +++ b/apps/api/src/controllers/v2/batch-scrape.ts @@ -130,12 +130,15 @@ export async function batchScrapeController( }; if (!req.body.appendToId) { - await crawlGroup.addGroup(id, [ - { - queue: scrapeQueue, - maxConcurrency: sc.maxConcurrency, - }, - ]); + await crawlGroup.addGroup(id, { + concurrency: [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ], + ownerId: sc.team_id, + }); await saveCrawl(id, sc); await markCrawlActive(id); } diff --git a/apps/api/src/controllers/v2/crawl.ts b/apps/api/src/controllers/v2/crawl.ts index 37b172c6e..3825cc706 100644 --- a/apps/api/src/controllers/v2/crawl.ts +++ b/apps/api/src/controllers/v2/crawl.ts @@ -203,12 +203,15 @@ export async function crawlController( }); } - await crawlGroup.addGroup(id, [ - { - queue: scrapeQueue, - maxConcurrency: sc.maxConcurrency, - }, - ]); + await crawlGroup.addGroup(id, { + concurrency: [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency, + }, + ], + ownerId: sc.team_id, + }); await saveCrawl(id, sc); await markCrawlActive(id); diff --git a/apps/api/src/services/indexing/index-worker.ts b/apps/api/src/services/indexing/index-worker.ts index b474a0e40..22b21f62c 100644 --- a/apps/api/src/services/indexing/index-worker.ts +++ b/apps/api/src/services/indexing/index-worker.ts @@ -491,12 +491,15 @@ const processPrecrawlJob = async (token: string, job: Job) => { // }); // } - await crawlGroup.addGroup(crawlId, [ - { - queue: scrapeQueue, - maxConcurrency: sc.maxConcurrency ?? undefined, - }, - ]); + await crawlGroup.addGroup(crawlId, { + concurrency: [ + { + queue: scrapeQueue, + maxConcurrency: sc.maxConcurrency ?? undefined, + }, + ], + ownerId: sc.team_id, + }); await saveCrawl(crawlId, sc); diff --git a/apps/api/src/services/worker/nuq.ts b/apps/api/src/services/worker/nuq.ts index 3d35a7b68..201ae2426 100644 --- a/apps/api/src/services/worker/nuq.ts +++ b/apps/api/src/services/worker/nuq.ts @@ -50,25 +50,35 @@ type NuQJobOptions = { timesOutAt?: Date; }; +function normalizeOwnerId(ownerId: string | undefined | null) { + const bareOwnerId = ownerId ?? undefined; + const normalizedOwnerId = bareOwnerId + ? uuidValidate(bareOwnerId) + ? bareOwnerId + : uuidv5(bareOwnerId, "b208cbac-8bdf-4599-bf17-da78426e3f7c") // preview namespace + : null; + return normalizedOwnerId; +} + class NuQ { constructor( public readonly queueName: string, public readonly options: NuQOptions = { concurrencyLimit: false }, - ) {} + ) { } // === Listener private listener: | { - type: "postgres"; - client: Client; - } + type: "postgres"; + client: Client; + } | { - type: "rabbitmq"; - connection: amqp.ChannelModel; - channel: amqp.Channel; - queue: string; - } + type: "rabbitmq"; + connection: amqp.ChannelModel; + channel: amqp.Channel; + queue: string; + } | null = null; private listens: { [key: string]: ((status: "completed" | "failed") => void)[]; @@ -265,7 +275,7 @@ class NuQ { channel.on("close", () => { logger.info("NuQ sender channel closed", { module: "nuq/rabbitmq" }); - connection.close().catch(() => {}); + connection.close().catch(() => { }); this.sender = null; }); @@ -475,6 +485,52 @@ class NuQ { } } + public async getJobCountsOfGroup( + groupId: string, + _logger: Logger = logger, + ): Promise> { + const start = Date.now(); + try { + const stats = await nuqPool.query(` + SELECT status, COUNT(*) as count FROM nuq.queue_scrape WHERE group_id = $1 GROUP BY status; + `, [groupId]); + + return Object.fromEntries(stats.rows.map(x => [x.status, x.count])); + } finally { + _logger.info("nuqGetJobCountsOfGroup metrics", { + module: "nuq/metrics", + method: "nuqGetJobCountsOfGroup", + duration: Date.now() - start, + groupId, + }); + } + } + + public async getJobsOfGroupWithStatus( + groupId: string, + status: NuQJobStatus, + limit: number, + offset: number, + _logger: Logger = logger, + ): Promise[]> { + const start = Date.now(); + try { + return (await nuqPool.query(` + SELECT ${this.jobReturning.join(", ")} FROM ${this.queueName} + WHERE group_id = $1 AND status = $2 + ORDER BY finished_at ASC + LIMIT $3 OFFSET $4 + `, [groupId, status, limit, offset])).rows.map(x => this.rowToJob(x)!); + } finally { + _logger.info("nuqGetJobsOfGroupWithStatus metrics", { + module: "nuq/metrics", + method: "nuqGetJobsOfGroupWithStatus", + duration: Date.now() - start, + groupId, + }); + } + } + public async removeJob( id: string, _logger: Logger = logger, @@ -531,13 +587,6 @@ class NuQ { 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, @@ -557,7 +606,7 @@ class NuQ { data, options.priority ?? 0, options.listenable ? listenChannelId : null, - normalizedOwnerId ?? null, + normalizeOwnerId(options.ownerId) ?? null, options.groupId ?? null, options.timesOutAt ? options.timesOutAt.toISOString() : null, ], @@ -592,13 +641,6 @@ class NuQ { options: NuQJobOptions = {}, ): Promise> { return withSpan("nuq.addJob", 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, @@ -618,7 +660,7 @@ class NuQ { data, options.priority ?? 0, options.listenable ? listenChannelId : null, - normalizedOwnerId ?? null, + normalizeOwnerId(options.ownerId) ?? null, options.groupId ?? null, options.timesOutAt ? options.timesOutAt.toISOString() : null, ], @@ -675,20 +717,13 @@ class NuQ { 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); + ownerIds.push(normalizeOwnerId(job.options?.ownerId)); groupIds.push(job.options?.groupId ?? null); timesOutAts.push( job.options?.timesOutAt @@ -1853,12 +1888,18 @@ type NuQGroupConcurrencySettings = { maxConcurrency?: number; }; +type NuQGroupInstanceOptions = { + concurrency: NuQGroupConcurrencySettings[]; + ownerId: string; +} + type NuQGroupInstance = { id: string; status: "active" | "completed"; createdAt: Date; finishedAt?: Date; expiresAt?: Date; + ownerId?: string; }; class NuQGroup { @@ -1873,6 +1914,7 @@ class NuQGroup { "created_at", "finished_at", "expires_at", + "owner_id", ]; private rowToGroup(row: any): NuQGroupInstance | null { @@ -1883,12 +1925,13 @@ class NuQGroup { 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, + ownerId: row.owner_id ?? undefined, }; } public async addGroup( id: string, - maxConcurrency: NuQGroupConcurrencySettings[], + options: NuQGroupInstanceOptions, ): Promise { const client = await nuqPool.connect(); @@ -1896,18 +1939,16 @@ class NuQGroup { try { const insert = await client.query( - `INSERT INTO ${this.groupName} (id) VALUES ($1) RETURNING ${this.groupReturning.join(", ")};`, - [id], + `INSERT INTO ${this.groupName} (id, owner_id) VALUES ($1, $2) RETURNING ${this.groupReturning.join(", ")};`, + [id, normalizeOwnerId(options.ownerId) ?? null], ); - 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], - ); - } + for (const entry of options.concurrency) { + 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], + ); } } @@ -1932,6 +1973,31 @@ class NuQGroup { ).rows[0], ); } + + public async cancelGroup(id: string): Promise { + const client = await nuqPool.connect(); + + await client.query("BEGIN"); + + try { + const updateOp = await client.query(`UPDATE ${this.groupName} SET status = 'cancelled'::nuq.group_status WHERE id = $1 AND status = 'active'::nuq.group_status`, [id]); + + if (updateOp.rowCount === 0) { + client.release(); + return false; + } + + for (const queue of this.options.memberQueues) { + await client.query(`UPDATE ${queue.queueName} SET status = 'failed'::nuq.job_status, lock = null, locked_at = null, finished_at = now(), failedreason = 'CANCELLED' WHERE group_id = $1 AND status = 'queued'::nuq.job_status`, [id]); + } + + client.release(); + return true; + } catch (e) { + client.release(e); + throw e; + } + } } // === Utilities diff --git a/apps/nuq-postgres/nuq.sql b/apps/nuq-postgres/nuq.sql index d553d5b1d..916c16f77 100644 --- a/apps/nuq-postgres/nuq.sql +++ b/apps/nuq-postgres/nuq.sql @@ -245,6 +245,7 @@ CREATE TABLE IF NOT EXISTS nuq.group_crawl ( created_at timestamp with time zone NOT NULL DEFAULT now(), finished_at timestamp with time zone, expires_at timestamp with time zone, + owner_id uuid, CONSTRAINT group_crawl_pkey PRIMARY KEY (id) );