From f3cea21d26ae300b747e66344ddcd2fe102101bf Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 25 Sep 2026 21:11:21 +0000 Subject: [PATCH] fix(core): don't strand runs on throttled writes past the delivery ceiling or in QuickJS inline claims (#4421) * fix(core): don't strand a run on throttled writes past the delivery ceiling or in QuickJS inline claims * fix(core): gate throttled-step stepInput on the CBOR queue transport * fix(core): defer a replay for a throttled QuickJS inline step, like the node engine Signed-off-by: Nathan Rajlich --- ...max-deliveries-transient-terminal-write.md | 5 + packages/core/src/runtime.test.ts | 128 +++++++++++++++++- packages/core/src/runtime.ts | 20 +++ 3 files changed, 152 insertions(+), 1 deletion(-) create mode 100644 .changeset/max-deliveries-transient-terminal-write.md diff --git a/.changeset/max-deliveries-transient-terminal-write.md b/.changeset/max-deliveries-transient-terminal-write.md new file mode 100644 index 0000000000..a9207dbd96 --- /dev/null +++ b/.changeset/max-deliveries-transient-terminal-write.md @@ -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`. diff --git a/packages/core/src/runtime.test.ts b/packages/core/src/runtime.test.ts index 103fd2fd22..317704dfaf 100644 --- a/packages/core/src/runtime.test.ts +++ b/packages/core/src/runtime.test.ts @@ -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, @@ -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) { + 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) => + 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']); + }); +}); diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index e51e0f0b2f..37d9f7ae6d 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -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. ` +