Skip to content

Commit c71befe

Browse files
pranaygpclaudeTooTallNate
authored
fix(core): premature suspension when hooks have buffered payloads with concurrent pending entities (#1294)
* test: reproduce hook+sleep promiseQueue regression Add failing tests that reproduce a regression from #1246 where the sleep's WorkflowSuspension fires before all hook payloads are delivered when a hook and sleep run concurrently. Root cause: when the null event fires, the sleep queues a suspension through promiseQueue. After the first hook payload resolves, subsequent hook payload resolutions are queued AFTER the sleep suspension, causing the workflow to terminate prematurely via Promise.race. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * test: expand tests to isolate bug to hook + pending entity pattern Add control tests proving sequential steps are NOT affected by the promiseQueue regression, isolating the bug to hooks specifically: Failing (hook-based): - hook + sleep: all 3 payloads → step invocation - hook + sleep: 2 payloads → return - hook + incomplete step: 2 payloads → return Passing (step-based controls): - sleep + sequential steps: both step events exist - sleep + sequential steps: only 1st step completed - incomplete step + sequential steps: all step events exist - hook only (no concurrent entity): payloads + step The bug is: any entity that queues a suspension through promiseQueue at null-event time preempts hook payload delivery, because hooks buffer payloads in payloadsQueue and only resolve them one-at-a-time as the workflow code iterates. Steps are unaffected because each step has its own events consumed before null fires. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> * fix(core): move WorkflowSuspension calls back to setTimeout(0) from promiseQueue The promiseQueue refactor (#1246) moved ALL promise resolutions through the queue, including WorkflowSuspension calls. This caused premature termination when hooks had buffered payloads and a concurrent entity (sleep or incomplete step) was pending: 1. Null event fires → all subscribers run 2. Sleep/step queues WorkflowSuspension via promiseQueue.then() 3. Hook's next payload resolution is also queued via promiseQueue.then() 4. Suspension fires first → Promise.race terminates workflow 5. Hook payload never delivered → infinite retry loop Fix: move WorkflowSuspension calls back to setTimeout(0) (macrotask). Suspensions fire AFTER all microtask-based deliveries (promiseQueue resolve/reject for step results and hook payloads) have completed. Non-suspension resolve/reject calls remain on promiseQueue for deterministic ordering of data delivery. Affected paths: - step.ts: null event handler - sleep.ts: null event handler - hook.ts: null event handler, eventLogEmpty suspension, dispose suspension * fix(core): use pendingDeliveries counter + scheduleWhenIdle for suspensions Replace all nested setTimeout/promiseQueue suspension patterns with a clean idle-polling mechanism: 1. pendingDeliveries counter: incremented before async hydration (step results, hook payloads), decremented in finally block after delivery 2. scheduleWhenIdle(ctx, fn): polls via setTimeout(0) → check counter → if > 0, wait for promiseQueue.then() → repeat. Only fires fn when pendingDeliveries reaches 0. This correctly handles: - Sync deserialization (no encryption): counter is 0, fires immediately after first setTimeout(0) - Async deserialization (with encryption): waits for decryption to complete before firing - Multi-round hook payload delivery: each createHookPromise() call increments the counter, preventing premature suspension between delivery rounds Also adds payloadsQueue.length check to hook null handler — don't trigger suspension if buffered payloads can satisfy pending awaits. Tests now run in both sync and async modes (14 total = 7 scenarios x 2). All 164 tests pass. --------- Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com> Co-authored-by: Nathan Rajlich <n@n8.io>
1 parent 5e4ef65 commit c71befe

13 files changed

Lines changed: 886 additions & 11 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@workflow/core": patch
3+
---
4+
5+
Fix premature workflow suspension when hooks have buffered payloads and a concurrent sleep or incomplete step is pending

‎packages/core/e2e/e2e.test.ts‎

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1968,4 +1968,76 @@ describe('e2e', () => {
19681968
expect(returnValue.endTime - returnValue.startTime).toBeGreaterThan(9999);
19691969
});
19701970
});
1971+
1972+
test(
1973+
'hookWithSleepWorkflow - hook payloads delivered correctly with concurrent sleep',
1974+
{ timeout: 90_000 },
1975+
async () => {
1976+
// Regression test: when a hook and sleep run concurrently, multiple
1977+
// hook_received events should all be processed even though the sleep
1978+
// has no wait_completed event. Previously, the sleep's WorkflowSuspension
1979+
// would terminate the workflow before all hook payloads were delivered.
1980+
const token = Math.random().toString(36).slice(2);
1981+
1982+
const run = await start(await e2e('hookWithSleepWorkflow'), [token]);
1983+
1984+
// Wait for the hook to be registered
1985+
await new Promise((resolve) => setTimeout(resolve, 5_000));
1986+
1987+
// Send 3 payloads: two normal ones, then one with done=true
1988+
let hook = await getHookByToken(token);
1989+
expect(hook.runId).toBe(run.runId);
1990+
await resumeHook(hook, { type: 'subscribe', id: 1 });
1991+
1992+
// Wait for the first payload to be processed (step must complete)
1993+
await new Promise((resolve) => setTimeout(resolve, 3_000));
1994+
1995+
hook = await getHookByToken(token);
1996+
await resumeHook(hook, { type: 'subscribe', id: 2 });
1997+
1998+
await new Promise((resolve) => setTimeout(resolve, 3_000));
1999+
2000+
hook = await getHookByToken(token);
2001+
await resumeHook(hook, { type: 'done', done: true });
2002+
2003+
const returnValue = await run.returnValue;
2004+
expect(returnValue).toBeInstanceOf(Array);
2005+
expect(returnValue).toHaveLength(3);
2006+
expect(returnValue[0]).toMatchObject({
2007+
processed: true,
2008+
type: 'subscribe',
2009+
id: 1,
2010+
});
2011+
expect(returnValue[1]).toMatchObject({
2012+
processed: true,
2013+
type: 'subscribe',
2014+
id: 2,
2015+
});
2016+
expect(returnValue[2]).toMatchObject({ processed: true, type: 'done' });
2017+
2018+
const { json: runData } = await cliInspectJson(`runs ${run.runId}`);
2019+
expect(runData.status).toBe('completed');
2020+
}
2021+
);
2022+
2023+
test(
2024+
'sleepWithSequentialStepsWorkflow - sequential steps work with concurrent sleep (control)',
2025+
{ timeout: 60_000 },
2026+
async () => {
2027+
// Control test: proves that void sleep('1d').then() does NOT break
2028+
// sequential step execution. Steps have per-event consumption so the
2029+
// sleep's pending suspension doesn't interfere. This contrasts with
2030+
// hookWithSleepWorkflow where the bug manifests.
2031+
const run = await start(
2032+
await e2e('sleepWithSequentialStepsWorkflow'),
2033+
[]
2034+
);
2035+
2036+
const returnValue = await run.returnValue;
2037+
expect(returnValue).toEqual({ a: 3, b: 6, c: 10, shouldCancel: false });
2038+
2039+
const { json: runData } = await cliInspectJson(`runs ${run.runId}`);
2040+
expect(runData.status).toBe('completed');
2041+
}
2042+
);
19712043
});

‎packages/core/src/async-deserialization-ordering.test.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext {
4343
),
4444
onWorkflowError: vi.fn(),
4545
promiseQueue: Promise.resolve(),
46+
pendingDeliveries: 0,
4647
};
4748
}
4849

0 commit comments

Comments
 (0)