From d7290b291a29e5ccaa3978bb9ec48f54879bc0eb Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Thu, 13 Aug 2026 20:12:10 -0700 Subject: [PATCH 1/7] Resume partial replay event streams --- .changeset/resume-partial-replay-streams.md | 6 + packages/world-vercel/src/events-v4.test.ts | 215 ++++++++++++++++++++ packages/world-vercel/src/events-v4.ts | 144 +++++++++++-- packages/world-vercel/src/events.test.ts | 59 +++++- packages/world-vercel/src/events.ts | 12 +- packages/world/src/events.ts | 16 ++ 6 files changed, 417 insertions(+), 35 deletions(-) create mode 100644 .changeset/resume-partial-replay-streams.md diff --git a/.changeset/resume-partial-replay-streams.md b/.changeset/resume-partial-replay-streams.md new file mode 100644 index 0000000000..9f5d35e640 --- /dev/null +++ b/.changeset/resume-partial-replay-streams.md @@ -0,0 +1,6 @@ +--- +"@workflow/world": patch +"@workflow/world-vercel": patch +--- + +Resume interrupted or partial replay event streams after their last validated event and expose decoded events to streaming consumers. diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index 911e1a57e3..503e0a3475 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -891,6 +891,221 @@ describe('createWorkflowRunEventV4 over HTTP', () => { agent.assertNoPendingInterceptors(); }); + it('continues a truncated run_started stream after its last event', async () => { + const origin = + WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; + const agent = new MockAgent(); + agent.disableNetConnect(); + const observed: string[] = []; + + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events/run_started', + method: 'POST', + headers: { accept: V4_FRAME_CONTENT_TYPE }, + }) + .reply( + 200, + encodeFrame( + { + eventId: 'evnt_1', + runId: 'wrun_1', + eventType: 'run_created', + createdAt: CREATED_AT, + eventData: { + deploymentId: 'dpl_1', + workflowName: 'workflow', + input: null, + }, + }, + new Uint8Array() + ), + { + headers: { + 'content-type': V4_FRAME_CONTENT_TYPE, + 'x-wf-max-events': '10000', + }, + } + ); + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_1', + method: 'GET', + }) + .reply( + 200, + Buffer.concat([ + encodeFrame( + { + eventId: 'evnt_2', + runId: 'wrun_1', + eventType: 'run_started', + createdAt: CREATED_AT, + }, + new Uint8Array() + ), + encodeFrame( + { _end: 1, next: 'eid:evnt_2', hasMore: false }, + new Uint8Array() + ), + ]), + { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } + ); + + const result = await createWorkflowRunStartedEventV4( + { runId: 'wrun_1', specVersion: 5 }, + { token: 'test-token', dispatcher: agent }, + (event) => observed.push(event.eventId) + ); + + expect(result.events.map((event) => event.eventId)).toEqual([ + 'evnt_1', + 'evnt_2', + ]); + expect(result.hasMore).toBe(false); + expect(observed).toEqual(['evnt_1', 'evnt_2']); + agent.assertNoPendingInterceptors(); + }); + + it('does not treat an event observer failure as stream truncation', async () => { + const origin = + WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; + const agent = new MockAgent(); + agent.disableNetConnect(); + + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events/run_started', + method: 'POST', + headers: { accept: V4_FRAME_CONTENT_TYPE }, + }) + .reply( + 200, + Buffer.concat([ + encodeFrame( + { + eventId: 'evnt_1', + runId: 'wrun_1', + eventType: 'run_created', + createdAt: CREATED_AT, + eventData: { + deploymentId: 'dpl_1', + workflowName: 'workflow', + input: null, + }, + }, + new Uint8Array() + ), + encodeFrame( + { _end: 1, next: 'eid:evnt_1', hasMore: false }, + new Uint8Array() + ), + ]), + { + headers: { + 'content-type': V4_FRAME_CONTENT_TYPE, + 'x-wf-max-events': '10000', + }, + } + ); + + await expect( + createWorkflowRunStartedEventV4( + { runId: 'wrun_1', specVersion: 5 }, + { token: 'test-token', dispatcher: agent }, + () => { + throw new Error('observer failed'); + } + ) + ).rejects.toThrow('observer failed'); + agent.assertNoPendingInterceptors(); + }); + + it('continues a graceful partial run_started stream from its sentinel cursor', async () => { + const origin = + WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; + const agent = new MockAgent(); + agent.disableNetConnect(); + + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events/run_started', + method: 'POST', + headers: { accept: V4_FRAME_CONTENT_TYPE }, + }) + .reply( + 200, + Buffer.concat([ + encodeFrame( + { + eventId: 'evnt_1', + runId: 'wrun_1', + eventType: 'run_created', + createdAt: CREATED_AT, + eventData: { + deploymentId: 'dpl_1', + workflowName: 'workflow', + input: null, + }, + }, + new Uint8Array() + ), + encodeFrame( + { _end: 1, next: 'eid:evnt_1', hasMore: true }, + new Uint8Array() + ), + ]), + { + headers: { + 'content-type': V4_FRAME_CONTENT_TYPE, + 'x-wf-max-events': '10000', + }, + } + ); + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_1', + method: 'GET', + }) + .reply( + 200, + Buffer.concat([ + encodeFrame( + { + eventId: 'evnt_2', + runId: 'wrun_1', + eventType: 'run_started', + createdAt: CREATED_AT, + }, + new Uint8Array() + ), + encodeFrame( + { _end: 1, next: 'eid:evnt_2', hasMore: false }, + new Uint8Array() + ), + ]), + { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } + ); + + const result = await createWorkflowRunStartedEventV4( + { runId: 'wrun_1', specVersion: 5 }, + { token: 'test-token', dispatcher: agent } + ); + + expect(result.events.map((event) => event.eventId)).toEqual([ + 'evnt_1', + 'evnt_2', + ]); + expect(result.cursor).toBe('eid:evnt_2'); + expect(result.hasMore).toBe(false); + agent.assertNoPendingInterceptors(); + }); + it('requires the event-stream response requested by run_started', async () => { const origin = WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index a3586f95de..4fc599eb37 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -27,6 +27,7 @@ import { type Event, type EventResult, EventSchema, + type EventStreamObserver, type EventType, EventTypeSchema, getEventDataPayloadField, @@ -779,7 +780,8 @@ async function decodeCreateEventResponse( export async function createWorkflowRunStartedEventV4( input: CreateEventV4InputBase, - config?: APIConfig + config?: APIConfig, + onEvent?: EventStreamObserver ) { const response = await postWorkflowRunEventV4( { ...input, eventType: 'run_started' }, @@ -787,7 +789,14 @@ export async function createWorkflowRunStartedEventV4( config ); const events: Event[] = []; - const page = await consumeEventFrameStream(response, 'createEvent', events); + const page = await consumeReplayLogResponse({ + response, + runId: input.runId, + opName: 'createEvent', + events, + config, + onEvent, + }); assert(page.cursor, 'v4 createEvent: event stream missing cursor'); const maxEvents = MaxEventsHeaderSchema.safeParse( response.headers.get(MAX_EVENTS_HEADER) @@ -1225,13 +1234,14 @@ export type HookReceivedPreloadV4Result = * A server that supports the lazy-hook replay stream answers the consumer's * idempotent re-ensure with the run's complete replay log as v4 frames — * the same event-frame sequence LIST uses, ending with the `_end` sentinel. - * A truncated stream (EOF without the sentinel) throws; the write is - * deduplicated by the server's `(runId, resumeId)` constraint, so retrying - * the whole request is safe and converges on the same canonical event. + * A truncated stream resumes with a GET after its last validated event. If it + * fails before producing any event, the outer POST retry remains safe because + * the server deduplicates `(runId, resumeId)` and returns the canonical event. */ export async function createHookReceivedPreloadEventV4( input: CreateEventV4InputBase, - config?: APIConfig + config?: APIConfig, + onEvent?: EventStreamObserver ): Promise { const response = await postWorkflowRunEventV4( { ...input, eventType: 'hook_received' }, @@ -1248,7 +1258,14 @@ export async function createHookReceivedPreloadEventV4( } const events: Event[] = []; - const page = await consumeEventFrameStream(response, 'createEvent', events); + const page = await consumeReplayLogResponse({ + response, + runId: input.runId, + opName: 'createEvent', + events, + config, + onEvent, + }); const maxEvents = MaxEventsHeaderSchema.safeParse( response.headers.get(MAX_EVENTS_HEADER) ); @@ -1326,6 +1343,8 @@ export interface ListEventsV4Params extends PaginationOptions { * it. */ remoteRefBehavior?: 'resolve' | 'lazy'; + /** Called synchronously after each event frame validates. */ + onEvent?: EventStreamObserver; } export interface ListEventsV4Result { @@ -1339,7 +1358,8 @@ export interface ListEventsV4Result { async function consumeEventFrameStream( response: Response, opName: string, - events: Event[] + events: Event[], + onEvent?: EventStreamObserver ): Promise> { const contentType = response.headers.get('content-type'); if (!contentType?.startsWith(V4_FRAME_CONTENT_TYPE)) { @@ -1349,24 +1369,91 @@ async function consumeEventFrameStream( } const chunks = response.body as unknown as AsyncIterable; + const frames = decodeFrames(chunks)[Symbol.asyncIterator](); - for await (const frame of decodeFrames(chunks)) { - if (frame.meta._end === 1) { - const end = EventStreamEndSchema.parse(frame.meta); - return { cursor: end.next ?? null, hasMore: end.hasMore }; - } - if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { - throw new Error(`v4 ${opName}: unexpected control frame`); + try { + for (;;) { + let next: IteratorResult; + try { + next = await frames.next(); + } catch (cause) { + throw new WorkflowWorldError( + `v4 ${opName}: event frame stream failed after ${events.length} events`, + { code: 'TRANSPORT', cause } + ); + } + if (next.done) break; + + const frame = next.value; + if (frame.meta._end === 1) { + const end = EventStreamEndSchema.parse(frame.meta); + return { cursor: end.next ?? null, hasMore: end.hasMore }; + } + if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { + throw new Error(`v4 ${opName}: unexpected control frame`); + } + const event = decodeEventFrame(frame); + events.push(event); + // Deliberately outside the stream-read catch above: observer/application + // failures are not truncation and must never advance recovery past this + // event. + onEvent?.(event); } - events.push(decodeEventFrame(frame)); + } finally { + await frames.return?.(undefined); } - throw new Error( + throw new WorkflowWorldError( `v4 ${opName}: frame stream ended without the end-of-stream sentinel ` + - `(${events.length} events read) — truncated response?` + `(${events.length} events read) — truncated response?`, + { code: 'TRANSPORT' } ); } +/** + * Finish a replay-log response, reusing every validated prefix. A graceful + * `hasMore` sentinel and an interrupted body now converge on the same GET + * continuation instead of making run_started download its accepted prefix + * again. + */ +async function consumeReplayLogResponse({ + response, + runId, + opName, + events, + config, + onEvent, +}: { + response: Response; + runId: string; + opName: string; + events: Event[]; + config?: APIConfig; + onEvent?: EventStreamObserver; +}): Promise> { + let page: Pick; + try { + page = await consumeEventFrameStream(response, opName, events, onEvent); + } catch (error) { + if (!WorkflowWorldError.is(error) || error.code !== 'TRANSPORT') { + throw error; + } + const lastEvent = events.at(-1); + if (!lastEvent) throw error; + page = { cursor: `eid:${lastEvent.eventId}`, hasMore: true }; + } + + if (!page.hasMore) return page; + assert(page.cursor, `v4 ${opName}: partial event stream missing cursor`); + const suffix = await getWorkflowRunEventsV4( + runId, + { cursor: page.cursor, onEvent }, + config + ); + events.push(...suffix.events); + return { cursor: suffix.cursor, hasMore: suffix.hasMore }; +} + /** * Drive a v4 frame-stream list response into an in-memory page. Used by * both the by-runId and by-correlationId list endpoints — the wire @@ -1381,7 +1468,8 @@ async function consumeListFrameStream( headers: Headers, config: APIConfig | undefined, opName: string, - events: Event[] + events: Event[], + onEvent?: EventStreamObserver ): Promise> { const response = await fetchV4( url, @@ -1389,7 +1477,7 @@ async function consumeListFrameStream( config, opName ); - return consumeEventFrameStream(response, opName, events); + return consumeEventFrameStream(response, opName, events, onEvent); } /** @@ -1443,8 +1531,19 @@ export async function getWorkflowRunEventsV4( headers, config, 'listEvents', - events + events, + params.onEvent ); + if (params.limit === undefined && page.hasMore) { + if (!page.cursor || page.cursor === cursor) { + throw new WorkflowWorldError( + `v4 listEvents: partial event stream made no cursor progress for run ${runId}`, + { code: 'SCHEMA_VALIDATION' } + ); + } + cursor = page.cursor; + continue; + } return { events, ...page }; } catch (error) { const lastEvent = events.at(-1); @@ -1492,7 +1591,8 @@ export async function getEventsByCorrelationIdV4( headers, config, 'listEventsByCorrelationId', - events + events, + params.onEvent ); return { events, ...page }; } diff --git a/packages/world-vercel/src/events.test.ts b/packages/world-vercel/src/events.test.ts index 5fea9675a4..d072603000 100644 --- a/packages/world-vercel/src/events.test.ts +++ b/packages/world-vercel/src/events.test.ts @@ -1393,7 +1393,14 @@ describe('getWorkflowRunEvents legacy structured-error compatibility', () => { * fallback preserves their (correct, if slower) behavior. */ describe('getWorkflowRunEvents hasMore mapping', () => { - function mockListResponse(agent: MockAgent, sentinelMeta: object) { + function mockListResponse( + agent: MockAgent, + sentinelMeta: object, + query: Record = { + returnAll: 'true', + remoteRefBehavior: 'resolve', + } + ) { const frames = Buffer.concat([ encodeFrame( { @@ -1412,9 +1419,7 @@ describe('getWorkflowRunEvents hasMore mapping', () => { .intercept({ path: '/api/v4/runs/wrun_1/events', method: 'GET', - // These tests omit the limit and use the default resolveData - // ('all' → resolve); match both translated query params. - query: { returnAll: 'true', remoteRefBehavior: 'resolve' }, + query, }) .reply(200, frames, { headers: { 'content-type': V4_FRAME_CONTENT_TYPE }, @@ -1439,10 +1444,14 @@ describe('getWorkflowRunEvents hasMore mapping', () => { it('maps an explicit hasMore:true through', async () => { const agent = mockAgent(); - mockListResponse(agent, { _end: 1, next: 'cursor-2', hasMore: true }); + mockListResponse( + agent, + { _end: 1, next: 'cursor-2', hasMore: true }, + { limit: '500', remoteRefBehavior: 'resolve' } + ); const result = await getWorkflowRunEvents( - { runId: 'wrun_1' }, + { runId: 'wrun_1', pagination: { limit: 500 } }, { token: 'test-token', dispatcher: agent } ); @@ -1842,7 +1851,7 @@ describe('createWorkflowRunEvent hook_received replay preload', () => { agent.assertNoPendingInterceptors(); }); - it('rejects a truncated preload stream (no end sentinel)', async () => { + it('continues a truncated preload stream after its last event', async () => { const agent = mockAgent(); agent .get(ORIGIN) @@ -1866,15 +1875,43 @@ describe('createWorkflowRunEvent hook_received replay preload', () => { }, PAYLOAD ), + { + headers: { + 'content-type': V4_FRAME_CONTENT_TYPE, + 'x-wf-event-id': 'evnt_4', + }, + } + ); + + agent + .get(ORIGIN) + .intercept({ + path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_4', + method: 'GET', + }) + .reply( + 200, + encodeFrame( + { _end: 1, next: 'eid:evnt_4', hasMore: false }, + new Uint8Array() + ), { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } ); - await expect( - createWorkflowRunEvent('wrun_1', hookReceivedRequest(), preloadParams, { + const result = await createWorkflowRunEvent( + 'wrun_1', + hookReceivedRequest(), + preloadParams, + { token: 'test-token', dispatcher: agent, - }) - ).rejects.toThrow(/end-of-stream sentinel/); + } + ); + + expect(result.event?.eventId).toBe('evnt_4'); + expect(result.events).toHaveLength(1); + expect(result.cursor).toBe('eid:evnt_4'); + expect(result.hasMore).toBe(false); agent.assertNoPendingInterceptors(); }); diff --git a/packages/world-vercel/src/events.ts b/packages/world-vercel/src/events.ts index 1d0279229d..8f6308a4a7 100644 --- a/packages/world-vercel/src/events.ts +++ b/packages/world-vercel/src/events.ts @@ -446,6 +446,9 @@ export async function getWorkflowRunEvents( const listParams: ListEventsV4Params = { ...pagination, remoteRefBehavior: resolveData === 'none' ? 'lazy' : 'resolve', + ...('onEvent' in params && params.onEvent + ? { onEvent: params.onEvent } + : {}), }; const result = await ('correlationId' in params @@ -744,7 +747,11 @@ async function createWorkflowRunEventInner( }; if (data.eventType === 'run_started' && !params?.skipPreload) { - const result = await createWorkflowRunStartedEventV4(input, config); + const result = await createWorkflowRunStartedEventV4( + input, + config, + params?.onEvent + ); const runCreated = result.events.find( (event) => event.eventType === 'run_created' ); @@ -811,7 +818,8 @@ async function createWorkflowRunEventInner( // an S3-backed hook payload the runtime would discard anyway. const outcome = await createHookReceivedPreloadEventV4( { ...input, remoteRefBehavior: 'lazy' }, - config + config, + params.onEvent ); if (outcome.kind === 'materialized') { // Older server (or optimization declined): the write still succeeded diff --git a/packages/world/src/events.ts b/packages/world/src/events.ts index 6f0c55ba6a..7dab0a9e3a 100644 --- a/packages/world/src/events.ts +++ b/packages/world/src/events.ts @@ -724,6 +724,15 @@ export type HookCreatedEventRequest = EventRequestOfType<'hook_created'>; export type HookReceivedEvent = z.infer; export type HookConflictEvent = z.infer; +/** + * Local observer invoked as each event is decoded from a streamed response. + * Worlds without a streaming implementation may ignore it. The callback runs + * synchronously, so its work intentionally applies response-stream backpressure. + * Keep it bounded to the event just decoded; large payloads may proportionally + * delay the next frame. + */ +export type EventStreamObserver = (event: Event) => void; + /** * Union of all possible event request types. * @internal Use CreateEventRequest or RunCreatedEventRequest instead. @@ -941,6 +950,11 @@ export interface CreateEventParams { * `resumeHook()` must not set it. */ preloadEvents?: true; + /** + * Observe replay-preload events as their frames are decoded. This is a + * client-side delivery hook only; it is never serialized to a backend. + */ + onEvent?: EventStreamObserver; } /** @@ -1113,6 +1127,8 @@ export interface ListEventsParams { /** Omit `limit` to return every remaining event. */ pagination?: PaginationOptions; resolveData?: ResolveData; + /** Observe events as a streaming World decodes them. */ + onEvent?: EventStreamObserver; } export interface ListEventsByCorrelationIdParams { From 6eb309b1290cb2d4b2bbd500bd6b5a528abc4c76 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Thu, 13 Aug 2026 23:53:30 -0700 Subject: [PATCH 2/7] [world-vercel] Preserve event observer failures --- packages/world-vercel/src/events-v4.test.ts | 46 +++++++++++++++++++++ packages/world-vercel/src/events-v4.ts | 9 ++-- 2 files changed, 52 insertions(+), 3 deletions(-) diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index 503e0a3475..eb3342f145 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -316,6 +316,52 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent.assertNoPendingInterceptors(); }); + it('does not resume past a GET observer failure', async () => { + const origin = + WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; + const agent = new MockAgent(); + agent.disableNetConnect(); + + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events?returnAll=true', + method: 'GET', + }) + .reply( + 200, + Buffer.concat([ + encodeFrame( + { + eventId: 'evnt_1', + runId: 'wrun_1', + eventType: 'run_started', + createdAt: CREATED_AT, + }, + new Uint8Array() + ), + encodeFrame( + { _end: 1, next: 'eid:evnt_1', hasMore: false }, + new Uint8Array() + ), + ]), + { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } + ); + + await expect( + getWorkflowRunEventsV4( + 'wrun_1', + { + onEvent: () => { + throw new Error('observer failed'); + }, + }, + { token: 'test-token', dispatcher: agent } + ) + ).rejects.toThrow('observer failed'); + agent.assertNoPendingInterceptors(); + }); + it.each([ ['an unknown event type', { eventType: 'future_event', eventData: {} }], ['invalid event metadata', { eventType: 'run_created', eventData: {} }], diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index 4fc599eb37..92eef634a3 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -1355,6 +1355,10 @@ export interface ListEventsV4Result { hasMore: boolean; } +function isTransportError(error: unknown): error is WorkflowWorldError { + return WorkflowWorldError.is(error) && error.code === 'TRANSPORT'; +} + async function consumeEventFrameStream( response: Response, opName: string, @@ -1435,9 +1439,7 @@ async function consumeReplayLogResponse({ try { page = await consumeEventFrameStream(response, opName, events, onEvent); } catch (error) { - if (!WorkflowWorldError.is(error) || error.code !== 'TRANSPORT') { - throw error; - } + if (!isTransportError(error)) throw error; const lastEvent = events.at(-1); if (!lastEvent) throw error; page = { cursor: `eid:${lastEvent.eventId}`, hasMore: true }; @@ -1546,6 +1548,7 @@ export async function getWorkflowRunEventsV4( } return { events, ...page }; } catch (error) { + if (!isTransportError(error)) throw error; const lastEvent = events.at(-1); if ( params.limit !== undefined || From fe2842b9ee01af078606bc69e82af8410192034d Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 12:13:51 -0700 Subject: [PATCH 3/7] [world-vercel] Simplify replay stream handling --- packages/world-vercel/src/events-v4.test.ts | 14 ++-- packages/world-vercel/src/events-v4.ts | 72 ++++++++++----------- packages/world-vercel/src/events.test.ts | 3 +- packages/world-vercel/src/events.ts | 54 ++++++++++------ packages/world-vercel/src/storage.ts | 4 +- packages/world/src/events.ts | 11 ++-- 6 files changed, 90 insertions(+), 68 deletions(-) diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index eb3342f145..d63fc95d7e 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -348,17 +348,20 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } ); + const observerError = new WorkflowWorldError('observer failed', { + code: 'TRANSPORT', + }); await expect( getWorkflowRunEventsV4( 'wrun_1', { onEvent: () => { - throw new Error('observer failed'); + throw observerError; }, }, { token: 'test-token', dispatcher: agent } ) - ).rejects.toThrow('observer failed'); + ).rejects.toBe(observerError); agent.assertNoPendingInterceptors(); }); @@ -1058,15 +1061,18 @@ describe('createWorkflowRunEventV4 over HTTP', () => { } ); + const observerError = new WorkflowWorldError('observer failed', { + code: 'TRANSPORT', + }); await expect( createWorkflowRunStartedEventV4( { runId: 'wrun_1', specVersion: 5 }, { token: 'test-token', dispatcher: agent }, () => { - throw new Error('observer failed'); + throw observerError; } ) - ).rejects.toThrow('observer failed'); + ).rejects.toBe(observerError); agent.assertNoPendingInterceptors(); }); diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index 92eef634a3..46b7a72e4c 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -1355,8 +1355,25 @@ export interface ListEventsV4Result { hasMore: boolean; } -function isTransportError(error: unknown): error is WorkflowWorldError { - return WorkflowWorldError.is(error) && error.code === 'TRANSPORT'; +class PartialEventStreamError extends WorkflowWorldError { + constructor(message: string, cause?: unknown) { + super(message, { code: 'TRANSPORT', cause }); + } +} + +async function* readEventFrames( + chunks: AsyncIterable, + opName: string, + events: readonly Event[] +): AsyncGenerator { + try { + yield* decodeFrames(chunks); + } catch (cause) { + throw new PartialEventStreamError( + `v4 ${opName}: event frame stream failed after ${events.length} events`, + cause + ); + } } async function consumeEventFrameStream( @@ -1373,44 +1390,25 @@ async function consumeEventFrameStream( } const chunks = response.body as unknown as AsyncIterable; - const frames = decodeFrames(chunks)[Symbol.asyncIterator](); - - try { - for (;;) { - let next: IteratorResult; - try { - next = await frames.next(); - } catch (cause) { - throw new WorkflowWorldError( - `v4 ${opName}: event frame stream failed after ${events.length} events`, - { code: 'TRANSPORT', cause } - ); - } - if (next.done) break; - const frame = next.value; - if (frame.meta._end === 1) { - const end = EventStreamEndSchema.parse(frame.meta); - return { cursor: end.next ?? null, hasMore: end.hasMore }; - } - if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { - throw new Error(`v4 ${opName}: unexpected control frame`); - } - const event = decodeEventFrame(frame); - events.push(event); - // Deliberately outside the stream-read catch above: observer/application - // failures are not truncation and must never advance recovery past this - // event. - onEvent?.(event); + for await (const frame of readEventFrames(chunks, opName, events)) { + if (frame.meta._end === 1) { + const end = EventStreamEndSchema.parse(frame.meta); + return { cursor: end.next ?? null, hasMore: end.hasMore }; + } + if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { + throw new Error(`v4 ${opName}: unexpected control frame`); } - } finally { - await frames.return?.(undefined); + const event = decodeEventFrame(frame); + events.push(event); + // Observer/application failures remain outside the frame reader, so they + // are not mistaken for truncation or skipped by partial-stream recovery. + onEvent?.(event); } - throw new WorkflowWorldError( + throw new PartialEventStreamError( `v4 ${opName}: frame stream ended without the end-of-stream sentinel ` + - `(${events.length} events read) — truncated response?`, - { code: 'TRANSPORT' } + `(${events.length} events read) — truncated response?` ); } @@ -1439,7 +1437,7 @@ async function consumeReplayLogResponse({ try { page = await consumeEventFrameStream(response, opName, events, onEvent); } catch (error) { - if (!isTransportError(error)) throw error; + if (!(error instanceof PartialEventStreamError)) throw error; const lastEvent = events.at(-1); if (!lastEvent) throw error; page = { cursor: `eid:${lastEvent.eventId}`, hasMore: true }; @@ -1548,7 +1546,7 @@ export async function getWorkflowRunEventsV4( } return { events, ...page }; } catch (error) { - if (!isTransportError(error)) throw error; + if (!(error instanceof PartialEventStreamError)) throw error; const lastEvent = events.at(-1); if ( params.limit !== undefined || diff --git a/packages/world-vercel/src/events.test.ts b/packages/world-vercel/src/events.test.ts index d072603000..95bad148ab 100644 --- a/packages/world-vercel/src/events.test.ts +++ b/packages/world-vercel/src/events.test.ts @@ -8,6 +8,7 @@ import { describe, expect, it } from 'vitest'; import { createWorkflowRunEvent, getWorkflowRunEvents, + getWorkflowRunEventsByCorrelationId, splitEventDataForV4, } from './events.js'; import { encodeFrame, V4_FRAME_CONTENT_TYPE } from './frames.js'; @@ -1520,7 +1521,7 @@ describe('getWorkflowRunEvents by correlation id is scoped to the run', () => { headers: { 'content-type': V4_FRAME_CONTENT_TYPE }, }); - const result = await getWorkflowRunEvents( + const result = await getWorkflowRunEventsByCorrelationId( { correlationId: 'step_001', runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ); diff --git a/packages/world-vercel/src/events.ts b/packages/world-vercel/src/events.ts index 8f6308a4a7..fae7ec6811 100644 --- a/packages/world-vercel/src/events.ts +++ b/packages/world-vercel/src/events.ts @@ -48,6 +48,7 @@ import { getEventDataPayloadField, isHookEventRequiringExistence, type ListEventsByCorrelationIdParams, + type ListEventsOptions, type ListEventsParams, type PaginatedResponse, validateUlidTimestamp, @@ -435,30 +436,46 @@ export async function getEvent( ); } -export async function getWorkflowRunEvents( - params: ListEventsParams | ListEventsByCorrelationIdParams, - config?: APIConfig -): Promise> { +function getListEventsV4Params(params: ListEventsOptions): ListEventsV4Params { const { pagination, resolveData = DEFAULT_RESOLVE_DATA_OPTION } = params; // `resolveData: 'none'` leaves payload refs unresolved, so the backend can // skip reading and streaming their contents. The validated lazy descriptors // remain on the returned events. - const listParams: ListEventsV4Params = { + return { ...pagination, remoteRefBehavior: resolveData === 'none' ? 'lazy' : 'resolve', - ...('onEvent' in params && params.onEvent - ? { onEvent: params.onEvent } - : {}), }; +} + +export async function getWorkflowRunEvents( + params: ListEventsParams, + config?: APIConfig +): Promise> { + const result = await getWorkflowRunEventsV4( + params.runId, + { ...getListEventsV4Params(params), onEvent: params.onEvent }, + config + ); + + return { + data: result.events, + // The cursor is present even on the final page because it is also the + // incremental-load resume point. `hasMore` is the pagination signal. + cursor: result.cursor, + hasMore: result.hasMore, + }; +} - const result = await ('correlationId' in params - ? getEventsByCorrelationIdV4( - params.correlationId, - params.runId, - listParams, - config - ) - : getWorkflowRunEventsV4(params.runId, listParams, config)); +export async function getWorkflowRunEventsByCorrelationId( + params: ListEventsByCorrelationIdParams, + config?: APIConfig +): Promise> { + const result = await getEventsByCorrelationIdV4( + params.correlationId, + params.runId, + getListEventsV4Params(params), + config + ); // A correlation id is unique per run, not globally — a slot-numbered run // numbers its own steps, so `step_…001` names the first step of every such @@ -467,10 +484,7 @@ export async function getWorkflowRunEvents( // `hasMore`/`cursor` stay the backend's, so a page that filters down to // nothing is still followed by the next one. return { - data: - 'correlationId' in params - ? result.events.filter((event) => event.runId === params.runId) - : result.events, + data: result.events.filter((event) => event.runId === params.runId), // The cursor is present even on the final page because it is also the // incremental-load resume point. `hasMore` is the pagination signal. cursor: result.cursor, diff --git a/packages/world-vercel/src/storage.ts b/packages/world-vercel/src/storage.ts index 186c759805..72abe6c57b 100644 --- a/packages/world-vercel/src/storage.ts +++ b/packages/world-vercel/src/storage.ts @@ -8,6 +8,7 @@ import { createWorkflowRunEventBatch, getEvent, getWorkflowRunEvents, + getWorkflowRunEventsByCorrelationId, } from './events.js'; import { getHook, getHookByToken, listHooks } from './hooks.js'; import { instrumentObject } from './instrumentObject.js'; @@ -53,7 +54,8 @@ export function createStorage(config?: APIConfig): Storage { createWorkflowRunEventBatch(runId, events, params, config), get: (runId, eventId, params) => getEvent(runId, eventId, params, config), list: (params) => getWorkflowRunEvents(params, config), - listByCorrelationId: (params) => getWorkflowRunEvents(params, config), + listByCorrelationId: (params) => + getWorkflowRunEventsByCorrelationId(params, config), }, hooks: { get: (hookId, params) => getHook(hookId, params, config), diff --git a/packages/world/src/events.ts b/packages/world/src/events.ts index 7dab0a9e3a..511e27b45b 100644 --- a/packages/world/src/events.ts +++ b/packages/world/src/events.ts @@ -1122,16 +1122,19 @@ export interface GetEventParams { resolveData?: ResolveData; } -export interface ListEventsParams { - runId: string; +export interface ListEventsOptions { /** Omit `limit` to return every remaining event. */ pagination?: PaginationOptions; resolveData?: ResolveData; +} + +export interface ListEventsParams extends ListEventsOptions { + runId: string; /** Observe events as a streaming World decodes them. */ onEvent?: EventStreamObserver; } -export interface ListEventsByCorrelationIdParams { +export interface ListEventsByCorrelationIdParams extends ListEventsOptions { correlationId: string; /** * The run the correlation id belongs to. A correlation id is unique per @@ -1142,6 +1145,4 @@ export interface ListEventsByCorrelationIdParams { * event id alone is not. */ runId: string; - pagination?: PaginationOptions; - resolveData?: ResolveData; } From 780e375846ed23873c91cfc5495ed046917ddba7 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 15:36:32 -0700 Subject: [PATCH 4/7] [world-vercel] Reuse world event list types --- packages/world-vercel/src/events-v4.test.ts | 65 +++++----- packages/world-vercel/src/events-v4.ts | 120 ++++++++---------- packages/world-vercel/src/events.test.ts | 2 +- packages/world-vercel/src/events.ts | 42 +----- .../src/trace-propagation.test.ts | 5 +- packages/world/src/events.ts | 30 ++--- 6 files changed, 103 insertions(+), 161 deletions(-) diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index d63fc95d7e..a40526ca4f 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -294,7 +294,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply(200, frames, { @@ -302,13 +302,12 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); const result = await getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ); - expect(result.events).toHaveLength(1); - expect(result.events[0]).toMatchObject({ + expect(result.data).toHaveLength(1); + expect(result.data[0]).toMatchObject({ eventId: 'evnt_1', eventData: { input: body }, }); @@ -325,7 +324,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( @@ -353,8 +352,8 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); await expect( getWorkflowRunEventsV4( - 'wrun_1', { + runId: 'wrun_1', onEvent: () => { throw observerError; }, @@ -377,7 +376,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( @@ -391,8 +390,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { await expect( getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ) ).rejects.toThrow(); @@ -431,7 +429,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply(200, frames, { @@ -439,8 +437,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); const result = await getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ); @@ -462,7 +459,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply(200, frames, { @@ -471,8 +468,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { await expect( getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ) ).rejects.toThrow(); @@ -505,7 +501,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?limit=500', + path: '/api/v4/runs/wrun_1/events?limit=500&remoteRefBehavior=resolve', method: 'GET', }) .reply(200, frames, { @@ -514,8 +510,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { await expect( getWorkflowRunEventsV4( - 'wrun_1', - { limit: 500 }, + { runId: 'wrun_1', pagination: { limit: 500 } }, { token: 'test-token', dispatcher: agent } ) ).rejects.toThrow(/end-of-stream sentinel/); @@ -530,7 +525,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( @@ -554,7 +549,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_1', + path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( @@ -578,12 +573,11 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { ); const result = await getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ); - expect(result.events.map((event) => event.eventId)).toEqual([ + expect(result.data.map((event) => event.eventId)).toEqual([ 'evnt_1', 'evnt_2', ]); @@ -638,9 +632,11 @@ describe('getEventsByCorrelationIdV4 over HTTP', () => { }); const result = await getEventsByCorrelationIdV4( - 'step_001', - 'wrun_1', - { limit: 10 }, + { + correlationId: 'step_001', + runId: 'wrun_1', + pagination: { limit: 10 }, + }, { token: 'test-token', dispatcher: agent } ); @@ -652,8 +648,8 @@ describe('getEventsByCorrelationIdV4 over HTTP', () => { expect(query.get('limit')).toBe('10'); } - expect(result.events).toHaveLength(1); - expect(result.events[0].runId).toBe('wrun_1'); + expect(result.data).toHaveLength(1); + expect(result.data[0].runId).toBe('wrun_1'); agent.assertNoPendingInterceptors(); }); }); @@ -735,7 +731,7 @@ describe('v4 transport uses global fetch (observability)', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply(200, encodeFrame({ _end: 1, hasMore: false }, new Uint8Array(0)), { @@ -747,8 +743,7 @@ describe('v4 transport uses global fetch (observability)', () => { const fetchSpy = vi.spyOn(globalThis, 'fetch'); await getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ); @@ -980,7 +975,7 @@ describe('createWorkflowRunEventV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_1', + path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( @@ -1121,7 +1116,7 @@ describe('createWorkflowRunEventV4 over HTTP', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_1', + path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( @@ -1804,7 +1799,7 @@ describe('v4 transport reports failures to the events recycler', () => { for (let i = 0; i < EVENTS_RECYCLE_AFTER_CONSECUTIVE_FAILURES; i++) { await expect( - getWorkflowRunEventsV4('wrun_1', {}, { token: 'test-token' }) + getWorkflowRunEventsV4({ runId: 'wrun_1' }, { token: 'test-token' }) ).rejects.toThrow(); // Still the same pool until the threshold is reached. if (i < EVENTS_RECYCLE_AFTER_CONSECUTIVE_FAILURES - 1) { diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index 46b7a72e4c..13c7d2ecb7 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -27,12 +27,13 @@ import { type Event, type EventResult, EventSchema, - type EventStreamObserver, type EventType, EventTypeSchema, getEventDataPayloadField, HookSchema, - type PaginationOptions, + type ListEventsByCorrelationIdParams, + type ListEventsParams, + type PaginatedResponse, StructuredErrorSchema, WaitSchema, WorkflowRunSchema, @@ -781,7 +782,7 @@ async function decodeCreateEventResponse( export async function createWorkflowRunStartedEventV4( input: CreateEventV4InputBase, config?: APIConfig, - onEvent?: EventStreamObserver + onEvent?: (event: Event) => void ) { const response = await postWorkflowRunEventV4( { ...input, eventType: 'run_started' }, @@ -1205,8 +1206,11 @@ async function postEventFrameOverWs( */ export type HookReceivedPreloadV4Result = /** The server streamed the replay log back as v4 frames. */ - | (ListEventsV4Result & { + | { kind: 'stream'; + events: Event[]; + cursor: string | null; + hasMore: boolean; /** * The canonical event this write created or converged on (the resume * claim winner's — ours or the producer's), named by the @@ -1215,7 +1219,7 @@ export type HookReceivedPreloadV4Result = canonicalEventId: string | undefined; /** Per-run event ceiling from the response header, when present. */ maxEvents: number | undefined; - }) + } /** * The server answered with the normal materialized CBOR body instead — * an older server, or one that declined the optimization. The @@ -1241,7 +1245,7 @@ export type HookReceivedPreloadV4Result = export async function createHookReceivedPreloadEventV4( input: CreateEventV4InputBase, config?: APIConfig, - onEvent?: EventStreamObserver + onEvent?: (event: Event) => void ): Promise { const response = await postWorkflowRunEventV4( { ...input, eventType: 'hook_received' }, @@ -1334,27 +1338,6 @@ export async function getEventV4( throw new Error(`v4 getEvent: empty frame stream for ${eventId}`); } -export interface ListEventsV4Params extends PaginationOptions { - /** - * Whether the backend resolves payload bytes into each frame body. - * `resolve` (default) streams the bytes; `lazy` emits empty-body frames - * (the ref descriptor stays in the frame meta) — for metadata-only - * listings that would otherwise download every payload just to discard - * it. - */ - remoteRefBehavior?: 'resolve' | 'lazy'; - /** Called synchronously after each event frame validates. */ - onEvent?: EventStreamObserver; -} - -export interface ListEventsV4Result { - events: Event[]; - /** Trailing event-log cursor, or null when the stream contained no events. */ - cursor: string | null; - /** Explicit "another page of results exists" flag from the sentinel. */ - hasMore: boolean; -} - class PartialEventStreamError extends WorkflowWorldError { constructor(message: string, cause?: unknown) { super(message, { code: 'TRANSPORT', cause }); @@ -1380,8 +1363,8 @@ async function consumeEventFrameStream( response: Response, opName: string, events: Event[], - onEvent?: EventStreamObserver -): Promise> { + onEvent?: (event: Event) => void +): Promise<{ cursor: string | null; hasMore: boolean }> { const contentType = response.headers.get('content-type'); if (!contentType?.startsWith(V4_FRAME_CONTENT_TYPE)) { throw new Error( @@ -1431,9 +1414,9 @@ async function consumeReplayLogResponse({ opName: string; events: Event[]; config?: APIConfig; - onEvent?: EventStreamObserver; -}): Promise> { - let page: Pick; + onEvent?: (event: Event) => void; +}): Promise<{ cursor: string | null; hasMore: boolean }> { + let page: { cursor: string | null; hasMore: boolean }; try { page = await consumeEventFrameStream(response, opName, events, onEvent); } catch (error) { @@ -1446,11 +1429,14 @@ async function consumeReplayLogResponse({ if (!page.hasMore) return page; assert(page.cursor, `v4 ${opName}: partial event stream missing cursor`); const suffix = await getWorkflowRunEventsV4( - runId, - { cursor: page.cursor, onEvent }, + { + runId, + pagination: { cursor: page.cursor }, + onEvent, + }, config ); - events.push(...suffix.events); + events.push(...suffix.data); return { cursor: suffix.cursor, hasMore: suffix.hasMore }; } @@ -1469,8 +1455,8 @@ async function consumeListFrameStream( config: APIConfig | undefined, opName: string, events: Event[], - onEvent?: EventStreamObserver -): Promise> { + onEvent?: (event: Event) => void +): Promise<{ cursor: string | null; hasMore: boolean }> { const response = await fetchV4( url, { method: 'GET', headers }, @@ -1485,20 +1471,26 @@ async function consumeListFrameStream( * Shared by the runId and correlationId list query builders so both send * `remoteRefBehavior` identically. */ -function appendListParams(sp: URLSearchParams, params: ListEventsV4Params) { - if (params.cursor) sp.set('cursor', params.cursor); - if (params.limit !== undefined) sp.set('limit', String(params.limit)); - if (params.sortOrder) sp.set('sortOrder', params.sortOrder); - if (params.remoteRefBehavior) { - sp.set('remoteRefBehavior', params.remoteRefBehavior); - } +function appendListParams( + sp: URLSearchParams, + params: ListEventsParams | ListEventsByCorrelationIdParams, + cursor = params.pagination?.cursor +) { + const { limit, sortOrder } = params.pagination ?? {}; + if (cursor) sp.set('cursor', cursor); + if (limit !== undefined) sp.set('limit', String(limit)); + if (sortOrder) sp.set('sortOrder', sortOrder); + sp.set( + 'remoteRefBehavior', + params.resolveData === 'none' ? 'lazy' : 'resolve' + ); } -function paginationToQuery(params: ListEventsV4Params): string { +function paginationToQuery(params: ListEventsParams, cursor?: string): string { const sp = new URLSearchParams(); // The World API uses an omitted limit for a complete event log. - if (params.limit === undefined) sp.set('returnAll', 'true'); - appendListParams(sp, params); + if (params.pagination?.limit === undefined) sp.set('returnAll', 'true'); + appendListParams(sp, params, cursor); return `?${sp.toString()}`; } @@ -1513,18 +1505,17 @@ function paginationToQuery(params: ListEventsV4Params): string { * after its last validated event instead of downloading accepted frames again. */ export async function getWorkflowRunEventsV4( - runId: string, - params: ListEventsV4Params = {}, + params: ListEventsParams, config?: APIConfig -): Promise { +): Promise> { const { baseUrl, headers } = await getHttpConfig(config); const events: Event[] = []; - let cursor = params.cursor; + let cursor = params.pagination?.cursor; while (true) { const url = - `${baseUrl}/v4/runs/${encodeURIComponent(runId)}/events` + - paginationToQuery({ ...params, cursor }); + `${baseUrl}/v4/runs/${encodeURIComponent(params.runId)}/events` + + paginationToQuery(params, cursor); try { const page = await consumeListFrameStream( url, @@ -1534,22 +1525,22 @@ export async function getWorkflowRunEventsV4( events, params.onEvent ); - if (params.limit === undefined && page.hasMore) { + if (params.pagination?.limit === undefined && page.hasMore) { if (!page.cursor || page.cursor === cursor) { throw new WorkflowWorldError( - `v4 listEvents: partial event stream made no cursor progress for run ${runId}`, + `v4 listEvents: partial event stream made no cursor progress for run ${params.runId}`, { code: 'SCHEMA_VALIDATION' } ); } cursor = page.cursor; continue; } - return { events, ...page }; + return { data: events, ...page }; } catch (error) { if (!(error instanceof PartialEventStreamError)) throw error; const lastEvent = events.at(-1); if ( - params.limit !== undefined || + params.pagination?.limit !== undefined || !lastEvent || `eid:${lastEvent.eventId}` === cursor ) { @@ -1575,15 +1566,13 @@ export async function getWorkflowRunEventsV4( * the page by run id. */ export async function getEventsByCorrelationIdV4( - correlationId: string, - runId: string, - params: ListEventsV4Params = {}, + params: ListEventsByCorrelationIdParams, config?: APIConfig -): Promise { +): Promise> { const { baseUrl, headers } = await getHttpConfig(config); const sp = new URLSearchParams(); - sp.set('correlationId', correlationId); - sp.set('runId', runId); + sp.set('correlationId', params.correlationId); + sp.set('runId', params.runId); appendListParams(sp, params); const url = `${baseUrl}/v4/events?${sp.toString()}`; const events: Event[] = []; @@ -1592,8 +1581,7 @@ export async function getEventsByCorrelationIdV4( headers, config, 'listEventsByCorrelationId', - events, - params.onEvent + events ); - return { events, ...page }; + return { data: events, ...page }; } diff --git a/packages/world-vercel/src/events.test.ts b/packages/world-vercel/src/events.test.ts index 95bad148ab..76d0c6896d 100644 --- a/packages/world-vercel/src/events.test.ts +++ b/packages/world-vercel/src/events.test.ts @@ -1887,7 +1887,7 @@ describe('createWorkflowRunEvent hook_received replay preload', () => { agent .get(ORIGIN) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true&cursor=eid%3Aevnt_4', + path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_4&remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply( diff --git a/packages/world-vercel/src/events.ts b/packages/world-vercel/src/events.ts index fae7ec6811..cb7d170a3d 100644 --- a/packages/world-vercel/src/events.ts +++ b/packages/world-vercel/src/events.ts @@ -48,7 +48,6 @@ import { getEventDataPayloadField, isHookEventRequiringExistence, type ListEventsByCorrelationIdParams, - type ListEventsOptions, type ListEventsParams, type PaginatedResponse, validateUlidTimestamp, @@ -63,15 +62,10 @@ import { getEventsByCorrelationIdV4, getEventV4, getWorkflowRunEventsV4, - type ListEventsV4Params, } from './events-v4.js'; import { decode as decodeRunId } from './run-id/index.js'; import { cancelWorkflowRunV1, createWorkflowRunV1 } from './runs.js'; -import { - type APIConfig, - DEFAULT_RESOLVE_DATA_OPTION, - makeRequest, -} from './utils.js'; +import { type APIConfig, makeRequest } from './utils.js'; function validateWorkflowRunIdTimestamp(id: string): string | null { const raw = id.startsWith('wrun_') ? id.slice('wrun_'.length) : id; @@ -436,46 +430,18 @@ export async function getEvent( ); } -function getListEventsV4Params(params: ListEventsOptions): ListEventsV4Params { - const { pagination, resolveData = DEFAULT_RESOLVE_DATA_OPTION } = params; - // `resolveData: 'none'` leaves payload refs unresolved, so the backend can - // skip reading and streaming their contents. The validated lazy descriptors - // remain on the returned events. - return { - ...pagination, - remoteRefBehavior: resolveData === 'none' ? 'lazy' : 'resolve', - }; -} - export async function getWorkflowRunEvents( params: ListEventsParams, config?: APIConfig ): Promise> { - const result = await getWorkflowRunEventsV4( - params.runId, - { ...getListEventsV4Params(params), onEvent: params.onEvent }, - config - ); - - return { - data: result.events, - // The cursor is present even on the final page because it is also the - // incremental-load resume point. `hasMore` is the pagination signal. - cursor: result.cursor, - hasMore: result.hasMore, - }; + return getWorkflowRunEventsV4(params, config); } export async function getWorkflowRunEventsByCorrelationId( params: ListEventsByCorrelationIdParams, config?: APIConfig ): Promise> { - const result = await getEventsByCorrelationIdV4( - params.correlationId, - params.runId, - getListEventsV4Params(params), - config - ); + const result = await getEventsByCorrelationIdV4(params, config); // A correlation id is unique per run, not globally — a slot-numbered run // numbers its own steps, so `step_…001` names the first step of every such @@ -484,7 +450,7 @@ export async function getWorkflowRunEventsByCorrelationId( // `hasMore`/`cursor` stay the backend's, so a page that filters down to // nothing is still followed by the next one. return { - data: result.events.filter((event) => event.runId === params.runId), + data: result.data.filter((event) => event.runId === params.runId), // The cursor is present even on the final page because it is also the // incremental-load resume point. `hasMore` is the pagination signal. cursor: result.cursor, diff --git a/packages/world-vercel/src/trace-propagation.test.ts b/packages/world-vercel/src/trace-propagation.test.ts index b22b2d0ec6..a8155c6f5d 100644 --- a/packages/world-vercel/src/trace-propagation.test.ts +++ b/packages/world-vercel/src/trace-propagation.test.ts @@ -166,7 +166,7 @@ describe('v4 event requests (fetchV4) trace propagation', () => { agent .get(origin) .intercept({ - path: '/api/v4/runs/wrun_1/events?returnAll=true', + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', }) .reply(200, encodeFrame({ _end: 1, hasMore: false }, new Uint8Array(0)), { @@ -184,8 +184,7 @@ describe('v4 event requests (fetchV4) trace propagation', () => { traceId = span.spanContext().traceId; spanId = span.spanContext().spanId; await getWorkflowRunEventsV4( - 'wrun_1', - {}, + { runId: 'wrun_1' }, { token: 'test-token', dispatcher: agent } ); span.end(); diff --git a/packages/world/src/events.ts b/packages/world/src/events.ts index 511e27b45b..9ac9226fe0 100644 --- a/packages/world/src/events.ts +++ b/packages/world/src/events.ts @@ -724,15 +724,6 @@ export type HookCreatedEventRequest = EventRequestOfType<'hook_created'>; export type HookReceivedEvent = z.infer; export type HookConflictEvent = z.infer; -/** - * Local observer invoked as each event is decoded from a streamed response. - * Worlds without a streaming implementation may ignore it. The callback runs - * synchronously, so its work intentionally applies response-stream backpressure. - * Keep it bounded to the event just decoded; large payloads may proportionally - * delay the next frame. - */ -export type EventStreamObserver = (event: Event) => void; - /** * Union of all possible event request types. * @internal Use CreateEventRequest or RunCreatedEventRequest instead. @@ -954,7 +945,7 @@ export interface CreateEventParams { * Observe replay-preload events as their frames are decoded. This is a * client-side delivery hook only; it is never serialized to a backend. */ - onEvent?: EventStreamObserver; + onEvent?: (event: Event) => void; } /** @@ -1122,19 +1113,19 @@ export interface GetEventParams { resolveData?: ResolveData; } -export interface ListEventsOptions { +export interface ListEventsParams { + runId: string; /** Omit `limit` to return every remaining event. */ pagination?: PaginationOptions; resolveData?: ResolveData; + /** + * Observe events as a streaming World decodes them. The callback runs + * synchronously and therefore applies response-stream backpressure. + */ + onEvent?: (event: Event) => void; } -export interface ListEventsParams extends ListEventsOptions { - runId: string; - /** Observe events as a streaming World decodes them. */ - onEvent?: EventStreamObserver; -} - -export interface ListEventsByCorrelationIdParams extends ListEventsOptions { +export interface ListEventsByCorrelationIdParams { correlationId: string; /** * The run the correlation id belongs to. A correlation id is unique per @@ -1145,4 +1136,7 @@ export interface ListEventsByCorrelationIdParams extends ListEventsOptions { * event id alone is not. */ runId: string; + /** Omit `limit` to return every remaining event. */ + pagination?: PaginationOptions; + resolveData?: ResolveData; } From 12cacdef903fee1a4aff1a6b0001bb2bfcfa49a8 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Fri, 14 Aug 2026 21:05:22 -0700 Subject: [PATCH 5/7] fix(world-vercel): preserve replay recovery cursor --- packages/world-vercel/src/events-v4.test.ts | 71 +++++++++++++++++++++ packages/world-vercel/src/events-v4.ts | 2 +- 2 files changed, 72 insertions(+), 1 deletion(-) diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index a40526ca4f..bdd5bc82b0 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -1013,6 +1013,77 @@ describe('createWorkflowRunEventV4 over HTTP', () => { agent.assertNoPendingInterceptors(); }); + it('preserves the POST cursor when truncation recovery returns an empty suffix', async () => { + const origin = + WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; + const agent = new MockAgent(); + agent.disableNetConnect(); + + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events/run_started', + method: 'POST', + headers: { accept: V4_FRAME_CONTENT_TYPE }, + }) + .reply( + 200, + Buffer.concat([ + encodeFrame( + { + eventId: 'evnt_1', + runId: 'wrun_1', + eventType: 'run_created', + createdAt: CREATED_AT, + eventData: { + deploymentId: 'dpl_1', + workflowName: 'workflow', + input: null, + }, + }, + new Uint8Array() + ), + encodeFrame( + { + eventId: 'evnt_2', + runId: 'wrun_1', + eventType: 'run_started', + createdAt: CREATED_AT, + }, + new Uint8Array() + ), + ]), + { + headers: { + 'content-type': V4_FRAME_CONTENT_TYPE, + 'x-wf-max-events': '10000', + }, + } + ); + agent + .get(origin) + .intercept({ + path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_2&remoteRefBehavior=resolve&returnAll=true', + method: 'GET', + }) + .reply(200, encodeFrame({ _end: 1, hasMore: false }, new Uint8Array()), { + headers: { 'content-type': V4_FRAME_CONTENT_TYPE }, + }); + + const result = await createWorkflowRunStartedEventV4( + { runId: 'wrun_1', specVersion: 5 }, + { token: 'test-token', dispatcher: agent } + ); + + expect(result.events.map((event) => event.eventId)).toEqual([ + 'evnt_1', + 'evnt_2', + ]); + expect(result.cursor).toBe('eid:evnt_2'); + expect(result.hasMore).toBe(false); + agent.assertNoPendingInterceptors(); + }); + it('does not treat an event observer failure as stream truncation', async () => { const origin = WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index 13c7d2ecb7..b9d812fa5c 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -1437,7 +1437,7 @@ async function consumeReplayLogResponse({ config ); events.push(...suffix.data); - return { cursor: suffix.cursor, hasMore: suffix.hasMore }; + return { cursor: suffix.cursor ?? page.cursor, hasMore: suffix.hasMore }; } /** From 0fc449abe429a685cb9fc3bf07be78fa62c33e78 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Mon, 17 Aug 2026 10:47:23 -0700 Subject: [PATCH 6/7] refactor(world-vercel): simplify replay frame recovery --- packages/world-vercel/src/events-v4.ts | 182 +++++++++++-------------- 1 file changed, 81 insertions(+), 101 deletions(-) diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index b9d812fa5c..2b64833217 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -789,16 +789,12 @@ export async function createWorkflowRunStartedEventV4( 'event-stream', config ); - const events: Event[] = []; - const page = await consumeReplayLogResponse({ + const replay = await consumeReplayLogResponse( response, - runId: input.runId, - opName: 'createEvent', - events, - config, - onEvent, - }); - assert(page.cursor, 'v4 createEvent: event stream missing cursor'); + { runId: input.runId, onEvent }, + config + ); + assert(replay.cursor, 'v4 createEvent: event stream missing cursor'); const maxEvents = MaxEventsHeaderSchema.safeParse( response.headers.get(MAX_EVENTS_HEADER) ); @@ -809,7 +805,7 @@ export async function createWorkflowRunStartedEventV4( }); } - return { events, ...page, maxEvents: maxEvents.data }; + return { ...replay, maxEvents: maxEvents.data }; } /** One event of a v4 batch POST, index-aligned with the response results. */ @@ -1261,22 +1257,17 @@ export async function createHookReceivedPreloadEventV4( }; } - const events: Event[] = []; - const page = await consumeReplayLogResponse({ + const replay = await consumeReplayLogResponse( response, - runId: input.runId, - opName: 'createEvent', - events, - config, - onEvent, - }); + { runId: input.runId, onEvent }, + config + ); const maxEvents = MaxEventsHeaderSchema.safeParse( response.headers.get(MAX_EVENTS_HEADER) ); return { kind: 'stream', - events, - ...page, + ...replay, canonicalEventId: response.headers.get(EVENT_ID_HEADER) ?? undefined, maxEvents: maxEvents.success ? maxEvents.data : undefined, }; @@ -1344,21 +1335,6 @@ class PartialEventStreamError extends WorkflowWorldError { } } -async function* readEventFrames( - chunks: AsyncIterable, - opName: string, - events: readonly Event[] -): AsyncGenerator { - try { - yield* decodeFrames(chunks); - } catch (cause) { - throw new PartialEventStreamError( - `v4 ${opName}: event frame stream failed after ${events.length} events`, - cause - ); - } -} - async function consumeEventFrameStream( response: Response, opName: string, @@ -1373,20 +1349,35 @@ async function consumeEventFrameStream( } const chunks = response.body as unknown as AsyncIterable; + const frames = decodeFrames(chunks); - for await (const frame of readEventFrames(chunks, opName, events)) { - if (frame.meta._end === 1) { - const end = EventStreamEndSchema.parse(frame.meta); - return { cursor: end.next ?? null, hasMore: end.hasMore }; - } - if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { - throw new Error(`v4 ${opName}: unexpected control frame`); + try { + while (true) { + let next: IteratorResult; + try { + next = await frames.next(); + } catch (cause) { + throw new PartialEventStreamError( + `v4 ${opName}: event frame stream failed after ${events.length} events`, + cause + ); + } + if (next.done) break; + + const frame = next.value; + if (frame.meta._end === 1) { + const end = EventStreamEndSchema.parse(frame.meta); + return { cursor: end.next ?? null, hasMore: end.hasMore }; + } + if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { + throw new Error(`v4 ${opName}: unexpected control frame`); + } + const event = decodeEventFrame(frame); + events.push(event); + onEvent?.(event); } - const event = decodeEventFrame(frame); - events.push(event); - // Observer/application failures remain outside the frame reader, so they - // are not mistaken for truncation or skipped by partial-stream recovery. - onEvent?.(event); + } finally { + void frames.return(undefined); } throw new PartialEventStreamError( @@ -1401,24 +1392,24 @@ async function consumeEventFrameStream( * continuation instead of making run_started download its accepted prefix * again. */ -async function consumeReplayLogResponse({ - response, - runId, - opName, - events, - config, - onEvent, -}: { - response: Response; - runId: string; - opName: string; +async function consumeReplayLogResponse( + response: Response, + params: Pick, + config?: APIConfig +): Promise<{ events: Event[]; - config?: APIConfig; - onEvent?: (event: Event) => void; -}): Promise<{ cursor: string | null; hasMore: boolean }> { + cursor: string | null; + hasMore: boolean; +}> { + const events: Event[] = []; let page: { cursor: string | null; hasMore: boolean }; try { - page = await consumeEventFrameStream(response, opName, events, onEvent); + page = await consumeEventFrameStream( + response, + 'createEvent', + events, + params.onEvent + ); } catch (error) { if (!(error instanceof PartialEventStreamError)) throw error; const lastEvent = events.at(-1); @@ -1426,44 +1417,22 @@ async function consumeReplayLogResponse({ page = { cursor: `eid:${lastEvent.eventId}`, hasMore: true }; } - if (!page.hasMore) return page; - assert(page.cursor, `v4 ${opName}: partial event stream missing cursor`); + if (!page.hasMore) return { events, ...page }; + assert(page.cursor, 'v4 createEvent: partial event stream missing cursor'); const suffix = await getWorkflowRunEventsV4( { - runId, + runId: params.runId, pagination: { cursor: page.cursor }, - onEvent, + onEvent: params.onEvent, }, config ); events.push(...suffix.data); - return { cursor: suffix.cursor ?? page.cursor, hasMore: suffix.hasMore }; -} - -/** - * Drive a v4 frame-stream list response into an in-memory page. Used by - * both the by-runId and by-correlationId list endpoints — the wire - * shape is identical, only the URL differs. - * - * `headers` come from the caller's single getHttpConfig resolution (the - * same call that produced the baseUrl in `url`) so each LIST resolves - * auth exactly once. - */ -async function consumeListFrameStream( - url: string, - headers: Headers, - config: APIConfig | undefined, - opName: string, - events: Event[], - onEvent?: (event: Event) => void -): Promise<{ cursor: string | null; hasMore: boolean }> { - const response = await fetchV4( - url, - { method: 'GET', headers }, - config, - opName - ); - return consumeEventFrameStream(response, opName, events, onEvent); + return { + events, + cursor: suffix.cursor ?? page.cursor, + hasMore: suffix.hasMore, + }; } /** @@ -1474,7 +1443,7 @@ async function consumeListFrameStream( function appendListParams( sp: URLSearchParams, params: ListEventsParams | ListEventsByCorrelationIdParams, - cursor = params.pagination?.cursor + cursor: string | null ) { const { limit, sortOrder } = params.pagination ?? {}; if (cursor) sp.set('cursor', cursor); @@ -1486,7 +1455,10 @@ function appendListParams( ); } -function paginationToQuery(params: ListEventsParams, cursor?: string): string { +function paginationToQuery( + params: ListEventsParams, + cursor: string | null +): string { const sp = new URLSearchParams(); // The World API uses an omitted limit for a complete event log. if (params.pagination?.limit === undefined) sp.set('returnAll', 'true'); @@ -1510,17 +1482,21 @@ export async function getWorkflowRunEventsV4( ): Promise> { const { baseUrl, headers } = await getHttpConfig(config); const events: Event[] = []; - let cursor = params.pagination?.cursor; + let cursor = params.pagination?.cursor ?? null; while (true) { const url = `${baseUrl}/v4/runs/${encodeURIComponent(params.runId)}/events` + paginationToQuery(params, cursor); try { - const page = await consumeListFrameStream( + const response = await fetchV4( url, - headers, + { method: 'GET', headers }, config, + 'listEvents' + ); + const page = await consumeEventFrameStream( + response, 'listEvents', events, params.onEvent @@ -1573,13 +1549,17 @@ export async function getEventsByCorrelationIdV4( const sp = new URLSearchParams(); sp.set('correlationId', params.correlationId); sp.set('runId', params.runId); - appendListParams(sp, params); + appendListParams(sp, params, params.pagination?.cursor ?? null); const url = `${baseUrl}/v4/events?${sp.toString()}`; const events: Event[] = []; - const page = await consumeListFrameStream( + const response = await fetchV4( url, - headers, + { method: 'GET', headers }, config, + 'listEventsByCorrelationId' + ); + const page = await consumeEventFrameStream( + response, 'listEventsByCorrelationId', events ); From 8eb05f8a0a8e5c2042a4074264e2f1a64298ca33 Mon Sep 17 00:00:00 2001 From: Nathan Colosimo <110621881+NathanColosimo@users.noreply.github.com> Date: Tue, 18 Aug 2026 14:11:36 -0700 Subject: [PATCH 7/7] fix(world-vercel): bound replay stream recovery --- packages/world-vercel/src/events-v4.test.ts | 231 ++++++++++---------- packages/world-vercel/src/events-v4.ts | 225 ++++++++++++------- packages/world-vercel/src/events.test.ts | 41 +++- packages/world-vercel/src/events.ts | 51 +++-- packages/world-vercel/src/http-core.ts | 7 +- 5 files changed, 328 insertions(+), 227 deletions(-) diff --git a/packages/world-vercel/src/events-v4.test.ts b/packages/world-vercel/src/events-v4.test.ts index bdd5bc82b0..a9e966efe1 100644 --- a/packages/world-vercel/src/events-v4.test.ts +++ b/packages/world-vercel/src/events-v4.test.ts @@ -29,6 +29,13 @@ import { import { WORKFLOW_SERVER_URL_OVERRIDE } from './utils.js'; const CREATED_AT = '2026-06-10T00:00:00.000Z'; +const ORIGIN = WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; + +function mockAgent() { + const agent = new MockAgent(); + agent.disableNetConnect(); + return agent; +} function createEventBody( event: AnyEventRequest, @@ -264,10 +271,7 @@ describe('throwForErrorResponse', () => { */ describe('getWorkflowRunEventsV4 over HTTP', () => { it('parses a frame stream fetched via a custom dispatcher', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); const body = new TextEncoder().encode('payload-bytes'); const frames = Buffer.concat([ @@ -292,7 +296,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { ]); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -316,13 +320,10 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); it('does not resume past a GET observer failure', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -368,13 +369,10 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { ['an unknown event type', { eventType: 'future_event', eventData: {} }], ['invalid event metadata', { eventType: 'run_created', eventData: {} }], ])('rejects %s', async (_description, meta) => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -398,10 +396,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); it('captures an explicit hasMore from the sentinel, independent of next', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); // The regression shape: a final page still carries a trailing `next` // cursor (incremental-load resume point) but hasMore is false. @@ -427,7 +422,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { ]); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -446,10 +441,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); it('rejects an end frame without hasMore', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); const frames = encodeFrame( { _end: 1, next: 'cursor-2' }, @@ -457,7 +449,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { ); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -475,10 +467,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); it('throws when the stream ends without the end sentinel (truncated response)', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); // A complete event frame but NO `{_end: 1}` sentinel — what a response // truncated on a frame boundary looks like. Returning this as a @@ -499,7 +488,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { ); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?limit=500&remoteRefBehavior=resolve', method: 'GET', @@ -517,13 +506,10 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { }); it('resumes a truncated full stream after its last accepted event', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -547,7 +533,7 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } ); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -585,6 +571,74 @@ describe('getWorkflowRunEventsV4 over HTTP', () => { expect(result.hasMore).toBe(false); agent.assertNoPendingInterceptors(); }); + + it('stops after three partial-stream recovery retries', async () => { + const agent = mockAgent(); + const paths = [ + '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', + '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', + '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_2&remoteRefBehavior=resolve&returnAll=true', + '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_3&remoteRefBehavior=resolve&returnAll=true', + ]; + + for (const [index, path] of paths.entries()) { + agent + .get(ORIGIN) + .intercept({ path, method: 'GET' }) + .reply( + 200, + encodeFrame( + { + eventId: `evnt_${index + 1}`, + runId: 'wrun_1', + eventType: 'run_created', + createdAt: CREATED_AT, + eventData: { + deploymentId: 'dpl_1', + workflowName: 'workflow', + input: null, + }, + }, + new Uint8Array() + ), + { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } + ); + } + + await expect( + getWorkflowRunEventsV4( + { runId: 'wrun_1' }, + { token: 'test-token', dispatcher: agent } + ) + ).rejects.toThrow(/end-of-stream sentinel/); + agent.assertNoPendingInterceptors(); + }); + + it('rejects an unexpected clean pagination response', async () => { + const agent = mockAgent(); + agent + .get(ORIGIN) + .intercept({ + path: '/api/v4/runs/wrun_1/events?remoteRefBehavior=resolve&returnAll=true', + method: 'GET', + }) + .reply( + 200, + encodeFrame( + { _end: 1, next: 'eid:evnt_1', hasMore: true }, + new Uint8Array() + ), + { headers: { 'content-type': V4_FRAME_CONTENT_TYPE } } + ); + + await expect( + getWorkflowRunEventsV4( + { runId: 'wrun_1' }, + { token: 'test-token', dispatcher: agent } + ) + ).rejects.toThrow(/returnAll response was unexpectedly paginated/); + agent.assertNoPendingInterceptors(); + }); }); /** @@ -936,14 +990,11 @@ describe('createWorkflowRunEventV4 over HTTP', () => { }); it('continues a truncated run_started stream after its last event', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); const observed: string[] = []; agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events/run_started', method: 'POST', @@ -973,7 +1024,7 @@ describe('createWorkflowRunEventV4 over HTTP', () => { } ); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -1014,13 +1065,10 @@ describe('createWorkflowRunEventV4 over HTTP', () => { }); it('preserves the POST cursor when truncation recovery returns an empty suffix', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events/run_started', method: 'POST', @@ -1061,7 +1109,7 @@ describe('createWorkflowRunEventV4 over HTTP', () => { } ); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_2&remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -1084,72 +1132,11 @@ describe('createWorkflowRunEventV4 over HTTP', () => { agent.assertNoPendingInterceptors(); }); - it('does not treat an event observer failure as stream truncation', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); - - agent - .get(origin) - .intercept({ - path: '/api/v4/runs/wrun_1/events/run_started', - method: 'POST', - headers: { accept: V4_FRAME_CONTENT_TYPE }, - }) - .reply( - 200, - Buffer.concat([ - encodeFrame( - { - eventId: 'evnt_1', - runId: 'wrun_1', - eventType: 'run_created', - createdAt: CREATED_AT, - eventData: { - deploymentId: 'dpl_1', - workflowName: 'workflow', - input: null, - }, - }, - new Uint8Array() - ), - encodeFrame( - { _end: 1, next: 'eid:evnt_1', hasMore: false }, - new Uint8Array() - ), - ]), - { - headers: { - 'content-type': V4_FRAME_CONTENT_TYPE, - 'x-wf-max-events': '10000', - }, - } - ); - - const observerError = new WorkflowWorldError('observer failed', { - code: 'TRANSPORT', - }); - await expect( - createWorkflowRunStartedEventV4( - { runId: 'wrun_1', specVersion: 5 }, - { token: 'test-token', dispatcher: agent }, - () => { - throw observerError; - } - ) - ).rejects.toBe(observerError); - agent.assertNoPendingInterceptors(); - }); - it('continues a graceful partial run_started stream from its sentinel cursor', async () => { - const origin = - WORKFLOW_SERVER_URL_OVERRIDE || 'https://vercel-workflow.com'; - const agent = new MockAgent(); - agent.disableNetConnect(); + const agent = mockAgent(); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events/run_started', method: 'POST', @@ -1185,7 +1172,7 @@ describe('createWorkflowRunEventV4 over HTTP', () => { } ); agent - .get(origin) + .get(ORIGIN) .intercept({ path: '/api/v4/runs/wrun_1/events?cursor=eid%3Aevnt_1&remoteRefBehavior=resolve&returnAll=true', method: 'GET', @@ -1834,10 +1821,9 @@ describe('v4 POST frame meta forwards every field the splitter produces', () => /** * The recycler in http-client only sees transport failures the v4 client - * reports to it. This covers that wiring end to end: a `fetch()` that rejects - * the way a wedged HTTP/2 session does must retire the shared events pool once - * the failures reach the threshold. Without the `onTransportOutcome` hook in - * `fetchV4` the recycler is never told anything and the pool lives forever. + * reports to it. A streamed response resolves `fetch()` as soon as headers + * arrive, before its body can fail, so the body consumer must own the success + * report or that early success erases every later stream failure. */ describe('v4 transport reports failures to the events recycler', () => { // There is only an undici pool to retire while the adapter owns one: @@ -1861,8 +1847,17 @@ describe('v4 transport reports failures to the events recycler', () => { }), }); - it('rebuilds the shared pool after repeated stream timeouts', async () => { - vi.spyOn(globalThis, 'fetch').mockRejectedValue(wedgedSessionError()); + it('rebuilds the shared pool after repeated response-body timeouts', async () => { + vi.spyOn(globalThis, 'fetch').mockImplementation(async () => { + const body = new ReadableStream({ + start(controller) { + controller.error(wedgedSessionError()); + }, + }); + return new Response(body, { + headers: { 'content-type': V4_FRAME_CONTENT_TYPE }, + }); + }); // No `dispatcher` in the config: the request must resolve the shared one, // which is what the recycler owns. diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index 2b64833217..f2ef4a3a30 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -100,15 +100,23 @@ import { isWsEventsTransportEnabled } from './ws-transport-enabled.js'; * for a large run can legitimately take a while to drain — a whole-request * deadline would abort it mid-stream. */ +interface V4Response { + response: Response; + reportTransportOutcome(error?: unknown): void; +} + async function fetchV4( url: string, init: { method: string; headers: Headers; body?: Uint8Array }, config: APIConfig | undefined, opName: string, + transportSuccess: 'headers' | 'body', attributes?: Record -): Promise { +): Promise { const dispatcher = getEventsDispatcher(config); - return instrumentedFetch({ + const reportTransportOutcome = (error?: unknown) => + noteEventsTransportOutcome(dispatcher, error); + const response = await instrumentedFetch({ method: init.method, url, headers: init.headers, @@ -119,8 +127,8 @@ async function fetchV4( // request builds a fresh one. undici keeps a black-holed HTTP/2 session in // service indefinitely, so without this every request routed onto it fails // until the compute instance is recycled — see noteEventsTransportOutcome. - onTransportOutcome: (error) => - noteEventsTransportOutcome(dispatcher, error), + onTransportOutcome: reportTransportOutcome, + deferTransportSuccess: transportSuccess === 'body', timeoutMs: null, logLabel: opName, // Read the body as bytes, not text: a CBOR error body (the fence 412 @@ -134,6 +142,7 @@ async function fetchV4( url ), }); + return { response, reportTransportOutcome }; } const EVENT_ID_HEADER = 'x-wf-event-id'; @@ -712,6 +721,7 @@ async function postWorkflowRunEventV4( { method: 'POST', headers, body: frame }, config, 'createEvent', + responseType === 'event-stream' ? 'body' : 'headers', { ...WorkflowEventsTransport('http'), ...WorkflowEventType(input.eventType), @@ -741,7 +751,11 @@ export async function createWorkflowRunEventV4( if (reply) return decodeCreateEventResponse(reply, input.eventType); } - const response = await postWorkflowRunEventV4(input, 'materialized', config); + const { response } = await postWorkflowRunEventV4( + input, + 'materialized', + config + ); const contentType = response.headers.get('content-type'); if (contentType?.startsWith(V4_FRAME_CONTENT_TYPE)) { @@ -784,19 +798,19 @@ export async function createWorkflowRunStartedEventV4( config?: APIConfig, onEvent?: (event: Event) => void ) { - const response = await postWorkflowRunEventV4( + const stream = await postWorkflowRunEventV4( { ...input, eventType: 'run_started' }, 'event-stream', config ); const replay = await consumeReplayLogResponse( - response, + stream, { runId: input.runId, onEvent }, config ); assert(replay.cursor, 'v4 createEvent: event stream missing cursor'); const maxEvents = MaxEventsHeaderSchema.safeParse( - response.headers.get(MAX_EVENTS_HEADER) + stream.response.headers.get(MAX_EVENTS_HEADER) ); if (!maxEvents.success) { throw new WorkflowWorldError('v4 createEvent: invalid max-events header', { @@ -895,11 +909,12 @@ export async function createWorkflowRunEventsBatchV4( } const url = `${baseUrl}/v4/runs/${encodeURIComponent(input.runId)}/events/batch`; - const response = await fetchV4( + const { response } = await fetchV4( url, { method: 'POST', headers, body }, config, 'createEventBatch', + 'headers', { ...WorkflowEventsTransport('http'), ...WorkflowEventType(input.events[0].eventType), @@ -1196,17 +1211,20 @@ async function postEventFrameOverWs( ); } +export interface ReplayLogResult { + events: Event[]; + cursor: string | null; + hasMore: boolean; +} + /** * Result of a `hook_received` POST that opted into the replay-log preload, * discriminated on `kind` (keyed on the response content type). */ export type HookReceivedPreloadV4Result = /** The server streamed the replay log back as v4 frames. */ - | { + | (ReplayLogResult & { kind: 'stream'; - events: Event[]; - cursor: string | null; - hasMore: boolean; /** * The canonical event this write created or converged on (the resume * claim winner's — ours or the producer's), named by the @@ -1215,7 +1233,7 @@ export type HookReceivedPreloadV4Result = canonicalEventId: string | undefined; /** Per-run event ceiling from the response header, when present. */ maxEvents: number | undefined; - } + }) /** * The server answered with the normal materialized CBOR body instead — * an older server, or one that declined the optimization. The @@ -1243,14 +1261,16 @@ export async function createHookReceivedPreloadEventV4( config?: APIConfig, onEvent?: (event: Event) => void ): Promise { - const response = await postWorkflowRunEventV4( + const stream = await postWorkflowRunEventV4( { ...input, eventType: 'hook_received' }, 'event-stream', config ); + const { response } = stream; const contentType = response.headers.get('content-type'); if (!contentType?.startsWith(V4_FRAME_CONTENT_TYPE)) { + stream.reportTransportOutcome(); return { kind: 'materialized', result: await decodeCreateEventResponse(response, 'hook_received'), @@ -1258,7 +1278,7 @@ export async function createHookReceivedPreloadEventV4( } const replay = await consumeReplayLogResponse( - response, + stream, { runId: input.runId, onEvent }, config ); @@ -1300,14 +1320,17 @@ export async function getEventV4( const url = `${baseUrl}/v4/runs/${encodeURIComponent(runId)}/events/${encodeURIComponent(eventId)}` + `?remoteRefBehavior=${remoteRefBehavior}`; - const response = await fetchV4( + const stream = await fetchV4( url, { method: 'GET', headers }, config, - 'getEvent' + 'getEvent', + 'body' ); + const { response } = stream; const contentType = response.headers.get('content-type'); if (!contentType?.startsWith(V4_FRAME_CONTENT_TYPE)) { + stream.reportTransportOutcome(); throw new Error( `v4 getEvent: expected ${V4_FRAME_CONTENT_TYPE}, got ${contentType ?? '(none)'}` ); @@ -1323,8 +1346,15 @@ export async function getEventV4( // GET emits a single frame (no sentinel); decodeFrames returns at EOF // after yielding it. - for await (const frame of decodeFrames(chunks)) { - return decodeEventFrame(frame); + try { + for await (const frame of decodeFrames(chunks)) { + stream.reportTransportOutcome(); + return decodeEventFrame(frame); + } + stream.reportTransportOutcome(); + } catch (error) { + stream.reportTransportOutcome(error); + throw error; } throw new Error(`v4 getEvent: empty frame stream for ${eventId}`); } @@ -1335,14 +1365,29 @@ class PartialEventStreamError extends WorkflowWorldError { } } +const MAX_PARTIAL_EVENT_STREAM_RETRIES = 3; + +type EventFrameStreamResult = + | { + kind: 'complete'; + cursor: string | null; + hasMore: boolean; + } + | { + kind: 'partial'; + error: PartialEventStreamError; + }; + async function consumeEventFrameStream( - response: Response, + stream: V4Response, opName: string, events: Event[], onEvent?: (event: Event) => void -): Promise<{ cursor: string | null; hasMore: boolean }> { +): Promise { + const { response } = stream; const contentType = response.headers.get('content-type'); if (!contentType?.startsWith(V4_FRAME_CONTENT_TYPE)) { + stream.reportTransportOutcome(); throw new Error( `v4 ${opName}: expected ${V4_FRAME_CONTENT_TYPE}, got ${contentType ?? '(none)'}` ); @@ -1357,33 +1402,47 @@ async function consumeEventFrameStream( try { next = await frames.next(); } catch (cause) { - throw new PartialEventStreamError( + const error = new PartialEventStreamError( `v4 ${opName}: event frame stream failed after ${events.length} events`, cause ); + stream.reportTransportOutcome(error); + return { kind: 'partial', error }; } if (next.done) break; - const frame = next.value; - if (frame.meta._end === 1) { - const end = EventStreamEndSchema.parse(frame.meta); - return { cursor: end.next ?? null, hasMore: end.hasMore }; - } - if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { - throw new Error(`v4 ${opName}: unexpected control frame`); + try { + const frame = next.value; + if (frame.meta._end === 1) { + const end = EventStreamEndSchema.parse(frame.meta); + stream.reportTransportOutcome(); + return { + kind: 'complete', + cursor: end.next ?? null, + hasMore: end.hasMore, + }; + } + if (Object.keys(frame.meta).some((key) => key.startsWith('_'))) { + throw new Error(`v4 ${opName}: unexpected control frame`); + } + const event = decodeEventFrame(frame); + events.push(event); + onEvent?.(event); + } catch (error) { + stream.reportTransportOutcome(); + throw error; } - const event = decodeEventFrame(frame); - events.push(event); - onEvent?.(event); } } finally { void frames.return(undefined); } - throw new PartialEventStreamError( + const error = new PartialEventStreamError( `v4 ${opName}: frame stream ended without the end-of-stream sentinel ` + `(${events.length} events read) — truncated response?` ); + stream.reportTransportOutcome(error); + return { kind: 'partial', error }; } /** @@ -1393,28 +1452,24 @@ async function consumeEventFrameStream( * again. */ async function consumeReplayLogResponse( - response: Response, + stream: V4Response, params: Pick, config?: APIConfig -): Promise<{ - events: Event[]; - cursor: string | null; - hasMore: boolean; -}> { +): Promise { const events: Event[] = []; let page: { cursor: string | null; hasMore: boolean }; - try { - page = await consumeEventFrameStream( - response, - 'createEvent', - events, - params.onEvent - ); - } catch (error) { - if (!(error instanceof PartialEventStreamError)) throw error; + const consumed = await consumeEventFrameStream( + stream, + 'createEvent', + events, + params.onEvent + ); + if (consumed.kind === 'partial') { const lastEvent = events.at(-1); - if (!lastEvent) throw error; + if (!lastEvent) throw consumed.error; page = { cursor: `eid:${lastEvent.eventId}`, hasMore: true }; + } else { + page = consumed; } if (!page.hasMore) return { events, ...page }; @@ -1428,6 +1483,10 @@ async function consumeReplayLogResponse( config ); events.push(...suffix.data); + assert( + suffix.data.length === 0 || suffix.cursor, + 'v4 createEvent: non-empty continuation missing cursor' + ); return { events, cursor: suffix.cursor ?? page.cursor, @@ -1484,46 +1543,44 @@ export async function getWorkflowRunEventsV4( const events: Event[] = []; let cursor = params.pagination?.cursor ?? null; - while (true) { + for (let partialRetries = 0; ; partialRetries++) { const url = `${baseUrl}/v4/runs/${encodeURIComponent(params.runId)}/events` + paginationToQuery(params, cursor); - try { - const response = await fetchV4( - url, - { method: 'GET', headers }, - config, - 'listEvents' - ); - const page = await consumeEventFrameStream( - response, - 'listEvents', - events, - params.onEvent - ); - if (params.pagination?.limit === undefined && page.hasMore) { - if (!page.cursor || page.cursor === cursor) { - throw new WorkflowWorldError( - `v4 listEvents: partial event stream made no cursor progress for run ${params.runId}`, - { code: 'SCHEMA_VALIDATION' } - ); - } - cursor = page.cursor; - continue; - } - return { data: events, ...page }; - } catch (error) { - if (!(error instanceof PartialEventStreamError)) throw error; + const stream = await fetchV4( + url, + { method: 'GET', headers }, + config, + 'listEvents', + 'body' + ); + const result = await consumeEventFrameStream( + stream, + 'listEvents', + events, + params.onEvent + ); + if (result.kind === 'partial') { const lastEvent = events.at(-1); if ( + partialRetries === MAX_PARTIAL_EVENT_STREAM_RETRIES || params.pagination?.limit !== undefined || !lastEvent || `eid:${lastEvent.eventId}` === cursor ) { - throw error; + throw result.error; } cursor = `eid:${lastEvent.eventId}`; + continue; } + + if (params.pagination?.limit === undefined && result.hasMore) { + throw new WorkflowWorldError( + `v4 listEvents: returnAll response was unexpectedly paginated for run ${params.runId}`, + { code: 'SCHEMA_VALIDATION' } + ); + } + return { data: events, cursor: result.cursor, hasMore: result.hasMore }; } } @@ -1552,16 +1609,18 @@ export async function getEventsByCorrelationIdV4( appendListParams(sp, params, params.pagination?.cursor ?? null); const url = `${baseUrl}/v4/events?${sp.toString()}`; const events: Event[] = []; - const response = await fetchV4( + const stream = await fetchV4( url, { method: 'GET', headers }, config, - 'listEventsByCorrelationId' + 'listEventsByCorrelationId', + 'body' ); - const page = await consumeEventFrameStream( - response, + const result = await consumeEventFrameStream( + stream, 'listEventsByCorrelationId', events ); - return { data: events, ...page }; + if (result.kind === 'partial') throw result.error; + return { data: events, cursor: result.cursor, hasMore: result.hasMore }; } diff --git a/packages/world-vercel/src/events.test.ts b/packages/world-vercel/src/events.test.ts index 76d0c6896d..87ba906319 100644 --- a/packages/world-vercel/src/events.test.ts +++ b/packages/world-vercel/src/events.test.ts @@ -1,10 +1,11 @@ import { Buffer } from 'node:buffer'; import { gzipSync } from 'node:zlib'; +import { WorkflowWorldError } from '@workflow/errors'; import type { AnyEventRequest, CreateEventParams } from '@workflow/world'; import { decode, encode } from 'cbor-x'; import { ulid } from 'ulid'; import { MockAgent } from 'undici'; -import { describe, expect, it } from 'vitest'; +import { describe, expect, it, vi } from 'vitest'; import { createWorkflowRunEvent, getWorkflowRunEvents, @@ -374,6 +375,44 @@ describe('createWorkflowRunEvent result contract', () => { ).rejects.toMatchObject(error); agent.assertNoPendingInterceptors(); }); + + it('does not retry an observer failure that looks like a transport error', async () => { + const agent = mockAgent(); + agent + .get(ORIGIN) + .intercept({ + path: '/api/v4/runs/wrun_1/events/run_started', + method: 'POST', + }) + .reply(200, runStartedResponse(), { + headers: { + 'content-type': V4_FRAME_CONTENT_TYPE, + 'x-wf-event-id': 'evnt_1', + 'x-wf-run-id': 'wrun_1', + 'x-wf-created-at': STARTED_AT.toISOString(), + 'x-wf-max-events': '10000', + }, + }); + + const error = new WorkflowWorldError('observer failed', { + code: 'TRANSPORT', + }); + const onEvent = vi.fn(() => { + throw error; + }); + + await expect( + createWorkflowRunEvent( + 'wrun_1', + { eventType: 'run_started', specVersion: 2 } as AnyEventRequest, + { onEvent }, + { token: 'test-token', dispatcher: agent } + ) + ).rejects.toBe(error); + + expect(onEvent).toHaveBeenCalledOnce(); + agent.assertNoPendingInterceptors(); + }); }); /** POSTs a v4 step_started with `params` and returns the decoded frame meta. */ diff --git a/packages/world-vercel/src/events.ts b/packages/world-vercel/src/events.ts index cb7d170a3d..e4f10b6ae0 100644 --- a/packages/world-vercel/src/events.ts +++ b/packages/world-vercel/src/events.ts @@ -31,6 +31,7 @@ * the v3 path. */ +import assert from 'node:assert/strict'; import { HookNotFoundError, WorkflowWorldError } from '@workflow/errors'; import { type AnyEventRequest, @@ -565,12 +566,33 @@ export async function createWorkflowRunEventBatch( }; } +class EventObserverError extends Error { + constructor(readonly error: unknown) { + super('event observer failed'); + } +} + export async function createWorkflowRunEvent( id: string | null, data: T, params?: CreateEventParams, config?: APIConfig ): Promise> { + const onEvent = params?.onEvent; + const requestParams = + onEvent === undefined + ? params + : { + ...params, + onEvent(event: Event) { + try { + onEvent(event); + } catch (error) { + throw new EventObserverError(error); + } + }, + }; + try { // Retry transient transport failures (UND_ERR_REQ_RETRY, ECONNRESET, // socket/headers timeouts, transient 5xx) in-process for event types that @@ -581,7 +603,7 @@ export async function createWorkflowRunEvent( // types (step_started, step_retrying, hook_received) run once. See // ./event-retry for the validated per-event classification. const result = await withEventPostRetry( - () => createWorkflowRunEventInner(id, data, params, config), + () => createWorkflowRunEventInner(id, data, requestParams, config), data.eventType, { // The atomic lazy-resume shape is deduplicated server-side by the @@ -613,6 +635,7 @@ export async function createWorkflowRunEvent( } return result as EventResult; } catch (err) { + if (err instanceof EventObserverError) throw err.error; // 404 on hook_disposed / hook_received → already-disposed hook. if ( isHookEventRequiringExistence(data.eventType) && @@ -748,32 +771,12 @@ async function createWorkflowRunEventInner( 'v4 createEvent: run_started stream is missing run_started' ); } - - let attributes = runCreated.eventData.attributes ?? {}; - let updatedAt = runStarted.createdAt; - for (const event of result.events) { - if (event.eventType === 'attr_set') { - attributes = applyAttributeChanges(attributes, event.eventData.changes); - updatedAt = event.createdAt; - } - } + const run = reconstructRunFromReplayEvents(result.events); + assert(run); return { event: runStarted, - run: { - runId: runCreated.runId, - status: 'running', - deploymentId: runCreated.eventData.deploymentId, - workflowName: runCreated.eventData.workflowName, - specVersion: runCreated.specVersion, - executionContext: runCreated.eventData.executionContext, - input: runCreated.eventData.input, - attributes, - encryptionPublicKey: runCreated.eventData.encryptionPublicKey, - startedAt: runStarted.createdAt, - createdAt: runCreated.createdAt, - updatedAt, - }, + run, events: result.events, cursor: result.cursor, hasMore: result.hasMore, diff --git a/packages/world-vercel/src/http-core.ts b/packages/world-vercel/src/http-core.ts index b8122bb7fd..c0729fbeb1 100644 --- a/packages/world-vercel/src/http-core.ts +++ b/packages/world-vercel/src/http-core.ts @@ -487,6 +487,8 @@ export interface InstrumentedFetchOptions extends HttpClientSpanOptions { * connections stop delivering (see noteEventsTransportOutcome). */ onTransportOutcome?: (error?: unknown) => void; + /** Let a streaming body consumer report success after it finishes. */ + deferTransportSuccess?: boolean; } /** @@ -520,6 +522,7 @@ export async function instrumentedFetch( attributes, durationAttribute, onTransportOutcome, + deferTransportSuccess = false, } = opts; const label = logLabel ?? url; @@ -601,7 +604,9 @@ export async function instrumentedFetch( throw error; } const ms = Date.now() - start; - onTransportOutcome?.(); + if (!deferTransportSuccess || !response.ok) { + onTransportOutcome?.(); + } httpLog(method, label, response, ms); recordClientSpanStatus(span, response.status);