Skip to content

Commit e7ef9d8

Browse files
pranaygpclaude
andauthored
perf(core): lazy inline step start (save one world round-trip per step) (#2478)
* perf(core): lazy inline step start to save a world round-trip per step The owned-inline runtime path used to write step_created (suspension handler) and then step_started (executeStep) as two separate world round-trips for a step it already owns and is about to run inline. This defers the step_created write: executeStep sends a single step_started carrying the step input, and the world creates the step on the fly (materializing the step entity plus a synthetic step_created event so replay still observes it). Mirrors the existing resilient run_started -> run_created pattern. Exactly-one ownership is preserved by the world's atomic create-claim: the loser of a concurrent lazy step_started gets EntityConflictError, which executeStep maps to `skipped`, so it never runs the body. A lazy step_started is only ever sent for a brand-new step (the suspension handler defers only steps with no prior step_created), so crash recovery still re-runs a `running` step via the normal non-lazy step_started. Worlds updated: world-local, world-postgres (implicit create + synthetic step_created event), world-vercel (routes the input as the v4 frame payload and threads the server's stepCreated flag). @workflow/world adds optional `input` to step_started and a `stepCreated` EventResult signal. Rollout: server-first. The matching workflow-server change must deploy before this ships; the Vercel world targets a single Vercel-operated backend (server always >= SDK). For local/postgres the world ships in the same package as the runtime, so there is no version skew. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * fix(core): materialize deferred step before failing unregistered step on lazy inline path The lazy inline step-start optimization defers a step's step_created write, expecting executeStep to materialize the step via a lazy step_started carrying its input. For an UNREGISTERED step, executeStep bails out before sending that step_started and writes step_failed directly — but the step entity was never created, so the world's "step must exist" ordering guard rejects the step_failed and the run wedges (times out). This regressed the StepNotRegisteredError e2e tests uniformly across every framework/world (the ghost step never reached `failed`). Fix: on the lazy path, send the lazy step_started first to materialize the step (entity + synthetic step_created, keeping replay correct), then write step_failed. The lazy step_started's atomic create-claim preserves exactly-one-owner: a concurrent winner makes ours reject with EntityConflictError → skipped, so the failure is never written twice. Adds world-level regression tests (world-local, world-postgres) asserting a lazy step_started followed by step_failed marks the step failed. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
1 parent 2074f91 commit e7ef9d8

13 files changed

Lines changed: 968 additions & 106 deletions

File tree

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
---
2+
'workflow': patch
3+
'@workflow/core': patch
4+
'@workflow/world': patch
5+
'@workflow/world-local': patch
6+
'@workflow/world-postgres': patch
7+
'@workflow/world-vercel': patch
8+
---
9+
10+
Lazy inline step start: the owned-inline runtime path now sends a single `step_started` carrying the step input, letting the world create the step on the fly and saving one round-trip per inline step.
11+
12+
`@workflow/world`: `step_started` event data accepts an optional `input`, and `EventResult` gains a `stepCreated` ownership signal.
13+
14+
`@workflow/world-local`: `step_started` with input atomically creates the step plus a synthetic `step_created` event; a lazy `step_started` for an already-existing step throws `EntityConflictError` so concurrent losers skip (exactly-once).
15+
16+
`@workflow/world-postgres`: same lazy-create + exactly-once create-claim for the Postgres backend.
17+
18+
`@workflow/world-vercel`: sends the step input on `step_started` over the v4 wire and threads the server's `stepCreated` signal into `EventResult`.

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

Lines changed: 91 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,11 @@
1-
import {
2-
EntityConflictError,
3-
RUN_ERROR_CODES,
4-
WorkflowWorldError,
5-
} from '@workflow/errors';
1+
import { RUN_ERROR_CODES, WorkflowWorldError } from '@workflow/errors';
62
import {
73
type Event,
84
SPEC_VERSION_CURRENT,
95
type WorkflowRun,
106
} from '@workflow/world';
117
import { afterEach, describe, expect, it, vi } from 'vitest';
8+
import { registerStepFunction } from './private.js';
129
import { REPLAY_DIVERGENCE_MAX_RETRIES } from './runtime/constants.js';
1310
import { setWorld } from './runtime/world.js';
1411
import { workflowEntrypoint } from './runtime.js';
@@ -629,7 +626,13 @@ describe('workflowEntrypoint replay guards', () => {
629626
expect(firstAttemptEvents).toContainEqual(
630627
expect.objectContaining({ eventType: 'attr_set' })
631628
);
632-
expect(firstAttemptEvents).toContainEqual(
629+
// Under lazy inline start the step that loses the attribute race is NOT
630+
// eagerly created: its step_created is deferred for a lazy step_started
631+
// that never fires, because the attr_set event triggers an immediate
632+
// re-invocation before any step executes. So no step_created is written on
633+
// this attempt — strictly less event-log garbage than the eager model,
634+
// and correct because the step loses the race and is abandoned on replay.
635+
expect(firstAttemptEvents).not.toContainEqual(
633636
expect.objectContaining({ eventType: 'step_created' })
634637
);
635638
expect(firstAttemptMessages).toEqual([]);
@@ -774,7 +777,7 @@ describe('workflowEntrypoint replay guards', () => {
774777
);
775778
});
776779

777-
it('propagates transient step_created failures to the queue handler without an unhandled rejection', async () => {
780+
it('propagates transient step-creation failures (lazy step_started) to the queue handler without an unhandled rejection', async () => {
778781
const createdEvents: unknown[] = [];
779782
const workflowRun: WorkflowRun = {
780783
runId: 'wrun_step_created_parse',
@@ -792,7 +795,10 @@ describe('workflowEntrypoint replay guards', () => {
792795
deploymentId: 'test-deployment',
793796
};
794797
// Simulates a transient network failure on POST /runs/{id}/events
795-
// (e.g. the connection terminated mid-response-body).
798+
// (e.g. the connection terminated mid-response-body). Under lazy inline
799+
// start the step is created on the fly by its step_started, so the
800+
// transient failure surfaces there (the standalone step_created round-trip
801+
// no longer exists on this path).
796802
const parseError = new WorkflowWorldError(
797803
'Failed to parse response body for POST /v3/runs/wrun_step_created_parse/events (Content-Type: application/cbor):\n\nTypeError: terminated',
798804
{ code: 'PARSE_ERROR' }
@@ -801,7 +807,7 @@ describe('workflowEntrypoint replay guards', () => {
801807
if (data.eventType === 'run_started') {
802808
return { run: workflowRun, events: [] };
803809
}
804-
if (data.eventType === 'step_created') {
810+
if (data.eventType === 'step_started') {
805811
throw parseError;
806812
}
807813
createdEvents.push(data);
@@ -893,21 +899,28 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => {
893899
`;globalThis.__private_workflows = new Map();
894900
globalThis.__private_workflows.set(${JSON.stringify(workflowName)}, ${workflowName});`;
895901

896-
// A workflow that suspends on a step AND a sleep. The harness's
897-
// events.create answers step_created with EntityConflictError (see
898-
// driveHandler), so this handler observes the pending step without owning
899-
// it — the crash-recovery shape where a prior handler wrote step_created
900-
// but died before queueing. That forces the unified dispatch to QUEUE the
901-
// step (inline execution is ownership-gated), exercising the
902-
// progress-critical step-dispatch queue() send that must complete before
903-
// the orchestrator message is acked.
902+
// A workflow that suspends on TWO parallel steps and a sleep. Under the
903+
// lazy-inline-start model exactly one pending step is run inline (its
904+
// step_created is deferred and folded into a lazy step_started); every other
905+
// pending step keeps its eager step_created and is QUEUED via the unified
906+
// dispatch. So the second step here is always queued, exercising the
907+
// progress-critical step-dispatch queue() send that must complete before the
908+
// orchestrator message is acked — independent of which step the runtime
909+
// happens to pick for inline execution.
904910
const stepWithSleepWorkflow = `const add = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("add");
911+
const addB = globalThis[Symbol.for("WORKFLOW_USE_STEP")]("addB");
905912
const sleep = globalThis[Symbol.for("WORKFLOW_SLEEP")];
906913
async function workflow() {
907-
const [a] = await Promise.all([add(1, 2), sleep('1h')]);
914+
const [a] = await Promise.all([add(1, 2), addB(3, 4), sleep('1h')]);
908915
return a;
909916
}${getWorkflowTransformCode('workflow')}`;
910917

918+
// Register the two steps so the one chosen for inline execution actually
919+
// runs (and completes) instead of failing as unregistered; the other is
920+
// queued. Both are no-ops — these tests only assert dispatch/ack ordering.
921+
registerStepFunction('add', async () => undefined);
922+
registerStepFunction('addB', async () => undefined);
923+
911924
async function makeRunningRun(runId: string): Promise<WorkflowRun> {
912925
return {
913926
runId,
@@ -946,25 +959,68 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => {
946959
// `.only`, not just the afterEach reset between this suite's tests.
947960
waitUntilPromises.length = 0;
948961

962+
// Stateful event log so replay converges instead of re-suspending forever:
963+
// the inline step's events and the queued step's eager step_created are
964+
// recorded here and returned by `list`, so a later loop iteration observes
965+
// the inline step as done and the queued step as already-created (and thus
966+
// not re-run/re-inlined).
967+
let eventSeq = 0;
968+
const durableEvents: Event[] = [];
969+
const recordEvent = (data: any): Event => {
970+
eventSeq += 1;
971+
const created = {
972+
eventId: `event-${eventSeq}`,
973+
runId: workflowRun.runId,
974+
createdAt: new Date(),
975+
...data,
976+
} as Event;
977+
durableEvents.push(created);
978+
return created;
979+
};
980+
949981
const eventsCreate = vi.fn(async (_runId: string, data: any) => {
950982
if (data.eventType === 'run_started') {
951983
return { run: workflowRun, events: [] as Event[] };
952984
}
953985
if (data.eventType === 'step_created') {
986+
// Eager step_created for the QUEUED step (the one not run inline).
987+
// It must be durably created before its dispatch send — the ordering
988+
// assertion below checks step_created precedes queue_dispatch_start.
954989
order.push('step_created');
955-
// Concurrent-handler simulation: the step_created event already
956-
// exists, so this handler doesn't own the step and must queue it
957-
// rather than execute it inline.
958-
throw new EntityConflictError('step already exists');
990+
return { event: recordEvent(data) };
959991
}
960-
return {
961-
event: {
962-
eventId: `event-${order.length}`,
963-
runId: workflowRun.runId,
964-
createdAt: new Date(),
965-
...data,
966-
},
967-
};
992+
if (data.eventType === 'step_started') {
993+
// The inline step's lazy step_started creates the step on the fly:
994+
// record a synthetic step_created so replay observes it, then the
995+
// step_started, and return a running step so executeStep can run the
996+
// (registered, no-op) body to completion.
997+
const lazy = data.eventData as { stepName?: string; input?: unknown };
998+
if (lazy?.input !== undefined) {
999+
recordEvent({
1000+
eventType: 'step_created',
1001+
specVersion: SPEC_VERSION_CURRENT,
1002+
correlationId: data.correlationId,
1003+
eventData: { stepName: lazy.stepName, input: lazy.input },
1004+
});
1005+
}
1006+
const created = recordEvent(data);
1007+
return {
1008+
event: created,
1009+
step: {
1010+
runId: workflowRun.runId,
1011+
stepId: data.correlationId,
1012+
stepName: lazy?.stepName,
1013+
status: 'running' as const,
1014+
attempt: 1,
1015+
input: lazy?.input,
1016+
startedAt: new Date(),
1017+
createdAt: new Date(),
1018+
updatedAt: new Date(),
1019+
},
1020+
...(lazy?.input !== undefined ? { stepCreated: true } : {}),
1021+
};
1022+
}
1023+
return { event: recordEvent(data) };
9681024
});
9691025

9701026
const queue = vi.fn(async (queueName: string, message: any) => {
@@ -1004,8 +1060,12 @@ describe('workflowEntrypoint step-dispatch ack ordering', () => {
10041060
),
10051061
events: {
10061062
create: eventsCreate,
1063+
// Return the accumulated event log so replay converges: a later loop
1064+
// iteration sees the inline step completed and the queued step already
1065+
// created (so neither is re-run), and the handler returns instead of
1066+
// re-suspending forever.
10071067
list: vi.fn(async () => ({
1008-
data: [] as Event[],
1068+
data: [...durableEvents],
10091069
hasMore: false,
10101070
cursor: 'cursor_test',
10111071
})),

‎packages/core/src/runtime.ts‎

Lines changed: 43 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -1245,41 +1245,42 @@ export function workflowEntrypoint(
12451245

12461246
const pendingSteps = suspensionResult.pendingSteps;
12471247

1248-
// Inline execution is gated on ownership: only the
1249-
// handler that actually wrote the step_created event
1250-
// may run the step body inline. The world-level
1251-
// step_created is atomic per-correlationId, so
1252-
// exactly one handler owns each step — concurrent
1253-
// handlers can't race on step execution.
1254-
const ownedPendingSteps = pendingSteps.filter((s) =>
1255-
suspensionResult.createdStepCorrelationIds.has(
1256-
s.correlationId
1257-
)
1258-
);
1259-
1260-
// Pick one owned step to execute inline (if any).
1261-
// The rest of the pending steps, plus any wait
1262-
// timer, are queued below in a single parallel
1263-
// batch.
1248+
// Inline execution is gated on ownership. The
1249+
// suspension handler deferred the step_created write for
1250+
// exactly one step (`lazyInlineStep`) so we can run it
1251+
// inline via a lazy `step_started` that creates the step
1252+
// on the fly — saving one world round-trip. Ownership is
1253+
// still atomic and exactly-one: the world's
1254+
// create-claim inside that step_started returns
1255+
// `EntityConflictError` (→ executeStep `skipped`) to any
1256+
// concurrent loser, so only one handler ever runs the
1257+
// body. Every other pending step keeps its eager
1258+
// step_created (in `createdStepCorrelationIds`) and is
1259+
// queued below.
12641260
//
1265-
// Skip inline execution when a created hook has a
1266-
// `hook.getConflict()` awaiter: an inline `await
1267-
// executeStep(...)` blocks this handler for the
1268-
// full step duration, so the awaiter's continuation
1269-
// (which only advances on the next replay) would be
1270-
// serialized behind the step — defeating work the
1271-
// workflow expressed as parallel (e.g.
1272-
// `hook.getConflict().then(() => stepB())` racing
1273-
// `await stepA()`). Queue every step and re-invoke
1274-
// immediately instead: the re-invocation replays
1275-
// over the just-committed hook_created and resolves
1276-
// the awaiter while the queued steps execute in
1261+
// The suspension handler only designates a
1262+
// `lazyInlineStep` when no `hook.getConflict()` awaiter
1263+
// is present. That awaiter case must execute nothing
1264+
// inline: an inline `await executeStep(...)` blocks this
1265+
// handler for the full step duration, so the awaiter's
1266+
// continuation (which only advances on the next replay)
1267+
// would be serialized behind the step — defeating work
1268+
// the workflow expressed as parallel (e.g.
1269+
// `hook.getConflict().then(() => stepB())` racing `await
1270+
// stepA()`). In that case `lazyInlineStep` is undefined
1271+
// and every step is queued for re-invocation, which
1272+
// replays over the just-committed hook_created and
1273+
// resolves the awaiter while queued steps run in
12771274
// parallel invocations.
1275+
const lazyInline = suspensionResult.lazyInlineStep;
12781276
const inlineStep:
12791277
| (typeof pendingSteps)[number]
1280-
| undefined = suspensionResult.hasAwaitedHookCreation
1281-
? undefined
1282-
: ownedPendingSteps[0];
1278+
| undefined = lazyInline
1279+
? pendingSteps.find(
1280+
(s) =>
1281+
s.correlationId === lazyInline.correlationId
1282+
)
1283+
: undefined;
12831284

12841285
// Unified queue dispatch for everything we are NOT
12851286
// inline-executing. Steps are queued with stepId so
@@ -1389,9 +1390,11 @@ export function workflowEntrypoint(
13891390
// (`preInlineWriteCursor`; a World may return none on
13901391
// the initial load).
13911392
// - This is the clean single-step sequential case:
1392-
// this suspension created exactly one step and no
1393-
// hooks or waits (`err.{step,hook,wait}Count`), and
1394-
// that one pending step is owned by us (no parallel
1393+
// this suspension produced exactly one step and no
1394+
// hooks or waits (`err.{step,hook,wait}Count`), that
1395+
// one step is the lone pending step
1396+
// (`pendingSteps.length === 1`), and it is the one we
1397+
// are running inline (`inlineStep` — no parallel
13951398
// siblings queued to background handlers, which would
13961399
// write their own events out of band).
13971400
// - No pending wait timer from THIS suspension.
@@ -1413,7 +1416,7 @@ export function workflowEntrypoint(
14131416
err.hookCount === 0 &&
14141417
err.waitCount === 0 &&
14151418
pendingSteps.length === 1 &&
1416-
ownedPendingSteps.length === 1 &&
1419+
inlineStep !== undefined &&
14171420
!suspensionResult.waitTimeout &&
14181421
!hasOpenHookOrWait(cachedEvents ?? []);
14191422

@@ -1429,6 +1432,11 @@ export function workflowEntrypoint(
14291432
stepId: inlineStep.correlationId,
14301433
stepName: inlineStep.stepName,
14311434
runSpecVersion: workflowRun.specVersion,
1435+
// Lazy inline start: send the deferred step's input
1436+
// on step_started so the world creates the step on
1437+
// the fly. lazyInline is defined whenever inlineStep
1438+
// is (both derive from suspensionResult.lazyInlineStep).
1439+
lazyStepInput: lazyInline?.dehydratedInput,
14321440
...(requestInlineDelta && preInlineWriteCursor
14331441
? { inlineDeltaSinceCursor: preInlineWriteCursor }
14341442
: {}),

0 commit comments

Comments
 (0)