From 40b880b95dde8b95c27c9ccf58063ad2f38289ee Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Fri, 25 Sep 2026 15:33:04 -0700 Subject: [PATCH 1/6] [core] Report executed workflow step IDs per request Signed-off-by: Karthik Kalyanaraman --- .changeset/tidy-pears-report.md | 5 + packages/core/src/runtime.test.ts | 19 +++ packages/core/src/runtime.ts | 108 +++++++++--------- .../src/runtime/invocation-step-ids.test.ts | 51 +++++++++ .../core/src/runtime/invocation-step-ids.ts | 38 ++++++ packages/core/src/runtime/step-executor.ts | 2 + 6 files changed, 172 insertions(+), 51 deletions(-) create mode 100644 .changeset/tidy-pears-report.md create mode 100644 packages/core/src/runtime/invocation-step-ids.test.ts create mode 100644 packages/core/src/runtime/invocation-step-ids.ts diff --git a/.changeset/tidy-pears-report.md b/.changeset/tidy-pears-report.md new file mode 100644 index 0000000000..c72f1c3ee8 --- /dev/null +++ b/.changeset/tidy-pears-report.md @@ -0,0 +1,5 @@ +--- +'@workflow/core': minor +--- + +Report the workflow step IDs executed by each handler invocation to the platform runtime. diff --git a/packages/core/src/runtime.test.ts b/packages/core/src/runtime.test.ts index 20909c0641..4c6ef4be8b 100644 --- a/packages/core/src/runtime.test.ts +++ b/packages/core/src/runtime.test.ts @@ -1682,6 +1682,7 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { handlerPromise, order, queue, + durableEvents, eventsList, stepIdSends, createdEventParams, @@ -1713,6 +1714,24 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { expect(queue).toHaveBeenCalled(); }); + it('reports every step whose user code ran during the invocation', async () => { + const { handlerPromise, durableEvents } = await driveHandler({ + runId: 'wrun_invocation_step_ids', + queueImpl: async () => ({ messageId: null }), + }); + + const res = (await handlerPromise) as Response; + const startedStepIds = durableEvents + .filter((event) => event.eventType === 'step_started') + .map((event) => event.correlationId); + + expect( + JSON.parse( + res.headers.get('x-vercel-internal-workflow-step-ids') ?? 'null' + ) + ).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.ts b/packages/core/src/runtime.ts index 349c4e59cc..842481f49f 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -94,6 +94,10 @@ import { withHealthCheck, } from './runtime/helpers.js'; import { withRunInputs } from './runtime/invocations.js'; +import { + attachInvocationStepIds, + withInvocationStepIds, +} from './runtime/invocation-step-ids.js'; import { dispatchRunCompletedHooks, dispatchRunFailedHooks, @@ -5508,59 +5512,61 @@ export function workflowEntrypoint( ? entrypointCreatedAt - options.routeModuleBodyStartedAt : undefined; - return withHealthCheck(async (req) => { - invocationCount += 1; - const handlerCached = cachedHandler !== undefined; - const spanKind = await getSpanKind('SERVER'); + return withHealthCheck((req) => + withInvocationStepIds(async () => { + invocationCount += 1; + const handlerCached = cachedHandler !== undefined; + const spanKind = await getSpanKind('SERVER'); - return trace( - 'workflow.route.flow', - { - kind: spanKind, - attributes: { - ...Attribute.WorkflowRouteType('flow'), - ...Attribute.FaasInstance(COMPUTE_INSTANCE_ID), - ...Attribute.WorkflowRouteHandlerCached(handlerCached), - ...Attribute.WorkflowRouteInvocationCount(invocationCount), - ...Attribute.WorkflowRouteEntrypointAgeMs( - Date.now() - entrypointCreatedAt - ), - ...(routeModuleBodyInitMs === undefined - ? {} - : Attribute.WorkflowRouteModuleBodyInitMs(routeModuleBodyInitMs)), - ...Attribute.HttpRequestMethod(req.method), - ...Attribute.HttpRoute('/.well-known/workflow/v1/flow'), + return trace( + 'workflow.route.flow', + { + kind: spanKind, + attributes: { + ...Attribute.WorkflowRouteType('flow'), + ...Attribute.FaasInstance(COMPUTE_INSTANCE_ID), + ...Attribute.WorkflowRouteHandlerCached(handlerCached), + ...Attribute.WorkflowRouteInvocationCount(invocationCount), + ...Attribute.WorkflowRouteEntrypointAgeMs( + Date.now() - entrypointCreatedAt + ), + ...(routeModuleBodyInitMs === undefined + ? {} + : Attribute.WorkflowRouteModuleBodyInitMs(routeModuleBodyInitMs)), + ...Attribute.HttpRequestMethod(req.method), + ...Attribute.HttpRoute('/.well-known/workflow/v1/flow'), + }, }, - }, - async (span) => { - if (!cachedHandler) { - cachedHandler = await trace('workflow.route.init', async () => { - // The full runtime World, not `getWorldHandlers()`. That accessor - // owns a second, build-time-safe cache, so calling it here built a - // second World in the same process: duplicate connection pools and - // queue workers for a stateful World, plus a second copy of that - // world package's modules once it is bundled, which is what - // silently demoted the events WebSocket transport to HTTP. #3665. - // - // The span keeps its original name. It is a distinct span from the - // per-request `workflow.route.get_world` at the top of the flow - // route, and renaming it would collide with that one. - const worldHandlers = await trace( - 'workflow.route.get_world_handlers', - async () => getWorld() - ); - return handler(worldHandlers); - }); - } + async (span) => { + if (!cachedHandler) { + cachedHandler = await trace('workflow.route.init', async () => { + // The full runtime World, not `getWorldHandlers()`. That accessor + // owns a second, build-time-safe cache, so calling it here built a + // second World in the same process: duplicate connection pools and + // queue workers for a stateful World, plus a second copy of that + // world package's modules once it is bundled, which is what + // silently demoted the events WebSocket transport to HTTP. #3665. + // + // The span keeps its original name. It is a distinct span from the + // per-request `workflow.route.get_world` at the top of the flow + // route, and renaming it would collide with that one. + const worldHandlers = await trace( + 'workflow.route.get_world_handlers', + async () => getWorld() + ); + return handler(worldHandlers); + }); + } - const response = await cachedHandler(req); - if (response instanceof Response) { - span?.setAttributes( - Attribute.HttpResponseStatusCode(response.status) - ); + const response = await cachedHandler!(req); + if (response instanceof Response) { + span?.setAttributes( + Attribute.HttpResponseStatusCode(response.status) + ); + } + return attachInvocationStepIds(response); } - return response; - } - ); - }); + ); + }) + ); } diff --git a/packages/core/src/runtime/invocation-step-ids.test.ts b/packages/core/src/runtime/invocation-step-ids.test.ts new file mode 100644 index 0000000000..4032924914 --- /dev/null +++ b/packages/core/src/runtime/invocation-step-ids.test.ts @@ -0,0 +1,51 @@ +import { describe, expect, it } from 'vitest'; +import { + attachInvocationStepIds, + recordInvocationStepId, + WORKFLOW_STEP_IDS_HEADER, + withInvocationStepIds, +} from './invocation-step-ids.js'; + +describe('invocation step IDs', () => { + it('reports unique step IDs in execution order', async () => { + await withInvocationStepIds(async () => { + recordInvocationStepId('step_a'); + recordInvocationStepId('step_b'); + recordInvocationStepId('step_a'); + + const response = attachInvocationStepIds(new Response(null)); + + expect( + JSON.parse(response.headers.get(WORKFLOW_STEP_IDS_HEADER)!) + ).toEqual(['step_a', 'step_b']); + }); + }); + + it('isolates concurrent invocations', async () => { + let releaseFirst!: () => void; + const firstBlocked = new Promise((resolve) => { + releaseFirst = resolve; + }); + + const first = withInvocationStepIds(async () => { + recordInvocationStepId('step_first'); + await firstBlocked; + return attachInvocationStepIds(new Response(null)); + }); + + const second = withInvocationStepIds(async () => { + recordInvocationStepId('step_second'); + return attachInvocationStepIds(new Response(null)); + }); + + releaseFirst(); + const [firstResponse, secondResponse] = await Promise.all([first, second]); + + expect( + JSON.parse(firstResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)!) + ).toEqual(['step_first']); + expect( + JSON.parse(secondResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)!) + ).toEqual(['step_second']); + }); +}); diff --git a/packages/core/src/runtime/invocation-step-ids.ts b/packages/core/src/runtime/invocation-step-ids.ts new file mode 100644 index 0000000000..d9f52648bd --- /dev/null +++ b/packages/core/src/runtime/invocation-step-ids.ts @@ -0,0 +1,38 @@ +import { AsyncLocalStorage } from 'node:async_hooks'; + +export const WORKFLOW_STEP_IDS_HEADER = 'x-vercel-internal-workflow-step-ids'; + +type InvocationStepIds = Set; + +const INVOCATION_STEP_IDS_STORAGE_SYMBOL = Symbol.for( + 'WORKFLOW_INVOCATION_STEP_IDS_STORAGE' +); + +const invocationStepIdsStorage: AsyncLocalStorage = (() => { + const store = globalThis as typeof globalThis & { + [INVOCATION_STEP_IDS_STORAGE_SYMBOL]?: AsyncLocalStorage; + }; + store[INVOCATION_STEP_IDS_STORAGE_SYMBOL] ??= + new AsyncLocalStorage(); + return store[INVOCATION_STEP_IDS_STORAGE_SYMBOL]; +})(); + +export function withInvocationStepIds(callback: () => T): T { + return invocationStepIdsStorage.run(new Set(), callback); +} + +export function recordInvocationStepId(stepId: string): void { + invocationStepIdsStorage.getStore()?.add(stepId); +} + +export function attachInvocationStepIds(response: Response): Response { + const stepIds = Array.from(invocationStepIdsStorage.getStore() ?? []); + const headers = new Headers(response.headers); + headers.set(WORKFLOW_STEP_IDS_HEADER, JSON.stringify(stepIds)); + + return new Response(response.body, { + headers, + status: response.status, + statusText: response.statusText, + }); +} diff --git a/packages/core/src/runtime/step-executor.ts b/packages/core/src/runtime/step-executor.ts index 8b26fc5899..8a2f0abf2e 100644 --- a/packages/core/src/runtime/step-executor.ts +++ b/packages/core/src/runtime/step-executor.ts @@ -57,6 +57,7 @@ import { } from './constants.js'; import { getPortLazy } from './get-port-lazy.js'; import { memoizeEncryptionKey, withReplayDelta } from './helpers.js'; +import { recordInvocationStepId } from './invocation-step-ids.js'; import { ReplayRecoveryReporter } from './replay-recovery-reporter.js'; import { computeResumeTtrAttributes, @@ -1147,6 +1148,7 @@ export async function executeStep( () => { // The last instant before user code: T7 of the resume window. reportResumeTtr(); + recordInvocationStepId(stepId); return stepFn.apply(thisVal, args); } ); From afffb231b7e50d0fbfe7f3236c88aea147147383 Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Fri, 25 Sep 2026 17:59:36 -0700 Subject: [PATCH 2/6] [core] Handle non-Response route results Signed-off-by: Karthik Kalyanaraman --- packages/core/src/runtime.ts | 19 ++++++++++----- .../src/runtime/invocation-step-ids.test.ts | 24 +++++++++---------- 2 files changed, 24 insertions(+), 19 deletions(-) diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 842481f49f..0726a164a4 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -93,11 +93,11 @@ import { stepDispatchIdempotencyKey, withHealthCheck, } from './runtime/helpers.js'; -import { withRunInputs } from './runtime/invocations.js'; import { attachInvocationStepIds, withInvocationStepIds, } from './runtime/invocation-step-ids.js'; +import { withRunInputs } from './runtime/invocations.js'; import { dispatchRunCompletedHooks, dispatchRunFailedHooks, @@ -5558,12 +5558,19 @@ export function workflowEntrypoint( }); } - const response = await cachedHandler!(req); - if (response instanceof Response) { - span?.setAttributes( - Attribute.HttpResponseStatusCode(response.status) - ); + const activeHandler = cachedHandler; + if (!activeHandler) { + throw new Error('Workflow route handler was not initialized'); } + + const response = await activeHandler(req); + if (!(response instanceof Response)) { + return response; + } + + span?.setAttributes( + Attribute.HttpResponseStatusCode(response.status) + ); return attachInvocationStepIds(response); } ); diff --git a/packages/core/src/runtime/invocation-step-ids.test.ts b/packages/core/src/runtime/invocation-step-ids.test.ts index 4032924914..1ed5b6f2ae 100644 --- a/packages/core/src/runtime/invocation-step-ids.test.ts +++ b/packages/core/src/runtime/invocation-step-ids.test.ts @@ -15,17 +15,15 @@ describe('invocation step IDs', () => { const response = attachInvocationStepIds(new Response(null)); - expect( - JSON.parse(response.headers.get(WORKFLOW_STEP_IDS_HEADER)!) - ).toEqual(['step_a', 'step_b']); + expect(response.headers.get(WORKFLOW_STEP_IDS_HEADER)).toBe( + JSON.stringify(['step_a', 'step_b']) + ); }); }); it('isolates concurrent invocations', async () => { - let releaseFirst!: () => void; - const firstBlocked = new Promise((resolve) => { - releaseFirst = resolve; - }); + const { promise: firstBlocked, resolve: releaseFirst } = + Promise.withResolvers(); const first = withInvocationStepIds(async () => { recordInvocationStepId('step_first'); @@ -41,11 +39,11 @@ describe('invocation step IDs', () => { releaseFirst(); const [firstResponse, secondResponse] = await Promise.all([first, second]); - expect( - JSON.parse(firstResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)!) - ).toEqual(['step_first']); - expect( - JSON.parse(secondResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)!) - ).toEqual(['step_second']); + expect(firstResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)).toBe( + JSON.stringify(['step_first']) + ); + expect(secondResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)).toBe( + JSON.stringify(['step_second']) + ); }); }); From 43f297ca0a6197e1ea84ac66703b6380d21085be Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Fri, 25 Sep 2026 18:18:50 -0700 Subject: [PATCH 3/6] [world-vercel] Scope workflow step headers to Vercel Signed-off-by: Karthik Kalyanaraman --- .changeset/tidy-pears-report.md | 4 +- packages/core/src/runtime.test.ts | 22 ++-- packages/core/src/runtime.ts | 109 ++++++++--------- .../src/runtime/invocation-step-ids.test.ts | 49 -------- .../core/src/runtime/invocation-step-ids.ts | 38 ------ packages/core/src/runtime/step-executor.ts | 3 +- packages/world-vercel/src/index.ts | 3 +- packages/world-vercel/src/queue.test.ts | 111 +++++++++++++++++- packages/world-vercel/src/queue.ts | 76 ++++++++++-- packages/world/src/interfaces.ts | 8 ++ 10 files changed, 248 insertions(+), 175 deletions(-) delete mode 100644 packages/core/src/runtime/invocation-step-ids.test.ts delete mode 100644 packages/core/src/runtime/invocation-step-ids.ts diff --git a/.changeset/tidy-pears-report.md b/.changeset/tidy-pears-report.md index c72f1c3ee8..414accaee3 100644 --- a/.changeset/tidy-pears-report.md +++ b/.changeset/tidy-pears-report.md @@ -1,5 +1,7 @@ --- '@workflow/core': minor +'@workflow/world': minor +'@workflow/world-vercel': minor --- -Report the workflow step IDs executed by each handler invocation to the platform runtime. +Let Vercel correlate flow requests with every workflow step executed inline. diff --git a/packages/core/src/runtime.test.ts b/packages/core/src/runtime.test.ts index 4c6ef4be8b..a5837bc231 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, + recordStepExecution, getDeploymentId: vi.fn(async () => workflowRun.deploymentId), createQueueHandler: vi.fn( ( @@ -1682,6 +1684,7 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { handlerPromise, order, queue, + recordStepExecution, durableEvents, eventsList, stepIdSends, @@ -1715,21 +1718,20 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { }); it('reports every step whose user code ran during the invocation', async () => { - const { handlerPromise, durableEvents } = await driveHandler({ - runId: 'wrun_invocation_step_ids', - queueImpl: async () => ({ messageId: null }), - }); + const { handlerPromise, durableEvents, recordStepExecution } = + await driveHandler({ + runId: 'wrun_invocation_step_ids', + queueImpl: async () => ({ messageId: null }), + }); - const res = (await handlerPromise) as Response; + await handlerPromise; const startedStepIds = durableEvents .filter((event) => event.eventType === 'step_started') .map((event) => event.correlationId); - expect( - JSON.parse( - res.headers.get('x-vercel-internal-workflow-step-ids') ?? 'null' - ) - ).toEqual(startedStepIds); + expect(recordStepExecution.mock.calls.map(([stepId]) => stepId)).toEqual( + startedStepIds + ); }); it('does not ack while the step-dispatch send is still in flight', async () => { diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 0726a164a4..349c4e59cc 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -93,10 +93,6 @@ import { stepDispatchIdempotencyKey, withHealthCheck, } from './runtime/helpers.js'; -import { - attachInvocationStepIds, - withInvocationStepIds, -} from './runtime/invocation-step-ids.js'; import { withRunInputs } from './runtime/invocations.js'; import { dispatchRunCompletedHooks, @@ -5512,68 +5508,59 @@ export function workflowEntrypoint( ? entrypointCreatedAt - options.routeModuleBodyStartedAt : undefined; - return withHealthCheck((req) => - withInvocationStepIds(async () => { - invocationCount += 1; - const handlerCached = cachedHandler !== undefined; - const spanKind = await getSpanKind('SERVER'); + return withHealthCheck(async (req) => { + invocationCount += 1; + const handlerCached = cachedHandler !== undefined; + const spanKind = await getSpanKind('SERVER'); - return trace( - 'workflow.route.flow', - { - kind: spanKind, - attributes: { - ...Attribute.WorkflowRouteType('flow'), - ...Attribute.FaasInstance(COMPUTE_INSTANCE_ID), - ...Attribute.WorkflowRouteHandlerCached(handlerCached), - ...Attribute.WorkflowRouteInvocationCount(invocationCount), - ...Attribute.WorkflowRouteEntrypointAgeMs( - Date.now() - entrypointCreatedAt - ), - ...(routeModuleBodyInitMs === undefined - ? {} - : Attribute.WorkflowRouteModuleBodyInitMs(routeModuleBodyInitMs)), - ...Attribute.HttpRequestMethod(req.method), - ...Attribute.HttpRoute('/.well-known/workflow/v1/flow'), - }, + return trace( + 'workflow.route.flow', + { + kind: spanKind, + attributes: { + ...Attribute.WorkflowRouteType('flow'), + ...Attribute.FaasInstance(COMPUTE_INSTANCE_ID), + ...Attribute.WorkflowRouteHandlerCached(handlerCached), + ...Attribute.WorkflowRouteInvocationCount(invocationCount), + ...Attribute.WorkflowRouteEntrypointAgeMs( + Date.now() - entrypointCreatedAt + ), + ...(routeModuleBodyInitMs === undefined + ? {} + : Attribute.WorkflowRouteModuleBodyInitMs(routeModuleBodyInitMs)), + ...Attribute.HttpRequestMethod(req.method), + ...Attribute.HttpRoute('/.well-known/workflow/v1/flow'), }, - async (span) => { - if (!cachedHandler) { - cachedHandler = await trace('workflow.route.init', async () => { - // The full runtime World, not `getWorldHandlers()`. That accessor - // owns a second, build-time-safe cache, so calling it here built a - // second World in the same process: duplicate connection pools and - // queue workers for a stateful World, plus a second copy of that - // world package's modules once it is bundled, which is what - // silently demoted the events WebSocket transport to HTTP. #3665. - // - // The span keeps its original name. It is a distinct span from the - // per-request `workflow.route.get_world` at the top of the flow - // route, and renaming it would collide with that one. - const worldHandlers = await trace( - 'workflow.route.get_world_handlers', - async () => getWorld() - ); - return handler(worldHandlers); - }); - } - - const activeHandler = cachedHandler; - if (!activeHandler) { - throw new Error('Workflow route handler was not initialized'); - } - - const response = await activeHandler(req); - if (!(response instanceof Response)) { - return response; - } + }, + async (span) => { + if (!cachedHandler) { + cachedHandler = await trace('workflow.route.init', async () => { + // The full runtime World, not `getWorldHandlers()`. That accessor + // owns a second, build-time-safe cache, so calling it here built a + // second World in the same process: duplicate connection pools and + // queue workers for a stateful World, plus a second copy of that + // world package's modules once it is bundled, which is what + // silently demoted the events WebSocket transport to HTTP. #3665. + // + // The span keeps its original name. It is a distinct span from the + // per-request `workflow.route.get_world` at the top of the flow + // route, and renaming it would collide with that one. + const worldHandlers = await trace( + 'workflow.route.get_world_handlers', + async () => getWorld() + ); + return handler(worldHandlers); + }); + } + const response = await cachedHandler(req); + if (response instanceof Response) { span?.setAttributes( Attribute.HttpResponseStatusCode(response.status) ); - return attachInvocationStepIds(response); } - ); - }) - ); + return response; + } + ); + }); } diff --git a/packages/core/src/runtime/invocation-step-ids.test.ts b/packages/core/src/runtime/invocation-step-ids.test.ts deleted file mode 100644 index 1ed5b6f2ae..0000000000 --- a/packages/core/src/runtime/invocation-step-ids.test.ts +++ /dev/null @@ -1,49 +0,0 @@ -import { describe, expect, it } from 'vitest'; -import { - attachInvocationStepIds, - recordInvocationStepId, - WORKFLOW_STEP_IDS_HEADER, - withInvocationStepIds, -} from './invocation-step-ids.js'; - -describe('invocation step IDs', () => { - it('reports unique step IDs in execution order', async () => { - await withInvocationStepIds(async () => { - recordInvocationStepId('step_a'); - recordInvocationStepId('step_b'); - recordInvocationStepId('step_a'); - - const response = attachInvocationStepIds(new Response(null)); - - expect(response.headers.get(WORKFLOW_STEP_IDS_HEADER)).toBe( - JSON.stringify(['step_a', 'step_b']) - ); - }); - }); - - it('isolates concurrent invocations', async () => { - const { promise: firstBlocked, resolve: releaseFirst } = - Promise.withResolvers(); - - const first = withInvocationStepIds(async () => { - recordInvocationStepId('step_first'); - await firstBlocked; - return attachInvocationStepIds(new Response(null)); - }); - - const second = withInvocationStepIds(async () => { - recordInvocationStepId('step_second'); - return attachInvocationStepIds(new Response(null)); - }); - - releaseFirst(); - const [firstResponse, secondResponse] = await Promise.all([first, second]); - - expect(firstResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)).toBe( - JSON.stringify(['step_first']) - ); - expect(secondResponse.headers.get(WORKFLOW_STEP_IDS_HEADER)).toBe( - JSON.stringify(['step_second']) - ); - }); -}); diff --git a/packages/core/src/runtime/invocation-step-ids.ts b/packages/core/src/runtime/invocation-step-ids.ts deleted file mode 100644 index d9f52648bd..0000000000 --- a/packages/core/src/runtime/invocation-step-ids.ts +++ /dev/null @@ -1,38 +0,0 @@ -import { AsyncLocalStorage } from 'node:async_hooks'; - -export const WORKFLOW_STEP_IDS_HEADER = 'x-vercel-internal-workflow-step-ids'; - -type InvocationStepIds = Set; - -const INVOCATION_STEP_IDS_STORAGE_SYMBOL = Symbol.for( - 'WORKFLOW_INVOCATION_STEP_IDS_STORAGE' -); - -const invocationStepIdsStorage: AsyncLocalStorage = (() => { - const store = globalThis as typeof globalThis & { - [INVOCATION_STEP_IDS_STORAGE_SYMBOL]?: AsyncLocalStorage; - }; - store[INVOCATION_STEP_IDS_STORAGE_SYMBOL] ??= - new AsyncLocalStorage(); - return store[INVOCATION_STEP_IDS_STORAGE_SYMBOL]; -})(); - -export function withInvocationStepIds(callback: () => T): T { - return invocationStepIdsStorage.run(new Set(), callback); -} - -export function recordInvocationStepId(stepId: string): void { - invocationStepIdsStorage.getStore()?.add(stepId); -} - -export function attachInvocationStepIds(response: Response): Response { - const stepIds = Array.from(invocationStepIdsStorage.getStore() ?? []); - const headers = new Headers(response.headers); - headers.set(WORKFLOW_STEP_IDS_HEADER, JSON.stringify(stepIds)); - - return new Response(response.body, { - headers, - status: response.status, - statusText: response.statusText, - }); -} diff --git a/packages/core/src/runtime/step-executor.ts b/packages/core/src/runtime/step-executor.ts index 8a2f0abf2e..64c4f50b13 100644 --- a/packages/core/src/runtime/step-executor.ts +++ b/packages/core/src/runtime/step-executor.ts @@ -57,7 +57,6 @@ import { } from './constants.js'; import { getPortLazy } from './get-port-lazy.js'; import { memoizeEncryptionKey, withReplayDelta } from './helpers.js'; -import { recordInvocationStepId } from './invocation-step-ids.js'; import { ReplayRecoveryReporter } from './replay-recovery-reporter.js'; import { computeResumeTtrAttributes, @@ -1148,7 +1147,7 @@ export async function executeStep( () => { // The last instant before user code: T7 of the resume window. reportResumeTtr(); - recordInvocationStepId(stepId); + world.recordStepExecution?.(stepId); return stepFn.apply(thisVal, args); } ); diff --git a/packages/world-vercel/src/index.ts b/packages/world-vercel/src/index.ts index 933bc1b836..fa0fde6d85 100644 --- a/packages/world-vercel/src/index.ts +++ b/packages/world-vercel/src/index.ts @@ -5,7 +5,7 @@ import { createRunId, describeRun } from './create-run-id.js'; import { createGetEncryptionKeyForRun } from './encryption.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'; @@ -70,6 +70,7 @@ export function createWorld(config?: APIConfig): World { // immediately, without a redeploy of this adapter. }, getRuntimeDeadline: getDeadline, + recordStepExecution, ...createQueue(config), ...createStorage(config), // Analytics list reads are served from an eventually-ingested store. 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 47234ff7c6..1bfacc29ab 100644 --- a/packages/world/src/interfaces.ts +++ b/packages/world/src/interfaces.ts @@ -613,6 +613,14 @@ export interface WorldCapabilities { * The "World" interface represents how Workflows are able to communicate with the outside world. */ export interface World extends Queue, Streamer, Storage { + /** + * 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; + /** * Optional analytics read namespace for observability surfaces. * From 2e48533278e36e9a9316ce3eaf3a909bc8f5fa4c Mon Sep 17 00:00:00 2001 From: Karthik Kalyan <105607645+karthikscale3@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:23:50 -0700 Subject: [PATCH 4/6] Update .changeset/tidy-pears-report.md Co-authored-by: Peter Wielander Signed-off-by: Karthik Kalyan <105607645+karthikscale3@users.noreply.github.com> --- .changeset/tidy-pears-report.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changeset/tidy-pears-report.md b/.changeset/tidy-pears-report.md index 414accaee3..3ba7aaaa65 100644 --- a/.changeset/tidy-pears-report.md +++ b/.changeset/tidy-pears-report.md @@ -4,4 +4,4 @@ '@workflow/world-vercel': minor --- -Let Vercel correlate flow requests with every workflow step executed inline. +Add `world.recordStepExecution?` interface to let Worlds correlate flow requests with inline-executed workflow steps. From 08dbaf2481a8402c52892aa4f41e06c002ccea97 Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Tue, 29 Sep 2026 10:37:16 -0700 Subject: [PATCH 5/6] refactor: namespace world telemetry hooks --- packages/core/src/runtime.test.ts | 2 +- packages/core/src/runtime/step-executor.ts | 2 +- packages/world-vercel/src/index.ts | 2 +- packages/world/src/interfaces.ts | 23 ++++++++++++++-------- 4 files changed, 18 insertions(+), 11 deletions(-) diff --git a/packages/core/src/runtime.test.ts b/packages/core/src/runtime.test.ts index a5837bc231..447dc4325c 100644 --- a/packages/core/src/runtime.test.ts +++ b/packages/core/src/runtime.test.ts @@ -1635,7 +1635,7 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => { setWorld({ specVersion: SPEC_VERSION_CURRENT, - recordStepExecution, + telemetry: { recordStepExecution }, getDeploymentId: vi.fn(async () => workflowRun.deploymentId), createQueueHandler: vi.fn( ( diff --git a/packages/core/src/runtime/step-executor.ts b/packages/core/src/runtime/step-executor.ts index 64c4f50b13..9263f1e7a6 100644 --- a/packages/core/src/runtime/step-executor.ts +++ b/packages/core/src/runtime/step-executor.ts @@ -1147,7 +1147,7 @@ export async function executeStep( () => { // The last instant before user code: T7 of the resume window. reportResumeTtr(); - world.recordStepExecution?.(stepId); + 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 fa0fde6d85..a54dbdba4b 100644 --- a/packages/world-vercel/src/index.ts +++ b/packages/world-vercel/src/index.ts @@ -70,7 +70,6 @@ export function createWorld(config?: APIConfig): World { // immediately, without a redeploy of this adapter. }, getRuntimeDeadline: getDeadline, - recordStepExecution, ...createQueue(config), ...createStorage(config), // Analytics list reads are served from an eventually-ingested store. @@ -95,5 +94,6 @@ export function createWorld(config?: APIConfig): World { config?.dispatcher ), resolveLatestDeploymentId: createResolveLatestDeploymentId(config), + telemetry: { recordStepExecution }, }; } diff --git a/packages/world/src/interfaces.ts b/packages/world/src/interfaces.ts index 1bfacc29ab..5bab0ecd6e 100644 --- a/packages/world/src/interfaces.ts +++ b/packages/world/src/interfaces.ts @@ -613,14 +613,6 @@ export interface WorldCapabilities { * The "World" interface represents how Workflows are able to communicate with the outside world. */ export interface World extends Queue, Streamer, Storage { - /** - * 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; - /** * Optional analytics read namespace for observability surfaces. * @@ -789,4 +781,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; } From 685d976ccae576dcbfda1bee285a3c064ca0a330 Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Tue, 29 Sep 2026 10:48:11 -0700 Subject: [PATCH 6/6] docs: update telemetry changeset --- .changeset/tidy-pears-report.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changeset/tidy-pears-report.md b/.changeset/tidy-pears-report.md index 3ba7aaaa65..95a88d6216 100644 --- a/.changeset/tidy-pears-report.md +++ b/.changeset/tidy-pears-report.md @@ -4,4 +4,4 @@ '@workflow/world-vercel': minor --- -Add `world.recordStepExecution?` interface to let Worlds correlate flow requests with inline-executed workflow steps. +Add the optional `world.telemetry.recordStepExecution` hook so Worlds can correlate flow requests with inline-executed workflow steps.