Skip to content
Closed
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
5 changes: 5 additions & 0 deletions .changeset/max-deliveries-transient-terminal-write.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/core': patch
---

Retry the max-deliveries `run_failed` write through queue redelivery when it fails transiently (429, 5xx, transport) instead of acking and leaving the run stuck `running`.
128 changes: 127 additions & 1 deletion packages/core/src/runtime.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,9 @@
import { RUN_ERROR_CODES, WorkflowWorldError } from '@workflow/errors';
import {
EntityConflictError,
RUN_ERROR_CODES,
ThrottleError,
WorkflowWorldError,
} from '@workflow/errors';
import {
type Event,
SPEC_VERSION_CURRENT,
Expand Down Expand Up @@ -1075,3 +1080,124 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => {
expect(await anyWaitUntilPromiseRejected()).toBe(false);
});
});

describe('workflowEntrypoint max-deliveries terminal write', () => {
afterEach(() => {
setWorld(undefined);
vi.clearAllMocks();
});

const workflowCode = `async function workflow() {
throw new Error('workflow code must not execute');
}
;globalThis.__private_workflows = new Map();
globalThis.__private_workflows.set("workflow", workflow);`;

/**
* Deliver one message past the delivery ceiling, with `run_failed` writes
* answered by `runFailed` (resolve for success, throw for failure).
*/
async function deliverPastCeiling(runFailed: () => Promise<unknown>) {
const runId = 'wrun_max_deliveries';
const created: string[] = [];
const eventsCreate = vi.fn(async (_runId: string, data: any) => {
created.push(data.eventType);
if (data.eventType === 'run_failed') return runFailed();
return { event: { ...data, eventId: 'evnt_1', runId } };
});
const queue = vi.fn(async () => ({ messageId: null }));
const runsGet = vi.fn();
const eventsList = vi.fn();
setWorld({
specVersion: SPEC_VERSION_CURRENT,
getDeploymentId: vi.fn(async () => 'dpl_test'),
createQueueHandler: vi.fn(
(_p: string, handler: (m: unknown, md: unknown) => Promise<unknown>) =>
async () => {
await handler(
{ runId, requestedAt: new Date() },
{
requestId: 'req',
// One past the default MAX_QUEUE_DELIVERIES (48).
attempt: 49,
queueName: '__wkf_workflow_workflow',
messageId: 'msg',
}
);
return new Response(null, { status: 204 });
}
),
events: { create: eventsCreate, list: eventsList },
runs: { get: runsGet },
queue,
getEncryptionKeyForRun: vi.fn(async () => undefined),
} as any);

const result = workflowEntrypoint(workflowCode)(
new Request('https://example.test')
);
return { result, created, queue, runsGet, eventsList };
}

it('fails the run and acks once the terminal write lands', async () => {
const { result, created, queue, runsGet, eventsList } =
await deliverPastCeiling(async () => ({ event: {} }));

await expect(result).resolves.toBeInstanceOf(Response);
expect(created).toEqual(['run_failed']);
// Nothing past the ceiling touches the replay path.
expect(runsGet).not.toHaveBeenCalled();
expect(eventsList).not.toHaveBeenCalled();
expect(queue).not.toHaveBeenCalled();
});

it.each([
[
'ThrottleError (429)',
new ThrottleError('rate limited', { retryAfter: 5 }),
],
[
'WorkflowWorldError (5xx)',
new WorkflowWorldError('server error', { status: 503 }),
],
[
'WorkflowWorldError (TRANSPORT)',
new WorkflowWorldError('socket hang up', { code: 'TRANSPORT' }),
],
])('does not ack (strand the run) when the terminal write hits a transient %s', async (_label, transientError) => {
// Regression: the give-up path used to log "the run will remain in its
// current state" and ack on ANY failure of this write. A 429 or 5xx here
// left the run `running` with no queue message left to drive it. The
// handler must throw so the queue redelivers; the redelivery is still past
// the ceiling, so it retries only this write.
const { result, created, runsGet, eventsList } = await deliverPastCeiling(
async () => {
throw transientError;
}
);

await expect(result).rejects.toBe(transientError);
expect(created).toEqual(['run_failed']);
expect(runsGet).not.toHaveBeenCalled();
expect(eventsList).not.toHaveBeenCalled();
});

it('acks when the run already reached a terminal state', async () => {
const { result } = await deliverPastCeiling(async () => {
throw new EntityConflictError('run already failed');
});

await expect(result).resolves.toBeInstanceOf(Response);
});

it('still acks on a non-transient terminal write failure', async () => {
// A definitive rejection gets the same verdict on every redelivery, so
// retrying it would only hot-loop the handler. Keep the logged give-up.
const { result, created } = await deliverPastCeiling(async () => {
throw new WorkflowWorldError('bad request', { status: 400 });
});

await expect(result).resolves.toBeInstanceOf(Response);
expect(created).toEqual(['run_failed']);
});
});
20 changes: 20 additions & 0 deletions packages/core/src/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -229,6 +229,26 @@ export function workflowEntrypoint(
// Run already finished, consume the message silently
return;
}
// A transient backend failure (429 / 5xx / transport) must not
// abandon the run: acking here leaves it `running` with no message
// left to drive it. Throw so the queue redelivers with its backoff
// (honoring a 429's Retry-After); the redelivery is still past the
// ceiling, so it only retries this terminal write, never the replay.
// This relies on the World redelivering past the ceiling: VQS and
// world-local do, while world-postgres currently caps its jobs at
// exactly this delivery (#4427, fixed by #4428).
if (isRetryableWorldError(err)) {
runtimeLogger.warn(
'Transient error marking run as failed after max deliveries, retrying via queue redelivery',
{
workflowRunId: runId,
attempt: metadata.attempt,
errorName: err instanceof Error ? err.name : 'UnknownError',
errorMessage: err instanceof Error ? err.message : String(err),
}
);
throw err;
}
runtimeLogger.error(
`Failed to mark run as failed after ${metadata.attempt} delivery attempts. ` +
`A persistent error is preventing the run from being terminated. ` +
Expand Down