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/force-claim-wake-recovery.md
Original file line number Diff line number Diff line change
@@ -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, ensuring resilience against crashes
Original file line number Diff line number Diff line change
Expand Up @@ -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.

<Callout type="info">
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.
Expand Down
6 changes: 6 additions & 0 deletions packages/core/src/runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string>();
const replayRecoveryReporter = replayDivergence
? new ReplayRecoveryReporter(replayDivergence.count)
: ReplayRecoveryReporter.inert();
Expand Down Expand Up @@ -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
Expand Down
182 changes: 115 additions & 67 deletions packages/core/src/runtime/hook-wake.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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-<hookId>`, 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<string>; nowMs?: number } = {}
): Promise<void> {
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,
}
);
})
);
}
23 changes: 12 additions & 11 deletions packages/core/src/runtime/quickjs-entrypoint.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ import {
} from './helpers.js';
import {
publishForceClaimVictimWake,
republishOwedForceClaimVictimWake,
republishOwedForceClaimVictimWakes,
} from './hook-wake.js';
import {
dispatchRunCompletedHooks,
Expand Down Expand Up @@ -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));
}
Expand Down Expand Up @@ -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) => {
Expand Down
Loading
Loading