Skip to content

Commit cafa190

Browse files
committed
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.
1 parent acaa5ec commit cafa190

5 files changed

Lines changed: 226 additions & 1 deletion

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: fix
4+
---
5+
6+
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.

internal-packages/run-engine/src/engine/systems/runAttemptSystem.ts

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1126,6 +1126,7 @@ export class RunAttemptSystem {
11261126
orgId: env.organizationId,
11271127
projectId: env.project.id,
11281128
timestamp: retryAt.getTime(),
1129+
resetQueueAttempts: true,
11291130
error: {
11301131
type: "INTERNAL_ERROR",
11311132
code: "TASK_RUN_DEQUEUED_MAX_RETRIES",
@@ -1249,6 +1250,7 @@ export class RunAttemptSystem {
12491250
checkpointId,
12501251
completedWaitpoints,
12511252
batchId,
1253+
resetQueueAttempts = false,
12521254
tx,
12531255
}: {
12541256
run: { id: string };
@@ -1269,6 +1271,11 @@ export class RunAttemptSystem {
12691271
index?: number;
12701272
}[];
12711273
batchId?: string;
1274+
/**
1275+
* Pass when the run is being requeued after an attempt actually executed (a task retry), so
1276+
* the queue's "dequeued but never started" budget is reset rather than consumed.
1277+
*/
1278+
resetQueueAttempts?: boolean;
12721279
}): Promise<{ wasRequeued: boolean } & ExecutionResult> {
12731280
const prisma = tx ?? this.$.prisma;
12741281

@@ -1278,6 +1285,7 @@ export class RunAttemptSystem {
12781285
orgId,
12791286
messageId: run.id,
12801287
retryAt: timestamp,
1288+
resetAttemptCount: resetQueueAttempts,
12811289
});
12821290

12831291
if (!gotRequeued) {

internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts

Lines changed: 127 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -491,6 +491,133 @@ describe("RunEngine attempt failures", () => {
491491
}
492492
});
493493

494+
containerTest(
495+
"task retries routed through the queue do not consume the queue's nack budget",
496+
async ({ prisma, redisOptions }) => {
497+
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");
498+
499+
const engine = new RunEngine({
500+
prisma,
501+
worker: {
502+
redis: redisOptions,
503+
workers: 1,
504+
tasksPerWorker: 10,
505+
pollIntervalMs: 100,
506+
},
507+
queue: {
508+
redis: redisOptions,
509+
retryOptions: {
510+
maxAttempts: 2,
511+
},
512+
masterQueueConsumersDisabled: true,
513+
processWorkerQueueDebounceMs: 50,
514+
},
515+
runLock: {
516+
redis: redisOptions,
517+
},
518+
machines: {
519+
defaultMachine: "small-1x",
520+
machines: {
521+
"small-1x": {
522+
name: "small-1x" as const,
523+
cpu: 0.5,
524+
memory: 0.5,
525+
centsPerMs: 0.0001,
526+
},
527+
},
528+
baseCostInCents: 0.0001,
529+
},
530+
retryWarmStartThresholdMs: 0,
531+
tracer: trace.getTracer("test", "0.0.0"),
532+
});
533+
534+
try {
535+
const taskIdentifier = "test-task";
536+
const taskMaxAttempts = 4;
537+
538+
await setupBackgroundWorker(engine, authenticatedEnvironment, taskIdentifier, undefined, {
539+
maxAttempts: taskMaxAttempts,
540+
factor: 1,
541+
minTimeoutInMs: 100,
542+
maxTimeoutInMs: 100,
543+
randomize: false,
544+
});
545+
546+
const run = await engine.trigger(
547+
{
548+
number: 1,
549+
friendlyId: "run_1234",
550+
environment: authenticatedEnvironment,
551+
taskIdentifier,
552+
payload: "{}",
553+
payloadType: "application/json",
554+
context: {},
555+
traceContext: {},
556+
traceId: "t12345",
557+
spanId: "s12345",
558+
workerQueue: "main",
559+
queue: "task/test-task",
560+
isTest: false,
561+
tags: [],
562+
},
563+
prisma
564+
);
565+
566+
const error = {
567+
type: "BUILT_IN_ERROR" as const,
568+
name: "Error",
569+
message: "boom",
570+
stackTrace: "Error: boom",
571+
};
572+
573+
for (let attempt = 1; attempt <= taskMaxAttempts; attempt++) {
574+
await setTimeout(500);
575+
const dequeued = await engine.dequeueFromWorkerQueue({
576+
consumerId: "test_12345",
577+
workerQueue: "main",
578+
});
579+
expect(dequeued.length).toBe(1);
580+
581+
const attemptResult = await engine.startRunAttempt({
582+
runId: dequeued[0].run.id,
583+
snapshotId: dequeued[0].snapshot.id,
584+
});
585+
expect(attemptResult.run.attemptNumber).toBe(attempt);
586+
587+
const result = await engine.completeRunAttempt({
588+
runId: dequeued[0].run.id,
589+
snapshotId: attemptResult.snapshot.id,
590+
completion: {
591+
ok: false,
592+
id: dequeued[0].run.id,
593+
error,
594+
retry: {
595+
timestamp: Date.now() + 100,
596+
delay: 100,
597+
},
598+
},
599+
});
600+
601+
if (attempt < taskMaxAttempts) {
602+
expect(result.attemptStatus).toBe("RETRY_QUEUED");
603+
expect(result.run.status).toBe("PENDING");
604+
} else {
605+
expect(result.attemptStatus).toBe("RUN_FINISHED");
606+
expect(result.run.status).toBe("COMPLETED_WITH_ERRORS");
607+
}
608+
}
609+
610+
const executionData = await engine.getRunExecutionData({ runId: run.id });
611+
assertNonNullable(executionData);
612+
expect(executionData.run.attemptNumber).toBe(taskMaxAttempts);
613+
expect(executionData.run.status).toBe("COMPLETED_WITH_ERRORS");
614+
expect(await engine.runQueue.lengthOfDeadLetterQueue(authenticatedEnvironment)).toBe(0);
615+
} finally {
616+
await engine.quit();
617+
}
618+
}
619+
);
620+
494621
containerTest("OOM retry on larger machine", async ({ prisma, redisOptions }) => {
495622
//create environment
496623
const authenticatedEnvironment = await setupAuthenticatedEnvironment(prisma, "PRODUCTION");

internal-packages/run-engine/src/run-queue/index.ts

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1116,12 +1116,19 @@ export class RunQueue {
11161116
messageId,
11171117
retryAt,
11181118
incrementAttemptCount = true,
1119+
resetAttemptCount = false,
11191120
skipDequeueProcessing = false,
11201121
}: {
11211122
orgId: string;
11221123
messageId: string;
11231124
retryAt?: number;
11241125
incrementAttemptCount?: boolean;
1126+
/**
1127+
* Zero the message's attempt counter instead of incrementing it. The counter is the budget
1128+
* for dequeues that never reach execution; a caller that knows an attempt did execute passes
1129+
* this so an ordinary task retry cannot exhaust it and dead-letter the run.
1130+
*/
1131+
resetAttemptCount?: boolean;
11251132
skipDequeueProcessing?: boolean;
11261133
}) {
11271134
return this.#trace(
@@ -1148,7 +1155,9 @@ export class RunQueue {
11481155
[SemanticAttributes.WORKER_QUEUE]: this.#getWorkerQueueFromMessage(message),
11491156
});
11501157

1151-
if (incrementAttemptCount) {
1158+
if (resetAttemptCount) {
1159+
message.attempt = 0;
1160+
} else if (incrementAttemptCount) {
11521161
message.attempt = message.attempt + 1;
11531162
if (message.attempt >= maxAttempts) {
11541163
await this.#callMoveToDeadLetterQueue({ message });

internal-packages/run-engine/src/run-queue/tests/nack.test.ts

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -130,6 +130,81 @@ describe("RunQueue.nackMessage", () => {
130130
}
131131
});
132132

133+
redisTest(
134+
"nacking with resetAttemptCount zeroes the counter instead of dead-lettering",
135+
async ({ redisContainer }) => {
136+
const queue = new RunQueue({
137+
...testOptions,
138+
retryOptions: {
139+
...testOptions.retryOptions,
140+
maxAttempts: 2,
141+
},
142+
queueSelectionStrategy: new FairQueueSelectionStrategy({
143+
redis: {
144+
keyPrefix: "runqueue:test:",
145+
host: redisContainer.getHost(),
146+
port: redisContainer.getPort(),
147+
},
148+
keys: testOptions.keys,
149+
}),
150+
redis: {
151+
keyPrefix: "runqueue:test:",
152+
host: redisContainer.getHost(),
153+
port: redisContainer.getPort(),
154+
},
155+
});
156+
157+
try {
158+
await queue.enqueueMessage({
159+
env: authenticatedEnvDev,
160+
message: messageDev,
161+
workerQueue: authenticatedEnvDev.id,
162+
});
163+
164+
await setTimeout(1000);
165+
166+
const dequeued = await queue.dequeueMessageFromWorkerQueue(
167+
"test_12345",
168+
authenticatedEnvDev.id
169+
);
170+
assertNonNullable(dequeued);
171+
172+
const first = await queue.nackMessage({
173+
orgId: messageDev.orgId,
174+
messageId: messageDev.runId,
175+
});
176+
expect(first).toBe(true);
177+
178+
const afterFirst = await queue.readMessage(messageDev.orgId, messageDev.runId);
179+
expect(afterFirst?.attempt).toBe(1);
180+
181+
await setTimeout(1000);
182+
183+
const dequeued2 = await queue.dequeueMessageFromWorkerQueue(
184+
"test_12345",
185+
authenticatedEnvDev.id
186+
);
187+
assertNonNullable(dequeued2);
188+
189+
// A plain nack here would hit maxAttempts and dead-letter the run
190+
const second = await queue.nackMessage({
191+
orgId: messageDev.orgId,
192+
messageId: messageDev.runId,
193+
resetAttemptCount: true,
194+
});
195+
expect(second).toBe(true);
196+
197+
const afterReset = await queue.readMessage(messageDev.orgId, messageDev.runId);
198+
expect(afterReset?.attempt).toBe(0);
199+
200+
expect(await queue.lengthOfEnvQueue(authenticatedEnvDev)).toBe(1);
201+
expect(await queue.lengthOfDeadLetterQueue(authenticatedEnvDev)).toBe(0);
202+
} finally {
203+
await queue.quit();
204+
}
205+
}
206+
);
207+
133208
redisTest(
134209
"nacking a message with maxAttempts reached should be moved to dead letter queue",
135210
async ({ redisContainer }) => {

0 commit comments

Comments
 (0)