Skip to content

Commit fe2fd8c

Browse files
authored
Classify Workflow stream failures (#3850)
## Summary & Motivation Stream infrastructure failures (HTTP/2 session wedges, transport timeouts, non-2xx stream responses) surfaced as plain `Error`, so terminal classification attributed them to customer code as `USER_ERROR`. They now carry a catchable `StreamError` with a `STREAM_ERROR` run error code, attributed to the SDK and retried when transport-level or 5xx. The v4 events response body is wrapped so a post-header stream failure is classified and reported to the dispatcher recycler — a response header arriving is not yet a successful streamed request. ## Test Plan Unit tests added across classification, serialization round-trip, the streamer, and the v4 transport; 331 `@workflow/core` and 123 `@workflow/world-vercel` focused tests pass.
1 parent 63143e5 commit fe2fd8c

25 files changed

Lines changed: 763 additions & 119 deletions
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
---
2+
'@workflow/errors': minor
3+
'@workflow/core': patch
4+
'@workflow/cli': patch
5+
'@workflow/web-shared': patch
6+
'@workflow/world-vercel': patch
7+
'workflow': minor
8+
---
9+
10+
Add a catchable `StreamError` and classify Workflow stream infrastructure failures as `STREAM_ERROR` instead of `USER_ERROR`.

‎packages/cli/src/lib/inspect/hydration.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -157,6 +157,7 @@ const ERROR_REVIVER_KEYS = [
157157
'ReferenceError',
158158
'RetryableError',
159159
'RuntimeDecryptionError',
160+
'StreamError',
160161
'SyntaxError',
161162
'TypeError',
162163
'URIError',

‎packages/core/src/classify-error.test.ts‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import {
55
ReplayDivergenceError,
66
RUN_ERROR_CODES,
77
RuntimeDecryptionError,
8+
StreamError,
89
ThrottleError,
910
TooEarlyError,
1011
WorkflowDeploymentMismatchError,
@@ -153,6 +154,12 @@ describe('classifyRunError', () => {
153154
).toBe(RUN_ERROR_CODES.WORLD_CONTRACT_ERROR);
154155
});
155156

157+
it('classifies StreamError as STREAM_ERROR', () => {
158+
expect(classifyRunError(new StreamError('stream write failed'))).toBe(
159+
RUN_ERROR_CODES.STREAM_ERROR
160+
);
161+
});
162+
156163
it('classifies a TRANSPORT error as WORLD_CONTRACT_ERROR (backend fault, not USER_ERROR)', () => {
157164
// Transport blips are normally redelivered via the queue (see
158165
// isRetryableWorldError); if one ever reaches terminal classification it is
@@ -194,6 +201,22 @@ describe('isRetryableWorldError', () => {
194201
).toBe(true);
195202
});
196203

204+
it('retries transport and server-side StreamErrors, but not terminal 4xx responses', () => {
205+
expect(isRetryableWorldError(new StreamError('stream read failed'))).toBe(
206+
true
207+
);
208+
expect(
209+
isRetryableWorldError(
210+
new StreamError('stream service failed', { status: 503 })
211+
)
212+
).toBe(true);
213+
expect(
214+
isRetryableWorldError(
215+
new StreamError('stream request was invalid', { status: 400 })
216+
)
217+
).toBe(false);
218+
});
219+
197220
it('treats TRANSPORT and TIMEOUT codes as retryable', () => {
198221
expect(
199222
isRetryableWorldError(

‎packages/core/src/classify-error.ts‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import {
66
type RunErrorCode,
77
RuntimeDecryptionError,
88
StepNotRegisteredError,
9+
StreamError,
910
ThrottleError,
1011
WorkflowDeploymentMismatchError,
1112
WorkflowNotRegisteredError,
@@ -29,7 +30,11 @@ const WORLD_CONTRACT_ERROR_CODES = new Set([
2930
* from. Kept distinct from `WORLD_CONTRACT_ERROR_CODES` so a transport blip is
3031
* never misclassified as the server returning a malformed response.
3132
*/
32-
const RETRYABLE_WORLD_ERROR_CODES = new Set(['TRANSPORT', 'TIMEOUT']);
33+
const RETRYABLE_WORLD_ERROR_CODES = new Set([
34+
'TRANSPORT',
35+
'TIMEOUT',
36+
RUN_ERROR_CODES.STREAM_ERROR,
37+
]);
3338

3439
/**
3540
* Set of error names that should classify as generic `RUNTIME_ERROR`. Each
@@ -100,6 +105,9 @@ export function isRetryableWorldError(err: unknown): boolean {
100105
if (ThrottleError.is(err)) {
101106
return true;
102107
}
108+
if (StreamError.is(err)) {
109+
return err.status === undefined || err.status >= 500;
110+
}
103111
if (!WorkflowWorldError.is(err)) {
104112
return false;
105113
}
@@ -133,6 +141,10 @@ export function classifyRunError(err: unknown): RunErrorCode {
133141
// attribute an outage correctly. Note the retryable variants are normally
134142
// redelivered via the queue (see `isRetryableWorldError`) and only reach this
135143
// terminal classification if the run ultimately gives up.
144+
if (StreamError.is(err)) {
145+
return RUN_ERROR_CODES.STREAM_ERROR;
146+
}
147+
136148
if (isWorldContractError(err) || isRetryableWorldError(err)) {
137149
return RUN_ERROR_CODES.WORLD_CONTRACT_ERROR;
138150
}

‎packages/core/src/describe-error.ts‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,8 @@ const MAX_DELIVERIES_HINT =
8585
'The workflow queue exceeded its max-delivery budget. This usually indicates a persistent runtime failure — check the most recent stack traces for the underlying cause.';
8686
const MAX_EVENTS_HINT =
8787
'The workflow exceeded the maximum number of events per run. This usually means unbounded work in the workflow function — e.g. a loop that keeps creating steps without terminating. Break long-running workflows into child workflows to stay under the limit.';
88+
const STREAM_ERROR_HINT =
89+
'Workflow stream infrastructure failed while reading or writing data. This is an SDK or backend failure, not an error in your workflow code; please retry the run and report persistent failures with the runId.';
8890
const WORLD_CONTRACT_HINT =
8991
'The workflow backend returned data that violated the SDK contract. This is not retryable; please report it with the stack trace and runId.';
9092
const DEPLOYMENT_MISMATCH_HINT =
@@ -144,6 +146,13 @@ export function describeRunError(
144146
hint: CORRUPTED_EVENT_LOG_HINT,
145147
};
146148
}
149+
if (errorCode === RUN_ERROR_CODES.STREAM_ERROR) {
150+
return {
151+
attribution: 'sdk',
152+
errorCode,
153+
hint: STREAM_ERROR_HINT,
154+
};
155+
}
147156
if (errorCode === RUN_ERROR_CODES.WORLD_CONTRACT_ERROR) {
148157
return {
149158
attribution: 'sdk',
@@ -264,6 +273,14 @@ export function describeError(
264273
};
265274
}
266275

276+
if (effectiveCode === RUN_ERROR_CODES.STREAM_ERROR) {
277+
return {
278+
attribution: 'sdk',
279+
errorCode: effectiveCode,
280+
hint: STREAM_ERROR_HINT,
281+
};
282+
}
283+
267284
if (effectiveCode === RUN_ERROR_CODES.WORLD_CONTRACT_ERROR) {
268285
return {
269286
attribution: 'sdk',

‎packages/core/src/reconnecting-framed-stream.test.ts‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { StreamExpiredError } from '@workflow/errors';
1+
import { StreamError, StreamExpiredError } from '@workflow/errors';
22
import { SPEC_VERSION_CURRENT, type World } from '@workflow/world';
33
import { afterEach, describe, expect, it, vi } from 'vitest';
44

@@ -660,8 +660,11 @@ describe('createReconnectingFramedStream', () => {
660660
setWorld(world);
661661

662662
const stream = createReconnectingFramedStream(RUN_ID, 's', 0);
663-
await expect(readAll(stream)).rejects.toThrow(
664-
/exceeded maximum reconnection attempts/
663+
const error = await readAll(stream).catch((cause: unknown) => cause);
664+
expect(StreamError.is(error)).toBe(true);
665+
expect(error).toHaveProperty(
666+
'message',
667+
expect.stringMatching(/exceeded maximum reconnection attempts/)
665668
);
666669
// Initial connect + one connect per allowed reconnect; the following
667670
// reconnect throws before opening another stream.

‎packages/core/src/runtime/quickjs-serde.test.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import fs from 'node:fs';
2121
import { createRequire } from 'node:module';
22+
import { StreamError } from '@workflow/errors';
2223
import { QuickJS } from 'quickjs-wasi';
2324
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
2425
import {
@@ -162,6 +163,21 @@ describe('wire parity: guest serialize matches the reference codec', () => {
162163
return e;
163164
},
164165
],
166+
[
167+
'StreamError',
168+
'(() => { const cause = new Error("socket closed"); cause.stack = "cause-stack"; const e = new Error("stream failed", { cause }); e.name = "StreamError"; e.stack = "stream-stack"; e.status = 503; e.url = "https://workflow.example/stream"; return e; })()',
169+
() => {
170+
const cause = new Error('socket closed');
171+
cause.stack = 'cause-stack';
172+
const error = new StreamError('stream failed', {
173+
cause,
174+
status: 503,
175+
url: 'https://workflow.example/stream',
176+
});
177+
error.stack = 'stream-stack';
178+
return error;
179+
},
180+
],
165181
[
166182
'Error with cause',
167183
'(() => { const c = new Error("cause"); c.stack = "cs"; const e = new Error("outer", { cause: c }); e.stack = "os"; return e; })()',

‎packages/core/src/runtime/quickjs-serde.ts‎

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -256,6 +256,7 @@ const SYMBOL_NAMES = [
256256
'@workflow/errors//HookConflictError',
257257
'@workflow/errors//RetryableError',
258258
'@workflow/errors//RuntimeDecryptionError',
259+
'@workflow/errors//StreamError',
259260
] as const;
260261
type SymbolName = (typeof SYMBOL_NAMES)[number];
261262

@@ -1167,6 +1168,23 @@ export function createQuickJSSerde(
11671168
if (Object.hasOwn(shape, 'cause')) reduced.cause = shape.cause;
11681169
return reduced;
11691170
},
1171+
StreamError: (value) => {
1172+
if (!isHandle(value) || !value.isError) return false;
1173+
if (chainedString(value, 'name') !== 'StreamError') return false;
1174+
const shape = reduceErrorShape(value) as Record<string, unknown>;
1175+
const reduced: Record<string, unknown> = {
1176+
message: shape.message,
1177+
stack: shape.stack,
1178+
};
1179+
if (Object.hasOwn(shape, 'cause')) reduced.cause = shape.cause;
1180+
const status = own(value, 'status');
1181+
if (status && !status.isUndefined) reduced.status = status;
1182+
else status?.dispose();
1183+
const url = own(value, 'url');
1184+
if (url && !url.isUndefined) reduced.url = url;
1185+
else url?.dispose();
1186+
return reduced;
1187+
},
11701188
SyntaxError: namedErrorSubclassReducer('SyntaxError'),
11711189
TypeError: namedErrorSubclassReducer('TypeError'),
11721190
URIError: namedErrorSubclassReducer('URIError'),
@@ -1849,6 +1867,37 @@ export function createQuickJSSerde(
18491867
}
18501868
return error;
18511869
},
1870+
StreamError: (value: JSValueHandle) => {
1871+
const cls = registeredErrorClass('@workflow/errors//StreamError');
1872+
let error: JSValueHandle;
1873+
if (cls) {
1874+
const options = vm.newObject();
1875+
for (const property of ['cause', 'status', 'url'] as const) {
1876+
if (!guestHasOwn(value, property)) continue;
1877+
const propertyValue = own(value, property) ?? vm.undefined;
1878+
options.setProp(property, propertyValue);
1879+
if (propertyValue !== vm.undefined) propertyValue.dispose();
1880+
}
1881+
const message = own(value, 'message') ?? vm.undefined;
1882+
error = vm.construct(cls, message, options);
1883+
if (message !== vm.undefined) message.dispose();
1884+
options.dispose();
1885+
const stack = own(value, 'stack');
1886+
if (stack && !stack.isUndefined) define(error, 'stack', stack);
1887+
stack?.dispose();
1888+
cls.dispose();
1889+
} else {
1890+
error = buildError(i.Error, value, { name: 'StreamError' });
1891+
for (const property of ['status', 'url'] as const) {
1892+
const propertyValue = own(value, property);
1893+
if (propertyValue && !propertyValue.isUndefined) {
1894+
define(error, property, propertyValue);
1895+
}
1896+
propertyValue?.dispose();
1897+
}
1898+
}
1899+
return error;
1900+
},
18521901
SyntaxError: namedErrorSubclassReviver('SyntaxError'),
18531902
TypeError: namedErrorSubclassReviver('TypeError'),
18541903
URIError: namedErrorSubclassReviver('URIError'),

‎packages/core/src/runtime/runs.ts‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { EntityConflictError } from '@workflow/errors';
1+
import { EntityConflictError, StreamError } from '@workflow/errors';
22
import {
33
BULK_CANCEL_MAX_RUN_IDS,
44
type BulkCancelWorkflowRunResult,
@@ -409,7 +409,8 @@ export async function readStream(
409409
try {
410410
return await world.streams.get(runId, streamId, options?.startIndex);
411411
} catch (err) {
412-
throw new Error(
412+
if (StreamError.is(err)) throw err;
413+
throw new StreamError(
413414
`Failed to read stream ${streamId}: ${err instanceof Error ? err.message : String(err)}`,
414415
{ cause: err }
415416
);
@@ -426,7 +427,8 @@ export async function listStreams(
426427
try {
427428
return await world.streams.list(runId);
428429
} catch (err) {
429-
throw new Error(
430+
if (StreamError.is(err)) throw err;
431+
throw new StreamError(
430432
`Failed to list streams for run ${runId}: ${err instanceof Error ? err.message : String(err)}`,
431433
{ cause: err }
432434
);

‎packages/core/src/serialization.test.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,9 @@ import {
44
FatalError,
55
HookConflictError,
66
RetryableError,
7+
RUN_ERROR_CODES,
78
RuntimeDecryptionError,
9+
StreamError,
810
} from '@workflow/errors';
911
import { WORKFLOW_DESERIALIZE, WORKFLOW_SERIALIZE } from '@workflow/serde';
1012
import { beforeAll, describe, expect, it, vi } from 'vitest';
@@ -4714,6 +4716,34 @@ describe('dehydrate/hydrateRunError', () => {
47144716
expect(cause?.name).toBe('OperationError');
47154717
});
47164718

4719+
it('should round-trip StreamError as a catchable error type', async () => {
4720+
const original = new StreamError('stream write failed', {
4721+
cause: new Error('HTTP 500'),
4722+
status: 503,
4723+
url: 'https://workflow.example/stream',
4724+
});
4725+
const serialized = await dehydrateRunError(
4726+
original,
4727+
mockRunId,
4728+
noEncryptionKey
4729+
);
4730+
const hydrated = (await hydrateRunError(
4731+
serialized,
4732+
mockRunId,
4733+
noEncryptionKey
4734+
)) as StreamError;
4735+
4736+
expect(StreamError.is(hydrated)).toBe(true);
4737+
expect(hydrated).toBeInstanceOf(StreamError);
4738+
expect(hydrated).toMatchObject({
4739+
message: 'stream write failed',
4740+
status: 503,
4741+
url: 'https://workflow.example/stream',
4742+
code: RUN_ERROR_CODES.STREAM_ERROR,
4743+
});
4744+
expect((hydrated.cause as Error).message).toBe('HTTP 500');
4745+
});
4746+
47174747
it('should produce DEVALUE_V1-prefixed binary output', async () => {
47184748
const serialized = await dehydrateRunError(
47194749
new Error('x'),

0 commit comments

Comments
 (0)