diff --git a/.changeset/terminal-run-errors-are-fatal.md b/.changeset/terminal-run-errors-are-fatal.md new file mode 100644 index 0000000000..d6cd3a7600 --- /dev/null +++ b/.changeset/terminal-run-errors-are-fatal.md @@ -0,0 +1,6 @@ +--- +'@workflow/errors': patch +'@workflow/core': patch +--- + +Mark `WorkflowRunFailedError` and `WorkflowRunCancelledError` as non-retryable, and make `FatalError.is()` honor the `fatal` marker, so a step that reads a terminal run's `returnValue` fails on its first attempt with the error intact instead of exhausting its retry budget first. diff --git a/docs/content/docs/api-reference/workflow-errors/workflow-run-cancelled-error.mdx b/docs/content/docs/api-reference/workflow-errors/workflow-run-cancelled-error.mdx index 3c3a6bc5ed..e78158d8e9 100644 --- a/docs/content/docs/api-reference/workflow-errors/workflow-run-cancelled-error.mdx +++ b/docs/content/docs/api-reference/workflow-errors/workflow-run-cancelled-error.mdx @@ -12,6 +12,8 @@ related: You can check for cancellation before awaiting by inspecting `run.status`. +A canceled run is terminal, so this error is non-retryable (`fatal: true`). Inside a workflow, `await run.returnValue` runs as a step, and that step fails on its first attempt instead of spending its retry budget re-reading a run that cannot change. Errors from *failing to read* the run, such as a transport blip, stay retryable. + ```typescript lineNumbers import { WorkflowRunCancelledError } from "workflow/errors" declare const run: { status: Promise; returnValue: Promise }; // @setup @@ -34,6 +36,11 @@ definition={` interface WorkflowRunCancelledError { /** The ID of the cancelled run. */ runId: string; + /** + * Always \`true\`. A canceled run is terminal, so a step that reads one is + * not retried. + */ + fatal: true; /** The error message. */ message: string; } diff --git a/docs/content/docs/api-reference/workflow-errors/workflow-run-failed-error.mdx b/docs/content/docs/api-reference/workflow-errors/workflow-run-failed-error.mdx index f4136b4a65..655489b9e7 100644 --- a/docs/content/docs/api-reference/workflow-errors/workflow-run-failed-error.mdx +++ b/docs/content/docs/api-reference/workflow-errors/workflow-run-failed-error.mdx @@ -13,6 +13,8 @@ related: The `cause` property contains the underlying error with its message, stack trace, and optional error code. +A failed run is terminal, so this error is non-retryable (`fatal: true`). Inside a workflow, `await run.returnValue` runs as a step, and that step fails on its first attempt instead of spending its retry budget re-reading a run that cannot change: the remote failure reaches the caller immediately, and the caller catches a `WorkflowRunFailedError` rather than a retry-exhaustion wrapper. Errors from *failing to read* the run, such as a transport blip, stay retryable. + ```typescript lineNumbers import { WorkflowRunFailedError } from "workflow/errors" declare const run: { status: Promise; returnValue: Promise }; // @setup @@ -40,6 +42,11 @@ interface WorkflowRunFailedError { runId: string; /** The underlying error that caused the failure. */ cause: Error & { code?: string }; + /** + * Always \`true\`. A failed run is terminal, so a step that reads one is not + * retried. + */ + fatal: true; /** The error message. */ message: string; } diff --git a/packages/core/src/runtime/run-return-value-terminal-retry.test.ts b/packages/core/src/runtime/run-return-value-terminal-retry.test.ts new file mode 100644 index 0000000000..cae27c2a57 --- /dev/null +++ b/packages/core/src/runtime/run-return-value-terminal-retry.test.ts @@ -0,0 +1,135 @@ +import { mkdtemp, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { + FatalError, + WorkflowRunCancelledError, + WorkflowRunFailedError, +} from '@workflow/errors'; +import type { World } from '@workflow/world'; +import { SPEC_VERSION_CURRENT } from '@workflow/world'; +import { createLocalWorld } from '@workflow/world-local'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +// Mock version module to avoid missing generated file +vi.mock('../version.js', () => ({ version: '0.0.0-test' })); + +import { + dehydrateStepArguments, + dehydrateWorkflowReturnValue, +} from '../serialization.js'; +import { getRun } from './run.js'; +import { setWorld } from './world.js'; + +/** + * `await run.returnValue` against a run that is already terminal. + * + * Inside a workflow the accessor is a step, so whatever it throws is classified + * by the step handler's retry policy — `FatalError.is(err)` is that gate (see + * `step-handler.ts`). A terminal run is immutable: once the accessor has + * *successfully read* one, re-running the body re-reads the same record and + * throws the same error. Retrying only delays the failure reaching the caller + * and, once the budget is spent, replaces the error the caller is documented to + * catch (`WorkflowRunFailedError`) with the retry-exhaustion wrapper. + * + * Errors from *failing to read* the run (transport blips, + * `WorkflowRunNotFoundError` for a resilient start) are a different case and + * stay retryable. + * + * See vercel/workflow#4288. + */ +describe('run.returnValue on a terminal run is not retried', () => { + let dir: string; + let world: World; + + beforeEach(async () => { + dir = await mkdtemp(join(tmpdir(), 'returnvalue-terminal-')); + world = createLocalWorld({ dataDir: dir }) as unknown as World; + setWorld(world); + }); + + afterEach(async () => { + setWorld(undefined); + await rm(dir, { recursive: true, force: true }); + }); + + /** A run in `running`, created through the event log like a real one. */ + async function startRun(workflowName: string): Promise { + const input = await dehydrateStepArguments([], 'run', undefined); + const created = await world.events.create(null, { + eventType: 'run_created', + specVersion: SPEC_VERSION_CURRENT, + eventData: { deploymentId: 'dpl_test', workflowName, input }, + }); + const runId = created.run?.runId; + if (!runId) throw new Error('expected the run to be created'); + await world.events.create(runId, { + eventType: 'run_started', + specVersion: SPEC_VERSION_CURRENT, + eventData: {}, + }); + return runId; + } + + it('throws a non-retryable WorkflowRunFailedError for a failed run', async () => { + const runId = await startRun('target'); + await world.events.create(runId, { + eventType: 'run_failed', + specVersion: SPEC_VERSION_CURRENT, + eventData: { + error: { message: 'target failed permanently' }, + errorCode: 'USER_ERROR', + }, + }); + + const error = await getRun(runId).returnValue.then( + () => { + throw new Error('expected returnValue to reject'); + }, + (err: unknown) => err + ); + + expect(WorkflowRunFailedError.is(error)).toBe(true); + // The accessor read the run successfully and the run cannot change, so the + // step fails here rather than burning its retry budget first — and the + // caller catches the error the docs point it at, not the retry-exhaustion + // wrapper that replaces it once the budget runs out. + expect(FatalError.is(error)).toBe(true); + expect((error as WorkflowRunFailedError).runId).toBe(runId); + expect((error as WorkflowRunFailedError).cause.message).toBe( + 'target failed permanently' + ); + }, 30_000); + + it('throws a non-retryable WorkflowRunCancelledError for a cancelled run', async () => { + const runId = await startRun('target'); + await world.events.create(runId, { + eventType: 'run_cancelled', + specVersion: SPEC_VERSION_CURRENT, + }); + + const error = await getRun(runId).returnValue.then( + () => { + throw new Error('expected returnValue to reject'); + }, + (err: unknown) => err + ); + + expect(WorkflowRunCancelledError.is(error)).toBe(true); + expect(FatalError.is(error)).toBe(true); + expect((error as WorkflowRunCancelledError).runId).toBe(runId); + }, 30_000); + + it('still resolves normally when the run succeeded', async () => { + const runId = await startRun('target'); + await world.events.create(runId, { + eventType: 'run_completed', + specVersion: SPEC_VERSION_CURRENT, + eventData: { + output: await dehydrateWorkflowReturnValue('done', runId, undefined), + }, + }); + + await expect(getRun(runId).returnValue).resolves.toBe('done'); + }, 30_000); +}); diff --git a/packages/errors/src/index.ts b/packages/errors/src/index.ts index c7ad587bd9..457161c5fb 100644 --- a/packages/errors/src/index.ts +++ b/packages/errors/src/index.ts @@ -151,6 +151,17 @@ export class WorkflowWorldError extends WorkflowError { * ``` */ export class WorkflowRunFailedError extends WorkflowError { + /** + * `failed` is terminal, and a run's terminal state is immutable. This error + * is only ever thrown after a *successful* read of such a run, so re-running + * the read returns the same record and throws the same error. Marking it + * non-retryable is what stops the step executor from spending a retry budget + * on that when the read happens inside a step — a parent awaiting a child's + * `returnValue` — and then replacing this error with its retry-exhaustion + * wrapper. A read that *fails* throws something else and stays retryable. + * See `FatalError.is()`. + */ + fatal = true; runId: string; declare cause: Error & { code?: string }; @@ -644,6 +655,8 @@ export class PreconditionFailedError extends WorkflowWorldError { * ``` */ export class WorkflowRunCancelledError extends WorkflowError { + /** Terminal and immutable, for the same reason as {@link WorkflowRunFailedError.fatal}. */ + fatal = true; runId: string; constructor(runId: string) { @@ -716,7 +729,12 @@ export class FatalError extends Error { } static is(value: unknown): value is FatalError { - return isError(value) && value.name === 'FatalError'; + if (!isError(value)) return false; + if (value.name === 'FatalError') return true; + // Other error classes opt out of retries by carrying the same marker, so + // the retry gate has to read the flag and not just the name. See + // `WorkflowRunFailedError.fatal`. + return (value as { fatal?: unknown }).fatal === true; } } diff --git a/packages/errors/src/terminal-run-error.test.ts b/packages/errors/src/terminal-run-error.test.ts new file mode 100644 index 0000000000..4fbeafa147 --- /dev/null +++ b/packages/errors/src/terminal-run-error.test.ts @@ -0,0 +1,50 @@ +import { describe, expect, it } from 'vitest'; +import { + FatalError, + WorkflowRunCancelledError, + WorkflowRunFailedError, + WorkflowRunNotCompletedError, + WorkflowRunNotFoundError, +} from './index.js'; + +/** + * `failed` and `cancelled` are terminal, and a run's terminal state is + * immutable. The errors that report them are only ever thrown after a run has + * been read successfully, so the read that produced them cannot come back + * different. `FatalError.is()` is the step executor's non-retry gate, and it + * has to say so — otherwise a parent awaiting a child's `returnValue` spends + * its whole retry budget re-reading the same record, and the error it finally + * surfaces is the executor's retry-exhaustion wrapper rather than the one + * callers are documented to catch. See vercel/workflow#4288. + */ +describe('terminal run errors are non-retryable', () => { + it('marks WorkflowRunFailedError fatal', () => { + const error = new WorkflowRunFailedError('wrun_1', { + message: 'boom', + code: 'USER_ERROR', + }); + expect(FatalError.is(error)).toBe(true); + // The marker is additive: identity and payload are unchanged. + expect(WorkflowRunFailedError.is(error)).toBe(true); + expect(error.runId).toBe('wrun_1'); + expect(error.cause.code).toBe('USER_ERROR'); + expect(error.cause.message).toBe('boom'); + }); + + it('marks WorkflowRunCancelledError fatal', () => { + const error = new WorkflowRunCancelledError('wrun_1'); + expect(FatalError.is(error)).toBe(true); + expect(WorkflowRunCancelledError.is(error)).toBe(true); + expect(error.runId).toBe('wrun_1'); + }); + + it('leaves the non-terminal run errors retryable', () => { + // Neither describes a settled run: `running` can still finish, and a run + // that is missing now can exist a moment later (a resilient start races + // its own `run_created`). Retrying either can produce a different answer. + expect( + FatalError.is(new WorkflowRunNotCompletedError('wrun_1', 'running')) + ).toBe(false); + expect(FatalError.is(new WorkflowRunNotFoundError('wrun_1'))).toBe(false); + }); +});