From a166367387d8892be37045cda4e95c6c3af1fd04 Mon Sep 17 00:00:00 2001 From: Nathan Rajlich Date: Fri, 25 Sep 2026 10:14:57 -0700 Subject: [PATCH] fix(core): don't strand a run when the max-deliveries terminal write is throttled --- ...max-deliveries-transient-terminal-write.md | 5 + packages/core/src/runtime.test.ts | 98 ++++++++++++++++++- packages/core/src/runtime.ts | 15 +++ .../core/src/runtime/step-handler.test.ts | 78 ++++++++++++++- packages/core/src/runtime/step-handler.ts | 84 ++++++++++------ 5 files changed, 249 insertions(+), 31 deletions(-) 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..d1b8917dac --- /dev/null +++ b/.changeset/max-deliveries-transient-terminal-write.md @@ -0,0 +1,5 @@ +--- +'@workflow/core': patch +--- + +Retry the max-deliveries `run_failed`/`step_failed` write and the step handler's workflow re-queue through queue redelivery when they fail 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..ab4b2fb894 100644 --- a/packages/core/src/runtime.test.ts +++ b/packages/core/src/runtime.test.ts @@ -1,4 +1,8 @@ -import { RUN_ERROR_CODES, WorkflowWorldError } from '@workflow/errors'; +import { + RUN_ERROR_CODES, + ThrottleError, + WorkflowWorldError, +} from '@workflow/errors'; import { type Event, SPEC_VERSION_CURRENT, @@ -1075,3 +1079,95 @@ 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);`; + + async function deliverPastCeiling(runFailed: () => Promise) { + const runId = 'wrun_max_deliveries'; + const created: string[] = []; + const runsGet = vi.fn(); + const eventsList = vi.fn(); + setWorld({ + specVersion: SPEC_VERSION_CURRENT, + createQueueHandler: vi.fn( + (_p: string, handler: (m: unknown, md: unknown) => Promise) => + async () => { + await handler( + { runId, requestedAt: new Date() }, + { + requestId: 'req', + // One past MAX_QUEUE_DELIVERIES (48). + attempt: 49, + queueName: '__wkf_workflow_workflow', + messageId: 'msg', + } + ); + return new Response(null, { status: 204 }); + } + ), + events: { + create: vi.fn(async (_runId: string, data: any) => { + created.push(data.eventType); + if (data.eventType === 'run_failed') return runFailed(); + return { event: data }; + }), + list: eventsList, + }, + runs: { get: runsGet }, + queue: vi.fn(async () => ({ messageId: null })), + getEncryptionKeyForRun: vi.fn(async () => undefined), + } as any); + + const result = workflowEntrypoint(workflowCode)( + new Request('https://example.test') + ); + return { result, created, runsGet, eventsList }; + } + + it('fails the run and acks once the terminal write lands', async () => { + const { result, created, runsGet } = await deliverPastCeiling(async () => ({ + event: {}, + })); + await expect(result).resolves.toBeInstanceOf(Response); + expect(created).toEqual(['run_failed']); + expect(runsGet).not.toHaveBeenCalled(); + }); + + it.each([ + [ + 'ThrottleError (429)', + new ThrottleError('rate limited', { retryAfter: 5 }), + ], + [ + 'WorkflowWorldError (5xx)', + new WorkflowWorldError('server error', { status: 503 }), + ], + ])('does not ack (strand the run) when the terminal write hits a transient %s', async (_label, transientError) => { + 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('still acks on a non-transient terminal write failure', async () => { + const { result } = await deliverPastCeiling(async () => { + throw new WorkflowWorldError('bad request', { status: 400 }); + }); + await expect(result).resolves.toBeInstanceOf(Response); + }); +}); diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index e51e0f0b2f..daf5278123 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -229,6 +229,21 @@ 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; the redelivery is + // still past the ceiling, so it only retries this terminal write. + if (isRetryableWorldError(err)) { + runtimeLogger.warn( + 'Transient error marking run as failed after max deliveries, retrying via queue redelivery', + { + workflowRunId: runId, + attempt: metadata.attempt, + error: 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. ` + diff --git a/packages/core/src/runtime/step-handler.test.ts b/packages/core/src/runtime/step-handler.test.ts index 2b0830d60a..be1aa8fea6 100644 --- a/packages/core/src/runtime/step-handler.test.ts +++ b/packages/core/src/runtime/step-handler.test.ts @@ -1,4 +1,9 @@ -import { EntityConflictError, WorkflowWorldError } from '@workflow/errors'; +import { + EntityConflictError, + RunExpiredError, + ThrottleError, + WorkflowWorldError, +} from '@workflow/errors'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; // Use vi.hoisted so these are available in mock factories @@ -574,7 +579,9 @@ describe('step-handler max deliveries', () => { ); }); - it('should consume message silently when step_failed fails with EntityConflictError', async () => { + it('should re-queue the workflow when step_failed fails with EntityConflictError', async () => { + // The step may have been failed by this message's own earlier delivery, + // whose workflow re-queue then failed: wake the workflow regardless. mockEventsCreate.mockRejectedValue( new EntityConflictError('Step already completed') ); @@ -586,6 +593,73 @@ describe('step-handler max deliveries', () => { expect(result).toBeUndefined(); expect(mockStepFn).not.toHaveBeenCalled(); + expect(mockQueueMessage).toHaveBeenCalledTimes(1); + }); + + it('should consume message silently when the run has already finished', async () => { + mockEventsCreate.mockRejectedValue(new RunExpiredError('Run completed')); + + const result = await capturedHandler(createMessage(), { + ...createMetadata('myStep'), + attempt: MAX_QUEUE_DELIVERIES + 1, + }); + + expect(result).toBeUndefined(); + expect(mockQueueMessage).not.toHaveBeenCalled(); + }); + + it.each([ + [ + 'ThrottleError (429)', + new ThrottleError('rate limited', { retryAfter: 5 }), + ], + [ + 'WorkflowWorldError (5xx)', + new WorkflowWorldError('server error', { status: 503 }), + ], + ])('should throw (not strand the run) when step_failed hits a transient %s', async (_label, transientError) => { + // Regression: acking on any step_failed failure left the step `running` + // and the run with no message left to drive it. + mockEventsCreate.mockRejectedValue(transientError); + + await expect( + capturedHandler(createMessage(), { + ...createMetadata('myStep'), + attempt: MAX_QUEUE_DELIVERIES + 1, + }) + ).rejects.toBe(transientError); + expect(mockStepFn).not.toHaveBeenCalled(); + expect(mockQueueMessage).not.toHaveBeenCalled(); + }); + + it('should throw when re-queueing the workflow fails', async () => { + const sendError = new Error('queue send failed'); + mockQueueMessage.mockRejectedValue(sendError); + + await expect( + capturedHandler(createMessage(), { + ...createMetadata('myStep'), + attempt: MAX_QUEUE_DELIVERIES + 1, + }) + ).rejects.toBe(sendError); + }); + + it('should consume the message on a definitive step_failed rejection', async () => { + mockEventsCreate.mockRejectedValue( + new WorkflowWorldError('bad request', { status: 400 }) + ); + + const result = await capturedHandler(createMessage(), { + ...createMetadata('myStep'), + attempt: MAX_QUEUE_DELIVERIES + 1, + }); + + expect(result).toBeUndefined(); + expect(mockQueueMessage).not.toHaveBeenCalled(); + expect(mockRuntimeLogger.error).toHaveBeenCalledWith( + expect.stringContaining('Failed to mark step as failed'), + expect.anything() + ); }); it('should not trigger max deliveries check when under limit', async () => { diff --git a/packages/core/src/runtime/step-handler.ts b/packages/core/src/runtime/step-handler.ts index db13fcfa01..b17f493e74 100644 --- a/packages/core/src/runtime/step-handler.ts +++ b/packages/core/src/runtime/step-handler.ts @@ -17,6 +17,7 @@ import { type Step, StepInvokePayloadSchema, } from '@workflow/world'; +import { isRetryableWorldError } from '../classify-error.js'; import { importKey } from '../encryption.js'; import { runtimeLogger, stepLogger } from '../logger.js'; import { getStepFunction } from '../private.js'; @@ -90,8 +91,9 @@ const stepHandler = createQueueHandler( // This prevents runaway steps from consuming infinite queue deliveries. // At this point, we want to do the minimal amount of work (no fetching // of the step details, etc. We simply attempt to mark the step as failed - // and enqueue the workflow once, and if either of those fails, the message - // is still consumed but with adequate logging that an error occurred. + // and enqueue the workflow once. A transient failure of either throws so + // the queue retries it; a definitive step_failed rejection consumes the + // message with adequate logging that an error occurred. if (metadata.attempt > MAX_QUEUE_DELIVERIES) { runtimeLogger.error( `Step handler exceeded max deliveries (${metadata.attempt}/${MAX_QUEUE_DELIVERIES})`, @@ -102,8 +104,8 @@ const stepHandler = createQueueHandler( attempt: metadata.attempt, } ); + const world = getWorld(); try { - const world = getWorld(); await world.events.create( workflowRunId, { @@ -117,36 +119,62 @@ const stepHandler = createQueueHandler( }, { requestId } ); - // Re-queue the workflow to handle the failed step - await queueMessage( - world, - getWorkflowQueueName(workflowName, stepNamespace), - { - runId: workflowRunId, - traceCarrier: await serializeTraceCarrier(), - requestedAt: new Date(), - } - ); } catch (err) { - if (EntityConflictError.is(err) || RunExpiredError.is(err)) { + if (RunExpiredError.is(err)) { return; } - // Can't even mark the step as failed. Consume the message to stop - // further retries. The run will remain in its current state. - runtimeLogger.error( - `Failed to mark step as failed after ${metadata.attempt} delivery attempts. ` + - `A persistent error is preventing the step from being terminated. ` + - `The run will remain in its current state until manually resolved. ` + - `This is most likely due to a persistent outage of the workflow backend ` + - `or a bug in the workflow runtime and should be reported to the Workflow team.`, - { - workflowRunId, - stepId, - attempt: metadata.attempt, - error: err instanceof Error ? err.message : String(err), + // EntityConflictError: the step is already terminal. That may be this + // message's own earlier delivery, which wrote step_failed and then + // failed to re-queue the workflow below, so fall through and re-queue + // it. An extra wake only costs one replay. + if (!EntityConflictError.is(err)) { + // 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. The redelivery + // is still past the ceiling, so it only retries this write and + // never re-runs the step body. + if (isRetryableWorldError(err)) { + runtimeLogger.warn( + 'Transient error marking step as failed after max deliveries, retrying via queue redelivery', + { + workflowRunId, + stepId, + attempt: metadata.attempt, + error: err instanceof Error ? err.message : String(err), + } + ); + throw err; } - ); + // Can't even mark the step as failed. Consume the message to stop + // further retries. The run will remain in its current state. + runtimeLogger.error( + `Failed to mark step as failed after ${metadata.attempt} delivery attempts. ` + + `A persistent error is preventing the step from being terminated. ` + + `The run will remain in its current state until manually resolved. ` + + `This is most likely due to a persistent outage of the workflow backend ` + + `or a bug in the workflow runtime and should be reported to the Workflow team.`, + { + workflowRunId, + stepId, + attempt: metadata.attempt, + error: err instanceof Error ? err.message : String(err), + } + ); + return; + } } + // Re-queue the workflow to handle the failed step. A failure here + // throws so the queue redelivers: acking would leave the step failed + // with no message left to wake the workflow. + await queueMessage( + world, + getWorkflowQueueName(workflowName, stepNamespace), + { + runId: workflowRunId, + traceCarrier: await serializeTraceCarrier(), + requestedAt: new Date(), + } + ); return; }