diff --git a/.changeset/tidy-pears-report.md b/.changeset/tidy-pears-report.md new file mode 100644 index 0000000000..95a88d6216 --- /dev/null +++ b/.changeset/tidy-pears-report.md @@ -0,0 +1,7 @@ +--- +'@workflow/core': minor +'@workflow/world': minor +'@workflow/world-vercel': minor +--- + +Add the optional `world.telemetry.recordStepExecution` hook so Worlds can correlate flow requests with inline-executed workflow steps. diff --git a/packages/core/src/runtime.test.ts b/packages/core/src/runtime.test.ts index 20909c0641..447dc4325c 100644 --- a/packages/core/src/runtime.test.ts +++ b/packages/core/src/runtime.test.ts @@ -1631,9 +1631,11 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { hasMore: false, cursor: 'cursor_test', })); + const recordStepExecution = vi.fn(); setWorld({ specVersion: SPEC_VERSION_CURRENT, + telemetry: { recordStepExecution }, getDeploymentId: vi.fn(async () => workflowRun.deploymentId), createQueueHandler: vi.fn( ( @@ -1682,6 +1684,8 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { handlerPromise, order, queue, + recordStepExecution, + durableEvents, eventsList, stepIdSends, createdEventParams, @@ -1713,6 +1717,23 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { expect(queue).toHaveBeenCalled(); }); + it('reports every step whose user code ran during the invocation', async () => { + const { handlerPromise, durableEvents, recordStepExecution } = + await driveHandler({ + runId: 'wrun_invocation_step_ids', + queueImpl: async () => ({ messageId: null }), + }); + + await handlerPromise; + const startedStepIds = durableEvents + .filter((event) => event.eventType === 'step_started') + .map((event) => event.correlationId); + + expect(recordStepExecution.mock.calls.map(([stepId]) => stepId)).toEqual( + startedStepIds + ); + }); + it('does not ack while the step-dispatch send is still in flight', async () => { let releaseSend!: () => void; const sendGate = new Promise((resolve) => { diff --git a/packages/core/src/runtime/step-executor.ts b/packages/core/src/runtime/step-executor.ts index 8b26fc5899..9263f1e7a6 100644 --- a/packages/core/src/runtime/step-executor.ts +++ b/packages/core/src/runtime/step-executor.ts @@ -1147,6 +1147,7 @@ export async function executeStep( () => { // The last instant before user code: T7 of the resume window. reportResumeTtr(); + world.telemetry?.recordStepExecution?.(stepId); return stepFn.apply(thisVal, args); } ); diff --git a/packages/world-vercel/src/index.ts b/packages/world-vercel/src/index.ts index a18a5900ff..4f7eb8834c 100644 --- a/packages/world-vercel/src/index.ts +++ b/packages/world-vercel/src/index.ts @@ -7,7 +7,7 @@ import { createGetEncryptionKeyForRun } from './encryption.js'; import { validateRunExecutionContext } from './execution-context.js'; import { getDeadline } from './get-deadline.js'; import { instrumentObject } from './instrumentObject.js'; -import { createQueue } from './queue.js'; +import { createQueue, recordStepExecution } from './queue.js'; import { createResolveLatestDeploymentId } from './resolve-latest-deployment.js'; import { createStorage } from './storage.js'; import { createStreamer } from './streamer.js'; @@ -119,5 +119,6 @@ export function createWorld(config?: APIConfig): World { config?.dispatcher ), resolveLatestDeploymentId: createResolveLatestDeploymentId(config), + telemetry: { recordStepExecution }, }; } diff --git a/packages/world-vercel/src/queue.test.ts b/packages/world-vercel/src/queue.test.ts index 4b9c1069de..81c429f3ae 100644 --- a/packages/world-vercel/src/queue.test.ts +++ b/packages/world-vercel/src/queue.test.ts @@ -73,7 +73,7 @@ vi.mock('./utils.js', () => ({ })); import { missingDeploymentIdMessage } from './deployment-id.js'; -import { createQueue } from './queue.js'; +import { createQueue, recordStepExecution } from './queue.js'; import { getHttpUrl } from './utils.js'; describe('createQueue', () => { @@ -1134,6 +1134,115 @@ describe('createQueue', () => { expect(capturedMeta.requestId).toBeUndefined(); }); + it('reports the step IDs executed by a flow request', async () => { + mockHandleCallback.mockImplementation((handler) => { + return async () => { + await handler( + { + payload: { runId: 'run-123' }, + queueName: '__wkf_workflow_test', + }, + { messageId: 'msg-123', deliveryCount: 1 } + ); + return new Response('ok'); + }; + }); + + const routeHandler = createQueue().createQueueHandler( + '__wkf_workflow_', + async () => { + recordStepExecution('step-a'); + recordStepExecution('step-b'); + recordStepExecution('step-a'); + } + ); + + const response = await routeHandler(new Request('http://localhost')); + + expect(response.headers.get('x-vercel-internal-workflow-step-ids')).toBe( + JSON.stringify(['step-a', 'step-b']) + ); + }); + + it('isolates step IDs between concurrent flow requests', async () => { + mockHandleCallback.mockImplementation((handler) => { + return async (request: Request) => { + const runId = new URL(request.url).pathname.slice(1); + await handler( + { + payload: { runId }, + queueName: '__wkf_workflow_test', + }, + { messageId: `msg-${runId}`, deliveryCount: 1 } + ); + return new Response('ok'); + }; + }); + + const firstStarted = Promise.withResolvers(); + const releaseFirst = Promise.withResolvers(); + const routeHandler = createQueue().createQueueHandler( + '__wkf_workflow_', + async (payload) => { + if ('runId' in payload && payload.runId === 'run-first') { + recordStepExecution('step-first'); + firstStarted.resolve(); + await releaseFirst.promise; + } else { + recordStepExecution('step-second'); + } + } + ); + + const firstResponse = routeHandler( + new Request('http://localhost/run-first') + ); + await firstStarted.promise; + const secondResponse = routeHandler( + new Request('http://localhost/run-second') + ); + releaseFirst.resolve(); + + const [first, second] = await Promise.all([ + firstResponse, + secondResponse, + ]); + expect(first.headers.get('x-vercel-internal-workflow-step-ids')).toBe( + JSON.stringify(['step-first']) + ); + expect(second.headers.get('x-vercel-internal-workflow-step-ids')).toBe( + JSON.stringify(['step-second']) + ); + }); + + it('keeps the scalar step delivery path unchanged', async () => { + mockHandleCallback.mockImplementation((handler) => { + return async () => { + await handler( + { + payload: { runId: 'run-123', stepId: 'step-a' }, + queueName: '__wkf_workflow_test', + }, + { messageId: 'msg-123', deliveryCount: 1 } + ); + return new Response('ok'); + }; + }); + + const routeHandler = createQueue().createQueueHandler( + '__wkf_workflow_', + async () => { + recordStepExecution('step-a'); + } + ); + + const response = await routeHandler(new Request('http://localhost')); + + expect( + response.headers.get('x-vercel-internal-workflow-step-ids') + ).toBeNull(); + }); + it('should re-enqueue inline step payloads correctly', async () => { mockSend.mockResolvedValue({ messageId: 'new-msg-123' }); diff --git a/packages/world-vercel/src/queue.ts b/packages/world-vercel/src/queue.ts index 0a292da95b..baeb683933 100644 --- a/packages/world-vercel/src/queue.ts +++ b/packages/world-vercel/src/queue.ts @@ -152,10 +152,43 @@ class DualTransport implements Transport { } } -// per-copy-ok: both ends of this store live in the same `createQueueHandler` -// closure: the `run()` wrapper and the `getStore()` read always come from the -// same module copy, so the context never has to cross a copy boundary. -const requestIdStorage = new AsyncLocalStorage(); +interface QueueInvocationContext { + collectStepIds: boolean; + requestId?: string; + stepIds: Set; +} + +// per-copy-ok: the route wrapper and the World hook exported from this module +// share the same module copy through the World instance, so the context never +// has to cross a copy boundary. +const invocationStorage = new AsyncLocalStorage(); + +const WORKFLOW_STEP_IDS_HEADER = 'x-vercel-internal-workflow-step-ids'; + +export function recordStepExecution(stepId: string): void { + const invocation = invocationStorage.getStore(); + if (invocation?.collectStepIds) { + invocation.stepIds.add(stepId); + } +} + +function attachStepIds( + response: Response, + invocation: QueueInvocationContext +): Response { + if (invocation.stepIds.size === 0) return response; + + const headers = new Headers(response.headers); + headers.set( + WORKFLOW_STEP_IDS_HEADER, + JSON.stringify(Array.from(invocation.stepIds)) + ); + return new Response(response.body, { + headers, + status: response.status, + statusText: response.statusText, + }); +} const MessageWrapper = z.compile( z.object({ @@ -731,7 +764,8 @@ export function createQueue(config?: APIConfig): Queue { return; } - const requestId = requestIdStorage.getStore(); + const invocation = invocationStorage.getStore(); + const requestId = invocation?.requestId; // The CborTransport handles CBOR decoding inside deserialize(), // so message is already a plain object with Uint8Array values intact. const { payload, queueName, deploymentId } = @@ -745,15 +779,25 @@ export function createQueue(config?: APIConfig): Queue { getRunIdFromPayload(payload), config ); + const collectStepIds = !( + 'stepId' in payload && typeof payload.stepId === 'string' + ); wsEvents.open(); try { - const result = await handler(payload, { - queueName, - messageId: MessageId.parse(metadata.messageId), - attempt: metadata.deliveryCount, - requestId, - }); + const runHandler = () => + handler(payload, { + queueName, + messageId: MessageId.parse(metadata.messageId), + attempt: metadata.deliveryCount, + requestId, + }); + const result = await (invocation + ? invocationStorage.run( + { ...invocation, collectStepIds }, + runHandler + ) + : runHandler()); if ( !('invoke' in payload && payload.invoke === true) && @@ -806,7 +850,15 @@ export function createQueue(config?: APIConfig): Queue { return async (req: Request) => { const rawId = req.headers.get('x-vercel-id'); const requestId = rawId?.trim() || undefined; - return requestIdStorage.run(requestId, () => vqsHandler(req)); + const invocation: QueueInvocationContext = { + collectStepIds: false, + requestId, + stepIds: new Set(), + }; + const response = await invocationStorage.run(invocation, () => + vqsHandler(req) + ); + return attachStepIds(response, invocation); }; }; diff --git a/packages/world/src/interfaces.ts b/packages/world/src/interfaces.ts index 82323b5e6d..8be1d6e6bd 100644 --- a/packages/world/src/interfaces.ts +++ b/packages/world/src/interfaces.ts @@ -838,4 +838,19 @@ export interface World extends Queue, Streamer, Storage { | Record | null | Promise | null>; + + /** + * Optional telemetry write namespace for non-critical observability signals. + */ + telemetry?: Telemetry; +} + +export interface Telemetry { + /** + * Called immediately before a step's user code begins executing. + * + * Worlds may use this synchronous hook to correlate step execution with the + * current platform invocation. Implementations must not throw. + */ + recordStepExecution?(stepId: string): void; }