Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-step-vs-wait-race.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@workflow/core": patch
---

Fix `Promise.race(step, sleep)` always blocking until step completed
10 changes: 9 additions & 1 deletion docs/content/docs/changelog/eager-processing.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -230,7 +230,15 @@ When an inline step fails with retries remaining:

### Mixed Suspensions

A suspension may contain steps, hooks, and waits simultaneously. The handler creates events for all, executes any pending step inline, and returns with the wait timeout if applicable. The workflow will re-suspend on next replay for the still-pending hooks/waits.
A suspension may contain steps, hooks, and waits simultaneously. The handler creates events for all, then chooses between inline execution and queue dispatch:

- **Steps only** (no waits): one owned step is executed inline; the rest are queued. The loop continues after the inline step completes.
- **Steps + at least one wait**: every step is queued (no inline execution). The handler returns with the wait timeout. Whichever lands first — a step's continuation or the wait timer — drives the next replay.
- **Hooks / waits only**: handler returns with the wait timeout (or no timeout, for hook-only suspensions). The next continuation is driven by external resume or the wait timer.

The "no inline when there's a wait" carve-out is necessary to preserve `Promise.race(step, sleep)` semantics. Inline `await executeStep(...)` blocks the handler for the full step duration, and `wait_completed` events are only created on the *next* loop iteration's "complete elapsed waits" pass — so a longer-running step would always swallow the shorter sleep and `Promise.race` would resolve incorrectly. Queueing the step in this case lets the wait timer drive a continuation in parallel, matching V1's behavior where each step ran in a separate function invocation.

Pure step suspensions (without waits) still benefit from inline execution; the carve-out only costs an extra queue roundtrip when a step and a sleep coexist.

### Hook Conflicts

Expand Down
16 changes: 16 additions & 0 deletions packages/core/e2e/e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -581,6 +581,22 @@ describe('e2e', () => {
expect(elapsed).toBeLessThan(25_000);
});

test('sleepWinsRaceWorkflow', { timeout: 60_000 }, async () => {
const run = await start(await e2e('sleepWinsRaceWorkflow'), []);
const returnValue = await run.returnValue;
expect(returnValue.winner).toBe('sleep');
// Sleep is 1s; step would take 10s. Should resolve in ~1s, well under 5s.
expect(returnValue.durationMs).toBeLessThan(5_000);
});

test('stepWinsRaceWorkflow', { timeout: 60_000 }, async () => {
const run = await start(await e2e('stepWinsRaceWorkflow'), []);
const returnValue = await run.returnValue;
expect(returnValue.winner).toBe('step');
// Step is 1s; sleep would take 10s. Should resolve in ~1s, well under 5s.
expect(returnValue.durationMs).toBeLessThan(5_000);
});

test('nullByteWorkflow', { timeout: 60_000 }, async () => {
const run = await start(await e2e('nullByteWorkflow'), []);
const returnValue = await run.returnValue;
Expand Down
19 changes: 18 additions & 1 deletion packages/core/src/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -882,9 +882,26 @@ export function workflowEntrypoint(

// Pick one owned step to execute inline (if any).
// The rest of the pending steps are queued below.
//
// Skip inline execution entirely when the suspension
// also has a pending wait (sleep): an inline `await
// executeStep(...)` blocks the handler for the full
// step duration, so the wait timer never has a chance
// to fire on time. That defeats `Promise.race(step,
// sleep)` semantics — if the sleep is shorter than
// the step, replay still picks the step because
// wait_completed is only created on the *next* loop
// iteration, which doesn't run until the step
// finishes. Queueing every step in this case lets
// the wait timeout drive a continuation in parallel,
// matching V1's behavior where each step ran in a
// separate function invocation.
const inlineStep:
| (typeof pendingSteps)[number]
| undefined = ownedPendingSteps[0];
| undefined =
suspensionResult.timeoutSeconds === undefined
? ownedPendingSteps[0]
: undefined;

// Queue every pending step except the one we're
// executing inline. This mirrors V1's unconditional
Expand Down
30 changes: 30 additions & 0 deletions workbench/example/workflows/99_e2e.ts
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,36 @@ export async function parallelSleepWorkflow() {

//////////////////////////////////////////////////////////

async function delayMsStep(ms: number, label: string) {
'use step';
await new Promise((resolve) => setTimeout(resolve, ms));
return label;
}

export async function sleepWinsRaceWorkflow() {
'use workflow';
const startTime = Date.now();
const winner = await Promise.race([
delayMsStep(10_000, 'step'),
sleep('1s').then(() => 'sleep'),
]);
const endTime = Date.now();
return { winner, durationMs: endTime - startTime };
}

export async function stepWinsRaceWorkflow() {
'use workflow';
const startTime = Date.now();
const winner = await Promise.race([
delayMsStep(1_000, 'step'),
sleep('10s').then(() => 'sleep'),
]);
const endTime = Date.now();
return { winner, durationMs: endTime - startTime };
}

//////////////////////////////////////////////////////////

async function nullByteStep() {
'use step';
return 'null byte \0';
Expand Down
Loading