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
6 changes: 6 additions & 0 deletions .changeset/terminal-run-errors-are-fatal.md
Original file line number Diff line number Diff line change
@@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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<string>; returnValue: Promise<any> }; // @setup
Expand All @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<string>; returnValue: Promise<any> }; // @setup
Expand Down Expand Up @@ -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;
}
Expand Down
135 changes: 135 additions & 0 deletions packages/core/src/runtime/run-return-value-terminal-retry.test.ts
Original file line number Diff line number Diff line change
@@ -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<string> {
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);
});
20 changes: 19 additions & 1 deletion packages/errors/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 };

Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
}
}

Expand Down
50 changes: 50 additions & 0 deletions packages/errors/src/terminal-run-error.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
Loading