Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .server-changes/retry-requeue-nack-budget.md
Original file line number Diff line number Diff line change
@@ -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.
Comment thread
matt-aitken marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -1126,6 +1126,7 @@ export class RunAttemptSystem {
orgId: env.organizationId,
projectId: env.project.id,
timestamp: retryAt.getTime(),
resetQueueAttempts: !forceRequeue,
Comment thread
matt-aitken marked this conversation as resolved.
error: {
type: "INTERNAL_ERROR",
code: "TASK_RUN_DEQUEUED_MAX_RETRIES",
Expand Down Expand Up @@ -1249,6 +1250,7 @@ export class RunAttemptSystem {
checkpointId,
completedWaitpoints,
batchId,
resetQueueAttempts = false,
tx,
}: {
run: { id: string };
Expand All @@ -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;

Expand All @@ -1278,6 +1286,7 @@ export class RunAttemptSystem {
orgId,
messageId: run.id,
retryAt: timestamp,
resetAttemptCount: resetQueueAttempts,
});

if (!gotRequeued) {
Expand Down
128 changes: 128 additions & 0 deletions internal-packages/run-engine/src/engine/tests/attemptFailures.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
},
Comment thread
coderabbitai[bot] marked this conversation as resolved.
},
});

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");
Expand Down
11 changes: 10 additions & 1 deletion internal-packages/run-engine/src/run-queue/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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 });
Expand Down
75 changes: 75 additions & 0 deletions internal-packages/run-engine/src/run-queue/tests/nack.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }) => {
Expand Down