From cafa1909180aefec07cfd874f42c0ab5dbb5c2db Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Thu, 27 Aug 2026 18:43:51 -0700 Subject: [PATCH 1/3] fix(run-engine): stop task retries consuming the queue nack budget A task retry whose delay is long enough to go back through the queue was counted as a failed dequeue, so after enough retries the run was dead-lettered and failed with TASK_RUN_DEQUEUED_MAX_RETRIES even though every attempt had actually executed. The queue attempt counter is now reset on that path, since a completed attempt proves the run can start. --- .server-changes/retry-requeue-nack-budget.md | 6 + .../src/engine/systems/runAttemptSystem.ts | 8 ++ .../src/engine/tests/attemptFailures.test.ts | 127 ++++++++++++++++++ .../run-engine/src/run-queue/index.ts | 11 +- .../src/run-queue/tests/nack.test.ts | 75 +++++++++++ 5 files changed, 226 insertions(+), 1 deletion(-) create mode 100644 .server-changes/retry-requeue-nack-budget.md 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..7f9e5f8a972 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: true, 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,11 @@ export class RunAttemptSystem { index?: number; }[]; batchId?: string; + /** + * Pass when the run is being requeued after an attempt actually executed (a task retry), so + * the queue's "dequeued but never started" budget is reset rather than consumed. + */ + resetQueueAttempts?: boolean; }): Promise<{ wasRequeued: boolean } & ExecutionResult> { const prisma = tx ?? this.$.prisma; @@ -1278,6 +1285,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..2a475a88949 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,133 @@ 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); + 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 }) => { From fb264d370bca1d842e6e5f67be4fd921d2d4a54e Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Fri, 28 Aug 2026 17:28:18 +0100 Subject: [PATCH 2/3] test(run-engine): process the master queue before dequeuing each retry The requeue-processing job is debounced to fire before the retry becomes due and no consumer runs after it in this test, so a dequeue that happened to run first found nothing. --- .../run-engine/src/engine/tests/attemptFailures.test.ts | 1 + 1 file changed, 1 insertion(+) 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 2a475a88949..146df29b4fd 100644 --- a/internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts +++ b/internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts @@ -572,6 +572,7 @@ describe("RunEngine attempt failures", () => { 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", From c98465b3c3678b395fac05128694c00ba53de5af Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Fri, 28 Aug 2026 17:49:36 +0100 Subject: [PATCH 3/3] fix(run-engine): keep the queue nack budget for engine-detected stalls Only a failure the worker reported itself resets the counter. A heartbeat timeout is requeued through the same path, and a run that keeps stalling still has to be bounded by the queue redelivery limit. --- .../run-engine/src/engine/systems/runAttemptSystem.ts | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts b/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts index 7f9e5f8a972..5b4f5a5ace3 100644 --- a/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts +++ b/internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts @@ -1126,7 +1126,7 @@ export class RunAttemptSystem { orgId: env.organizationId, projectId: env.project.id, timestamp: retryAt.getTime(), - resetQueueAttempts: true, + resetQueueAttempts: !forceRequeue, error: { type: "INTERNAL_ERROR", code: "TASK_RUN_DEQUEUED_MAX_RETRIES", @@ -1272,8 +1272,9 @@ export class RunAttemptSystem { }[]; batchId?: string; /** - * Pass when the run is being requeued after an attempt actually executed (a task retry), so - * the queue's "dequeued but never started" budget is reset rather than consumed. + * 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> {