From e4a61b2efa70f7b197e99de205dc6fede66263fe Mon Sep 17 00:00:00 2001 From: "vercel[bot]" <35613825+vercel[bot]@users.noreply.github.com> Date: Fri, 25 Sep 2026 00:09:13 +0000 Subject: [PATCH 1/3] fix(core): repay a force-claim victim's wake independent of the claimer's log tail A replay republished the victim's wake only while the forced hook_created was the last row the claimer itself wrote, so any own row landing between the creation and the wake (a step or wait terminal from another invocation, or, since #4392, a row the same suspension writes alongside the creation) hid the debt, and a victim parked only on `await hook` never woke. Every replay now republishes for each forced hook_created in the loaded log whose createdAt is within 24 hours (the queue's message retention and idempotency window), under the existing `hook-force-claim-` key. The node:vm handler sends it alongside the suspension's writes and at most once per hook per invocation; QuickJS once per invocation on load. Both use the shared rule in hook-wake.ts. Closes #4393. Co-Authored-By: Claude Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> --- .changeset/force-claim-wake-recovery.md | 5 + .../v5/api-reference/workflow/create-hook.mdx | 2 +- packages/core/src/runtime.ts | 6 + packages/core/src/runtime/hook-wake.ts | 182 ++++++---- .../core/src/runtime/quickjs-entrypoint.ts | 23 +- .../runtime/quickjs-force-claim-wake.test.ts | 67 +++- .../src/runtime/suspension-handler.test.ts | 320 +++++++++++++----- .../core/src/runtime/suspension-handler.ts | 61 +++- 8 files changed, 467 insertions(+), 199 deletions(-) create mode 100644 .changeset/force-claim-wake-recovery.md diff --git a/.changeset/force-claim-wake-recovery.md b/.changeset/force-claim-wake-recovery.md new file mode 100644 index 0000000000..1ad21c1717 --- /dev/null +++ b/.changeset/force-claim-wake-recovery.md @@ -0,0 +1,5 @@ +--- +'@workflow/core': patch +--- + +Republish a force-claimed hook's victim wake on every replay within 24 hours of the takeover instead of only while the forced `hook_created` is the claimer's last own event, so a crash before the wake is repaid even when the claimer wrote other events after the creation. diff --git a/docs/content/docs/v5/api-reference/workflow/create-hook.mdx b/docs/content/docs/v5/api-reference/workflow/create-hook.mdx index e58e41109c..30b7ea04fe 100644 --- a/docs/content/docs/v5/api-reference/workflow/create-hook.mdx +++ b/docs/content/docs/v5/api-reference/workflow/create-hook.mdx @@ -251,7 +251,7 @@ With `experimental_force`, this run always ends up owning the token: - Any number of runs forcing the same token at the same time converge on a single owner. The takeovers form a chain: each run that loses the token gets `HookForceClaimedError`, exactly one run ends up owning it, and none of them can get stuck. Which run wins among simultaneous claimers is not defined; if the order matters, start them in order. - A finished run that still holds the token under [`experimental_minRetention`](#keep-a-token-unavailable-after-the-run-ends) is taken over silently, since there is nothing left to wake. A run can also take over a token held by its own earlier Hook. -The takeover is durable. If either run's compute fails partway through, the next request for the token completes it, so the token never ends up held by nobody or by both runs. The previous owner's wake is repaired on a best-effort basis: if the new owner's compute fails between registering the Hook and waking the previous owner, the new owner's next invocation republishes the wake (idempotently), as long as the new owner has not recorded any other event since. If it has, for example a step it started alongside the Hook, the previous owner still sees the takeover, but only the next time something else invokes it. +The takeover is durable. If either run's compute fails partway through, the next request for the token completes it, so the token never ends up held by nobody or by both runs. The previous owner's wake is durable too: if the new owner's compute fails between registering the Hook and waking the previous owner, the new owner's next invocation republishes the wake, whatever else the new owner has recorded since (a step it started alongside the Hook, for example). Every invocation of the new owner within 24 hours of the takeover republishes it under the same idempotency key, which collapses the repeats into one wake; a repeat that does get through only replays the previous owner, which finds nothing new. A token can only be taken from a run whose runtime understands being taken from. Runs started at a Workflow spec version below 8, which includes every run started by an older SDK release, a Python SDK run, or a deployment with `WORKFLOW_SEALED_LOG=0`, would never learn that their Hook was disposed. The World declines to take their token and the forced Hook rejects with the ordinary [`HookConflictError`](/docs/api-reference/workflow-errors/hook-conflict-error) instead, exactly as if `experimental_force` had not been set. Finished runs holding a retained token are taken over at any version. diff --git a/packages/core/src/runtime.ts b/packages/core/src/runtime.ts index 4149cfe410..0f90c6e67c 100644 --- a/packages/core/src/runtime.ts +++ b/packages/core/src/runtime.ts @@ -967,6 +967,11 @@ export function workflowEntrypoint( // than when the wait's own timer would have fired. let eventLogFromInlineDelta = false; let loopIteration = 0; + // Hooks whose force-claim victim wake this invocation has + // already sent (its own forced creations, and the replay's + // republishes), so each suspension of the loop below does + // not send them again. See `forcedCreationsOwingWake`. + const forceClaimVictimWakes = new Set(); const replayRecoveryReporter = replayDivergence ? new ReplayRecoveryReporter(replayDivergence.count) : ReplayRecoveryReporter.inert(); @@ -3544,6 +3549,7 @@ export function workflowEntrypoint( eventLog, runReadyBarrier, replayRecoveryReporter, + forceClaimVictimWakes, // Resilient step dispatch: lets eligible newly // created steps publish their step-execution // message (carrying `stepInput`) in parallel with diff --git a/packages/core/src/runtime/hook-wake.ts b/packages/core/src/runtime/hook-wake.ts index 7936d44444..1019f578a3 100644 --- a/packages/core/src/runtime/hook-wake.ts +++ b/packages/core/src/runtime/hook-wake.ts @@ -100,9 +100,11 @@ export async function publishHookWakeWithRetry( * Same durability contract as `resumeHook()`'s wake: the row is durable * before this runs, the publish is retried on transport-shaped failures, and * a publish that still fails is logged rather than failing the claimer — - * nothing of the claimer's is wrong, and the victim reads the row on its - * next invocation for any reason. The idempotency key is the claimer's hook - * id, so a claimer retrying its creation republishes at most one wake. + * nothing of the claimer's is wrong, every later replay of the claimer inside + * the republish window tries again ({@link forcedCreationsOwingWake}), and + * the victim reads the row on its next invocation for any reason. The + * idempotency key is the claimer's hook id, so those republishes collapse + * into one wake. * * Skipped when the victim is the claimer itself (a run taking over its own * earlier hook is already running) and when the World recorded no @@ -154,86 +156,132 @@ export async function publishForceClaimVictimWake( } /** - * The forced hook creation whose victim wake this run still owes, if any: the - * forced `hook_created` is the last event the run's own replay appended. + * How long after a forced `hook_created` a replay of the claimer keeps + * republishing its victim's wake: 24 hours, measured from the row's + * `createdAt`. + * + * The bound has to outlast every redelivery of the invocation that journaled + * the creation, because that redelivery is the replay that must repay a wake + * the invocation died before publishing. A queue message is retained for 24 + * hours from its send and the creation is written after the send, so any such + * redelivery arrives within 24 hours of the creation. The same 24 hours is the + * Vercel queue's idempotency window (`min(retention, 24h)`), so every + * republish inside it collapses, under `hook-force-claim-`, into the + * one wake that was (or now is) delivered. Past it a republish would be a + * genuinely new message, which is what the bound saves. world-postgres + * remembers a completed key in-process to the same effect; world-local + * dedupes a key only while its message is in flight, so there a republish can + * deliver the victim one more replay, which reads nothing new. + */ +export const FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS = 24 * 60 * 60 * 1000; + +/** + * The forced hook creations whose victim wake this run may still owe: every + * forced `hook_created` (one carrying `forceClaimedFrom`) in the log whose + * `createdAt` is within {@link FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS} of + * `nowMs`. * * A forced creation is followed by a wake of the run it took the token from. * If the invocation died between the two, the creation is in the log and the - * victim was never told; the row itself is the durable record of that debt. - * As long as it is the last event THIS RUN wrote, the run has made no progress - * since, so the invocation that should have woken the victim did not finish, - * and the replay republishes (under the hook's idempotency key, so a wake that - * did go out is not duplicated). The first event the run appends after it - * ends the republishing. + * victim was never told; the row itself is the durable record of that debt, + * and nothing records that the wake went out. So the replay does not try to + * infer it: it republishes for every recent forced creation, and the hook's + * idempotency key collapses a wake that did go out (a duplicate that slips + * past a World's dedupe is one harmless replay of the victim). * - * "This run wrote" matters: a delivery appends `hook_received` to this log - * from another request, a sealed-log World appends `noop`, and a LATER - * claimer taking the token from this run appends - * `hook_disposed{forceClaimedBy}`. None is progress of this run — the model - * (`ForceWakeOnce.cfg`'s sibling trace) has a delivery land between the crash - * and the retry, and a rule that looked at the bare tail would then never - * wake the victim. The foreign disposal is the chain case: this run took the - * token from A, died before waking A, and was itself taken from by C. Its own - * wake (from C) is the very invocation that must repay A's — the run's own - * `hook_disposed` (a `dispose()` in its code) IS its progress and still ends - * the debt, but a row another run put here does not. Both engines call this - * on the log they loaded for the invocation, before writing anything. + * The rule reads nothing written after the creation, which is what makes it + * sound. Any row can land between the creation and the wake: a step, wait, + * attribute or other hook row the same suspension writes concurrently + * (neither engine holds those for the wake), a step or wait terminal from + * another invocation, a delivery's `hook_received`, a sealed-log `noop`, or a + * later claimer's `hook_disposed{forceClaimedBy}` taking the token from this + * run. None of them says the wake was published. The last is the chain case: + * this run took the token from A, died before waking A, and was taken from by + * C; C's wake of this run is the invocation that repays A's. * - * Reading only the last own row is best-effort. Any row this run writes after - * the forced creation and before the wake goes out ends the debt as if the - * wake had been published: a step, wait, attribute or other hook row the same - * suspension writes concurrently (neither engine holds those for the wake), or - * a step or wait terminal from another invocation. If the invocation then - * dies before publishing, the victim reads its disposal only on its next - * invocation for any other reason — the same outcome as a wake whose publish - * fails outright. Making the recovery independent of the log's tail is - * vercel/workflow#4393. + * The time is the row's `createdAt`, never one decoded from its event id (a + * slot-numbered id carries none). Under slot identity it is the writer's + * client clock, clamped by the World to at most an hour ahead of its own, so + * skew moves the window's edge by at most that much; a creation dated ahead of + * `nowMs` counts as recent, and one whose time cannot be read is treated as + * recent too, erring toward a wake rather than a stranded victim. */ -export function forcedCreationOwingWake( - events: readonly Event[] | undefined -): (Event & { eventType: 'hook_created' }) | undefined { - if (!events) return undefined; - for (let i = events.length - 1; i >= 0; i--) { - const event = events[i]; +export function forcedCreationsOwingWake( + events: readonly Event[] | undefined, + nowMs: number = Date.now() +): (Event & { eventType: 'hook_created' })[] { + const owed: (Event & { eventType: 'hook_created' })[] = []; + if (!events) return owed; + for (const event of events) { if ( - event.eventType === 'hook_received' || - event.eventType === 'noop' || - (event.eventType === 'hook_disposed' && - event.eventData?.forceClaimedBy !== undefined) + event.eventType !== 'hook_created' || + event.eventData?.forceClaimedFrom === undefined ) { continue; } - return event.eventType === 'hook_created' && - event.eventData.forceClaimedFrom !== undefined - ? event - : undefined; + const createdAtMs = new Date(event.createdAt).getTime(); + if ( + Number.isNaN(createdAtMs) || + nowMs - createdAtMs < FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS + ) { + owed.push(event); + } } - return undefined; + return owed; } /** - * Republish the wake {@link forcedCreationOwingWake} says is owed. Shared by - * the node:vm suspension handler and the QuickJS entrypoint so the two engines - * cannot drift on the durability contract. + * Republish the wakes {@link forcedCreationsOwingWake} says may be owed. + * Shared by the node:vm suspension handler and the QuickJS entrypoint so the + * two engines cannot drift on the durability contract. + * + * `alreadyWoken` is the invocation's record of the hooks whose victim wake it + * has already published or attempted, the forced creations it made itself + * included: an engine that replays more than once per invocation passes the + * same set every time, so each hook costs one send per invocation. Hooks this + * call publishes are added to it. + * + * Self-claims and victims with no recorded `workflowName` are skipped + * silently: the creating invocation already logged the latter, and a replay + * repeating it would say nothing new. Never rejects; a publish that fails + * after its retries is logged, and the next replay inside the window tries + * again. */ -export async function republishOwedForceClaimVictimWake( +export async function republishOwedForceClaimVictimWakes( world: World, runId: string, - events: readonly Event[] | undefined + events: readonly Event[] | undefined, + options: { alreadyWoken?: Set; nowMs?: number } = {} ): Promise { - const owed = forcedCreationOwingWake(events); - if (!owed) return; - const claimedFrom = owed.eventData.forceClaimedFrom as HookClaimedFrom; - const outcome = await publishForceClaimVictimWake(world, runId, { - hookId: owed.correlationId, - claimedFrom, - }); - if (outcome !== 'skipped') { - runtimeLogger.info('Republished the wake of a force-claimed hook victim', { - workflowRunId: runId, - hookId: owed.correlationId, - victimRunId: claimedFrom.runId, - victimWake: outcome, - }); - } + const { alreadyWoken } = options; + const owed = forcedCreationsOwingWake(events, options.nowMs).filter( + (event) => { + const from = event.eventData.forceClaimedFrom as HookClaimedFrom; + return ( + from.runId !== runId && + from.workflowName !== undefined && + !alreadyWoken?.has(event.correlationId) + ); + } + ); + if (owed.length === 0) return; + await Promise.all( + owed.map(async (event) => { + alreadyWoken?.add(event.correlationId); + const claimedFrom = event.eventData.forceClaimedFrom as HookClaimedFrom; + const outcome = await publishForceClaimVictimWake(world, runId, { + hookId: event.correlationId, + claimedFrom, + }); + runtimeLogger.debug( + 'Republished the wake of a force-claimed hook victim', + { + workflowRunId: runId, + hookId: event.correlationId, + victimRunId: claimedFrom.runId, + victimWake: outcome, + } + ); + }) + ); } diff --git a/packages/core/src/runtime/quickjs-entrypoint.ts b/packages/core/src/runtime/quickjs-entrypoint.ts index 381ff2e1a4..7719625206 100644 --- a/packages/core/src/runtime/quickjs-entrypoint.ts +++ b/packages/core/src/runtime/quickjs-entrypoint.ts @@ -67,7 +67,7 @@ import { } from './helpers.js'; import { publishForceClaimVictimWake, - republishOwedForceClaimVictimWake, + republishOwedForceClaimVictimWakes, } from './hook-wake.js'; import { dispatchRunCompletedHooks, @@ -594,10 +594,10 @@ async function dispatchPendingOps(params: { }; // Token groups run in parallel with every other op, forced creations // included. A forced creation publishes its victim's wake before its group's - // next write, but nothing else waits for it, so a crash before the wake can - // leave another row as the log's last and hide the owed wake from the - // replay's `forcedCreationOwingWake`. Making that recovery independent of - // the log's tail is tracked in vercel/workflow#4393. + // next write, but nothing else waits for it, and nothing needs to: a crash + // before the wake is repaid by the next replay from the forced + // `hook_created` itself, which `forcedCreationsOwingWake` finds wherever it + // sits in the log, so no row written after it can hide the debt. for (const group of hookOpsByToken.values()) { opsPromises.push(runHookGroup(group)); } @@ -1129,12 +1129,13 @@ export async function runWorkflowWithQuickJS(params: { // handed back on a write that the VM has not been given yet. Every write // made from this view goes through `createEvent` below so it names the // position it was decided against and its response is queued here. - // Same durability contract as the node:vm suspension handler: a forced - // hook creation that is still the last event this run wrote owes its - // victim a wake, because the invocation that created it died before - // publishing one. Repaid here, on the log as loaded, before this - // invocation writes anything. - await republishOwedForceClaimVictimWake(world, runId, events); + // Same durability contract as the node:vm suspension handler: every + // recent forced hook creation in the log may still owe its victim a wake, + // because the invocation that created it may have died before publishing + // one, so it is republished under the hook's idempotency key (see + // `forcedCreationsOwingWake`). Once per invocation, on the log as loaded; + // the forced creations this invocation makes publish their own. + await republishOwedForceClaimVictimWakes(world, runId, events); const logView = new QuickJSLogView(events, loadedCursor); const createEvent: EventCreator = async (data, eventParams) => { diff --git a/packages/core/src/runtime/quickjs-force-claim-wake.test.ts b/packages/core/src/runtime/quickjs-force-claim-wake.test.ts index c53a2d21ab..ef6787d68c 100644 --- a/packages/core/src/runtime/quickjs-force-claim-wake.test.ts +++ b/packages/core/src/runtime/quickjs-force-claim-wake.test.ts @@ -4,11 +4,11 @@ * * A forced `hook_created` is followed by a wake of the run it took the token * from. If the invocation dies between the two, the creation is in the log and - * the victim was never told. On the next invocation the row is still the last - * event this run wrote, so the entrypoint republishes the wake — before the VM - * runs and before anything is written — under the hook's idempotency key. - * `forcedCreationOwingWake` is the shared rule; this test drives the QuickJS - * entrypoint through it from a committed log, with the VM mocked. + * the victim was never told. Every later invocation inside the republish + * window republishes the wake under the hook's idempotency key, whatever the + * run wrote after the creation. `forcedCreationsOwingWake` is the shared rule; + * this test drives the QuickJS entrypoint through it from a committed log, + * with the VM mocked. */ import { type CreateEventRequest, @@ -19,6 +19,7 @@ import { type World, } from '@workflow/world'; import { describe, expect, it, vi } from 'vitest'; +import { FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS } from './hook-wake.js'; import { setWorld } from './world.js'; vi.mock('@vercel/functions', () => ({ waitUntil: vi.fn() })); @@ -49,19 +50,20 @@ const event = ( slot: number, eventType: Event['eventType'], eventData?: unknown, - correlationId?: string + correlationId?: string, + createdAt: Date = new Date() ): Event => ({ eventType, eventId: slotToEventId(slot), runId, correlationId, - createdAt: startedAt, + createdAt, specVersion: SPEC_VERSION_CURRENT, eventData, }) as Event; -const forcedCreation = (slot: number) => +const forcedCreation = (slot: number, createdAt: Date = new Date()) => event( slot, 'hook_created', @@ -76,7 +78,8 @@ const forcedCreation = (slot: number) => runSpecVersion: SPEC_VERSION_CURRENT, }, }, - 'hook_claimer' + 'hook_claimer', + createdAt ); /** Replay the entrypoint over a committed log; the VM suspends with nothing pending. */ @@ -112,7 +115,7 @@ async function replayWith(events: Event[]) { } describe('QuickJS force-claim victim wake on replay', () => { - it('republishes the wake when the forced creation is the last event this run wrote', async () => { + it('republishes the wake for a recent forced creation', async () => { const queue = await replayWith([ event(1, 'run_created', { workflowName: 'workflow', input: [] }), event(2, 'run_started'), @@ -128,7 +131,7 @@ describe('QuickJS force-claim victim wake on replay', () => { }); }); - it("still republishes past a delivery's hook_received, which is not this run's progress", async () => { + it("republishes past a delivery's hook_received", async () => { const queue = await replayWith([ event(1, 'run_created', { workflowName: 'workflow', input: [] }), event(2, 'run_started'), @@ -143,20 +146,56 @@ describe('QuickJS force-claim victim wake on replay', () => { expect(queue).toHaveBeenCalledTimes(1); }); - it('stops once this run has written anything after the creation', async () => { + it("repays the wake even when the run's own step and wait rows landed after the forced creation", async () => { + // vercel/workflow#4393: rows this run wrote after the creation (a step or + // wait terminal from another invocation, or a row written alongside the + // creation) do not say the wake was published. const queue = await replayWith([ event(1, 'run_created', { workflowName: 'workflow', input: [] }), event(2, 'run_started'), forcedCreation(3), event(4, 'step_created', { stepName: 'after', input: [] }, 'step_after'), + event(5, 'step_completed', { result: [] }, 'step_after'), + event(6, 'wait_completed', {}, 'wait_1'), + ]); + expect(queue).toHaveBeenCalledTimes(1); + expect(queue.mock.calls[0][2]).toMatchObject({ + idempotencyKey: 'hook-force-claim-hook_claimer', + }); + }); + + it("repays the victim's wake when a later claimer took the token from this run (the chain)", async () => { + const queue = await replayWith([ + event(1, 'run_created', { workflowName: 'workflow', input: [] }), + event(2, 'run_started'), + forcedCreation(3), + event( + 4, + 'hook_disposed', + { forceClaimedBy: { runId: 'wrun_third', hookId: 'hook_third' } }, + 'hook_claimer' + ), + ]); + expect(queue).toHaveBeenCalledTimes(1); + expect(queue.mock.calls[0][1]).toEqual({ runId: 'wrun_victim' }); + }); + + it('stops republishing once the forced creation is outside the window', async () => { + const queue = await replayWith([ + event(1, 'run_created', { workflowName: 'workflow', input: [] }), + event(2, 'run_started'), + forcedCreation( + 3, + new Date(Date.now() - FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS - 1_000) + ), ]); expect(queue).not.toHaveBeenCalled(); }); it("does not hold the other writes for a forced creation's victim wake", async () => { // The wake still goes out, but the sibling hook and the wait are written - // while it is in flight rather than after it (vercel/workflow#4393 tracks - // the recovery gap that leaves after a crash). + // while it is in flight rather than after it. A crash before it goes out + // is repaid by the next replay from the forced creation itself. const order: string[] = []; const queue = vi.fn( async (_queueName: string, message: { runId: string }) => { diff --git a/packages/core/src/runtime/suspension-handler.test.ts b/packages/core/src/runtime/suspension-handler.test.ts index be8573bbf3..38aa035615 100644 --- a/packages/core/src/runtime/suspension-handler.test.ts +++ b/packages/core/src/runtime/suspension-handler.test.ts @@ -19,6 +19,7 @@ import { type QueueItem, WorkflowSuspension } from '../global.js'; import { hydrateStepArguments, hydrateStepError } from '../serialization.js'; import { COMPUTE_INSTANCE_ID } from './compute-instance.js'; import { maxEventSlot, stepDispatchIdempotencyKey } from './helpers.js'; +import { FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS } from './hook-wake.js'; import { ReplayRecoveryReporter } from './replay-recovery-reporter.js'; import { handleSuspension } from './suspension-handler.js'; import { isUnserializableStepInputPlaceholder } from './unserializable-step.js'; @@ -378,49 +379,74 @@ describe('handleSuspension', () => { deploymentId: 'dpl_victim', runSpecVersion: SPEC_VERSION_CURRENT, }; - const forcedCreation = (slot: number): Event => + const forcedCreation = ( + slot: number, + { + createdAt = new Date(), + hookId = 'hook_claimer', + from = claimedFrom, + }: { + createdAt?: Date; + hookId?: string; + from?: Record; + } = {} + ): Event => ({ eventType: 'hook_created', eventId: slotToEventId(slot), runId: run.runId, - correlationId: 'hook_claimer', - createdAt: new Date(), + correlationId: hookId, + createdAt, specVersion: SPEC_VERSION_CURRENT, eventData: { - token: 'channel:1', + token: `channel:${hookId}`, force: true, - forceClaimedFrom: claimedFrom, + forceClaimedFrom: from, }, }) as Event; - const stepCreated = (slot: number): Event => + const ownRow = ( + slot: number, + eventType: Event['eventType'], + correlationId: string, + eventData: unknown = {} + ): Event => ({ - eventType: 'step_created', + eventType, eventId: slotToEventId(slot), runId: run.runId, - correlationId: 'step_after', + correlationId, createdAt: new Date(), specVersion: SPEC_VERSION_CURRENT, - eventData: { stepName: 'after', input: [] }, + eventData, }) as Event; - const worldWithQueue = (queue: ReturnType): World => + const worldWithQueue = ( + queue: ReturnType, + eventsCreate: ReturnType = vi.fn() + ): World => ({ - events: { create: vi.fn() }, + events: { create: eventsCreate }, getEncryptionKeyForRun: vi.fn().mockResolvedValue(undefined), queue, }) as unknown as World; - - it('republishes the victim wake while the forced creation is the last event in the log', async () => { - // The invocation that created the hook died before waking the victim: - // its replay finds the creation as the log's tail and republishes, - // under the hook's idempotency key, so a wake that did go out is not - // duplicated. - const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); - await handleSuspension({ + const suspendOver = ( + queue: ReturnType, + events: Event[], + extra: { forceClaimVictimWakes?: Set } = {} + ) => + handleSuspension({ suspension: new WorkflowSuspension(new Map(), globalThis), world: worldWithQueue(queue), run, - eventLog: { events: [forcedCreation(3)], cursor: null }, + eventLog: { events, cursor: null }, + ...extra, }); + + it('republishes the victim wake for a recent forced creation', async () => { + // The invocation that created the hook may have died before waking the + // victim: its replay republishes, under the hook's idempotency key, so a + // wake that did go out is not duplicated. + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + await suspendOver(queue, [forcedCreation(3)]); expect(queue).toHaveBeenCalledTimes(1); const [queueName, message, options] = queue.mock.calls[0]; expect(queueName).toContain('victim-workflow'); @@ -431,93 +457,209 @@ describe('handleSuspension', () => { }); }); - it("republishes past rows other actors appended: a delivery's hook_received is not this run's progress", async () => { - // The trace TLC found: the claimer dies after journaling, a delivery - // lands `hook_received` in its log before it comes back. The bare tail - // is no longer the creation, but the run itself has written nothing - // since, so the wake is still owed. + it("repays the wake even when the run's own step and wait rows landed after the forced creation", async () => { + // vercel/workflow#4393: a step or wait terminal from another invocation, + // or a row this suspension wrote alongside the creation, can land before + // the wake goes out. None of them says the wake was published, so none + // of them may end the republishing. const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); - const received: Event = { - eventType: 'hook_received', - eventId: slotToEventId(4), - runId: run.runId, - correlationId: 'hook_claimer', - createdAt: new Date(), - specVersion: SPEC_VERSION_CURRENT, - eventData: { token: 'channel:1', payload: { n: 1 } as never }, - } as Event; - await handleSuspension({ - suspension: new WorkflowSuspension(new Map(), globalThis), - world: worldWithQueue(queue), - run, - eventLog: { events: [forcedCreation(3), received], cursor: null }, + await suspendOver(queue, [ + forcedCreation(3), + ownRow(4, 'step_created', 'step_after', { + stepName: 'after', + input: [], + }), + ownRow(5, 'step_completed', 'step_after', { result: [] }), + ownRow(6, 'wait_completed', 'wait_1'), + ]); + expect(queue).toHaveBeenCalledTimes(1); + expect(queue.mock.calls[0][2]).toMatchObject({ + idempotencyKey: 'hook-force-claim-hook_claimer', }); + }); + + it("republishes past a delivery's hook_received", async () => { + // The trace TLC found for the tail rule: the claimer dies after + // journaling and a delivery lands before it comes back. The window rule + // never looks past the creation, so it is unaffected. + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + await suspendOver(queue, [ + forcedCreation(3), + ownRow(4, 'hook_received', 'hook_claimer', { + token: 'channel:1', + payload: { n: 1 }, + }), + ]); expect(queue).toHaveBeenCalledTimes(1); }); - it("republishes past a later claimer's hook_disposed{forceClaimedBy}: being taken from is not this run's progress", async () => { - // The chain: this run took the token from the victim, died before - // waking it, and a third run then took the token from THIS run — - // appending `hook_disposed{forceClaimedBy}` to this log and waking it. - // That wake is the invocation that must repay the victim's; the foreign - // row must not read as "moved on" or the victim is never invoked. + it("repays the victim's wake when a later claimer took the token from this run (the chain)", async () => { + // This run took the token from the victim, died before waking it, and a + // third run then took the token from THIS run, appending + // `hook_disposed{forceClaimedBy}` here and waking it. That wake is the + // invocation that must repay the victim's. const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); - const takenFrom: Event = { - eventType: 'hook_disposed', - eventId: slotToEventId(4), - runId: run.runId, - correlationId: 'hook_claimer', - createdAt: new Date(), - specVersion: SPEC_VERSION_CURRENT, - eventData: { + await suspendOver(queue, [ + forcedCreation(3), + ownRow(4, 'hook_disposed', 'hook_claimer', { forceClaimedBy: { runId: 'wrun_third', hookId: 'hook_third' }, - }, - } as Event; - await handleSuspension({ - suspension: new WorkflowSuspension(new Map(), globalThis), - world: worldWithQueue(queue), - run, - eventLog: { events: [forcedCreation(3), takenFrom], cursor: null }, - }); + }), + ownRow(5, 'step_completed', 'step_after', { result: [] }), + ]); expect(queue).toHaveBeenCalledTimes(1); expect(queue.mock.calls[0][1]).toEqual({ runId: 'wrun_victim' }); }); - it("stops republishing after the run's OWN hook_disposed: a dispose() in its code is its progress", async () => { + it("still republishes after the run's own hook_disposed", async () => { + // Disposing the hook says nothing about whether its victim was woken. const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); - const ownDisposal: Event = { - eventType: 'hook_disposed', - eventId: slotToEventId(4), - runId: run.runId, - correlationId: 'hook_claimer', - createdAt: new Date(), - specVersion: SPEC_VERSION_CURRENT, - eventData: {}, - } as Event; + await suspendOver(queue, [ + forcedCreation(3), + ownRow(4, 'hook_disposed', 'hook_claimer'), + ]); + expect(queue).toHaveBeenCalledTimes(1); + }); + + it('stops republishing once the forced creation is outside the window', async () => { + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + await suspendOver(queue, [ + forcedCreation(3, { + createdAt: new Date( + Date.now() - FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS - 1_000 + ), + }), + ]); + expect(queue).not.toHaveBeenCalled(); + }); + + it('republishes every recent forced creation in the log, each under its own key', async () => { + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + await suspendOver(queue, [ + forcedCreation(3, { + hookId: 'hook_old', + createdAt: new Date( + Date.now() - FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS - 1_000 + ), + }), + forcedCreation(4, { hookId: 'hook_a' }), + forcedCreation(5, { + hookId: 'hook_b', + from: { ...claimedFrom, runId: 'wrun_victim_b' }, + }), + ]); + expect( + queue.mock.calls.map((call) => call[2].idempotencyKey).sort() + ).toEqual(['hook-force-claim-hook_a', 'hook-force-claim-hook_b']); + }); + + it('skips a self-claim and a victim with no recorded workflowName', async () => { + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + await suspendOver(queue, [ + forcedCreation(3, { + hookId: 'hook_self', + from: { ...claimedFrom, runId: run.runId }, + }), + forcedCreation(4, { + hookId: 'hook_legacy', + from: { runId: 'wrun_legacy', hookId: 'hook_legacy_victim' }, + }), + ]); + expect(queue).not.toHaveBeenCalled(); + }); + + it('sends each hook once per invocation across its suspensions', async () => { + // The caller passes one set for the whole invocation, so a run that + // suspends more than once does not resend the wake on every pass. + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + const forceClaimVictimWakes = new Set(); + const events = [forcedCreation(3)]; + await suspendOver(queue, events, { forceClaimVictimWakes }); + await suspendOver(queue, events, { forceClaimVictimWakes }); + expect(queue).toHaveBeenCalledTimes(1); + expect([...forceClaimVictimWakes]).toEqual(['hook_claimer']); + }); + + it('does not resend the wake of a forced creation this invocation made itself', async () => { + const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + const eventsCreate = vi.fn(async (_runId, event) => ({ + event: { + ...event, + eventId: slotToEventId(3), + createdAt: new Date(), + }, + hook: { hookId: 'hook_claimer', claimedFrom }, + })); + const forceClaimVictimWakes = new Set(); + const eventLog = { events: [] as Event[], cursor: null }; + await handleSuspension({ + suspension: new WorkflowSuspension( + new Map([ + [ + 'hook_claimer', + { + type: 'hook', + correlationId: 'hook_claimer', + token: 'channel:1', + force: true, + }, + ], + ]), + globalThis + ), + world: worldWithQueue(queue, eventsCreate), + run, + eventLog, + forceClaimVictimWakes, + }); + // The next pass of the same invocation replays over the creation. await handleSuspension({ suspension: new WorkflowSuspension(new Map(), globalThis), - world: worldWithQueue(queue), + world: worldWithQueue(queue, eventsCreate), run, - eventLog: { events: [forcedCreation(3), ownDisposal], cursor: null }, + eventLog: { events: [forcedCreation(3)], cursor: null }, + forceClaimVictimWakes, }); - expect(queue).not.toHaveBeenCalled(); + expect(queue).toHaveBeenCalledTimes(1); }); - it('stops republishing once the run has appended anything after the creation', async () => { - // Progress after the creation means the invocation that made it - // finished, wake included. Republishing on every later replay would be - // a wasted wake of the victim for the rest of this run's life. - const queue = vi.fn().mockResolvedValue({ messageId: 'msg_wake' }); + it("does not hold the suspension's writes for the republish", async () => { + // The rule reads nothing the suspension writes, so the republish rides + // alongside the writes; the queue here only answers once a write has + // been issued, which a republish-first order would never let happen. + let releaseQueue!: () => void; + const queueHeld = new Promise((resolve) => { + releaseQueue = resolve; + }); + const order: string[] = []; + const queue = vi.fn(async () => { + await queueHeld; + order.push('wake'); + return { messageId: 'msg_wake' }; + }); + const eventsCreate = vi.fn(async (_runId, event) => { + order.push(event.eventType); + releaseQueue(); + return { event }; + }); await handleSuspension({ - suspension: new WorkflowSuspension(new Map(), globalThis), - world: worldWithQueue(queue), + suspension: new WorkflowSuspension( + new Map([ + [ + 'w1', + { + type: 'wait', + correlationId: 'w1', + resumeAt: new Date(Date.now() + 60_000), + }, + ], + ]), + globalThis + ), + world: worldWithQueue(queue, eventsCreate), run, - eventLog: { - events: [forcedCreation(3), stepCreated(4)], - cursor: null, - }, + eventLog: { events: [forcedCreation(3)], cursor: null }, }); - expect(queue).not.toHaveBeenCalled(); + expect(order).toEqual(['wait_created', 'wake']); }); }); @@ -605,9 +747,9 @@ describe('handleSuspension', () => { it("does not hold other writes for a forced creation's victim wake", async () => { // The wake still goes out, once, but the sibling hook and the step are - // written while it is in flight rather than after it - // (vercel/workflow#4393 tracks the recovery gap that leaves after a - // crash). + // written while it is in flight rather than after it. A crash before + // it goes out is repaid by the next replay from the forced creation + // itself, wherever those rows land. const order: string[] = []; const eventsCreate = vi.fn(async (_runId, event) => { order.push(`${event.eventType}:${event.correlationId}`); diff --git a/packages/core/src/runtime/suspension-handler.ts b/packages/core/src/runtime/suspension-handler.ts index 777d49361d..8968f93a41 100644 --- a/packages/core/src/runtime/suspension-handler.ts +++ b/packages/core/src/runtime/suspension-handler.ts @@ -67,7 +67,7 @@ import { } from './helpers.js'; import { publishForceClaimVictimWake, - republishOwedForceClaimVictimWake, + republishOwedForceClaimVictimWakes, } from './hook-wake.js'; import { ReplayRecoveryReporter } from './replay-recovery-reporter.js'; import type { PreclaimedInlineStart } from './step-executor.js'; @@ -152,6 +152,15 @@ export interface SuspensionHandlerParams { * keep the everything-durable-at-return behavior. */ allowDeferredBatchWork?: boolean; + /** + * The invocation's record of hooks whose force-claim victim wake it has + * already published or attempted, whether as the forced creation's own wake + * or as a replay's republish (see `republishOwedForceClaimVictimWakes`). The + * caller passes one set for the whole invocation so a run that suspends more + * than once per invocation sends each hook's wake once, not once per + * suspension. Omitted, every suspension republishes on its own. + */ + forceClaimVictimWakes?: Set; } /** @@ -358,6 +367,7 @@ async function createHookEvent({ sinceCursor, createEvent, world, + forceClaimVictimWakes, }: { runId: string; hookEvent: CreateEventRequest; @@ -368,6 +378,8 @@ async function createHookEvent({ * `publishForceClaimVictimWake`. */ world: World; + /** See {@link SuspensionHandlerParams.forceClaimVictimWakes}. */ + forceClaimVictimWakes?: Set; /** * Cursor to ask the World for the event-log delta against, or undefined to * not ask. See `hookDeltaCursor` in {@link handleSuspension} for when it is @@ -404,12 +416,14 @@ async function createHookEvent({ // A forced creation that took the token over: the World journaled the // victim's `hook_disposed{forceClaimedBy}` and recorded the victim on the // hook. The victim only reads that row when something invokes it, and the - // World has no queue, so the wake is ours to publish — before the hook - // phase is considered done, so a claimer that dies here re-posts and - // republishes (the World answers a completed takeover with the same - // `claimedFrom`, on the adoption path). See `publishForceClaimVictimWake` - // for why a wake that still fails does not fail the claimer. + // World has no queue, so the wake is ours to publish. A claimer that dies + // before it goes out has left the forced `hook_created` in its log, and + // every replay inside the republish window repays it + // (`forcedCreationsOwingWake`), whatever else this suspension wrote. See + // `publishForceClaimVictimWake` for why a wake that still fails does not + // fail the claimer. if (result.hook?.claimedFrom) { + forceClaimVictimWakes?.add(result.hook.hookId); const outcome = await publishForceClaimVictimWake( world, runId, @@ -491,14 +505,23 @@ export async function handleSuspension({ stepDispatch, ownerMessageId, allowDeferredBatchWork, + forceClaimVictimWakes, }: SuspensionHandlerParams): Promise { const runId = run.runId; - // A forced creation whose victim wake this run still owes (the invocation - // that created it died before publishing) is repaid before anything else; - // see `forcedCreationOwingWake` for the rule and why it reads the log as - // loaded, before this suspension's writes. - await republishOwedForceClaimVictimWake(world, runId, eventLog?.events); + // Every recent forced creation in the loaded log may still owe its victim a + // wake (the invocation that created it may have died before publishing), so + // it is republished under the hook's idempotency key; see + // `forcedCreationsOwingWake` for the rule and its window. The rule reads + // nothing this suspension writes, so the republish goes out alongside the + // writes below instead of ahead of them, and is joined before returning. + // It never rejects. + const owedVictimWakes = republishOwedForceClaimVictimWakes( + world, + runId, + eventLog?.events, + { alreadyWoken: forceClaimVictimWakes } + ); // Turbo mode: hold every world write below until the backgrounded // `run_started` has *settled*, so we never write a step/hook/wait event for a @@ -820,6 +843,7 @@ export async function handleSuspension({ sinceCursor: hookDeltaCursor, createEvent: createGuarded, world, + forceClaimVictimWakes, }); if (result.hasHookConflict) { hookConflictCorrelationIds.push(queueItem.correlationId); @@ -912,11 +936,10 @@ export async function handleSuspension({ // // That includes forced creations. A forced creation publishes its victim's // wake before its token group's next write, but other groups and the step - // writes are not held for it, so a crash before the wake can leave one of - // them as the log's last row and hide the owed wake from the replay's - // `forcedCreationOwingWake`. That check was never airtight (an earlier - // suspension's step can finish in the same window), and making the recovery - // independent of the log's tail is tracked in vercel/workflow#4393. + // writes are not held for it. They need not be: a crash before the wake is + // repaid by the next replay from the forced `hook_created` itself, which + // `forcedCreationsOwingWake` finds wherever it sits in the log, so no row + // written after it can hide the debt. const hookGroups = [...hookItemsByToken.values()]; const hooksNeedingAbort = allHookItems.filter( (item) => item.abortRequested && !item.disposed @@ -2048,7 +2071,11 @@ export async function handleSuspension({ const nonHookOpsSettled = Promise.allSettled( ops.filter((op) => op !== hookOp) ).then(() => Date.now()); - await settlePhase(ops); + try { + await settlePhase(ops); + } finally { + await owedVictimWakes; + } // The hook writes' share of this suspension's wall time: only the stretch // they held it after every other write had settled, since until then the From 2db67e31e87bbc6a1bdf97153bd892b5288c345e Mon Sep 17 00:00:00 2001 From: "vercel[bot]" <35613825+vercel[bot]@users.noreply.github.com> Date: Fri, 25 Sep 2026 00:09:14 +0000 Subject: [PATCH 2/3] test(core): pin concurrent creation of forced hooks on different tokens Both engines used to create forced hooks one token at a time, each group waiting on its victim wake before the next started. That barrier is gone (#4392); these tests hold the first forced create until the second is issued, so a serial implementation deadlocks, and check that each victim is still woken once under its own key. Co-Authored-By: Claude Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> --- .../runtime/quickjs-force-claim-wake.test.ts | 82 +++++++++++++++++++ .../src/runtime/suspension-handler.test.ts | 60 ++++++++++++++ 2 files changed, 142 insertions(+) diff --git a/packages/core/src/runtime/quickjs-force-claim-wake.test.ts b/packages/core/src/runtime/quickjs-force-claim-wake.test.ts index ef6787d68c..dac7c06251 100644 --- a/packages/core/src/runtime/quickjs-force-claim-wake.test.ts +++ b/packages/core/src/runtime/quickjs-force-claim-wake.test.ts @@ -192,6 +192,88 @@ describe('QuickJS force-claim victim wake on replay', () => { expect(queue).not.toHaveBeenCalled(); }); + it('creates forced hooks on different tokens concurrently', async () => { + // The first forced create is held until the second has been issued, so + // creating one token at a time would deadlock here. + let releaseFirst!: () => void; + const firstHeld = new Promise((resolve) => { + releaseFirst = resolve; + }); + const wakes: string[] = []; + const queue = vi.fn( + async ( + _queueName: string, + message: { runId: string }, + opts?: { idempotencyKey?: string } + ) => { + if (message.runId !== runId) { + wakes.push(`${message.runId}:${opts?.idempotencyKey}`); + } + return { messageId: 'msg' }; + } + ); + setWorld({ + specVersion: SPEC_VERSION_CURRENT, + capabilities: { hookForceClaim: true }, + events: { + list: vi.fn(async () => ({ data: [], cursor: null, hasMore: false })), + create: vi.fn(async (_runId: string, request: CreateEventRequest) => { + const correlationId = request.correlationId as string; + if (correlationId === 'hook_forced_a') { + await firstHeld; + } else if (correlationId === 'hook_forced_b') { + releaseFirst(); + } + const letter = correlationId.slice(-1); + return { + event: { ...request, runId, eventId: `evnt_${correlationId}` }, + hook: { + hookId: correlationId, + claimedFrom: { + runId: `wrun_victim_${letter}`, + hookId: `hook_victim_${letter}`, + workflowName: 'victim-workflow', + }, + }, + }; + }), + }, + runs: { get: vi.fn(async () => workflowRun) }, + queue, + getEncryptionKeyForRun: vi.fn().mockResolvedValue(undefined), + } as unknown as World); + startQuickJSWorkflow.mockResolvedValue({ + result: { + suspended: { + pendingOperations: ['a', 'b'].map((letter) => ({ + type: 'hook', + correlationId: `hook_forced_${letter}`, + token: `forced-token-${letter}`, + isWebhook: false, + force: true, + hasCreatedEvent: false, + })), + }, + }, + continueWithEvents: vi.fn(), + dispose: vi.fn(), + }); + + const { runWorkflowWithQuickJS } = await import('./quickjs-entrypoint.js'); + await runWorkflowWithQuickJS({ + workflowCode: '// not evaluated: the VM is mocked', + workflowName: 'workflow', + workflowRun, + preloadedEvents: [], + }); + + // Each victim is still woken once, under its own claimer hook's key. + expect(wakes.sort()).toEqual([ + 'wrun_victim_a:hook-force-claim-hook_forced_a', + 'wrun_victim_b:hook-force-claim-hook_forced_b', + ]); + }); + it("does not hold the other writes for a forced creation's victim wake", async () => { // The wake still goes out, but the sibling hook and the wait are written // while it is in flight rather than after it. A crash before it goes out diff --git a/packages/core/src/runtime/suspension-handler.test.ts b/packages/core/src/runtime/suspension-handler.test.ts index 38aa035615..b1e32b8770 100644 --- a/packages/core/src/runtime/suspension-handler.test.ts +++ b/packages/core/src/runtime/suspension-handler.test.ts @@ -803,6 +803,66 @@ describe('handleSuspension', () => { }); }); + it('creates forced hooks on different tokens concurrently', async () => { + // The first forced create is held until the second has been issued, so + // creating one token at a time (each waiting on its victim wake before + // the next starts) would deadlock here. + let releaseFirst!: () => void; + const firstHeld = new Promise((resolve) => { + releaseFirst = resolve; + }); + const eventsCreate = vi.fn(async (_runId, event) => { + if (event.correlationId === 'hook_forced_a') { + await firstHeld; + } else if (event.correlationId === 'hook_forced_b') { + releaseFirst(); + } + const letter = event.correlationId.slice(-1); + return { + event, + hook: { + hookId: event.correlationId, + claimedFrom: { + runId: `wrun_victim_${letter}`, + hookId: `hook_victim_${letter}`, + workflowName: 'victim-workflow', + }, + }, + }; + }); + const queue = vi.fn(async () => ({ messageId: 'msg_wake' })); + const world = { + events: { create: eventsCreate }, + getEncryptionKeyForRun: vi.fn().mockResolvedValue(undefined), + queue, + } as unknown as World; + + await handleSuspension({ + suspension: new WorkflowSuspension( + new Map([ + hook('hook_forced_a', { force: true }), + hook('hook_forced_b', { force: true }), + ]), + globalThis + ), + world, + run, + }); + + // Each victim is still woken once, under its own claimer hook's key. + expect( + queue.mock.calls + .map(([, message, opts]) => [ + (message as { runId: string }).runId, + (opts as { idempotencyKey: string }).idempotencyKey, + ]) + .sort() + ).toEqual([ + ['wrun_victim_a', 'hook-force-claim-hook_forced_a'], + ['wrun_victim_b', 'hook-force-claim-hook_forced_b'], + ]); + }); + it('creates a hook before delivering its abort within one suspension', async () => { const order: string[] = []; const eventsCreate = vi.fn(async (_runId, event) => { From a866bbe21d043e33aa9ff6d6bfe700db13a3f9b7 Mon Sep 17 00:00:00 2001 From: Peter Wielander Date: Thu, 24 Sep 2026 17:20:02 -0700 Subject: [PATCH 3/3] Apply suggestion from @VaguelySerious Signed-off-by: Peter Wielander --- .changeset/force-claim-wake-recovery.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.changeset/force-claim-wake-recovery.md b/.changeset/force-claim-wake-recovery.md index 1ad21c1717..ef6afb529f 100644 --- a/.changeset/force-claim-wake-recovery.md +++ b/.changeset/force-claim-wake-recovery.md @@ -2,4 +2,4 @@ '@workflow/core': patch --- -Republish a force-claimed hook's victim wake on every replay within 24 hours of the takeover instead of only while the forced `hook_created` is the claimer's last own event, so a crash before the wake is repaid even when the claimer wrote other events after the creation. +Republish a force-claimed hook's victim wake on every replay within 24 hours of the takeover instead of only while the forced `hook_created` is the claimer's last own event, ensuring resilience against crashes