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
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`/`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`.
98 changes: 97 additions & 1 deletion packages/core/src/runtime.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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<unknown>) {
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<unknown>) =>
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);
});
});
15 changes: 15 additions & 0 deletions packages/core/src/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same caveat as #4421, for anyone reading this later: this relies on the queue redelivering after delivery 49. VQS (24 h TTL, no delivery cap) and world-local (256) do. On this branch, world-postgres jobs cap at 3 attempts, so they never reach the ceiling.

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. ` +
Expand Down
78 changes: 76 additions & 2 deletions packages/core/src/runtime/step-handler.test.ts
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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')
);
Expand All @@ -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 () => {
Expand Down
84 changes: 56 additions & 28 deletions packages/core/src/runtime/step-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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})`,
Expand All @@ -102,8 +104,8 @@ const stepHandler = createQueueHandler(
attempt: metadata.attempt,
}
);
const world = getWorld();
try {
const world = getWorld();
await world.events.create(
workflowRunId,
{
Expand All @@ -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(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is the one spot that throws on any error rather than checking isRetryableWorldError. A definitive publish failure would redeliver every ~15 min until the message TTL. It's bounded, and publish failures are almost always transient, so this is probably the right trade: acking here strands the run for certain. I'm pointing it out because it differs from the step_failed branch above.

world,
getWorkflowQueueName(workflowName, stepNamespace),
{
runId: workflowRunId,
traceCarrier: await serializeTraceCarrier(),
requestedAt: new Date(),
}
);
return;
}

Expand Down
Loading