Skip to content

Commit 3560218

Browse files
authored
fix(core): preserve order across adjacent hook deliveries (#3552)
Signed-off-by: Andrew Barba <barba@hey.com>
1 parent de2a86c commit 3560218

3 files changed

Lines changed: 85 additions & 9 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+
Preserve event-log order when adjacent hook payloads resume independent workflow branches.

‎packages/core/src/delivery-barrier-coverage.test.ts‎

Lines changed: 78 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -20,16 +20,20 @@
2020
* its barrier. They then skipped both the gate and
2121
* `awaitEarlierDeliveries`' macrotask yield.
2222
*
23-
* 3. ABORT deliveries. `_setAborted` fires the signal's listeners, and a
23+
* 3. HOOK behind HOOK. Adjacent payloads for independently awaited hooks
24+
* must reach their consumers in log order even when the earlier consumer
25+
* takes more microtask hops to act on its value.
26+
*
27+
* 4. ABORT deliveries. `_setAborted` fires the signal's listeners, and a
2428
* listener may invoke a step and draw a ULID, so an abort is as
2529
* branch-deciding as any other delivery — but it resolved straight off its
2630
* `promiseQueue` slot and registered no barrier.
2731
*
28-
* Cases 1-3 assert the same thing: the replay allocates its follow-up step
32+
* Cases 1-4 assert the same thing: the replay allocates its follow-up step
2933
* ULIDs in the order the committed log recorded. A regression surfaces as the
3034
* production `ReplayDivergenceError`.
3135
*
32-
* Section 4 asserts the registry's other job directly, over a registry built
36+
* Section 5 asserts the registry's other job directly, over a registry built
3337
* by hand rather than by a replay: which entries an idle check may ignore. Get
3438
* that wrong in one direction and a chain parked on an unclaimed hook payload
3539
* deadlocks; wrong in the other and a suspension preempts a batch of parked
@@ -406,7 +410,75 @@ describe('hook payload delivery ordering against an earlier step result', () =>
406410
}
407411
});
408412

409-
// ─── 3. abort deliveries ───────────────────────────────────────────────────
413+
// ─── 3. hook behind hook, across different consumer shapes ─────────────────
414+
describe('hook payload delivery ordering against an earlier hook payload', () => {
415+
it('keeps the recorded ULID allocation when the earlier hook uses an async iterator', async () => {
416+
const ops: Promise<unknown>[] = [];
417+
const [firstPayload, secondPayload] = await Promise.all([
418+
dehydrateStepReturnValue({ v: 1 }, 'wrun_test', undefined, ops),
419+
dehydrateStepReturnValue({ v: 2 }, 'wrun_test', undefined, ops),
420+
]);
421+
422+
const events: Event[] = [
423+
event('evnt_0', 'hook_created', `hook_${ULIDS[0]}`, {
424+
token: 'first',
425+
isWebhook: false,
426+
}),
427+
event('evnt_1', 'hook_created', `hook_${ULIDS[1]}`, {
428+
token: 'second',
429+
isWebhook: false,
430+
}),
431+
event('evnt_2', 'hook_received', `hook_${ULIDS[0]}`, {
432+
token: 'first',
433+
payload: firstPayload,
434+
}),
435+
event('evnt_3', 'hook_received', `hook_${ULIDS[1]}`, {
436+
token: 'second',
437+
payload: secondPayload,
438+
}),
439+
event('evnt_4', 'step_created', `step_${ULIDS[2]}`, {
440+
stepName: 'afterFirst',
441+
}),
442+
event('evnt_5', 'step_created', `step_${ULIDS[3]}`, {
443+
stepName: 'afterSecond',
444+
}),
445+
];
446+
447+
const ctx = setupWorkflowContext(events);
448+
const useStep = createUseStep(ctx);
449+
const createHook = createCreateHook(ctx);
450+
451+
const error = await replay(ctx, async () => {
452+
const first = createHook<{ v: number }>({ token: 'first' });
453+
const second = createHook<{ v: number }>({ token: 'second' });
454+
const afterFirst = useStep('afterFirst');
455+
const afterSecond = useStep('afterSecond');
456+
457+
await Promise.all([
458+
(async () => {
459+
for await (const payload of first) {
460+
void payload;
461+
// Model a layered hook consumer such as an async inbox merger:
462+
// the delivery order must survive its continuation depth.
463+
for (let i = 0; i < 16; i++) {
464+
await Promise.resolve();
465+
}
466+
await afterFirst();
467+
break;
468+
}
469+
})(),
470+
(async () => {
471+
await second;
472+
await afterSecond();
473+
})(),
474+
]);
475+
});
476+
477+
expectSuspendedWithPendingSteps(ctx, error, ['afterFirst', 'afterSecond']);
478+
});
479+
});
480+
481+
// ─── 4. abort deliveries ───────────────────────────────────────────────────
410482
//
411483
// evnt_4 wait_completed ← makes the step result defer
412484
// evnt_5 step_completed stepA ← deferred behind the wait
@@ -491,7 +563,7 @@ describe('abort delivery ordering against an earlier step result', () => {
491563
});
492564
});
493565

494-
// ─── 4. idle reachability over the barrier registry ────────────────────────
566+
// ─── 5. idle reachability over the barrier registry ────────────────────────
495567
//
496568
// `hasParkedCommittedDelivery` decides whether an idle check may observe idle,
497569
// and it is the only remaining caller of the recursive `resolvesOnItsOwn`
@@ -609,7 +681,7 @@ describe('delivery-barrier idle reachability', () => {
609681
});
610682
});
611683

612-
// ─── 5. suspension timing: idle must wait out parked deliveries ────────────
684+
// ─── 6. suspension timing: idle must wait out parked deliveries ────────────
613685
//
614686
// Field shape from vercel/workflow#3183: a fire-and-forget `sleep()` (a
615687
// watchdog — never awaited, never completing in-run) plus a parallel batch of

‎packages/core/src/private.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -245,8 +245,7 @@ interface DeliveryBarrierEntry {
245245
* Which earlier kinds each delivery kind defers behind. Chosen so that no kind
246246
* blocks on a peer it does not need to:
247247
*
248-
* - a hook defers behind earlier WAITS and STEPS — not earlier hooks, which
249-
* are sequential same-entity payloads and must not block one another;
248+
* - a hook defers behind earlier HOOKS, WAITS and STEPS;
250249
* - a wait defers behind earlier HOOKS and STEPS — not earlier waits, since a
251250
* wait never needs to queue behind another wait;
252251
* - a step defers behind earlier WAITS, HOOKS and STEPS.
@@ -269,7 +268,7 @@ interface DeliveryBarrierEntry {
269268
* wait-for graph can never contain a cycle.
270269
*/
271270
const DEFER_BEHIND: Record<DeliveryKind, readonly DeliveryKind[]> = {
272-
hook: ['wait', 'step'],
271+
hook: ['hook', 'wait', 'step'],
273272
wait: ['hook', 'step'],
274273
step: ['wait', 'hook', 'step'],
275274
};

0 commit comments

Comments
 (0)