diff --git a/apps/api/src/services/worker/nuq.ts b/apps/api/src/services/worker/nuq.ts index 05732737b..2e20925fc 100644 --- a/apps/api/src/services/worker/nuq.ts +++ b/apps/api/src/services/worker/nuq.ts @@ -39,15 +39,18 @@ class NuQ { // === Listener - private listener: { - type: "postgres"; - client: Client; - } | { - type: "rabbitmq"; - connection: amqp.ChannelModel; - channel: amqp.Channel; - queue: string; - } | null = null; + private listener: + | { + type: "postgres"; + client: Client; + } + | { + type: "rabbitmq"; + connection: amqp.ChannelModel; + channel: amqp.Channel; + queue: string; + } + | null = null; private listens: { [key: string]: ((status: "completed" | "failed") => void)[]; } = {}; @@ -60,11 +63,18 @@ class NuQ { const connection = await amqp.connect(process.env.NUQ_RABBITMQ_URL); const channel = await connection.createChannel(); await channel.prefetch(1); - const queue = await channel.assertQueue(this.queueName + ".listen." + listenChannelId, { - exclusive: true, - autoDelete: true, - durable: false, - }); + const queue = await channel.assertQueue( + this.queueName + ".listen." + listenChannelId, + { + exclusive: true, + autoDelete: true, + durable: false, + arguments: { + "x-queue-type": "classic", + "x-message-ttl": 60000, + }, + }, + ); this.listener = { type: "rabbitmq", @@ -73,43 +83,51 @@ class NuQ { queue: queue.queue, }; - await this.listener.channel.consume(this.listener.queue, (msg => { - if (msg === null) { - logger.info("NuQ listener channel closed", { module: "nuq/rabbitmq" }); - this.listener = null; + await this.listener.channel.consume( + this.listener.queue, + (msg => { + if (msg === null) { + logger.info("NuQ listener channel closed", { + module: "nuq/rabbitmq", + }); + this.listener = null; - setTimeout( - (() => { - this.startListener().catch(err => - logger.error("Error in NuQ listener reconnect", { - err, - module: "nuq/rabbitmq", - }), - ); - }).bind(this), - 250, - ); - return; - } + setTimeout( + (() => { + this.startListener().catch(err => + logger.error("Error in NuQ listener reconnect", { + err, + module: "nuq/rabbitmq", + }), + ); + }).bind(this), + 250, + ); + return; + } - logger.info("NuQ job received", { module: "nuq/rabbitmq", jobId: msg.properties.correlationId, status: msg.content.toString() }); + logger.info("NuQ job received", { + module: "nuq/rabbitmq", + jobId: msg.properties.correlationId, + status: msg.content.toString(), + }); - const jobId = msg.properties.correlationId as string; - const status = msg.content.toString() as "completed" | "failed"; + const jobId = msg.properties.correlationId as string; + const status = msg.content.toString() as "completed" | "failed"; - if (jobId in this.listens) { - this.listens[jobId].forEach(listener => - listener(status), - ); - } - delete this.listens[jobId]; + if (jobId in this.listens) { + this.listens[jobId].forEach(listener => listener(status)); + } + delete this.listens[jobId]; - if (this.listener && this.listener.type === "rabbitmq") { - this.listener.channel.ack(msg); - } - }).bind(this), { - noAck: false, - }); + if (this.listener && this.listener.type === "rabbitmq") { + this.listener.channel.ack(msg); + } + }).bind(this), + { + noAck: false, + }, + ); } else { this.listener = { type: "postgres", @@ -204,7 +222,7 @@ class NuQ { await channel.assertQueue(this.queueName + ".prefetch", { durable: true, arguments: { - "x-message-ttl": 30000, + "x-queue-type": "quorum", "x-max-length": 20000, }, }); @@ -228,27 +246,44 @@ class NuQ { } } - private async sendJobEnd(id: string, status: "completed" | "failed", listenChannelId: string, _logger: Logger = logger) { + private async sendJobEnd( + id: string, + status: "completed" | "failed", + listenChannelId: string, + _logger: Logger = logger, + ) { await this.startSender(); if (this.sender) { - await this.sender.channel.sendToQueue(this.queueName + ".listen." + listenChannelId, Buffer.from(status, "utf8"), { - correlationId: id, - }); + await this.sender.channel.sendToQueue( + this.queueName + ".listen." + listenChannelId, + Buffer.from(status, "utf8"), + { + correlationId: id, + }, + ); _logger.info("NuQ job sent", { module: "nuq/rabbitmq" }); } else { _logger.warn("NuQ sender not started", { module: "nuq/rabbitmq" }); } } - private async sendJobPrefetch(job: NuQJob, _logger: Logger = logger) { + private async sendJobPrefetch( + job: NuQJob, + _logger: Logger = logger, + ) { await this.startSender(); if (this.sender) { - await this.sender.channel.sendToQueue(this.queueName + ".prefetch", Buffer.from(JSON.stringify(job), "utf8"), { - correlationId: job.id, - persistent: true, - }); + await this.sender.channel.sendToQueue( + this.queueName + ".prefetch", + Buffer.from(JSON.stringify(job), "utf8"), + { + correlationId: job.id, + persistent: true, + expiration: "30000", + }, + ); _logger.info("NuQ job prefetch sent", { module: "nuq/rabbitmq" }); } else { _logger.warn("NuQ sender not started", { module: "nuq/rabbitmq" }); @@ -613,7 +648,7 @@ class NuQ { return result; }); } - + // === Prefetch public async prefetchJobs(_logger: Logger = logger): Promise { @@ -624,16 +659,25 @@ class NuQ { ` WITH next AS (SELECT id 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 500) 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(", ")}; - ` + `, ) ).rows.map(row => this.rowToJob(row)!); - + for (const job of jobs) { - await this.sendJobPrefetch(job, _logger.child({ jobId: job.id, zeroDataRetention: !!((job.data || {} as any).zeroDataRetention) })); + await this.sendJobPrefetch( + job, + _logger.child({ + jobId: job.id, + zeroDataRetention: !!(job.data || ({} as any)).zeroDataRetention, + }), + ); } - - _logger.info("Prefetched jobs", { module: "nuq/metrics", jobCount: jobs.length }); - + + _logger.info("Prefetched jobs", { + module: "nuq/metrics", + jobCount: jobs.length, + }); + return jobs.length; } finally { _logger.info("nuqPrefetchJobs metrics", { @@ -653,14 +697,19 @@ class NuQ { await this.startSender(); if (this.sender) { - const job = await this.sender.channel.get(this.queueName + ".prefetch", { noAck: true }); + const job = await this.sender.channel.get( + this.queueName + ".prefetch", + { noAck: true }, + ); if (job !== false) { return this.rowToJob(JSON.parse(job.content.toString())); } else { return null; } } else { - logger.warn("NuQ sender not started, falling back to postgres", { module: "nuq/rabbitmq" }); + logger.warn("NuQ sender not started, falling back to postgres", { + module: "nuq/rabbitmq", + }); } } @@ -732,12 +781,19 @@ class NuQ { if (success) { const job = result.rows[0]; if (this.nuqWaitMode === "listen" && !process.env.NUQ_RABBITMQ_URL) { - await nuqPool.query(`SELECT pg_notify('${this.queueName}', $1);`, [job.id + "|completed"]); + await nuqPool.query(`SELECT pg_notify('${this.queueName}', $1);`, [ + job.id + "|completed", + ]); } else if (process.env.NUQ_RABBITMQ_URL && job.listen_channel_id) { - await this.sendJobEnd(job.id, "completed", job.listen_channel_id, _logger); + await this.sendJobEnd( + job.id, + "completed", + job.listen_channel_id, + _logger, + ); } } - + setSpanAttributes(span, { "nuq.job_finished": success, }); @@ -773,7 +829,8 @@ class NuQ { const start = Date.now(); try { - const result = await nuqPool.query(`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;`, + const result = await nuqPool.query( + `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;`, [id, lock, failedReason], ); @@ -782,9 +839,16 @@ class NuQ { if (success) { const job = result.rows[0]; if (this.nuqWaitMode === "listen" && !process.env.NUQ_RABBITMQ_URL) { - await nuqPool.query(`SELECT pg_notify('${this.queueName}', $1);`, [job.id + "|failed"]); + await nuqPool.query(`SELECT pg_notify('${this.queueName}', $1);`, [ + job.id + "|failed", + ]); } else if (process.env.NUQ_RABBITMQ_URL && job.listen_channel_id) { - await this.sendJobEnd(job.id, "failed", job.listen_channel_id, _logger); + await this.sendJobEnd( + job.id, + "failed", + job.listen_channel_id, + _logger, + ); } }