diff --git a/.changeset/send-event-occurred-at.md b/.changeset/send-event-occurred-at.md new file mode 100644 index 0000000000..8d968fbc86 --- /dev/null +++ b/.changeset/send-event-occurred-at.md @@ -0,0 +1,7 @@ +--- +'@workflow/world': patch +'@workflow/world-vercel': patch +'@workflow/world-postgres': patch +--- + +Send optional client-side event occurrence timestamps through world event creation. diff --git a/packages/world-postgres/src/drizzle/schema.ts b/packages/world-postgres/src/drizzle/schema.ts index aef97831e8..4761299730 100644 --- a/packages/world-postgres/src/drizzle/schema.ts +++ b/packages/world-postgres/src/drizzle/schema.ts @@ -114,7 +114,7 @@ export const events = schema.table( eventData: Cbor()('payload_cbor'), specVersion: integer('spec_version'), } satisfies DrizzlishOfType< - Cborized + Cborized & { eventData?: undefined }, 'eventData'> >, (tb) => [ index().on(tb.runId), diff --git a/packages/world-vercel/src/events-v4.ts b/packages/world-vercel/src/events-v4.ts index d5202ea99a..a3c13d8e4a 100644 --- a/packages/world-vercel/src/events-v4.ts +++ b/packages/world-vercel/src/events-v4.ts @@ -54,6 +54,8 @@ export interface CreateEventV4Input { specVersion: number; correlationId?: string; vercelId?: string; + /** Client-side time at which the event occurred. */ + occurredAt?: Date; remoteRefBehavior?: 'resolve' | 'lazy'; deploymentId?: string; workflowName?: string; @@ -119,6 +121,7 @@ function buildPostFrameMeta( if (input.correlationId !== undefined) meta.correlationId = input.correlationId; if (input.vercelId !== undefined) meta.vercelId = input.vercelId; + if (input.occurredAt !== undefined) meta.occurredAt = input.occurredAt; if (input.remoteRefBehavior !== undefined) { meta.remoteRefBehavior = input.remoteRefBehavior; } @@ -277,6 +280,7 @@ export interface DecodedV4Event { eventType: string; correlationId?: string; createdAt: Date | string; + occurredAt?: Date | string; specVersion?: number; eventData?: Record; } diff --git a/packages/world-vercel/src/events.test.ts b/packages/world-vercel/src/events.test.ts index b499b8e607..7ec3db5683 100644 --- a/packages/world-vercel/src/events.test.ts +++ b/packages/world-vercel/src/events.test.ts @@ -1,5 +1,5 @@ import type { AnyEventRequest } from '@workflow/world'; -import { encode } from 'cbor-x'; +import { decode, encode } from 'cbor-x'; import { MockAgent } from 'undici'; import { describe, expect, it } from 'vitest'; import { @@ -17,6 +17,19 @@ function mockAgent() { return agent; } +function decodePostedMeta(rawBody: unknown): Record { + const bytes = + typeof rawBody === 'string' + ? new TextEncoder().encode(rawBody) + : new Uint8Array(rawBody as ArrayBufferLike); + const metaLen = new DataView( + bytes.buffer, + bytes.byteOffset, + bytes.byteLength + ).getUint32(0, false); + return decode(bytes.subarray(4, 4 + metaLen)) as Record; +} + /** * Legacy (spec-version-1) runs predate event sourcing: the runtime still * posts hook_received (resumeHook) and wait_completed (wakeUpRun) for them @@ -212,6 +225,52 @@ describe('splitEventDataForV4 hook fields', () => { }); describe('createWorkflowRunEvent response coercion', () => { + it('sends occurredAt in the v4 frame meta', async () => { + const agent = mockAgent(); + const occurredAt = new Date('2026-06-10T00:00:03.000Z'); + let capturedMeta: Record | undefined; + + agent + .get(ORIGIN) + .intercept({ + path: '/api/v4/runs/wrun_1/events/run_started', + method: 'POST', + }) + .reply( + 200, + (opts: { body?: unknown }) => { + capturedMeta = decodePostedMeta(opts.body); + return encode({ + run: { + runId: 'wrun_1', + status: 'running', + startedAt: new Date('2026-06-10T00:00:04.000Z'), + }, + }); + }, + { + headers: { + 'x-wf-event-id': 'evnt_1', + 'x-wf-run-id': 'wrun_1', + 'x-wf-created-at': '2026-06-10T00:00:04.000Z', + }, + } + ); + + await createWorkflowRunEvent( + 'wrun_1', + { eventType: 'run_started', specVersion: 2 } as AnyEventRequest, + { occurredAt }, + { token: 'test-token', dispatcher: agent } + ); + + expect(capturedMeta?.occurredAt).toBeInstanceOf(Date); + expect((capturedMeta?.occurredAt as Date).getTime()).toBe( + occurredAt.getTime() + ); + agent.assertNoPendingInterceptors(); + }); + it('coerces ISO-string dates in the returned event and preloaded events', async () => { // Persisted events store nested eventData dates as ISO strings // (the backend's entity layer converts Date → toISOString on write with @@ -239,6 +298,7 @@ describe('createWorkflowRunEvent response coercion', () => { runId: 'wrun_1', eventType: 'run_started', createdAt: '2026-06-10T00:00:01.000Z', + occurredAt: '2026-06-10T00:00:00.500Z', eventData: {}, }, events: [ @@ -248,6 +308,7 @@ describe('createWorkflowRunEvent response coercion', () => { eventType: 'wait_created', correlationId: 'wait_1', createdAt: '2026-06-10T00:00:02.000Z', + occurredAt: '2026-06-10T00:00:01.500Z', specVersion: 2, eventData: { resumeAt: '2026-06-10T01:00:00.000Z' }, }, @@ -272,11 +333,14 @@ describe('createWorkflowRunEvent response coercion', () => { ); expect(result.event?.createdAt).toBeInstanceOf(Date); + expect(result.event?.occurredAt).toBeInstanceOf(Date); const preloaded = result.events?.[0] as { createdAt: Date; + occurredAt: Date; eventData: { resumeAt: Date }; }; expect(preloaded.createdAt).toBeInstanceOf(Date); + expect(preloaded.occurredAt).toBeInstanceOf(Date); expect(preloaded.eventData.resumeAt).toBeInstanceOf(Date); expect(preloaded.eventData.resumeAt.getTime()).toBe( new Date('2026-06-10T01:00:00.000Z').getTime() diff --git a/packages/world-vercel/src/events.ts b/packages/world-vercel/src/events.ts index c01b34782f..1baec032fa 100644 --- a/packages/world-vercel/src/events.ts +++ b/packages/world-vercel/src/events.ts @@ -461,6 +461,14 @@ function buildEventFromV4( decoded.createdAt instanceof Date ? decoded.createdAt : new Date(decoded.createdAt), + ...(decoded.occurredAt !== undefined + ? { + occurredAt: + decoded.occurredAt instanceof Date + ? decoded.occurredAt + : new Date(decoded.occurredAt), + } + : {}), ...(decoded.correlationId ? { correlationId: decoded.correlationId } : {}), eventData, ...(decoded.specVersion !== undefined @@ -619,6 +627,7 @@ async function createWorkflowRunEventInner( specVersion: data.specVersion ?? 2, ...(data.correlationId ? { correlationId: data.correlationId } : {}), ...(params?.requestId ? { vercelId: params.requestId } : {}), + occurredAt: params?.occurredAt ?? new Date(), remoteRefBehavior, payload, ...meta, diff --git a/packages/world/src/events.ts b/packages/world/src/events.ts index 61b7b2e665..48bd70d22e 100644 --- a/packages/world/src/events.ts +++ b/packages/world/src/events.ts @@ -362,6 +362,7 @@ export const EventSchema = AllEventsSchema.and( runId: z.string(), eventId: z.string(), createdAt: z.coerce.date(), + occurredAt: z.coerce.date().optional(), specVersion: z.number().optional(), }) ); @@ -397,6 +398,12 @@ export interface CreateEventParams { resolveData?: ResolveData; /** Request ID (x-vercel-id when on Vercel) for correlating request logs with workflow events. */ requestId?: string; + /** + * Timestamp for when the event occurred on the client side. Worlds that + * support this can persist it separately from `createdAt`, which represents + * when the backing service accepted or stored the event. + */ + occurredAt?: Date; } /**