Skip to content

Commit 03bf52c

Browse files
committed
Address review: terminal-only lazy status gate, terminal-run start fence in local worlds, prologue telemetry, last-completer coverage
- The last completer's lazy runs.get result is now gated on isTerminalWorkflowRunStatus (with a debug log): a stale 'pending' read — a run with completed steps has necessarily started — no longer silently abandons the fan-out's continuation; it falls through to the inline replay, whose next entity write is fenced server-side if the run truly ended meanwhile. - world-local / world-postgres now reject step_started on terminal runs even when the step row still reads 'running' (a redelivered start a previous delivery claimed): starting work on a finished run is never valid, and previously the body re-ran with its outcome unconsumable. In-flight steps still write their terminal events unchanged. This closes the adapter gap behind the fetch-free prologue's reliance on the step_started claim as the run-liveness check, and the prologue comment now states the contract precisely. - workflow.step.dispatch_prologue span attribute ('run_context' | 'runs_get') makes fetch-free adoption and the saved round trip observable during version-skew windows. - Restated why the eager redelivery re-ensure survives on the fetch-free path (no run fetch to overlap; still cheaper than the in-band recovery's failed-start round trip). - New two-phase fan-out coverage: a real replay emits the queued step message (asserting the stamped runContext), then its redelivery runs as the LAST completer — zero reads before the step, exactly one lazy runs.get, run completed; plus the stale-'pending' fall-through and the genuinely-terminal skip.
1 parent 890f912 commit 03bf52c

6 files changed

Lines changed: 316 additions & 15 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
'@workflow/world-local': patch
3+
'@workflow/world-postgres': patch
4+
---
5+
6+
Reject `step_started` on terminal runs even when the step row still reads `running`: a redelivered start on a cancelled/completed run previously passed the claim and re-executed the step body whose outcome nothing would consume. In-flight steps can still write their terminal events (`step_completed`/`step_failed`) unchanged.

‎packages/core/src/runtime.test.ts‎

Lines changed: 225 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2194,6 +2194,231 @@ describe('workflowEntrypoint resilient step consumption (stepInput re-ensure)',
21942194
});
21952195
});
21962196

2197+
// The fan-out path this PR optimizes end to end: phase 1 drives a REAL replay
2198+
// of a two-step Promise.all fan-out (inline cap 1, so the runtime itself
2199+
// emits the queued step message with its seeded correlation id and stamped
2200+
// runContext); phase 2 redelivers that exact message against the shared event
2201+
// log, making it the LAST completer — exercising the fetch-free prologue, the
2202+
// single lazy runs.get before the inline replay, and the terminal-only status
2203+
// gate on its result.
2204+
describe('workflowEntrypoint fan-out last completer (fetch-free prologue)', () => {
2205+
beforeEach(() => {
2206+
process.env.WORKFLOW_MAX_INLINE_STEPS = '1';
2207+
});
2208+
afterEach(() => {
2209+
delete process.env.WORKFLOW_MAX_INLINE_STEPS;
2210+
setWorld(undefined);
2211+
vi.clearAllMocks();
2212+
});
2213+
2214+
const getWorkflowTransformCode = (workflowName: string) =>
2215+
`;globalThis.__private_workflows = new Map();
2216+
globalThis.__private_workflows.set(${JSON.stringify(workflowName)}, ${workflowName});`;
2217+
2218+
const fanOutWorkflow = `const fanA = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("fanA");
2219+
const fanB = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("fanB");
2220+
async function workflow() {
2221+
const [a, b] = await Promise.all([fanA(1, 2), fanB(3, 4)]);
2222+
return a + b;
2223+
}${getWorkflowTransformCode('workflow')}`;
2224+
2225+
registerStepFunction('fanA', async (a: number, b: number) => a + b);
2226+
registerStepFunction('fanB', async (a: number, b: number) => a + b);
2227+
2228+
async function driveFanOut(opts: {
2229+
runId: string;
2230+
/** Status the LAST completer's lazy runs.get returns. */
2231+
lazyRunStatus: 'running' | 'pending' | 'completed';
2232+
}) {
2233+
const workflowRun: WorkflowRun = {
2234+
runId: opts.runId,
2235+
workflowName: 'workflow',
2236+
status: 'running',
2237+
specVersion: SPEC_VERSION_CURRENT,
2238+
input: await dehydrateWorkflowArguments([], opts.runId, undefined, []),
2239+
createdAt: new Date('2024-01-01T00:00:00.000Z'),
2240+
updatedAt: new Date('2024-01-01T00:00:00.000Z'),
2241+
startedAt: new Date('2024-01-01T00:00:00.000Z'),
2242+
deploymentId: 'test-deployment',
2243+
};
2244+
2245+
// Shared durable state across the two deliveries.
2246+
let eventSeq = 0;
2247+
const durableEvents: Event[] = [
2248+
{
2249+
eventId: 'event-run-created',
2250+
runId: opts.runId,
2251+
createdAt: new Date('2024-01-01T00:00:00.000Z'),
2252+
eventType: 'run_created',
2253+
specVersion: SPEC_VERSION_CURRENT,
2254+
eventData: {
2255+
deploymentId: 'test-deployment',
2256+
workflowName: 'workflow',
2257+
input: workflowRun.input,
2258+
},
2259+
} as unknown as Event,
2260+
];
2261+
const recordEvent = (data: any): Event => {
2262+
eventSeq += 1;
2263+
const created = {
2264+
eventId: `event-${eventSeq}`,
2265+
runId: opts.runId,
2266+
createdAt: new Date(),
2267+
...data,
2268+
} as Event;
2269+
durableEvents.push(created);
2270+
return created;
2271+
};
2272+
// Step inputs recorded at creation, so a BARE start can return the entity
2273+
// with the input the body will hydrate (like a real World).
2274+
const stepInputsByCorrelationId = new Map<string, unknown>();
2275+
2276+
const eventsCreate = vi.fn(async (_runId: string, data: any) => {
2277+
if (data.eventType === 'run_started') {
2278+
return { run: workflowRun, events: [...durableEvents] };
2279+
}
2280+
if (data.eventType === 'step_created') {
2281+
stepInputsByCorrelationId.set(
2282+
data.correlationId,
2283+
data.eventData?.input
2284+
);
2285+
return { event: recordEvent(data) };
2286+
}
2287+
if (data.eventType === 'step_started') {
2288+
const lazy = data.eventData as { stepName?: string; input?: unknown };
2289+
if (lazy?.input !== undefined) {
2290+
// Lazy inline start: create the step on the fly.
2291+
stepInputsByCorrelationId.set(data.correlationId, lazy.input);
2292+
recordEvent({
2293+
eventType: 'step_created',
2294+
specVersion: SPEC_VERSION_CURRENT,
2295+
correlationId: data.correlationId,
2296+
eventData: { stepName: lazy.stepName, input: lazy.input },
2297+
});
2298+
}
2299+
const created = recordEvent(data);
2300+
return {
2301+
event: created,
2302+
step: {
2303+
runId: opts.runId,
2304+
stepId: data.correlationId,
2305+
stepName: lazy?.stepName,
2306+
status: 'running' as const,
2307+
attempt: 1,
2308+
input: stepInputsByCorrelationId.get(data.correlationId),
2309+
startedAt: new Date(),
2310+
createdAt: new Date(),
2311+
updatedAt: new Date(),
2312+
},
2313+
...(lazy?.input !== undefined ? { stepCreated: true } : {}),
2314+
};
2315+
}
2316+
return { event: recordEvent(data) };
2317+
});
2318+
2319+
const queuedMessages: any[] = [];
2320+
const queue = vi.fn(async (_queueName: string, message: any) => {
2321+
queuedMessages.push(message);
2322+
return { messageId: null };
2323+
});
2324+
const runsGet = vi.fn(async () => ({
2325+
...workflowRun,
2326+
status: opts.lazyRunStatus,
2327+
}));
2328+
2329+
// Capture the queue handler so both phases can deliver messages directly.
2330+
let handler!: (message: unknown, metadata: unknown) => Promise<unknown>;
2331+
setWorld({
2332+
specVersion: SPEC_VERSION_CURRENT,
2333+
createQueueHandler: vi.fn((_prefix: string, h: typeof handler) => {
2334+
handler = h;
2335+
return async () => new Response(null, { status: 204 });
2336+
}),
2337+
events: {
2338+
create: eventsCreate,
2339+
list: vi.fn(async () => ({
2340+
data: [...durableEvents],
2341+
hasMore: false,
2342+
cursor: 'cursor_test',
2343+
})),
2344+
},
2345+
runs: { get: runsGet },
2346+
queue,
2347+
getEncryptionKeyForRun: vi.fn(async () => undefined),
2348+
} as any);
2349+
2350+
const entry = workflowEntrypoint(fanOutWorkflow);
2351+
await entry(new Request('https://example.test')); // binds `handler`
2352+
2353+
// Phase 1: flow delivery. The real replay inline-executes one step (cap 1)
2354+
// and queues the other with its seeded correlation id and runContext.
2355+
await handler(
2356+
{ runId: opts.runId, requestedAt: new Date() },
2357+
{
2358+
requestId: 'req_1',
2359+
attempt: 1,
2360+
queueName: '__wkf_workflow_workflow',
2361+
messageId: 'msg_flow_1',
2362+
}
2363+
);
2364+
const stepMessage = queuedMessages.find((m) => m && m.stepId);
2365+
expect(stepMessage).toBeDefined();
2366+
expect(runsGet).not.toHaveBeenCalled();
2367+
2368+
// Phase 2: redeliver the runtime's own step message — the last completer.
2369+
await handler(stepMessage, {
2370+
requestId: 'req_2',
2371+
attempt: 1,
2372+
queueName: '__wkf_workflow_workflow',
2373+
messageId: 'msg_step_1',
2374+
});
2375+
2376+
return { stepMessage, durableEvents, runsGet };
2377+
}
2378+
2379+
it('completes the run via the fetch-free prologue with exactly one lazy runs.get', async () => {
2380+
const { stepMessage, durableEvents, runsGet } = await driveFanOut({
2381+
runId: 'wrun_fanout_last_completer',
2382+
lazyRunStatus: 'running',
2383+
});
2384+
2385+
// The producer stamped the run identity on its own message.
2386+
expect(stepMessage.runContext).toMatchObject({
2387+
deploymentId: 'test-deployment',
2388+
specVersion: SPEC_VERSION_CURRENT,
2389+
rootRunId: 'wrun_fanout_last_completer',
2390+
});
2391+
// Zero reads before the step; exactly one for the inline replay.
2392+
expect(runsGet).toHaveBeenCalledTimes(1);
2393+
// The last completer replayed inline and finished the run.
2394+
expect(durableEvents.map((e) => e.eventType)).toContain('run_completed');
2395+
});
2396+
2397+
it('falls through a stale `pending` read instead of silently abandoning the fan-out', async () => {
2398+
const { durableEvents, runsGet } = await driveFanOut({
2399+
runId: 'wrun_fanout_stale_pending',
2400+
lazyRunStatus: 'pending',
2401+
});
2402+
2403+
// A run with completed steps cannot truly be pending — the read is stale.
2404+
// The terminal-only gate lets the replay proceed, so the continuation is
2405+
// not dropped on the floor.
2406+
expect(runsGet).toHaveBeenCalledTimes(1);
2407+
expect(durableEvents.map((e) => e.eventType)).toContain('run_completed');
2408+
});
2409+
2410+
it('still skips the inline replay when the run is genuinely terminal', async () => {
2411+
const { durableEvents } = await driveFanOut({
2412+
runId: 'wrun_fanout_terminal',
2413+
lazyRunStatus: 'completed',
2414+
});
2415+
2416+
expect(durableEvents.map((e) => e.eventType)).not.toContain(
2417+
'run_completed'
2418+
);
2419+
});
2420+
});
2421+
21972422
describe('workflowEntrypoint turbo mode', () => {
21982423
const ORIG_TURBO = process.env.WORKFLOW_TURBO;
21992424
const ORIG_OPT = process.env.WORKFLOW_OPTIMISTIC_INLINE_START;

‎packages/core/src/runtime.ts‎

Lines changed: 48 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import {
2727
getQueueTopicPrefix,
2828
isLegacySpecVersion,
2929
isTerminalRunEventType,
30+
isTerminalWorkflowRunStatus,
3031
type RunInput,
3132
resolveQueueNamespace,
3233
SPEC_VERSION_CURRENT,
@@ -1415,11 +1416,14 @@ export function workflowEntrypoint(
14151416
// ~300s visibility-timeout redelivery — measured
14161417
// exactly so in the durabench parallel sweeps before
14171418
// this path existed.
1418-
// - EAGERLY on a genuine redelivery (attempt > 1),
1419-
// in parallel with the run fetch below — a
1419+
// - EAGERLY on a genuine redelivery (attempt > 1) — a
14201420
// redelivered dispatch already had its create race
1421-
// resolved either way, so this saves the failed
1422-
// start round-trip at no wall-time cost. First
1421+
// resolved either way, so ensuring up front saves the
1422+
// failed-start round trip the in-band recovery would
1423+
// otherwise pay. On the legacy prologue it overlaps
1424+
// the run fetch (no wall-time cost); on the
1425+
// fetch-free (runContext) prologue it is the sole
1426+
// pre-step write and still the cheaper trade. First
14231427
// deliveries skip it: the producer's write almost
14241428
// always lands, and an eager ensure would burn a
14251429
// conditional write per step.
@@ -1492,18 +1496,37 @@ export function workflowEntrypoint(
14921496
// entirely — one less round trip on the TTLS-critical
14931497
// path, and N fewer reads on the run's partition per
14941498
// fan-out (vercel/workflow#3456). The run-status early
1495-
// exit is not lost: a terminal run rejects the
1496-
// `step_started` claim server-side (RunExpired → gone,
1497-
// terminal step → skipped). Older messages without the
1498-
// field keep the legacy fetch.
1499+
// exit is not lost: every World rejects a `step_started`
1500+
// claim on a terminal run (RunExpired → gone, terminal
1501+
// step → skipped) — including a redelivered start whose
1502+
// step row still reads `running`, which world-local and
1503+
// world-postgres reject as of this change (previously
1504+
// only world-vercel's run-status fence covered that
1505+
// shape, and the body could re-run on a finished run).
1506+
// Older messages without the field keep the legacy
1507+
// fetch.
14991508
let bgRun: WorkflowRun | undefined;
15001509
let runIdentity: {
15011510
deploymentId: string;
15021511
specVersion: number;
15031512
startedAt?: number;
15041513
rootRunId?: string;
15051514
};
1515+
// Which prologue ran — makes adoption of the fetch-free
1516+
// path (and the round trip it saves) directly observable
1517+
// during version-skew windows.
1518+
span?.setAttributes(
1519+
Attribute.StepDispatchPrologue(
1520+
runContext ? 'run_context' : 'runs_get'
1521+
)
1522+
);
15061523
if (runContext) {
1524+
// The eager redelivery re-ensure has no run fetch to
1525+
// overlap with on this path — it is kept because a
1526+
// redelivered dispatch has already had its create race
1527+
// resolved, so one conditional write here is cheaper
1528+
// than letting the bare start fail and paying the
1529+
// in-band recovery's extra start round trip.
15071530
const ensureOutcome =
15081531
stepInput && metadata.attempt > 1
15091532
? await ensureStepFromMessage()
@@ -1782,7 +1805,23 @@ export function workflowEntrypoint(
17821805
(await world.runs.get(runId, {
17831806
resolveData: 'none',
17841807
}));
1785-
if (replayRunRow.status !== 'running') {
1808+
// Terminal statuses only: a `pending` read here is a
1809+
// stale row (this run has completed steps, so it has
1810+
// started), and under the fetch-free prologue this is
1811+
// the last completer's ONLY status read — returning on
1812+
// it would silently abandon the fan-out's continuation
1813+
// (the final step_completed is written, the inline
1814+
// replay never runs). Fall through instead: the replay
1815+
// synthesizes `running` and the next entity write is
1816+
// fenced server-side if the run truly ended meanwhile.
1817+
if (isTerminalWorkflowRunStatus(replayRunRow.status)) {
1818+
runtimeLogger.debug(
1819+
'Run already finished, skipping inline replay after background step',
1820+
{
1821+
workflowRunId: runId,
1822+
status: replayRunRow.status,
1823+
}
1824+
);
17861825
return;
17871826
}
17881827
const runCreatedEvent = loaded.events.find(

‎packages/core/src/telemetry/semantic-conventions.ts‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -491,6 +491,17 @@ export const StepResilientDispatchMaterialized = SemanticConvention<boolean>(
491491
'workflow.step.resilient_dispatch_materialized'
492492
);
493493

494+
/**
495+
* How the queued-step consumer resolved the run identity for this execution:
496+
* `run_context` — carried on the dispatch message, no `runs.get` before the
497+
* step (the fetch-free prologue); `runs_get` — the legacy blocking fetch
498+
* (message from an older producer). Distinguishes the two paths during
499+
* version-skew windows and makes the saved round trip directly measurable.
500+
*/
501+
export const StepDispatchPrologue = SemanticConvention<
502+
'run_context' | 'runs_get'
503+
>('workflow.step.dispatch_prologue');
504+
494505
// Webhook attributes
495506

496507
/** Number of webhook handlers triggered */

‎packages/world-local/src/storage/events-storage.ts‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1305,11 +1305,21 @@ export function createEventsStorage(
13051305
);
13061306
}
13071307

1308-
// On terminal runs: only allow completing/failing in-progress steps
1308+
// On terminal runs: only allow completing/failing in-progress
1309+
// steps. A step_started is never that — it begins work, and no
1310+
// work should begin on a finished run — so it is rejected even
1311+
// when the step row still reads `running` (a redelivery of a
1312+
// start a previous delivery already claimed). Without this, a
1313+
// redelivered start on a cancelled/completed run passes the
1314+
// claim and executes the step body whose outcome nothing will
1315+
// ever consume.
13091316
if (currentRun && isTerminalWorkflowRunStatus(currentRun.status)) {
1310-
if (validatedStep.status !== 'running') {
1317+
if (
1318+
validatedStep.status !== 'running' ||
1319+
data.eventType === 'step_started'
1320+
) {
13111321
throw new RunExpiredError(
1312-
`Cannot modify non-running step on run in terminal state "${currentRun.status}"`
1322+
`Cannot ${data.eventType === 'step_started' ? 'start' : 'modify non-running'} step on run in terminal state "${currentRun.status}"`
13131323
);
13141324
}
13151325
}

0 commit comments

Comments
 (0)