diff --git a/.changeset/transport-failure-not-user-error.md b/.changeset/transport-failure-not-user-error.md new file mode 100644 index 0000000000..5e4dc517d7 --- /dev/null +++ b/.changeset/transport-failure-not-user-error.md @@ -0,0 +1,6 @@ +--- +'@workflow/core': patch +'@workflow/world-vercel': patch +--- + +Route unrecognized backend connection and stream failures through existing retry policies, rebuilding shared event connections after repeated HTTP/2 failures. Keep invalid backend URLs, blocked ports, and unsupported request headers out of those retries. Include error cause chains in run-failure logs to expose underlying socket, DNS, and TLS failures. diff --git a/docs/content/docs/foundations/errors-and-retries.mdx b/docs/content/docs/foundations/errors-and-retries.mdx index fc6c309765..2cd766dc3d 100644 --- a/docs/content/docs/foundations/errors-and-retries.mdx +++ b/docs/content/docs/foundations/errors-and-retries.mdx @@ -164,6 +164,16 @@ export async function myWorkflow(input: unknown) { Uncaught, the run fails immediately with the `USER_ERROR` code — without retrying. See [serialization-failed](/docs/errors/serialization-failed) for common causes and fixes. +## Backend Connection Failures + +On Vercel, backend connection failures and interrupted event streams use the existing retry policies, even when their error codes are unrecognized. This lets workflows recover from network failures instead of immediately failing with `USER_ERROR`. Persistent failures can still exhaust the retry budget. + +The SDK also replaces its shared events connection pool after repeated HTTP/2 session failures. Invalid backend URLs, including unsupported protocols and embedded credentials, fail immediately. Fetch requests to blocked ports or with unsupported headers (such as `Expect`) also fail without retrying. + +A connection failure does not prove that the backend rejected a write: it may have accepted it before the response was lost. Continue to make step side effects [idempotent](/docs/foundations/idempotency). + +Failed-run logs include the underlying error causes and their codes, exposing socket, DNS, or TLS errors behind messages such as `TypeError: fetch failed`. If a cause cannot be read, the log includes `[unavailable cause]` and the run can still be recorded as failed. + ## Error Codes When a workflow run fails, the error may include a `code` that classifies the failure. You can access it programmatically via the `Run` class: diff --git a/packages/core/README.md b/packages/core/README.md index fc99bf5f99..2e47b71726 100644 --- a/packages/core/README.md +++ b/packages/core/README.md @@ -1,3 +1,5 @@ # @workflow/core Core runtime package for [Workflow SDK](https://useworkflow.dev). + +Failed-run logs include underlying error causes and codes to help diagnose failures such as socket, DNS, and TLS errors. Unreadable causes are marked without preventing the run from being recorded as failed. diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 251cdc9a47..e51e0f0b2f 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -56,7 +56,12 @@ import { withTraceContext, withWorkflowBaggage, } from './telemetry.js'; -import { getErrorName, getErrorStack, normalizeUnknownError } from './types.js'; +import { + formatErrorCauseChain, + getErrorName, + getErrorStack, + normalizeUnknownError, +} from './types.js'; import { buildWorkflowSuspensionMessage } from './util.js'; import { runWorkflow } from './workflow.js'; @@ -919,6 +924,15 @@ export function workflowEntrypoint( errorCode, errorName, errorStack, + // Neither the message nor the stack reaches a wrapped + // error's reason: `TypeError: fetch failed` carries an + // empty message by design and a stack of pure + // `node:internal/` frames, and the world layer's own + // wrappers name the request that failed rather than what + // failed about it. Undefined when there is no cause, so + // the field disappears for an ordinary user throw. + errorCause: + formatErrorCauseChain(terminalError) || undefined, }); // Fail the workflow run via event (event-sourced architecture) diff --git a/packages/core/src/types.test.ts b/packages/core/src/types.test.ts new file mode 100644 index 0000000000..0213bf61ba --- /dev/null +++ b/packages/core/src/types.test.ts @@ -0,0 +1,135 @@ +import { describe, expect, it } from 'vitest'; +import { formatErrorCauseChain } from './types.js'; + +describe('formatErrorCauseChain', () => { + it.each([ + 'cause', + 'name', + 'message', + 'code', + 'errors', + ])('tolerates a throwing %s getter in the cause chain', (property) => { + const cause = new Error('inner'); + Object.defineProperty(cause, property, { + get() { + throw new Error('getter failed'); + }, + }); + + expect(formatErrorCauseChain(new Error('outer', { cause }))).toContain( + '[unavailable cause]' + ); + }); + + it('preserves readable links before an inaccessible cause', () => { + const { proxy, revoke } = Proxy.revocable({}, {}); + revoke(); + const middle = new Error('middle', { cause: proxy }); + + expect(formatErrorCauseChain(new Error('outer', { cause: middle }))).toBe( + 'Error: middle\n[unavailable cause]' + ); + }); + + it('tolerates a throwing getter on the outer cause', () => { + const error = new Error('outer'); + Object.defineProperty(error, 'cause', { + get() { + throw new Error('getter failed'); + }, + }); + + expect(formatErrorCauseChain(error)).toBe('[unavailable cause]'); + }); + + it('returns an empty string when there is no cause', () => { + expect(formatErrorCauseChain(new Error('boom'))).toBe(''); + expect(formatErrorCauseChain('not an error')).toBe(''); + expect(formatErrorCauseChain(undefined)).toBe(''); + }); + + it('renders the wrapped reason a `fetch failed` hides', () => { + // The whole point: the outer error says nothing, the cause says + // everything. The outer link is skipped — the log already prints it. + const cause = Object.assign(new Error('other side closed'), { + name: 'SocketError', + code: 'UND_ERR_SOCKET', + }); + const wrapper = new TypeError('fetch failed', { cause }); + + expect(formatErrorCauseChain(wrapper)).toBe( + 'SocketError: other side closed (UND_ERR_SOCKET)' + ); + }); + + it('renders each link of a multi-level chain, outermost first', () => { + const inner = Object.assign( + new Error('getaddrinfo ENOTFOUND ai-gateway.vercel.sh'), + { code: 'ENOTFOUND' } + ); + const wrapper = new Error('POST /v4/… transport failure (ENOTFOUND)', { + cause: new TypeError('fetch failed', { cause: inner }), + }); + + expect(formatErrorCauseChain(wrapper)).toBe( + [ + 'TypeError: fetch failed', + // The code is already in the message, so it is not repeated. + 'Error: getaddrinfo ENOTFOUND ai-gateway.vercel.sh', + ].join('\n') + ); + }); + + it('summarizes the attempts an AggregateError collects', () => { + // A happy-eyeballs connect reports every address it tried on `errors` + // and leaves the AggregateError itself blank. + const aggregate = new AggregateError( + [ + Object.assign(new Error('connect ECONNREFUSED 10.0.0.1:443'), { + code: 'ECONNREFUSED', + }), + Object.assign(new Error('connect ECONNREFUSED [::1]:443'), { + code: 'ECONNREFUSED', + }), + ], + '' + ); + + expect( + formatErrorCauseChain(new TypeError('fetch failed', { cause: aggregate })) + ).toBe( + [ + 'AggregateError', + 'Error: connect ECONNREFUSED 10.0.0.1:443', + 'Error: connect ECONNREFUSED [::1]:443', + ].join('\n') + ); + }); + + it('caps a long chain', () => { + let error = new Error('innermost'); + for (let i = 0; i < 8; i++) { + error = new Error(`level ${i}`, { cause: error }); + } + + const lines = formatErrorCauseChain(error).split('\n'); + expect(lines).toHaveLength(5); + expect(lines.at(-1)).toBe('…'); + }); + + it('stops on a cyclic chain', () => { + const inner = new Error('inner') as Error & { cause?: unknown }; + const outer = new Error('outer', { cause: inner }); + inner.cause = outer; + + expect(formatErrorCauseChain(outer)).toBe( + ['Error: inner', 'Error: outer'].join('\n') + ); + }); + + it('renders a non-error cause', () => { + expect( + formatErrorCauseChain(new Error('boom', { cause: 'a string' })) + ).toBe('a string'); + }); +}); diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index b505cb65da..813a65be5e 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -14,6 +14,92 @@ export function getErrorStack(v: unknown): string { return ''; } +/** Upper bound on the links {@link formatErrorCauseChain} renders. */ +const MAX_CAUSE_LINKS = 4; + +/** One `Name: message (CODE)` line for a link in a cause chain. */ +function describeErrorLink(value: unknown): string { + if (typeof value !== 'object' || value === null) { + return String(value); + } + const { name, message, code } = value as { + name?: unknown; + message?: unknown; + code?: unknown; + }; + const label = typeof name === 'string' && name ? name : 'Error'; + const text = + typeof message === 'string' && message ? `${label}: ${message}` : label; + return typeof code === 'string' && code && !text.includes(code) + ? `${text} (${code})` + : text; +} + +/** + * Summarize the `cause` chain hanging off a thrown value, one link per line, + * outermost first. The value itself is skipped: whatever logs this already + * states it, in the message or in the stack header. + * + * `util.inspect` renders `[cause]` when Node prints an error, but the + * structured logs read `name` / `message` / `stack` and drop everything else + * — exactly the wrong half for errors that arrive pre-wrapped. + * `TypeError: fetch failed` is the canonical one: undici's wrapper says + * nothing on its own and its stack is all `node:internal/` frames, so the DNS, + * socket or TLS failure that actually happened is only readable one or two + * `cause` hops down. Same for the world layer's own wrapping, where the + * request that failed is on the wrapper and the reason it failed is on the + * cause. + * + * `AggregateError` also gets its `errors` summarized, because a + * happy-eyeballs connect reports every attempt there and leaves the + * `AggregateError` itself blank. + * + * Returns `''` when there is no cause, so callers can drop the field. + */ +export function formatErrorCauseChain(value: unknown): string { + const lines: string[] = []; + const seen = new Set(); + + try { + for ( + let current = causeOf(value); + current != null && lines.length <= MAX_CAUSE_LINKS; + current = causeOf(current) + ) { + if (typeof current !== 'object') { + lines.push(String(current)); + break; + } + // A cause chain can loop (`err.cause = err`) or repeat a shared error. + if (seen.has(current)) break; + seen.add(current); + lines.push(describeErrorLink(current), ...aggregatedLinks(current)); + } + } catch { + // Causes can contain getters or proxies that throw. Logging must not + // replace the original error or prevent the run_failed event from being written. + lines.push('[unavailable cause]'); + } + + return lines.length > MAX_CAUSE_LINKS + ? [...lines.slice(0, MAX_CAUSE_LINKS), '…'].join('\n') + : lines.join('\n'); +} + +function causeOf(value: unknown): unknown { + return typeof value === 'object' && value !== null + ? (value as { cause?: unknown }).cause + : undefined; +} + +/** The attempts an `AggregateError` collected, if this link is one. */ +function aggregatedLinks(value: object): string[] { + const errors = (value as { errors?: unknown }).errors; + return Array.isArray(errors) + ? errors.slice(0, MAX_CAUSE_LINKS).map(describeErrorLink) + : []; +} + export interface NormalizedUnknownError { name: string; message: string; diff --git a/packages/world-vercel/README.md b/packages/world-vercel/README.md index 7c7ac1fb1d..14e0269706 100644 --- a/packages/world-vercel/README.md +++ b/packages/world-vercel/README.md @@ -6,6 +6,12 @@ Integrates with Vercel's infrastructure for storage, queuing, and authentication Used by default for deployments on Vercel. Authentication and API endpoints are configured automatically in Vercel deployments. +## Connection failures + +Backend connection failures and interrupted event streams follow existing retry policies, including failures with unrecognized error codes. Repeated HTTP/2 session failures rebuild the shared events connection pool. Invalid backend URL protocols, embedded credentials, Fetch-blocked ports, and unsupported request headers fail immediately. + +See [Backend Connection Failures](https://useworkflow.dev/docs/foundations/errors-and-retries#backend-connection-failures) for retry behavior and diagnostics. + ## Custom dispatcher HTTP requests (including the queue) default to a shared undici `RetryAgent` that handles connection pooling and retries. Pass a custom `dispatcher` to override it — e.g. to tune undici on newer Node runtimes: diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index f07cfbc9f3..1d7ea9cff2 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -488,4 +488,192 @@ describe('v4 transport reports failures to the events recycler', () => { expect(getEventsDispatcher({ token: 'test-token' })).not.toBe(before); }); + + it.each([ + 'pre-header', + 'post-header', + ])('keeps %s HTTP/2 session failures retryable and rebuilds the shared pool', async (phase) => { + const now = Date.now(); + // The preceding tests rebuilt the process-global pool as late as + // `now + 60_000`; stay clear of its anti-thrash cooldown. + vi.spyOn(Date, 'now').mockReturnValue( + now + (phase === 'pre-header' ? 100_000 : 120_000) + ); + const error = new TypeError('fetch failed', { + cause: Object.assign(new Error('Session received GOAWAY'), { + code: 'ERR_HTTP2_GOAWAY_SESSION', + }), + }); + vi.spyOn(globalThis, 'fetch') + .mockResolvedValueOnce(new Response('', { status: 404 })) + .mockImplementation(async () => { + if (phase === 'pre-header') throw error; + return new Response( + new ReadableStream({ + pull(controller) { + controller.error(error); + }, + }), + { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } + ); + }); + + // A completed response resets any failure streak from earlier requests + // that used this process-wide pool, even when the HTTP status is an error. + await expect( + getWorkflowRunEventsV4('wrun_1', {}, { token: 'test-token' }) + ).rejects.toMatchObject({ status: 404 }); + + const before = getEventsDispatcher({ token: 'test-token' }); + for (let i = 0; i < EVENTS_RECYCLE_AFTER_CONSECUTIVE_FAILURES; i++) { + const rejection = await getWorkflowRunEventsV4( + 'wrun_1', + {}, + { token: 'test-token' } + ).catch((cause: unknown) => cause); + expect(rejection).toMatchObject({ + name: 'WorkflowWorldError', + code: 'TRANSPORT', + }); + expect(rejection).toHaveProperty('cause', error); + if (i < EVENTS_RECYCLE_AFTER_CONSECUTIVE_FAILURES - 1) { + expect(getEventsDispatcher({ token: 'test-token' })).toBe(before); + } + } + expect(getEventsDispatcher({ token: 'test-token' })).not.toBe(before); + }); +}); + +/** + * A rejection from `fetch` means no response was produced, which is a + * transport failure regardless of what the cause chain says. Left raw, a + * `TypeError: fetch failed` reaches the runtime as an ordinary throw: + * `classifyRunError` reads it as USER_ERROR and the queue never redelivers + * the run, so a backend blip fails the run and blames customer code. + */ +describe('v4 transport wraps pre-response failures the allowlist misses', () => { + beforeEach(() => { + vi.stubEnv(NODE_HTTP_ENV_VAR, '0'); + }); + + afterEach(() => { + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + }); + + it('maps a bare `TypeError: fetch failed` to a TRANSPORT failure', async () => { + vi.spyOn(globalThis, 'fetch').mockRejectedValue( + new TypeError('fetch failed') + ); + + const rejection = await getWorkflowRunEventsV4( + 'wrun_1', + {}, + { token: 'test-token', dispatcher: {} } + ).catch((e) => e); + + expect(rejection).toMatchObject({ + name: 'WorkflowWorldError', + code: 'TRANSPORT', + }); + expect(rejection.message).toContain('transport failure'); + }); + + it('rejects a credential-bearing backend URL without dispatch or retry', async () => { + vi.stubEnv('VERCEL_WORKFLOW_SERVER_URL', 'http://user:password@127.0.0.1'); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + + const rejection = await getWorkflowRunEventsV4( + 'wrun_1', + {}, + { token: 'test-token' } + ).catch((error: unknown) => error); + + expect(rejection).toMatchObject({ + name: 'TypeError', + message: 'HTTP(S) URLs with embedded credentials are unsupported', + }); + expect(WorkflowWorldError.is(rejection)).toBe(false); + expect(fetchSpy).not.toHaveBeenCalled(); + }); + + it('preserves unsupported headers from the backend configuration as non-retryable', async () => { + vi.stubEnv('VERCEL_WORKFLOW_SERVER_URL', 'http://127.0.0.1:12345'); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const rejection = await getWorkflowRunEventsV4( + 'wrun_1', + {}, + { + token: 'test-token', + headers: { Expect: '100-continue' }, + } + ).catch((error: unknown) => error); + expect(fetchSpy).toHaveBeenCalledTimes(1); + await expect(fetchSpy.mock.results[0].value).rejects.toBe(rejection); + expect(rejection).toMatchObject({ + name: 'TypeError', + cause: { code: 'UND_ERR_NOT_SUPPORTED' }, + }); + expect(WorkflowWorldError.is(rejection)).toBe(false); + }); + + it('maps an unrecognized post-header write failure to a TRANSPORT failure', async () => { + const sessionFailure = Object.assign( + new Error('The session has been destroyed'), + { code: 'ERR_HTTP2_GOAWAY_SESSION' } + ); + vi.spyOn(globalThis, 'fetch').mockResolvedValue( + new Response( + new ReadableStream({ + pull(controller) { + controller.error(sessionFailure); + }, + }), + { + status: 200, + headers: { + 'x-wf-event-id': 'evnt_1', + 'x-wf-run-id': 'wrun_1', + 'x-wf-created-at': '2026-06-10T00:00:00.000Z', + }, + } + ) + ); + + const rejection = await createWorkflowRunEventV4( + { + runId: 'wrun_1', + eventType: 'step_completed', + specVersion: 6, + correlationId: 'step_1', + }, + { token: 'test-token', dispatcher: {} } + ).catch((error: unknown) => error); + + expect(rejection).toMatchObject({ + name: 'WorkflowWorldError', + code: 'TRANSPORT', + cause: sessionFailure, + }); + }); + + it('rethrows a request-construction fault unchanged', async () => { + const constructionFault = Object.assign( + new TypeError('Failed to parse URL from nonsense'), + { + cause: Object.assign(new TypeError('Invalid URL'), { + code: 'ERR_INVALID_URL', + }), + } + ); + vi.spyOn(globalThis, 'fetch').mockRejectedValue(constructionFault); + + await expect( + getWorkflowRunEventsV4( + 'wrun_1', + {}, + { token: 'test-token', dispatcher: {} } + ) + ).rejects.toBe(constructionFault); + }); }); diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index 7641898e10..fb06a1c045 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -29,8 +29,8 @@ import { noteEventsTransportOutcome, } from './http-client.js'; import { + describeTransportFailure, errorForResponse, - getTransientTransportCode, instrumentedFetch, parseRetryAfter, } from './http-core.js'; @@ -109,7 +109,10 @@ async function fetchV4( } } catch (cause) { noteEventsTransportOutcome(dispatcher, cause); - const transportCode = getTransientTransportCode(cause); + // A body read can fail after response headers have arrived. Classify + // every such failure as transport unless it has the shape of a + // permanent request-construction fault, just like the pre-header path. + const transportCode = describeTransportFailure(cause); controller.error( transportCode ? new WorkflowWorldError( diff --git a/packages/world-vercel/src/http-client.test.ts b/packages/world-vercel/src/http-client.test.ts index bf894933fb..3ae716a36f 100644 --- a/packages/world-vercel/src/http-client.test.ts +++ b/packages/world-vercel/src/http-client.test.ts @@ -950,7 +950,14 @@ describe('dispatcher recycling accounting', () => { // and an abort is the caller's own doing. it('counts only transport failures a rebuild can fix', () => { expect(isRecyclableTransportError(h2StreamTimeout())).toBe(true); - for (const code of ['UND_ERR_HEADERS_TIMEOUT', 'UND_ERR_BODY_TIMEOUT']) { + for (const code of [ + 'UND_ERR_HEADERS_TIMEOUT', + 'UND_ERR_BODY_TIMEOUT', + 'ERR_HTTP2_GOAWAY_SESSION', + 'ERR_HTTP2_INVALID_SESSION', + 'ERR_HTTP2_SESSION_ERROR', + 'ERR_HTTP2_STREAM_ERROR', + ]) { expect( isRecyclableTransportError(Object.assign(new Error(code), { code })) ).toBe(true); @@ -961,6 +968,8 @@ describe('dispatcher recycling accounting', () => { 'ENOTFOUND', 'ECONNREFUSED', 'CERT_HAS_EXPIRED', + 'ERR_HTTP2_INVALID_HEADER_VALUE', + 'ERR_HTTP2_STREAM_CANCEL', ]) { expect( isRecyclableTransportError(Object.assign(new Error(code), { code })) diff --git a/packages/world-vercel/src/http-client.ts b/packages/world-vercel/src/http-client.ts index 064a2dbe5c..4f5e4ccb87 100644 --- a/packages/world-vercel/src/http-client.ts +++ b/packages/world-vercel/src/http-client.ts @@ -736,8 +736,8 @@ const RETIRED_CLOSE_DELAY_MS = 5_000; const RETIRED_DESTROY_DELAY_MS = 60_000; /** - * undici error codes that mean "no response arrived over a connection that was - * already established". These are the failures a rebuild can fix; DNS, connect + * Errors from an established connection or a session that can no longer accept + * requests. These are the failures a rebuild can fix; DNS, connect * and TLS errors are excluded because a new agent would hit the same wall, and an * abort is excluded because it is the caller's own doing. */ @@ -746,6 +746,14 @@ const RECYCLABLE_ERROR_CODES = new Set([ 'UND_ERR_INFO', 'UND_ERR_HEADERS_TIMEOUT', 'UND_ERR_BODY_TIMEOUT', + // Node can surface these directly instead of undici's UND_ERR_INFO. Once + // repeated, retire the pool even when classification used its fallback. + // Keep this scoped to session/stream failures: ERR_HTTP2_* also contains + // request-validation and caller-cancellation errors a new pool cannot fix. + 'ERR_HTTP2_GOAWAY_SESSION', + 'ERR_HTTP2_INVALID_SESSION', + 'ERR_HTTP2_SESSION_ERROR', + 'ERR_HTTP2_STREAM_ERROR', ]); /** Guard against a self-referential `cause` chain. */ diff --git a/packages/world-vercel/src/http-core.test.ts b/packages/world-vercel/src/http-core.test.ts index 9b7d7703c4..3bff56b0aa 100644 --- a/packages/world-vercel/src/http-core.test.ts +++ b/packages/world-vercel/src/http-core.test.ts @@ -5,12 +5,15 @@ import { TooEarlyError, WorkflowWorldError, } from '@workflow/errors'; +import { NODE_HTTP_ENV_VAR } from '@workflow/world'; import { afterEach, describe, expect, it, vi } from 'vitest'; import { + describeTransportFailure, errorForResponse, formatVercelDiagnostics, getTransientTransportCode, getVercelDiagnostics, + instrumentedFetch, parseRetryAfter, resolveVercelApiToken, } from './http-core.js'; @@ -77,6 +80,158 @@ describe('getTransientTransportCode', () => { }); }); +describe('instrumentedFetch URL validation', () => { + afterEach(() => { + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + }); + + it('preserves Fetch port-blocking errors instead of wrapping them as TRANSPORT', async () => { + vi.stubEnv(NODE_HTTP_ENV_VAR, '0'); + // Use real Fetch: Request construction accepts this URL, but Fetch + // rejects it locally with a code-less `Error: bad port` cause. + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const rejection = await instrumentedFetch({ + method: 'GET', + url: 'http://127.0.0.1:21/events', + headers: new Headers(), + dispatcher: undefined, + peerService: 'workflow-server', + }).catch((error: unknown) => error); + + expect(fetchSpy).toHaveBeenCalledTimes(1); + await expect(fetchSpy.mock.results[0].value).rejects.toBe(rejection); + expect(rejection).toMatchObject({ + name: 'TypeError', + message: 'fetch failed', + cause: { message: 'bad port' }, + }); + expect(WorkflowWorldError.is(rejection)).toBe(false); + }); + + it.each([ + [{ Expect: '100-continue' }, 'UND_ERR_NOT_SUPPORTED'], + [{ Connection: 'invalid' }, 'UND_ERR_INVALID_ARG'], + ] as const)('preserves request-validation errors for %j', async (headers, code) => { + vi.stubEnv(NODE_HTTP_ENV_VAR, '0'); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const rejection = await instrumentedFetch({ + method: 'GET', + url: 'http://127.0.0.1:12345/events', + headers: new Headers(headers), + dispatcher: undefined, + }).catch((error: unknown) => error); + expect(fetchSpy).toHaveBeenCalledTimes(1); + await expect(fetchSpy.mock.results[0].value).rejects.toBe(rejection); + expect(rejection).toMatchObject({ name: 'TypeError', cause: { code } }); + expect(WorkflowWorldError.is(rejection)).toBe(false); + }); + + it.each([ + '0', + '1', + ])('rejects unsupported protocols before dispatch (WORKFLOW_NODE_HTTP=%s)', async (mode) => { + vi.stubEnv(NODE_HTTP_ENV_VAR, mode); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const onTransportOutcome = vi.fn(); + + await expect( + instrumentedFetch({ + method: 'GET', + url: 'ftp://localhost/events', + headers: new Headers(), + dispatcher: undefined, + peerService: 'workflow-server', + onTransportOutcome, + }) + ).rejects.toThrow(TypeError); + expect(fetchSpy).not.toHaveBeenCalled(); + expect(onTransportOutcome).not.toHaveBeenCalled(); + }); + + it.each([ + '0', + '1', + ])('rejects embedded URL credentials before dispatch (WORKFLOW_NODE_HTTP=%s)', async (mode) => { + vi.stubEnv(NODE_HTTP_ENV_VAR, mode); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const onTransportOutcome = vi.fn(); + + await expect( + instrumentedFetch({ + method: 'GET', + url: 'http://user:password@localhost/events', + headers: new Headers(), + dispatcher: undefined, + peerService: 'workflow-server', + onTransportOutcome, + }) + ).rejects.toThrow('HTTP(S) URLs with embedded credentials are unsupported'); + expect(fetchSpy).not.toHaveBeenCalled(); + expect(onTransportOutcome).not.toHaveBeenCalled(); + }); +}); + +describe('describeTransportFailure', () => { + it('reports the allowlisted code when the cause chain carries one', () => { + // Known codes keep naming themselves, so messages and the + // `UND_ERR_CONNECT_TIMEOUT` retry branch in makeRequest are unchanged. + const cause = Object.assign(new Error('Request failed'), { + code: 'UND_ERR_REQ_RETRY', + }); + const err = Object.assign(new TypeError('fetch failed'), { cause }); + expect(describeTransportFailure(err)).toBe('UND_ERR_REQ_RETRY'); + expect(getTransientTransportCode(err)).toBe('UND_ERR_REQ_RETRY'); + }); + + it.each([ + ['an h2 session error', 'ERR_HTTP2_GOAWAY_SESSION'], + ['an unreachable network', 'ENETUNREACH'], + ['a TLS handshake failure', 'UNABLE_TO_VERIFY_LEAF_SIGNATURE'], + ])('reports %s the allowlist has never seen', (_label, code) => { + // These are the failures the allowlist misses today. Each one still means + // no response was produced, so each one is a transport failure. + const cause = Object.assign(new Error('nope'), { code }); + const err = Object.assign(new TypeError('fetch failed'), { cause }); + expect(getTransientTransportCode(err)).toBeUndefined(); + expect(describeTransportFailure(err)).toBe(code); + }); + + it('falls back to the innermost error name when no code is present', () => { + // A happy-eyeballs connect rejects with an AggregateError whose codes live + // on `errors[]`, out of reach of a `cause` walk. + const cause = new AggregateError([new Error('ECONNREFUSED')], ''); + const err = Object.assign(new TypeError('fetch failed'), { cause }); + expect(describeTransportFailure(err)).toBe('AggregateError'); + }); + + it('describes a bare `TypeError: fetch failed`', () => { + expect(describeTransportFailure(new TypeError('fetch failed'))).toBe( + 'TypeError' + ); + }); + + it.each([ + ['a malformed URL', 'ERR_INVALID_URL'], + ['an invalid header value', 'ERR_HTTP_INVALID_HEADER_VALUE'], + ['a bad argument', 'ERR_INVALID_ARG_TYPE'], + ])('returns undefined for %s', (_label, code) => { + // The request was never formed. Redelivering it just re-forms the same + // broken request until the run runs out of deliveries. + const cause = Object.assign(new TypeError('Invalid URL'), { code }); + const err = Object.assign(new TypeError('Failed to parse URL from x'), { + cause, + }); + expect(describeTransportFailure(err)).toBeUndefined(); + }); + + it('stops walking a self-referential cause chain', () => { + const err = new Error('loop') as Error & { cause?: unknown }; + err.cause = err; + expect(describeTransportFailure(err)).toBe('Error'); + }); +}); + describe('parseRetryAfter', () => { it('parses integer seconds', () => { expect(parseRetryAfter('30')).toBe(30); diff --git a/packages/world-vercel/src/http-core.ts b/packages/world-vercel/src/http-core.ts index 8fd21934cc..a069ea5cb5 100644 --- a/packages/world-vercel/src/http-core.ts +++ b/packages/world-vercel/src/http-core.ts @@ -101,6 +101,108 @@ export function getTransientTransportCode(error: unknown): string | undefined { return undefined; } +/** Reject invalid URLs before dispatch, where a failure would be retryable. */ +export function validateHttpUrl(url: string): void { + const { protocol, username, password } = new URL(url); + // Both fetch and nodeHttpFetch can reject unsupported schemes without an + // error code, so describeTransportFailure cannot identify these faults. + if (protocol !== 'http:' && protocol !== 'https:') { + throw new TypeError( + `Unsupported URL protocol ${protocol}; expected http: or https:` + ); + } + // Fetch rejects URL userinfo locally with a code-less TypeError. Keep that + // permanent configuration fault outside the transport classifier, and make + // the node:http and Fetch paths agree instead of allowing one to send it. + if (username || password) { + throw new TypeError( + 'HTTP(S) URLs with embedded credentials are unsupported' + ); + } +} + +/** + * Codes that mean the request was never *formed*, as opposed to formed and + * then failed on the wire. `fetch()` reports a malformed URL, an invalid + * header name/value, or a bad argument as a rejected `TypeError` that is + * structurally identical to the `TypeError: fetch failed` it raises for a dead + * socket, and the node:http path throws Node's own `ERR_*`. + * + * These faults are permanent — every redelivery re-forms the same broken + * request — so they must keep propagating raw rather than being classified as + * a retryable transport failure, which would spend the run's whole delivery + * budget before failing it with a less specific error than it started with. + */ +const REQUEST_CONSTRUCTION_ERROR_CODES = new Set([ + 'ERR_INVALID_URL', + 'ERR_INVALID_ARG_TYPE', + 'ERR_INVALID_ARG_VALUE', + 'ERR_INVALID_CHAR', + 'ERR_INVALID_HTTP_TOKEN', + 'ERR_HTTP_INVALID_HEADER_VALUE', + 'ERR_UNESCAPED_CHARACTERS', + // Undici validates request options and headers during dispatch, after + // Fetch has constructed the Request (e.g. unsupported Expect headers). + 'UND_ERR_INVALID_ARG', + 'UND_ERR_NOT_SUPPORTED', +]); + +/** + * Classify a rejection from `fetch()` / `nodeHttpFetch()`, calls that only + * settle once the response headers are in hand. + * + * A rejection leaves the request outcome unknown: it may have failed locally, + * or the backend may have applied it without a response reaching the caller. + * Preserve known request-construction faults; route other failures through + * the existing retry policies instead of attributing them to user code. + * + * {@link TRANSIENT_TRANSPORT_ERROR_CODES} alone could not hold that line, + * because it can only list failures someone has already seen. The ones it + * misses are not exotic: HTTP/2 session errors (`ERR_HTTP2_GOAWAY_SESSION` and + * friends — the shared events pool negotiates h2), TLS handshake failures, + * `ENETUNREACH` / `EHOSTUNREACH`, and the `AggregateError` a happy-eyeballs + * connect raises, which carries its codes on `errors[]` where a `cause` walk + * cannot see them. Each of those used to propagate raw, and a raw + * `TypeError: fetch failed` is indistinguishable from a user throw by the time + * it reaches `classifyRunError`: the run failed as `USER_ERROR`, attributing a + * backend outage to the customer, and the queue never redelivered it. + * + * Returns the most specific marker available to name the failure in the error + * message: the allowlisted code when there is one (so known failures keep + * reporting exactly what they reported before), otherwise the first `code` in + * the cause chain, otherwise the innermost error name. `undefined` means the + * request was never formed and the caller should rethrow as-is. + */ +export function describeTransportFailure(error: unknown): string | undefined { + const known = getTransientTransportCode(error); + if (known) return known; + + let firstCode: string | undefined; + let innermostName: string | undefined; + let current = error; + for (let depth = 0; current && depth < 8; depth++) { + const { code, name, message } = current as { + code?: unknown; + name?: unknown; + message?: unknown; + }; + // Node Fetch enforces the Fetch Standard's port blocking after Request + // construction. Its `TypeError: fetch failed` wraps a code-less + // `Error: bad port`; retrying cannot make that URL acceptable. Preserve + // the original rejection without duplicating Fetch's blocked-port list. + if (name === 'Error' && message === 'bad port' && code === undefined) { + return undefined; + } + if (typeof code === 'string' && code) { + if (REQUEST_CONSTRUCTION_ERROR_CODES.has(code)) return undefined; + firstCode ??= code; + } + if (typeof name === 'string' && name) innermostName = name; + current = (current as { cause?: unknown }).cause; + } + return firstCode ?? innermostName ?? 'unknown'; +} + /** * Lightweight debug logger toggle for HTTP requests. Activated when the DEBUG * env var contains "workflow:" or is "*". @@ -413,6 +515,7 @@ export async function instrumentedFetch( deferTransportSuccessUntilBody = false, } = opts; const label = logLabel ?? url; + validateHttpUrl(url); return trace( `http ${method}`, @@ -439,19 +542,24 @@ export async function instrumentedFetch( ? AbortSignal.any([callerSignal, timeoutSignal]) : (callerSignal ?? timeoutSignal); + // With no dispatcher to honor, `WORKFLOW_NODE_HTTP` takes the request + // off undici altogether rather than leaving it on the undici behind + // `fetch`. A dispatcher the caller supplied is an instruction to use + // undici, so it keeps the request on `fetch`. + // + // Resolved outside the try: the catch below reads everything it sees as + // a failure of the request on the wire, and picking an agent happens + // before there is one. + const nodeAgents = dispatcher ? undefined : getNodeHttpAgents(); + // Both transports issue the same span against the same URL, so this is + // the only thing that tells them apart in a trace. + span?.setAttributes({ + ...WorkflowHttpTransport(nodeAgents ? 'node-http' : 'undici'), + }); + const start = Date.now(); let response: Response; try { - // With no dispatcher to honour, `WORKFLOW_NODE_HTTP` takes the request - // off undici altogether rather than leaving it on the undici behind - // `fetch`. A dispatcher the caller supplied is an instruction to use - // undici, so it keeps the request on `fetch`. - const nodeAgents = dispatcher ? undefined : getNodeHttpAgents(); - // Both transports issue the same span against the same URL, so this is - // the only thing that tells them apart in a trace. - span?.setAttributes({ - ...WorkflowHttpTransport(nodeAgents ? 'node-http' : 'undici'), - }); response = nodeAgents ? await nodeHttpFetch(url, { method, @@ -493,7 +601,10 @@ export async function instrumentedFetch( span?.recordException?.(timeoutError); throw timeoutError; } - const transportCode = getTransientTransportCode(error); + // Nothing below this point saw a response, so anything that is not a + // request-construction fault is a transport failure — including codes + // the allowlist has never seen. See describeTransportFailure. + const transportCode = describeTransportFailure(error); if (transportCode) { const transportError = new WorkflowWorldError( `${method} ${label} transport failure after ${elapsed}ms (${transportCode})`, diff --git a/packages/world-vercel/src/utils.test.ts b/packages/world-vercel/src/utils.test.ts index 21ef03259e..ca80db6337 100644 --- a/packages/world-vercel/src/utils.test.ts +++ b/packages/world-vercel/src/utils.test.ts @@ -5,6 +5,7 @@ import { NODE_HTTP_ENV_VAR } from '@workflow/world'; import { encode } from 'cbor-x'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { z } from 'zod'; +import { isRetryableEventPostError } from './event-retry.js'; import { getHeaders, getHttpConfig, @@ -420,6 +421,60 @@ describe('makeRequest body-parse retry', () => { }); }); +describe('makeRequest URL validation', () => { + beforeEach(() => { + vi.restoreAllMocks(); + }); + + afterEach(() => { + vi.unstubAllEnvs(); + vi.restoreAllMocks(); + }); + + it.each([ + 'http', + 'https', + ])('preserves Fetch port-blocking errors for %s backend URLs', async (scheme) => { + vi.stubEnv(NODE_HTTP_ENV_VAR, '0'); + vi.stubEnv('VERCEL_WORKFLOW_SERVER_URL', `${scheme}://127.0.0.1:21`); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const rejection = await makeRequest({ + endpoint: '/v3/runs/wrun_test/events', + options: { method: 'GET' }, + schema: z.unknown(), + config: { token: 'test-token' }, + }).catch((error: unknown) => error); + + expect(fetchSpy).toHaveBeenCalledTimes(1); + await expect(fetchSpy.mock.results[0].value).rejects.toBe(rejection); + expect(rejection).toMatchObject({ + name: 'TypeError', + message: 'fetch failed', + cause: { message: 'bad port' }, + }); + expect(isRetryableEventPostError(rejection)).toBe(false); + }); + + it.each([ + '0', + '1', + ])('rejects unsupported backend protocols before dispatch (WORKFLOW_NODE_HTTP=%s)', async (mode) => { + vi.stubEnv(NODE_HTTP_ENV_VAR, mode); + vi.stubEnv('VERCEL_WORKFLOW_SERVER_URL', 'ftp://localhost'); + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + + await expect( + makeRequest({ + endpoint: '/v3/runs/wrun_test/events', + options: { method: 'GET' }, + schema: z.unknown(), + config: { token: 'test-token' }, + }) + ).rejects.toThrow(TypeError); + expect(fetchSpy).not.toHaveBeenCalled(); + }); +}); + describe('makeRequest transport errors', () => { const schema = z.object({ value: z.string() }); const originalEnv = process.env; @@ -494,8 +549,68 @@ describe('makeRequest transport errors', () => { expect(rejection.cause).toBe(fetchErr); }); - it('rethrows a non-transient fetch error unchanged', async () => { - const fetchErr = new Error('some unexpected non-network error'); + it('maps a fetch failure whose code the allowlist has never seen to TRANSPORT', async () => { + // A dead h2 session is the shape that motivated inverting the default: + // the events pool negotiates HTTP/2, `ERR_HTTP2_GOAWAY_SESSION` is not in + // TRANSIENT_TRANSPORT_ERROR_CODES, and no response was produced either + // way. The unknown code still names the failure in the message. + const cause = Object.assign(new Error('The session has been destroyed'), { + code: 'ERR_HTTP2_GOAWAY_SESSION', + }); + const fetchErr = Object.assign(new TypeError('fetch failed'), { cause }); + vi.stubGlobal('fetch', vi.fn().mockRejectedValue(fetchErr)); + + const rejection = await makeRequest({ + endpoint: '/v3/runs/wrun_test/events', + options: { method: 'GET' }, + schema, + }).catch((e) => e); + + expect(rejection).toMatchObject({ + name: 'WorkflowWorldError', + code: 'TRANSPORT', + }); + expect(rejection.message).toContain('ERR_HTTP2_GOAWAY_SESSION'); + expect(rejection.cause).toBe(fetchErr); + }); + + it('maps a bare `TypeError: fetch failed` to a retryable TRANSPORT error', async () => { + // undici's wrapper with nothing usable underneath (the AggregateError a + // happy-eyeballs connect raises hangs its codes off `errors[]`, where a + // `cause` walk cannot see them). Rethrown raw, this reached + // `classifyRunError` as an ordinary `TypeError` and failed the run as + // USER_ERROR — a backend outage billed to the customer — without the + // queue ever redelivering it. + const fetchErr = new TypeError('fetch failed'); + vi.stubGlobal('fetch', vi.fn().mockRejectedValue(fetchErr)); + + const rejection = await makeRequest({ + endpoint: '/v3/runs/wrun_test/events', + options: { method: 'GET' }, + schema, + }).catch((e) => e); + + expect(rejection).toMatchObject({ + name: 'WorkflowWorldError', + code: 'TRANSPORT', + }); + expect(rejection.cause).toBe(fetchErr); + // That code is what queue redelivery in `@workflow/core` keys on + // (isRetryableWorldError, which also classifies the terminal failure as + // WORLD_CONTRACT_ERROR rather than USER_ERROR). + }); + + it('rethrows a request-construction fault unchanged', async () => { + // A malformed URL or header is permanent: every redelivery re-forms the + // same broken request, so it must not be dressed up as retryable. + const fetchErr = Object.assign( + new TypeError('Failed to parse URL from nonsense'), + { + cause: Object.assign(new TypeError('Invalid URL'), { + code: 'ERR_INVALID_URL', + }), + } + ); vi.stubGlobal('fetch', vi.fn().mockRejectedValue(fetchErr)); await expect( diff --git a/packages/world-vercel/src/utils.ts b/packages/world-vercel/src/utils.ts index 8fd4736ef9..15a92e4cf4 100644 --- a/packages/world-vercel/src/utils.ts +++ b/packages/world-vercel/src/utils.ts @@ -12,15 +12,16 @@ import { getNodeHttpPhaseTimeouts, } from './http-client.js'; import { + describeTransportFailure, errorForResponse, formatVercelDiagnostics, - getTransientTransportCode, HTTP_DEBUG_ENABLED, httpClientSpanAttributes, httpLog, logCurlRepro, parseRetryAfter, REQUEST_TIMEOUT_MS, + validateHttpUrl, } from './http-core.js'; import { ErrorType, @@ -349,6 +350,7 @@ export async function makeRequest({ const method = (options.method || 'GET').toUpperCase(); const { baseUrl, headers } = await getHttpConfig(config); const url = `${baseUrl}${endpoint}`; + validateHttpUrl(url); // Standard OTEL span name for HTTP client: "{method}" // See: https://opentelemetry.io/docs/specs/semconv/http/http-spans/#name @@ -403,19 +405,28 @@ export async function makeRequest({ const signal = options.signal ? AbortSignal.any([options.signal, timeoutSignal]) : timeoutSignal; + // `WORKFLOW_NODE_HTTP` takes this request off undici entirely, rather + // than leaving it on the undici behind `fetch`. `getNodeHttpAgents` + // returns the pool only when no caller dispatcher was supplied, so + // an explicit `config.dispatcher` still keeps the request on `fetch`. + // + // Agent selection and `Request` construction (which validates the URL + // and the headers) sit outside the try on purpose: the catch below + // reads everything it sees as a failure of the request on the wire, + // and these run before there is one. + const nodeAgents = getNodeHttpAgents(config); + const undiciRequest = nodeAgents + ? undefined + : new Request(url, { ...options, body, headers, signal }); + const undiciDispatcher = nodeAgents ? undefined : getDispatcher(config); + // Both transports issue the same span against the same URL, so this + // is the only thing that tells them apart in a trace. + span?.setAttributes({ + ...WorkflowHttpTransport(nodeAgents ? 'node-http' : 'undici'), + }); const fetchStart = Date.now(); let response: Response; try { - // `WORKFLOW_NODE_HTTP` takes this request off undici entirely, rather - // than leaving it on the undici behind `fetch`. `getNodeHttpAgents` - // returns the pool only when no caller dispatcher was supplied, so - // an explicit `config.dispatcher` still keeps the request on `fetch`. - const nodeAgents = getNodeHttpAgents(config); - // Both transports issue the same span against the same URL, so this - // is the only thing that tells them apart in a trace. - span?.setAttributes({ - ...WorkflowHttpTransport(nodeAgents ? 'node-http' : 'undici'), - }); response = nodeAgents ? await nodeHttpFetch(url, { method, @@ -429,10 +440,10 @@ export async function makeRequest({ ...getNodeHttpPhaseTimeouts(), }) : await fetch( - new Request(url, { ...options, body, headers, signal }), + undiciRequest as Request, { // eslint-disable-next-line @typescript-eslint/no-explicit-any -- undici v7 dispatcher types don't match @types/node's RequestInit - dispatcher: getDispatcher(config), + dispatcher: undiciDispatcher, } as any ); } catch (error) { @@ -452,11 +463,14 @@ export async function makeRequest({ span?.recordException?.(timeoutError); throw timeoutError; } - // Transient transport failure (RetryAgent retries exhausted, socket - // reset, connect/DNS failure). Surface as a retryable - // WorkflowWorldError so the runtime redrives via the queue instead - // of failing the run. See TRANSIENT_TRANSPORT_ERROR_CODES. - const transportCode = getTransientTransportCode(error); + // The request produced no response (RetryAgent retries exhausted, + // socket reset, connect/DNS/TLS failure, dead h2 session, …), so + // this is a transport failure whether or not the code is one we + // have seen before. Surface it as a retryable WorkflowWorldError so + // the runtime redrives via the queue instead of failing the run + // with a backend outage attributed to user code. See + // describeTransportFailure. + const transportCode = describeTransportFailure(error); if (transportCode) { if ( retryConnectTimeout &&