Skip to content
Merged
7 changes: 7 additions & 0 deletions .changeset/tidy-pears-report.md
Original file line number Diff line number Diff line change
@@ -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.
21 changes: 21 additions & 0 deletions packages/core/src/runtime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
(
Expand Down Expand Up @@ -1682,6 +1684,8 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => {
handlerPromise,
order,
queue,
recordStepExecution,
durableEvents,
eventsList,
stepIdSends,
createdEventParams,
Expand Down Expand Up @@ -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<void>((resolve) => {
Expand Down
1 change: 1 addition & 0 deletions packages/core/src/runtime/step-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
);
Expand Down
3 changes: 2 additions & 1 deletion packages/world-vercel/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -119,5 +119,6 @@ export function createWorld(config?: APIConfig): World {
config?.dispatcher
),
resolveLatestDeploymentId: createResolveLatestDeploymentId(config),
telemetry: { recordStepExecution },
};
}
111 changes: 110 additions & 1 deletion packages/world-vercel/src/queue.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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', () => {
Expand Down Expand Up @@ -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<void>();
const releaseFirst = Promise.withResolvers<void>();
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' });

Expand Down
76 changes: 64 additions & 12 deletions packages/world-vercel/src/queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -152,10 +152,43 @@ class DualTransport implements Transport<unknown> {
}
}

// 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<string | undefined>();
interface QueueInvocationContext {
collectStepIds: boolean;
requestId?: string;
stepIds: Set<string>;
}

// 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<QueueInvocationContext>();

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({
Expand Down Expand Up @@ -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 } =
Expand All @@ -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) &&
Expand Down Expand Up @@ -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);
};
};

Expand Down
15 changes: 15 additions & 0 deletions packages/world/src/interfaces.ts
Original file line number Diff line number Diff line change
Expand Up @@ -838,4 +838,19 @@ export interface World extends Queue, Streamer, Storage {
| Record<string, string | null>
| null
| Promise<Record<string, string | null> | 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;
}
Loading