diff --git a/.vscode/launch.json b/.vscode/launch.json index 2bd25ee36..8242758d3 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -146,7 +146,7 @@ "type": "node-terminal", "request": "launch", "name": "Debug RunQueue tests", - "command": "pnpm run test ./src/run-queue/tests/dequeueMessageFromMasterQueue.test.ts", + "command": "pnpm run test ./src/run-queue/index.test.ts", "cwd": "${workspaceFolder}/internal-packages/run-engine", "sourceMaps": true } diff --git a/internal-packages/run-engine/src/run-queue/index.test.ts b/internal-packages/run-engine/src/run-queue/index.test.ts index 7fcae56ee..952df3b28 100644 --- a/internal-packages/run-engine/src/run-queue/index.test.ts +++ b/internal-packages/run-engine/src/run-queue/index.test.ts @@ -826,7 +826,7 @@ describe("RunQueue", () => { } ); - redisTest("Dead Letter Queue", { timeout: 8_000 }, async ({ redisContainer, redisOptions }) => { + redisTest("Dead Letter Queue", async ({ redisContainer, redisOptions }) => { const queue = new RunQueue({ ...testOptions, retryOptions: { @@ -891,36 +891,30 @@ describe("RunQueue", () => { expect(taskConcurrency2).toBe(0); //check the message is still there - const exists2 = await redis.exists(key); - expect(exists2).toBe(1); + const message = await queue.readMessage(messages[0].message.orgId, messages[0].messageId); + expect(message).toBeDefined(); - //check it's in the dlq - const dlqKey = "dlq"; - const dlqExists = await redis.exists(dlqKey); - expect(dlqExists).toBe(1); - const dlqMembers = await redis.zrange(dlqKey, 0, -1); - expect(dlqMembers).toContain(messageProd.runId); + const deadLetterQueueLengthBefore = await queue.lengthOfDeadLetterQueue(authenticatedEnvProd); + expect(deadLetterQueueLengthBefore).toBe(1); + + const existsInDlq = await queue.messageInDeadLetterQueue( + authenticatedEnvProd, + messageProd.runId + ); + expect(existsInDlq).toBe(true); //redrive - const redisClient = createRedisClient({ - host: redisContainer.getHost(), - port: redisContainer.getPort(), - password: redisContainer.getPassword(), - }); - - // Publish redrive message - await redisClient.publish( - "rq:redrive", - JSON.stringify({ runId: messageProd.runId, orgId: messageProd.orgId }) - ); + await queue.redriveMessage(authenticatedEnvProd, messageProd.runId); // Wait for the item to be redrived and processed await setTimeout(5_000); - await redisClient.quit(); //shouldn't be in the dlq now - const dlqMembersAfter = await redis.zrange(dlqKey, 0, -1); - expect(dlqMembersAfter).not.toContain(messageProd.runId); + const existsInDlqAfter = await queue.messageInDeadLetterQueue( + authenticatedEnvProd, + messageProd.runId + ); + expect(existsInDlqAfter).toBe(false); //dequeue const messages3 = await queue.dequeueMessageFromMasterQueue("test_12345", "main", 10); diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 64f466b85..4d1734753 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -161,6 +161,28 @@ export class RunQueue { return this.redis.zcard(this.keys.envQueueKey(env)); } + public async lengthOfDeadLetterQueue(env: MinimalAuthenticatedEnvironment) { + return this.redis.zcard(this.keys.deadLetterQueueKey(env)); + } + + public async messageInDeadLetterQueue(env: MinimalAuthenticatedEnvironment, messageId: string) { + const result = await this.redis.zscore(this.keys.deadLetterQueueKey(env), messageId); + return !!result; + } + + public async redriveMessage(env: MinimalAuthenticatedEnvironment, messageId: string) { + // Publish redrive message + await this.redis.publish( + "rq:redrive", + JSON.stringify({ + runId: messageId, + orgId: env.organization.id, + envId: env.id, + projectId: env.project.id, + }) + ); + } + public async oldestMessageInQueue( env: MinimalAuthenticatedEnvironment, queue: string, @@ -275,6 +297,10 @@ export class RunQueue { return this.redis.scard(this.keys.projectCurrentConcurrencyKey(env)); } + public async messageExists(orgId: string, messageId: string) { + return this.redis.exists(this.keys.messageKey(orgId, messageId)); + } + public async currentConcurrencyOfTask( env: MinimalAuthenticatedEnvironment, taskIdentifier: string @@ -282,6 +308,41 @@ export class RunQueue { return this.redis.scard(this.keys.taskIdentifierCurrentConcurrencyKey(env, taskIdentifier)); } + public async readMessage(orgId: string, messageId: string) { + return this.#trace( + "readMessage", + async (span) => { + const rawMessage = await this.redis.get(this.keys.messageKey(orgId, messageId)); + + if (!rawMessage) { + return; + } + + const message = OutputPayload.safeParse(JSON.parse(rawMessage)); + + if (!message.success) { + this.logger.error(`[${this.name}] Failed to parse message`, { + messageId, + error: message.error, + service: this.name, + }); + + return; + } + + return message.data; + }, + { + attributes: { + [SEMATTRS_MESSAGING_OPERATION]: "receive", + [SEMATTRS_MESSAGE_ID]: messageId, + [SEMATTRS_MESSAGING_SYSTEM]: "marqs", + [SemanticAttributes.RUN_ID]: messageId, + }, + } + ); + } + public async enqueueMessage({ env, message, @@ -434,7 +495,7 @@ export class RunQueue { return this.#trace( "acknowledgeMessage", async (span) => { - const message = await this.#readMessage(orgId, messageId); + const message = await this.readMessage(orgId, messageId); if (!message) { this.logger.log(`[${this.name}].acknowledgeMessage() message not found`, { @@ -452,18 +513,7 @@ export class RunQueue { }); await this.#callAcknowledgeMessage({ - messageId, - messageQueue: message.queue, - masterQueues: message.masterQueues, - messageKey: this.keys.messageKey(orgId, messageId), - concurrencyKey: this.keys.currentConcurrencyKeyFromQueue(message.queue), - envConcurrencyKey: this.keys.envCurrentConcurrencyKeyFromQueue(message.queue), - taskConcurrencyKey: this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue( - message.queue, - message.taskIdentifier - ), - envQueueKey: this.keys.envQueueKeyFromQueue(message.queue), - projectConcurrencyKey: this.keys.projectCurrentConcurrencyKeyFromQueue(message.queue), + message, }); }, { @@ -497,7 +547,7 @@ export class RunQueue { async (span) => { const maxAttempts = this.retryOptions.maxAttempts ?? defaultRetrySettings.maxAttempts; - const message = await this.#readMessage(orgId, messageId); + const message = await this.readMessage(orgId, messageId); if (!message) { this.logger.log(`[${this.name}].nackMessage() message not found`, { orgId, @@ -516,75 +566,16 @@ export class RunQueue { [SemanticAttributes.MASTER_QUEUES]: message.masterQueues.join(","), }); - const messageKey = this.keys.messageKey(orgId, messageId); - const messageQueue = message.queue; - const concurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue); - const envConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue); - const taskConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue( - message.queue, - message.taskIdentifier - ); - const projectConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue( - message.queue - ); - const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue); - if (incrementAttemptCount) { message.attempt = message.attempt + 1; if (message.attempt >= maxAttempts) { - await this.redis.moveToDeadLetterQueue( - messageKey, - messageQueue, - concurrencyKey, - envConcurrencyKey, - projectConcurrencyKey, - envQueueKey, - taskConcurrencyKey, - "dlq", - messageId, - messageQueue, - JSON.stringify(message.masterQueues), - this.options.redis.keyPrefix ?? "" - ); + await this.#callMoveToDeadLetterQueue({ message }); return false; } } - const nextRetryDelay = calculateNextRetryDelay(this.retryOptions, message.attempt); - const messageScore = retryAt ?? (nextRetryDelay ? Date.now() + nextRetryDelay : Date.now()); + await this.#callNackMessage({ message }); - this.logger.debug("Calling nackMessage", { - messageKey, - messageQueue, - masterQueues: message.masterQueues, - concurrencyKey, - envConcurrencyKey, - projectConcurrencyKey, - envQueueKey, - taskConcurrencyKey, - messageId, - messageScore, - attempt: message.attempt, - service: this.name, - }); - - await this.redis.nackMessage( - //keys - messageKey, - messageQueue, - concurrencyKey, - envConcurrencyKey, - projectConcurrencyKey, - envQueueKey, - taskConcurrencyKey, - //args - messageId, - messageQueue, - JSON.stringify(message), - String(messageScore), - JSON.stringify(message.masterQueues), - this.options.redis.keyPrefix ?? "" - ); return true; }, { @@ -606,7 +597,7 @@ export class RunQueue { return this.#trace( "releaseConcurrency", async (span) => { - const message = await this.#readMessage(orgId, messageId); + const message = await this.readMessage(orgId, messageId); if (!message) { this.logger.log(`[${this.name}].acknowledgeMessage() message not found`, { @@ -652,7 +643,7 @@ export class RunQueue { return this.#trace( "reacquireConcurrency", async (span) => { - const message = await this.#readMessage(orgId, messageId); + const message = await this.readMessage(orgId, messageId); if (!message) { this.logger.log(`[${this.name}].acknowledgeMessage() message not found`, { @@ -702,16 +693,21 @@ export class RunQueue { private async handleRedriveMessage(channel: string, message: string) { try { - const { runId, orgId } = JSON.parse(message) as any; - if (typeof orgId !== "string" || typeof runId !== "string") { + const { runId, envId, projectId, orgId } = JSON.parse(message) as any; + if ( + typeof orgId !== "string" || + typeof runId !== "string" || + typeof envId !== "string" || + typeof projectId !== "string" + ) { this.logger.error( - "handleRedriveMessage: invalid message format: runId and orgId must be strings", + "handleRedriveMessage: invalid message format: runId, envId, projectId and orgId must be strings", { message, channel } ); return; } - const data = await this.#readMessage(orgId, runId); + const data = await this.readMessage(orgId, runId); if (!data) { this.logger.error(`handleRedriveMessage: couldn't read message`, { orgId, runId, channel }); @@ -739,7 +735,10 @@ export class RunQueue { }); //remove from the dlq - const result = await this.redis.zrem("dlq", runId); + const result = await this.redis.zrem( + this.keys.deadLetterQueueKey({ envId, orgId, projectId }), + runId + ); if (result === 0) { this.logger.error(`handleRedriveMessage: couldn't remove message from dlq`, { @@ -800,41 +799,6 @@ export class RunQueue { this.subscriber.on("message", this.handleRedriveMessage.bind(this)); } - async #readMessage(orgId: string, messageId: string) { - return this.#trace( - "readMessage", - async (span) => { - const rawMessage = await this.redis.get(this.keys.messageKey(orgId, messageId)); - - if (!rawMessage) { - return; - } - - const message = OutputPayload.safeParse(JSON.parse(rawMessage)); - - if (!message.success) { - this.logger.error(`[${this.name}] Failed to parse message`, { - messageId, - error: message.error, - service: this.name, - }); - - return; - } - - return message.data; - }, - { - attributes: { - [SEMATTRS_MESSAGING_OPERATION]: "receive", - [SEMATTRS_MESSAGE_ID]: messageId, - [SEMATTRS_MESSAGING_SYSTEM]: "marqs", - [SemanticAttributes.RUN_ID]: messageId, - }, - } - ); - } - async #callEnqueueMessage( message: OutputPayload, masterQueues: string[], @@ -1085,35 +1049,30 @@ export class RunQueue { }; } - async #callAcknowledgeMessage({ - messageId, - masterQueues, - messageKey, - messageQueue, - concurrencyKey, - envConcurrencyKey, - taskConcurrencyKey, - envQueueKey, - projectConcurrencyKey, - }: { - masterQueues: string[]; - messageKey: string; - messageQueue: string; - concurrencyKey: string; - envConcurrencyKey: string; - taskConcurrencyKey: string; - envQueueKey: string; - projectConcurrencyKey: string; - messageId: string; - }) { + async #callAcknowledgeMessage({ message }: { message: OutputPayload }) { + const messageId = message.runId; + const messageKey = this.keys.messageKey(message.orgId, messageId); + const messageQueue = message.queue; + const queueCurrentConcurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue); + const envCurrentConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue); + const projectCurrentConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue( + message.queue + ); + const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue); + const taskCurrentConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue( + message.queue, + message.taskIdentifier + ); + const masterQueues = message.masterQueues; + this.logger.debug("Calling acknowledgeMessage", { messageKey, messageQueue, - concurrencyKey, - envConcurrencyKey, - projectConcurrencyKey, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + projectCurrentConcurrencyKey, envQueueKey, - taskConcurrencyKey, + taskCurrentConcurrencyKey, messageId, masterQueues, service: this.name, @@ -1125,11 +1084,11 @@ export class RunQueue { return this.redis.acknowledgeMessage( messageKey, messageQueue, - concurrencyKey, - envConcurrencyKey, - projectConcurrencyKey, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + projectCurrentConcurrencyKey, envQueueKey, - taskConcurrencyKey, + taskCurrentConcurrencyKey, queueReserveConcurrencyKey, envReserveConcurrencyKey, messageId, @@ -1139,6 +1098,90 @@ export class RunQueue { ); } + async #callNackMessage({ message, retryAt }: { message: OutputPayload; retryAt?: number }) { + const messageId = message.runId; + const messageKey = this.keys.messageKey(message.orgId, message.runId); + const messageQueue = message.queue; + const queueCurrentConcurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue); + const envCurrentConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue); + const projectCurrentConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue( + message.queue + ); + const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue); + const taskCurrentConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue( + message.queue, + message.taskIdentifier + ); + + const nextRetryDelay = calculateNextRetryDelay(this.retryOptions, message.attempt); + const messageScore = retryAt ?? (nextRetryDelay ? Date.now() + nextRetryDelay : Date.now()); + + this.logger.debug("Calling nackMessage", { + messageKey, + messageQueue, + masterQueues: message.masterQueues, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + projectCurrentConcurrencyKey, + envQueueKey, + taskCurrentConcurrencyKey, + messageId, + messageScore, + attempt: message.attempt, + service: this.name, + }); + + await this.redis.nackMessage( + //keys + messageKey, + messageQueue, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + projectCurrentConcurrencyKey, + envQueueKey, + taskCurrentConcurrencyKey, + //args + messageId, + messageQueue, + JSON.stringify(message), + String(messageScore), + JSON.stringify(message.masterQueues), + this.options.redis.keyPrefix ?? "" + ); + } + + async #callMoveToDeadLetterQueue({ message }: { message: OutputPayload }) { + const messageId = message.runId; + const messageKey = this.keys.messageKey(message.orgId, message.runId); + const messageQueue = message.queue; + const queueCurrentConcurrencyKey = this.keys.currentConcurrencyKeyFromQueue(message.queue); + const envCurrentConcurrencyKey = this.keys.envCurrentConcurrencyKeyFromQueue(message.queue); + const projectCurrentConcurrencyKey = this.keys.projectCurrentConcurrencyKeyFromQueue( + message.queue + ); + const envQueueKey = this.keys.envQueueKeyFromQueue(message.queue); + const taskCurrentConcurrencyKey = this.keys.taskIdentifierCurrentConcurrencyKeyFromQueue( + message.queue, + message.taskIdentifier + ); + const deadLetterQueueKey = this.keys.deadLetterQueueKeyFromQueue(message.queue); + + await this.redis.moveToDeadLetterQueue( + messageKey, + messageQueue, + queueCurrentConcurrencyKey, + envCurrentConcurrencyKey, + projectCurrentConcurrencyKey, + envQueueKey, + taskCurrentConcurrencyKey, + deadLetterQueueKey, + messageId, + messageQueue, + JSON.stringify(message.masterQueues), + this.options.redis.keyPrefix ?? "" + ); + } + #callUpdateGlobalConcurrencyLimits({ envConcurrencyLimitKey, envConcurrencyLimit, @@ -1501,11 +1544,11 @@ redis.call('SREM', envReserveConcurrencyKey, messageId) -- Keys: local messageKey = KEYS[1] local messageQueueKey = KEYS[2] -local concurrencyKey = KEYS[3] -local envConcurrencyKey = KEYS[4] -local projectConcurrencyKey = KEYS[5] +local queueCurrentConcurrencyKey = KEYS[3] +local envCurrentConcurrencyKey = KEYS[4] +local projectCurrentConcurrencyKey = KEYS[5] local envQueueKey = KEYS[6] -local taskConcurrencyKey = KEYS[7] +local taskCurrentConcurrencyKey = KEYS[7] -- Args: local messageId = ARGV[1] @@ -1519,10 +1562,10 @@ local keyPrefix = ARGV[6] redis.call('SET', messageKey, messageData) -- Update the concurrency keys -redis.call('SREM', concurrencyKey, messageId) -redis.call('SREM', envConcurrencyKey, messageId) -redis.call('SREM', projectConcurrencyKey, messageId) -redis.call('SREM', taskConcurrencyKey, messageId) +redis.call('SREM', queueCurrentConcurrencyKey, messageId) +redis.call('SREM', envCurrentConcurrencyKey, messageId) +redis.call('SREM', projectCurrentConcurrencyKey, messageId) +redis.call('SREM', taskCurrentConcurrencyKey, messageId) -- Enqueue the message into the queue redis.call('ZADD', messageQueueKey, messageScore, messageId) @@ -1547,11 +1590,11 @@ end -- Keys: local messageKey = KEYS[1] local messageQueue = KEYS[2] -local concurrencyKey = KEYS[3] +local queueCurrentConcurrencyKey = KEYS[3] local envCurrentConcurrencyKey = KEYS[4] local projectCurrentConcurrencyKey = KEYS[5] local envQueueKey = KEYS[6] -local taskCurrentConcurrencyKey = KEYS[7] +local taskCurrentConcurrencyKeyPrefix = KEYS[7] local deadLetterQueueKey = KEYS[8] -- Args: @@ -1579,10 +1622,10 @@ end redis.call('ZADD', deadLetterQueueKey, tonumber(redis.call('TIME')[1]), messageId) -- Update the concurrency keys -redis.call('SREM', concurrencyKey, messageId) +redis.call('SREM', queueCurrentConcurrencyKey, messageId) redis.call('SREM', envCurrentConcurrencyKey, messageId) redis.call('SREM', projectCurrentConcurrencyKey, messageId) -redis.call('SREM', taskCurrentConcurrencyKey, messageId) +redis.call('SREM', taskCurrentConcurrencyKeyPrefix, messageId) `, }); @@ -1760,11 +1803,11 @@ declare module "@internal/redis" { nackMessage( messageKey: string, messageQueue: string, - concurrencyKey: string, - envConcurrencyKey: string, - projectConcurrencyKey: string, + queueCurrentConcurrencyKey: string, + envCurrentConcurrencyKey: string, + projectCurrentConcurrencyKey: string, envQueueKey: string, - taskConcurrencyKey: string, + taskCurrentConcurrencyKey: string, messageId: string, messageQueueName: string, messageData: string, @@ -1777,11 +1820,11 @@ declare module "@internal/redis" { moveToDeadLetterQueue( messageKey: string, messageQueue: string, - concurrencyKey: string, - envConcurrencyKey: string, - projectConcurrencyKey: string, + queueCurrentConcurrencyKey: string, + envCurrentConcurrencyKey: string, + projectCurrentConcurrencyKey: string, envQueueKey: string, - taskConcurrencyKey: string, + taskCurrentConcurrencyKey: string, deadLetterQueueKey: string, messageId: string, messageQueueName: string, diff --git a/internal-packages/run-engine/src/run-queue/keyProducer.ts b/internal-packages/run-engine/src/run-queue/keyProducer.ts index 1feac4a7c..a24bf8dbd 100644 --- a/internal-packages/run-engine/src/run-queue/keyProducer.ts +++ b/internal-packages/run-engine/src/run-queue/keyProducer.ts @@ -13,6 +13,7 @@ const constants = { TASK_PART: "task", MESSAGE_PART: "message", RESERVE_CONCURRENCY_PART: "reserveConcurrency", + DEAD_LETTER_QUEUE_PART: "deadLetter", } as const; export class RunQueueFullKeyProducer implements RunQueueKeyProducer { @@ -256,6 +257,30 @@ export class RunQueueFullKeyProducer implements RunQueueKeyProducer { return this.envReserveConcurrencyKey(descriptor); } + deadLetterQueueKey(env: MinimalAuthenticatedEnvironment): string; + deadLetterQueueKey(env: EnvDescriptor): string; + deadLetterQueueKey(envOrDescriptor: EnvDescriptor | MinimalAuthenticatedEnvironment): string { + if ("id" in envOrDescriptor) { + return [ + this.orgKeySection(envOrDescriptor.organization.id), + this.projKeySection(envOrDescriptor.project.id), + this.envKeySection(envOrDescriptor.id), + constants.DEAD_LETTER_QUEUE_PART, + ].join(":"); + } else { + return [ + this.orgKeySection(envOrDescriptor.orgId), + this.projKeySection(envOrDescriptor.projectId), + this.envKeySection(envOrDescriptor.envId), + constants.DEAD_LETTER_QUEUE_PART, + ].join(":"); + } + } + deadLetterQueueKeyFromQueue(queue: string): string { + const descriptor = this.descriptorFromQueue(queue); + + return this.deadLetterQueueKey(descriptor); + } private queueReserveConcurrencyKeyFromDescriptor(descriptor: QueueDescriptor) { return [ diff --git a/internal-packages/run-engine/src/run-queue/tests/nack.test.ts b/internal-packages/run-engine/src/run-queue/tests/nack.test.ts new file mode 100644 index 000000000..67ec26b58 --- /dev/null +++ b/internal-packages/run-engine/src/run-queue/tests/nack.test.ts @@ -0,0 +1,244 @@ +import { redisTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { Logger } from "@trigger.dev/core/logger"; +import { describe } from "node:test"; +import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; +import { RunQueue } from "../index.js"; +import { RunQueueFullKeyProducer } from "../keyProducer.js"; +import { InputPayload } from "../types.js"; +import { setTimeout } from "node:timers/promises"; + +const testOptions = { + name: "rq", + tracer: trace.getTracer("rq"), + workers: 1, + defaultEnvConcurrency: 25, + logger: new Logger("RunQueue", "warn"), + retryOptions: { + maxAttempts: 5, + factor: 1.1, + minTimeoutInMs: 100, + maxTimeoutInMs: 1_000, + randomize: true, + }, + keys: new RunQueueFullKeyProducer(), +}; + +const authenticatedEnvDev = { + id: "e1234", + type: "DEVELOPMENT" as const, + maximumConcurrencyLimit: 10, + project: { id: "p1234" }, + organization: { id: "o1234" }, +}; + +const messageDev: InputPayload = { + runId: "r4321", + taskIdentifier: "task/my-task", + orgId: "o1234", + projectId: "p1234", + environmentId: "e4321", + environmentType: "DEVELOPMENT", + queue: "task/my-task", + timestamp: Date.now(), + attempt: 0, +}; + +vi.setConfig({ testTimeout: 60_000 }); + +describe("RunQueue.nackMessage", () => { + redisTest("nacking a message clears all concurrency", async ({ redisContainer }) => { + const queue = new RunQueue({ + ...testOptions, + queueSelectionStrategy: new FairQueueSelectionStrategy({ + redis: { + keyPrefix: "runqueue:test:", + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + keys: testOptions.keys, + }), + redis: { + keyPrefix: "runqueue:test:", + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + }); + + try { + const envMasterQueue = `env:${authenticatedEnvDev.id}`; + + // Enqueue message with reserve concurrency + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: messageDev, + masterQueues: ["main", envMasterQueue], + reserveConcurrency: { + messageId: messageDev.runId, + recursiveQueue: true, + }, + }); + + // Verify reserve concurrency is set + const queueReserveConcurrency = await queue.reserveConcurrencyOfQueue( + authenticatedEnvDev, + messageDev.queue + ); + expect(queueReserveConcurrency).toBe(1); + + const envReserveConcurrency = await queue.reserveConcurrencyOfEnvironment( + authenticatedEnvDev + ); + expect(envReserveConcurrency).toBe(1); + + // Dequeue message + const dequeued = await queue.dequeueMessageFromMasterQueue("test_12345", envMasterQueue, 10); + expect(dequeued.length).toBe(1); + + // Verify current concurrency is set and reserve is cleared + const queueCurrentConcurrency = await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + messageDev.queue + ); + expect(queueCurrentConcurrency).toBe(1); + + const envCurrentConcurrency = await queue.currentConcurrencyOfEnvironment( + authenticatedEnvDev + ); + expect(envCurrentConcurrency).toBe(1); + + // Nack the message + await queue.nackMessage({ + orgId: messageDev.orgId, + messageId: messageDev.runId, + }); + + // Verify all concurrency is cleared + const queueCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfQueue( + authenticatedEnvDev, + messageDev.queue + ); + expect(queueCurrentConcurrencyAfterNack).toBe(0); + + const envCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfEnvironment( + authenticatedEnvDev + ); + expect(envCurrentConcurrencyAfterNack).toBe(0); + + const projectCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfProject( + authenticatedEnvDev + ); + expect(projectCurrentConcurrencyAfterNack).toBe(0); + + const taskCurrentConcurrencyAfterNack = await queue.currentConcurrencyOfTask( + authenticatedEnvDev, + messageDev.taskIdentifier + ); + expect(taskCurrentConcurrencyAfterNack).toBe(0); + + const envQueueLength = await queue.lengthOfEnvQueue(authenticatedEnvDev); + expect(envQueueLength).toBe(1); + + const message = await queue.readMessage(messageDev.orgId, messageDev.runId); + expect(message?.attempt).toBe(1); + + //we need to wait because the default wait is 1 second + await setTimeout(300); + + // Now we should be able to dequeue it again + const dequeued2 = await queue.dequeueMessageFromMasterQueue("test_12345", envMasterQueue, 10); + expect(dequeued2.length).toBe(1); + } finally { + await queue.quit(); + } + }); + + redisTest( + "nacking a message with maxAttempts reached should be moved to dead letter queue", + async ({ redisContainer }) => { + const queue = new RunQueue({ + ...testOptions, + retryOptions: { + ...testOptions.retryOptions, + maxAttempts: 2, // Set lower for testing + }, + queueSelectionStrategy: new FairQueueSelectionStrategy({ + redis: { + keyPrefix: "runqueue:test:", + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + keys: testOptions.keys, + }), + redis: { + keyPrefix: "runqueue:test:", + host: redisContainer.getHost(), + port: redisContainer.getPort(), + }, + }); + + try { + const envMasterQueue = `env:${authenticatedEnvDev.id}`; + + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: messageDev, + masterQueues: ["main", envMasterQueue], + }); + + const dequeued = await queue.dequeueMessageFromMasterQueue( + "test_12345", + envMasterQueue, + 10 + ); + expect(dequeued.length).toBe(1); + + await queue.nackMessage({ + orgId: messageDev.orgId, + messageId: messageDev.runId, + }); + + // Wait for any requeue delay + await setTimeout(300); + + // Message should not be requeued as max attempts reached + const envQueueLength = await queue.lengthOfEnvQueue(authenticatedEnvDev); + expect(envQueueLength).toBe(1); + + const message = await queue.readMessage(messageDev.orgId, messageDev.runId); + expect(message?.attempt).toBe(1); + + // Now we dequeue and nack again, and it should be moved to dead letter queue + const dequeued3 = await queue.dequeueMessageFromMasterQueue( + "test_12345", + envMasterQueue, + 10 + ); + expect(dequeued3.length).toBe(1); + + const envQueueLengthDequeue = await queue.lengthOfEnvQueue(authenticatedEnvDev); + expect(envQueueLengthDequeue).toBe(0); + + const deadLetterQueueLengthBefore = await queue.lengthOfDeadLetterQueue( + authenticatedEnvDev + ); + expect(deadLetterQueueLengthBefore).toBe(0); + + await queue.nackMessage({ + orgId: messageDev.orgId, + messageId: messageDev.runId, + }); + + const envQueueLengthAfterNack = await queue.lengthOfEnvQueue(authenticatedEnvDev); + expect(envQueueLengthAfterNack).toBe(0); + + const deadLetterQueueLengthAfterNack = await queue.lengthOfDeadLetterQueue( + authenticatedEnvDev + ); + expect(deadLetterQueueLengthAfterNack).toBe(1); + } finally { + await queue.quit(); + } + } + ); +}); diff --git a/internal-packages/run-engine/src/run-queue/types.ts b/internal-packages/run-engine/src/run-queue/types.ts index 2f005b28b..563cececa 100644 --- a/internal-packages/run-engine/src/run-queue/types.ts +++ b/internal-packages/run-engine/src/run-queue/types.ts @@ -92,6 +92,10 @@ export interface RunQueueKeyProducer { reserveConcurrencyKeyFromQueue(queue: string): string; envReserveConcurrencyKeyFromQueue(queue: string): string; + + deadLetterQueueKey(env: MinimalAuthenticatedEnvironment): string; + deadLetterQueueKey(env: EnvDescriptor): string; + deadLetterQueueKeyFromQueue(queue: string): string; } export type EnvQueues = { diff --git a/internal-packages/run-engine/vitest.config.ts b/internal-packages/run-engine/vitest.config.ts index e10e77f70..044760d6b 100644 --- a/internal-packages/run-engine/vitest.config.ts +++ b/internal-packages/run-engine/vitest.config.ts @@ -12,5 +12,6 @@ export default defineConfig({ singleThread: true, }, }, + testTimeout: 60_000, }, }); diff --git a/internal-packages/run-queue/src/run-queue/tests/nack.test.ts b/internal-packages/run-queue/src/run-queue/tests/nack.test.ts new file mode 100644 index 000000000..c553fce41 --- /dev/null +++ b/internal-packages/run-queue/src/run-queue/tests/nack.test.ts @@ -0,0 +1,11 @@ +import { redisTest } from "@internal/testcontainers"; +import { trace } from "@internal/tracing"; +import { Logger } from "@trigger.dev/core/logger"; +import { describe } from "node:test"; +import { FairQueueSelectionStrategy } from "../fairQueueSelectionStrategy.js"; +import { RunQueue } from "../index.js"; +import { RunQueueFullKeyProducer } from "../keyProducer.js"; +import { InputPayload } from "../types.js"; +import { createRedisClient } from "@internal/redis"; + +// ... existing code ...