diff --git a/.changeset/late-hooks-rearm.md b/.changeset/late-hooks-rearm.md new file mode 100644 index 0000000000..9ca2387ad0 --- /dev/null +++ b/.changeset/late-hooks-rearm.md @@ -0,0 +1,5 @@ +--- +"@workflow/core": patch +--- + +Keep late-claimed buffered hook payloads from being preempted by workflow suspension in retained VMs. diff --git a/packages/core/src/async-deserialization-ordering.test.ts b/packages/core/src/async-deserialization-ordering.test.ts index e82360eb35..a7b552c92f 100644 --- a/packages/core/src/async-deserialization-ordering.test.ts +++ b/packages/core/src/async-deserialization-ordering.test.ts @@ -4,7 +4,7 @@ import * as nanoid from 'nanoid'; import { monotonicFactory } from 'ulid'; import { afterEach, describe, expect, it, vi } from 'vitest'; import { EventsConsumer } from './events-consumer.js'; -import type { WorkflowOrchestratorContext } from './private.js'; +import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js'; import { dehydrateStepReturnValue } from './serialization.js'; import { createUseStep } from './step.js'; import { createContext } from './vm/index.js'; @@ -685,4 +685,48 @@ describe('async deserialization ordering', () => { expect(ctx.pendingDeliveryBarriers?.size).toBe(0); }); + + it('should restore delivery protection when a buffered hook payload is claimed after its barrier retires', async () => { + const payload = await dehydrateStepReturnValue( + 'buffered', + 'wrun_test', + undefined + ); + const ctx = setupWorkflowContext([ + { + eventId: 'evnt_0', + runId: 'wrun_test', + eventType: 'hook_received', + correlationId: 'hook_01K11TFZ62YS0YYFDQ3E8B9YCV', + eventData: { payload }, + createdAt: new Date(), + }, + ]); + + // Model another payload hydrating in the same replay so the idle safety + // net cannot retire this hook's barrier before the test observes it. + ctx.pendingDeliveries = 1; + const createHook = createCreateHook(ctx); + const hook = createHook(); + + // The payload arrived before a consumer requested it, so it starts with + // an unarmed barrier and remains buffered after the idle safety net retires + // that barrier. + await vi.waitFor(() => { + expect(ctx.pendingDeliveryBarriers?.size).toBe(1); + }); + await ctx.promiseQueue; + expect(ctx.pendingDeliveryBarriers?.size).toBe(1); + ctx.pendingDeliveries = 0; + await vi.waitFor(() => { + expect(ctx.pendingDeliveryBarriers?.size).toBe(0); + }); + + // Claiming the buffered payload commits it to reaching this consumer. Its + // delivery must become non-idle again until the claim resolves. + const delivery = hook.then((value) => value); + expect(isDeliveryIdle(ctx)).toBe(false); + await expect(delivery).resolves.toBe('buffered'); + expect(isDeliveryIdle(ctx)).toBe(true); + }); }); diff --git a/packages/core/src/delivery-barrier-dispenser.test.ts b/packages/core/src/delivery-barrier-dispenser.test.ts index 226b07b7cf..63970846b5 100644 --- a/packages/core/src/delivery-barrier-dispenser.test.ts +++ b/packages/core/src/delivery-barrier-dispenser.test.ts @@ -117,6 +117,18 @@ const resumeAtA = new Date('2026-07-27T12:00:05.000Z'); const resumeAtB = new Date('2026-07-27T12:00:06.000Z'); describe('barrier safety-net dispenser', () => { + it('rejects a second barrier owner for the same event index', () => { + const ctx = setupWorkflowContext([]); + const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false }); + + expect(() => registerDeliveryBarrier(ctx, 0, 'step')).toThrowError( + 'Delivery barrier already registered at event index 0' + ); + + barrier.markDelivered(); + expect(isDeliveryIdle(ctx)).toBe(true); + }); + it('suspends only after deliveries parked behind an unclaimed payload have run', async () => { const ops: Promise[] = []; const payload = await dehydrateStepReturnValue( diff --git a/packages/core/src/private.ts b/packages/core/src/private.ts index ee421ae3f0..638a08e79e 100644 --- a/packages/core/src/private.ts +++ b/packages/core/src/private.ts @@ -2,6 +2,7 @@ * Utils used by the bundler when transforming code */ +import { WorkflowRuntimeError } from '@workflow/errors'; import { withResolvers } from '@workflow/utils'; import type { CryptoKey } from './encryption.js'; import type { EventsConsumer } from './events-consumer.js'; @@ -215,8 +216,8 @@ export type DeliveryKind = 'hook' | 'wait' | 'step'; interface DeliveryBarrierEntry { kind: DeliveryKind; - /** Resolves once this delivery has resolved to the workflow. */ - delivered: Promise; + /** Resolves once this delivery is handed to the workflow or retired. */ + released: Promise; /** * Whether this delivery is committed to reaching the workflow without any * further action by workflow code. True for wait completions and step @@ -229,11 +230,15 @@ interface DeliveryBarrierEntry { * once a consumer takes the payload. */ armed: boolean; + /** Whether this entry has been removed and its `released` promise settled. */ + retired: boolean; /** - * Retire this entry: resolve `delivered` and remove it from the registry, - * exactly as `markDelivered` would. Called only by the context's safety-net - * dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index - * entry at delivery idle. Idempotent. + * Retire this entry: resolve `released` and remove it from the registry, + * without marking the handle delivered to the workflow. A safety-retired + * buffered payload may therefore install a fresh entry if it is claimed by + * a retained VM later. Called only by the context's safety-net dispenser + * ({@link ensureBarrierSafetyNet}), and only on the lowest-index entry at + * delivery idle. Idempotent. */ retire: () => void; } @@ -396,7 +401,7 @@ export async function awaitEarlierDeliveries( ) { continue; } - earlier.push(entry.delivered); + earlier.push(entry.released); } if (earlier.length > 0) { await Promise.all(earlier); @@ -422,7 +427,7 @@ export async function awaitEarlierDeliveries( export interface DeliveryBarrier { /** * Mark this delivery as delivered to the workflow. Resolves its - * `delivered` promise so any later-in-log delivery gated on it (via + * `released` promise so any later-in-log delivery gated on it (via * {@link awaitEarlierDeliveries}) may proceed, and removes it from the * registry. Idempotent. */ @@ -468,41 +473,71 @@ export function registerDeliveryBarrier( return { markDelivered: () => {}, arm: () => {} }; } - let done = false; - const { promise, resolve } = withResolvers(); - - const finish = () => { - if (done) { - return; - } - done = true; - if (barriers.get(eventIndex) === entry) { - barriers.delete(eventIndex); + const install = (armed: boolean): DeliveryBarrierEntry => { + if (barriers.has(eventIndex)) { + throw new WorkflowRuntimeError( + `Delivery barrier already registered at event index ${eventIndex}` + ); } - resolve(); + const { promise, resolve } = withResolvers(); + const entry: DeliveryBarrierEntry = { + kind, + released: promise, + armed, + retired: false, + retire: () => { + if (entry.retired) { + return; + } + entry.retired = true; + if (barriers.get(eventIndex) === entry) { + barriers.delete(eventIndex); + } + resolve(); + }, + }; + barriers.set(eventIndex, entry); + + // Safety net: if this delivery is never delivered to the workflow (its + // branch was not taken / the run is suspending, or a buffered hook payload + // is only claimed after a later delivery the workflow is still waiting + // on), it is retired at idle so a later delivery gated on it cannot + // deadlock and the registry cannot leak an entry per abandoned delivery. + // Retirement goes through the context's single ordered dispenser rather + // than a per-barrier idle poll. See {@link ensureBarrierSafetyNet} for why + // the ORDER of these retirements is load-bearing. + ensureBarrierSafetyNet(ctx); + return entry; }; - const entry: DeliveryBarrierEntry = { - kind, - delivered: promise, - armed: options.armed ?? true, - retire: finish, - }; - barriers.set(eventIndex, entry); - - // Safety net: if this delivery is never delivered to the workflow (its - // branch was not taken / the run is suspending, or a buffered hook payload - // is only claimed after a later delivery the workflow is still waiting on), - // it is retired at idle so a later delivery gated on it cannot deadlock and - // the registry cannot leak an entry per abandoned delivery. Retirement goes - // through the context's single ordered dispenser rather than a per-barrier - // idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these - // retirements is load-bearing. - ensureBarrierSafetyNet(ctx); + let entry = install(options.armed ?? true); + let deliveredToWorkflow = false; return { - markDelivered: finish, + markDelivered: () => { + if (deliveredToWorkflow) { + return; + } + deliveredToWorkflow = true; + entry.retire(); + }, arm: () => { + if (deliveredToWorkflow) { + return; + } + // The idle safety net may retire an unclaimed buffered hook payload + // while a retained VM keeps its `claim()` closure alive. If workflow + // code later claims that payload, replace the settled entry so delivery + // remains non-idle until the claim reaches the workflow. + if (entry.retired) { + entry = install(true); + return; + } + if (barriers.get(eventIndex) !== entry) { + throw new WorkflowRuntimeError( + `Delivery barrier lost ownership of event index ${eventIndex}` + ); + } entry.armed = true; }, };