diff --git a/.server-changes/retry-requeue-nack-budget.md b/.server-changes/retry-requeue-nack-budget.md new file mode 100644 index 00000000000..c8cabee6661 --- /dev/null +++ b/.server-changes/retry-requeue-nack-budget.md @@ -0,0 +1,6 @@ +--- +area: webapp +type: fix +--- + +Task retries that wait in the queue no longer count against the queue's internal redelivery limit, so runs with many long-delay retries are not wrongly failed with TASK_RUN_DEQUEUED_MAX_RETRIES. diff --git a/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts b/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts index e999d35676d..5b4f5a5ace3 100644 --- a/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts +++ b/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts @@ -1126,6 +1126,7 @@ export class RunAttemptSystem { orgId: env.organizationId, projectId: env.project.id, timestamp: retryAt.getTime(), + resetQueueAttempts: !forceRequeue, error: { type: "INTERNAL_ERROR", code: "TASK_RUN_DEQUEUED_MAX_RETRIES", @@ -1249,6 +1250,7 @@ export class RunAttemptSystem { checkpointId, completedWaitpoints, batchId, + resetQueueAttempts = false, tx, }: { run: { id: string }; @@ -1269,6 +1271,12 @@ export class RunAttemptSystem { index?: number; }[]; batchId?: string; + /** + * Pass when the worker reported the attempt's failure itself (an ordinary task retry), so the + * queue's redelivery budget is reset rather than consumed. Engine-detected stalls and dequeue + * failures leave it unset so a run that never comes back healthy is still bounded. + */ + resetQueueAttempts?: boolean; }): Promise<{ wasRequeued: boolean } & ExecutionResult> { const prisma = tx ?? this.$.prisma; @@ -1278,6 +1286,7 @@ export class RunAttemptSystem { orgId, messageId: run.id, retryAt: timestamp, + resetAttemptCount: resetQueueAttempts, }); if (!gotRequeued) { diff --git a/internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts b/internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts index 36b04f768b2..146df29b4fd 100644 --- a/internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts +++ b/internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts @@ -491,6 +491,134 @@ describe("RunEngine attempt failures", () => { } }); + containerTest( + "task retries routed through the queue do not consume the queue's nack budget", + async ({ prisma, redisOptions }) => { + const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION"); + + const engine = new RunEngine({ + prisma, + worker: { + redis: redisOptions, + workers: 1, + tasksPerWorker: 10, + pollIntervalMs: 100, + }, + queue: { + redis: redisOptions, + retryOptions: { + maxAttempts: 2, + }, + masterQueueConsumersDisabled: true, + processWorkerQueueDebounceMs: 50, + }, + runLock: { + redis: redisOptions, + }, + machines: { + defaultMachine: "small-1x", + machines: { + "small-1x": { + name: "small-1x" as const, + cpu: 0.5, + memory: 0.5, + centsPerMs: 0.0001, + }, + }, + baseCostInCents: 0.0001, + }, + retryWarmStartThresholdMs: 0, + tracer: trace.getTracer("test", "0.0.0"), + }); + + try { + const taskIdentifier = "test-task"; + const taskMaxAttempts = 4; + + await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier, undefined, { + maxAttempts: taskMaxAttempts, + factor: 1, + minTimeoutInMs: 100, + maxTimeoutInMs: 100, + randomize: false, + }); + + const run = await engine.trigger( + { + number: 1, + friendlyId: "run_1234", + environment: authenticatedEnvironment, + taskIdentifier, + payload: "{}", + payloadType: "application/json", + context: {}, + traceContext: {}, + traceId: "t12345", + spanId: "s12345", + workerQueue: "main", + queue: "task/test-task", + isTest: false, + tags: [], + }, + prisma + ); + + const error = { + type: "BUILT_IN_ERROR" as const, + name: "Error", + message: "boom", + stackTrace: "Error: boom", + }; + + for (let attempt = 1; attempt <= taskMaxAttempts; attempt++) { + await setTimeout(500); + await engine.runQueue.processMasterQueueForEnvironment(authenticatedEnvironment.id); + const dequeued = await engine.dequeueFromWorkerQueue({ + consumerId: "test_12345", + workerQueue: "main", + }); + expect(dequeued.length).toBe(1); + + const attemptResult = await engine.startRunAttempt({ + runId: dequeued[0].run.id, + snapshotId: dequeued[0].snapshot.id, + }); + expect(attemptResult.run.attemptNumber).toBe(attempt); + + const result = await engine.completeRunAttempt({ + runId: dequeued[0].run.id, + snapshotId: attemptResult.snapshot.id, + completion: { + ok: false, + id: dequeued[0].run.id, + error, + retry: { + timestamp: Date.now() + 100, + delay: 100, + }, + }, + }); + + if (attempt < taskMaxAttempts) { + expect(result.attemptStatus).toBe("RETRY_QUEUED"); + expect(result.run.status).toBe("PENDING"); + } else { + expect(result.attemptStatus).toBe("RUN_FINISHED"); + expect(result.run.status).toBe("COMPLETED_WITH_ERRORS"); + } + } + + const executionData = await engine.getRunExecutionData({ runId: run.id }); + assertNonNullable(executionData); + expect(executionData.run.attemptNumber).toBe(taskMaxAttempts); + expect(executionData.run.status).toBe("COMPLETED_WITH_ERRORS"); + expect(await engine.runQueue.lengthOfDeadLetterQueue(authenticatedEnvironment)).toBe(0); + } finally { + await engine.quit(); + } + } + ); + containerTest("OOM retry on larger machine", async ({ prisma, redisOptions }) => { //create environment const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION"); diff --git a/internal-packages/run-engine/src/run-queue/index.ts b/internal-packages/run-engine/src/run-queue/index.ts index 48a84785134..4afa5b2bab9 100644 --- a/internal-packages/run-engine/src/run-queue/index.ts +++ b/internal-packages/run-engine/src/run-queue/index.ts @@ -1116,12 +1116,19 @@ export class RunQueue { messageId, retryAt, incrementAttemptCount = true, + resetAttemptCount = false, skipDequeueProcessing = false, }: { orgId: string; messageId: string; retryAt?: number; incrementAttemptCount?: boolean; + /** + * Zero the message's attempt counter instead of incrementing it. The counter is the budget + * for dequeues that never reach execution; a caller that knows an attempt did execute passes + * this so an ordinary task retry cannot exhaust it and dead-letter the run. + */ + resetAttemptCount?: boolean; skipDequeueProcessing?: boolean; }) { return this.#trace( @@ -1148,7 +1155,9 @@ export class RunQueue { [SemanticAttributes.WORKER_QUEUE]: this.#getWorkerQueueFromMessage(message), }); - if (incrementAttemptCount) { + if (resetAttemptCount) { + message.attempt = 0; + } else if (incrementAttemptCount) { message.attempt = message.attempt + 1; if (message.attempt >= maxAttempts) { await this.#callMoveToDeadLetterQueue({ message }); 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 index 8711c816b1b..889b32ec338 100644 --- a/internal-packages/run-engine/src/run-queue/tests/nack.test.ts +++ b/internal-packages/run-engine/src/run-queue/tests/nack.test.ts @@ -130,6 +130,81 @@ describe("RunQueue.nackMessage", () => { } }); + redisTest( + "nacking with resetAttemptCount zeroes the counter instead of dead-lettering", + async ({ redisContainer }) => { + const queue = new RunQueue({ + ...testOptions, + retryOptions: { + ...testOptions.retryOptions, + maxAttempts: 2, + }, + 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 { + await queue.enqueueMessage({ + env: authenticatedEnvDev, + message: messageDev, + workerQueue: authenticatedEnvDev.id, + }); + + await setTimeout(1000); + + const dequeued = await queue.dequeueMessageFromWorkerQueue( + "test_12345", + authenticatedEnvDev.id + ); + assertNonNullable(dequeued); + + const first = await queue.nackMessage({ + orgId: messageDev.orgId, + messageId: messageDev.runId, + }); + expect(first).toBe(true); + + const afterFirst = await queue.readMessage(messageDev.orgId, messageDev.runId); + expect(afterFirst?.attempt).toBe(1); + + await setTimeout(1000); + + const dequeued2 = await queue.dequeueMessageFromWorkerQueue( + "test_12345", + authenticatedEnvDev.id + ); + assertNonNullable(dequeued2); + + // A plain nack here would hit maxAttempts and dead-letter the run + const second = await queue.nackMessage({ + orgId: messageDev.orgId, + messageId: messageDev.runId, + resetAttemptCount: true, + }); + expect(second).toBe(true); + + const afterReset = await queue.readMessage(messageDev.orgId, messageDev.runId); + expect(afterReset?.attempt).toBe(0); + + expect(await queue.lengthOfEnvQueue(authenticatedEnvDev)).toBe(1); + expect(await queue.lengthOfDeadLetterQueue(authenticatedEnvDev)).toBe(0); + } finally { + await queue.quit(); + } + } + ); + redisTest( "nacking a message with maxAttempts reached should be moved to dead letter queue", async ({ redisContainer }) => {