Skip to content

Commit ff80d62

Browse files
pranaygpvercel[bot]VaguelySerious
authored
fix(core): repay a force-claim victim's wake independent of the claimer's log tail (#4398)
* 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-<hookId>` 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 <noreply@anthropic.com> Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> * 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 <noreply@anthropic.com> Co-Authored-By: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> * Apply suggestion from @VaguelySerious Signed-off-by: Peter Wielander <mittgfu@gmail.com> --------- Signed-off-by: Peter Wielander <mittgfu@gmail.com> Co-authored-by: vercel[bot] <35613825+vercel[bot]@users.noreply.github.com> Co-authored-by: Pranay Prakash <1797812+pranaygp@users.noreply.github.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com>
1 parent 35ebeb4 commit ff80d62

8 files changed

Lines changed: 609 additions & 199 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+
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

‎docs/content/docs/v5/api-reference/workflow/create-hook.mdx‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -251,7 +251,7 @@ With `experimental_force`, this run always ends up owning the token:
251251
- 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.
252252
- 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.
253253

254-
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.
254+
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.
255255

256256
<Callout type="info">
257257
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.

‎packages/core/src/runtime.ts‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -967,6 +967,11 @@ export function workflowEntrypoint(
967967
// than when the wait's own timer would have fired.
968968
let eventLogFromInlineDelta = false;
969969
let loopIteration = 0;
970+
// Hooks whose force-claim victim wake this invocation has
971+
// already sent (its own forced creations, and the replay's
972+
// republishes), so each suspension of the loop below does
973+
// not send them again. See `forcedCreationsOwingWake`.
974+
const forceClaimVictimWakes = new Set<string>();
970975
const replayRecoveryReporter = replayDivergence
971976
? new ReplayRecoveryReporter(replayDivergence.count)
972977
: ReplayRecoveryReporter.inert();
@@ -3544,6 +3549,7 @@ export function workflowEntrypoint(
35443549
eventLog,
35453550
runReadyBarrier,
35463551
replayRecoveryReporter,
3552+
forceClaimVictimWakes,
35473553
// Resilient step dispatch: lets eligible newly
35483554
// created steps publish their step-execution
35493555
// message (carrying `stepInput`) in parallel with

‎packages/core/src/runtime/hook-wake.ts‎

Lines changed: 115 additions & 67 deletions
Original file line numberDiff line numberDiff line change
@@ -100,9 +100,11 @@ export async function publishHookWakeWithRetry(
100100
* Same durability contract as `resumeHook()`'s wake: the row is durable
101101
* before this runs, the publish is retried on transport-shaped failures, and
102102
* a publish that still fails is logged rather than failing the claimer —
103-
* nothing of the claimer's is wrong, and the victim reads the row on its
104-
* next invocation for any reason. The idempotency key is the claimer's hook
105-
* id, so a claimer retrying its creation republishes at most one wake.
103+
* nothing of the claimer's is wrong, every later replay of the claimer inside
104+
* the republish window tries again ({@link forcedCreationsOwingWake}), and
105+
* the victim reads the row on its next invocation for any reason. The
106+
* idempotency key is the claimer's hook id, so those republishes collapse
107+
* into one wake.
106108
*
107109
* Skipped when the victim is the claimer itself (a run taking over its own
108110
* earlier hook is already running) and when the World recorded no
@@ -154,86 +156,132 @@ export async function publishForceClaimVictimWake(
154156
}
155157

156158
/**
157-
* The forced hook creation whose victim wake this run still owes, if any: the
158-
* forced `hook_created` is the last event the run's own replay appended.
159+
* How long after a forced `hook_created` a replay of the claimer keeps
160+
* republishing its victim's wake: 24 hours, measured from the row's
161+
* `createdAt`.
162+
*
163+
* The bound has to outlast every redelivery of the invocation that journaled
164+
* the creation, because that redelivery is the replay that must repay a wake
165+
* the invocation died before publishing. A queue message is retained for 24
166+
* hours from its send and the creation is written after the send, so any such
167+
* redelivery arrives within 24 hours of the creation. The same 24 hours is the
168+
* Vercel queue's idempotency window (`min(retention, 24h)`), so every
169+
* republish inside it collapses, under `hook-force-claim-<hookId>`, into the
170+
* one wake that was (or now is) delivered. Past it a republish would be a
171+
* genuinely new message, which is what the bound saves. world-postgres
172+
* remembers a completed key in-process to the same effect; world-local
173+
* dedupes a key only while its message is in flight, so there a republish can
174+
* deliver the victim one more replay, which reads nothing new.
175+
*/
176+
export const FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS = 24 * 60 * 60 * 1000;
177+
178+
/**
179+
* The forced hook creations whose victim wake this run may still owe: every
180+
* forced `hook_created` (one carrying `forceClaimedFrom`) in the log whose
181+
* `createdAt` is within {@link FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS} of
182+
* `nowMs`.
159183
*
160184
* A forced creation is followed by a wake of the run it took the token from.
161185
* If the invocation died between the two, the creation is in the log and the
162-
* victim was never told; the row itself is the durable record of that debt.
163-
* As long as it is the last event THIS RUN wrote, the run has made no progress
164-
* since, so the invocation that should have woken the victim did not finish,
165-
* and the replay republishes (under the hook's idempotency key, so a wake that
166-
* did go out is not duplicated). The first event the run appends after it
167-
* ends the republishing.
186+
* victim was never told; the row itself is the durable record of that debt,
187+
* and nothing records that the wake went out. So the replay does not try to
188+
* infer it: it republishes for every recent forced creation, and the hook's
189+
* idempotency key collapses a wake that did go out (a duplicate that slips
190+
* past a World's dedupe is one harmless replay of the victim).
168191
*
169-
* "This run wrote" matters: a delivery appends `hook_received` to this log
170-
* from another request, a sealed-log World appends `noop`, and a LATER
171-
* claimer taking the token from this run appends
172-
* `hook_disposed{forceClaimedBy}`. None is progress of this run — the model
173-
* (`ForceWakeOnce.cfg`'s sibling trace) has a delivery land between the crash
174-
* and the retry, and a rule that looked at the bare tail would then never
175-
* wake the victim. The foreign disposal is the chain case: this run took the
176-
* token from A, died before waking A, and was itself taken from by C. Its own
177-
* wake (from C) is the very invocation that must repay A's — the run's own
178-
* `hook_disposed` (a `dispose()` in its code) IS its progress and still ends
179-
* the debt, but a row another run put here does not. Both engines call this
180-
* on the log they loaded for the invocation, before writing anything.
192+
* The rule reads nothing written after the creation, which is what makes it
193+
* sound. Any row can land between the creation and the wake: a step, wait,
194+
* attribute or other hook row the same suspension writes concurrently
195+
* (neither engine holds those for the wake), a step or wait terminal from
196+
* another invocation, a delivery's `hook_received`, a sealed-log `noop`, or a
197+
* later claimer's `hook_disposed{forceClaimedBy}` taking the token from this
198+
* run. None of them says the wake was published. The last is the chain case:
199+
* this run took the token from A, died before waking A, and was taken from by
200+
* C; C's wake of this run is the invocation that repays A's.
181201
*
182-
* Reading only the last own row is best-effort. Any row this run writes after
183-
* the forced creation and before the wake goes out ends the debt as if the
184-
* wake had been published: a step, wait, attribute or other hook row the same
185-
* suspension writes concurrently (neither engine holds those for the wake), or
186-
* a step or wait terminal from another invocation. If the invocation then
187-
* dies before publishing, the victim reads its disposal only on its next
188-
* invocation for any other reason — the same outcome as a wake whose publish
189-
* fails outright. Making the recovery independent of the log's tail is
190-
* vercel/workflow#4393.
202+
* The time is the row's `createdAt`, never one decoded from its event id (a
203+
* slot-numbered id carries none). Under slot identity it is the writer's
204+
* client clock, clamped by the World to at most an hour ahead of its own, so
205+
* skew moves the window's edge by at most that much; a creation dated ahead of
206+
* `nowMs` counts as recent, and one whose time cannot be read is treated as
207+
* recent too, erring toward a wake rather than a stranded victim.
191208
*/
192-
export function forcedCreationOwingWake(
193-
events: readonly Event[] | undefined
194-
): (Event & { eventType: 'hook_created' }) | undefined {
195-
if (!events) return undefined;
196-
for (let i = events.length - 1; i >= 0; i--) {
197-
const event = events[i];
209+
export function forcedCreationsOwingWake(
210+
events: readonly Event[] | undefined,
211+
nowMs: number = Date.now()
212+
): (Event & { eventType: 'hook_created' })[] {
213+
const owed: (Event & { eventType: 'hook_created' })[] = [];
214+
if (!events) return owed;
215+
for (const event of events) {
198216
if (
199-
event.eventType === 'hook_received' ||
200-
event.eventType === 'noop' ||
201-
(event.eventType === 'hook_disposed' &&
202-
event.eventData?.forceClaimedBy !== undefined)
217+
event.eventType !== 'hook_created' ||
218+
event.eventData?.forceClaimedFrom === undefined
203219
) {
204220
continue;
205221
}
206-
return event.eventType === 'hook_created' &&
207-
event.eventData.forceClaimedFrom !== undefined
208-
? event
209-
: undefined;
222+
const createdAtMs = new Date(event.createdAt).getTime();
223+
if (
224+
Number.isNaN(createdAtMs) ||
225+
nowMs - createdAtMs < FORCE_CLAIM_WAKE_REPUBLISH_WINDOW_MS
226+
) {
227+
owed.push(event);
228+
}
210229
}
211-
return undefined;
230+
return owed;
212231
}
213232

214233
/**
215-
* Republish the wake {@link forcedCreationOwingWake} says is owed. Shared by
216-
* the node:vm suspension handler and the QuickJS entrypoint so the two engines
217-
* cannot drift on the durability contract.
234+
* Republish the wakes {@link forcedCreationsOwingWake} says may be owed.
235+
* Shared by the node:vm suspension handler and the QuickJS entrypoint so the
236+
* two engines cannot drift on the durability contract.
237+
*
238+
* `alreadyWoken` is the invocation's record of the hooks whose victim wake it
239+
* has already published or attempted, the forced creations it made itself
240+
* included: an engine that replays more than once per invocation passes the
241+
* same set every time, so each hook costs one send per invocation. Hooks this
242+
* call publishes are added to it.
243+
*
244+
* Self-claims and victims with no recorded `workflowName` are skipped
245+
* silently: the creating invocation already logged the latter, and a replay
246+
* repeating it would say nothing new. Never rejects; a publish that fails
247+
* after its retries is logged, and the next replay inside the window tries
248+
* again.
218249
*/
219-
export async function republishOwedForceClaimVictimWake(
250+
export async function republishOwedForceClaimVictimWakes(
220251
world: World,
221252
runId: string,
222-
events: readonly Event[] | undefined
253+
events: readonly Event[] | undefined,
254+
options: { alreadyWoken?: Set<string>; nowMs?: number } = {}
223255
): Promise<void> {
224-
const owed = forcedCreationOwingWake(events);
225-
if (!owed) return;
226-
const claimedFrom = owed.eventData.forceClaimedFrom as HookClaimedFrom;
227-
const outcome = await publishForceClaimVictimWake(world, runId, {
228-
hookId: owed.correlationId,
229-
claimedFrom,
230-
});
231-
if (outcome !== 'skipped') {
232-
runtimeLogger.info('Republished the wake of a force-claimed hook victim', {
233-
workflowRunId: runId,
234-
hookId: owed.correlationId,
235-
victimRunId: claimedFrom.runId,
236-
victimWake: outcome,
237-
});
238-
}
256+
const { alreadyWoken } = options;
257+
const owed = forcedCreationsOwingWake(events, options.nowMs).filter(
258+
(event) => {
259+
const from = event.eventData.forceClaimedFrom as HookClaimedFrom;
260+
return (
261+
from.runId !== runId &&
262+
from.workflowName !== undefined &&
263+
!alreadyWoken?.has(event.correlationId)
264+
);
265+
}
266+
);
267+
if (owed.length === 0) return;
268+
await Promise.all(
269+
owed.map(async (event) => {
270+
alreadyWoken?.add(event.correlationId);
271+
const claimedFrom = event.eventData.forceClaimedFrom as HookClaimedFrom;
272+
const outcome = await publishForceClaimVictimWake(world, runId, {
273+
hookId: event.correlationId,
274+
claimedFrom,
275+
});
276+
runtimeLogger.debug(
277+
'Republished the wake of a force-claimed hook victim',
278+
{
279+
workflowRunId: runId,
280+
hookId: event.correlationId,
281+
victimRunId: claimedFrom.runId,
282+
victimWake: outcome,
283+
}
284+
);
285+
})
286+
);
239287
}

‎packages/core/src/runtime/quickjs-entrypoint.ts‎

Lines changed: 12 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -67,7 +67,7 @@ import {
6767
} from './helpers.js';
6868
import {
6969
publishForceClaimVictimWake,
70-
republishOwedForceClaimVictimWake,
70+
republishOwedForceClaimVictimWakes,
7171
} from './hook-wake.js';
7272
import {
7373
dispatchRunCompletedHooks,
@@ -594,10 +594,10 @@ async function dispatchPendingOps(params: {
594594
};
595595
// Token groups run in parallel with every other op, forced creations
596596
// included. A forced creation publishes its victim's wake before its group's
597-
// next write, but nothing else waits for it, so a crash before the wake can
598-
// leave another row as the log's last and hide the owed wake from the
599-
// replay's `forcedCreationOwingWake`. Making that recovery independent of
600-
// the log's tail is tracked in vercel/workflow#4393.
597+
// next write, but nothing else waits for it, and nothing needs to: a crash
598+
// before the wake is repaid by the next replay from the forced
599+
// `hook_created` itself, which `forcedCreationsOwingWake` finds wherever it
600+
// sits in the log, so no row written after it can hide the debt.
601601
for (const group of hookOpsByToken.values()) {
602602
opsPromises.push(runHookGroup(group));
603603
}
@@ -1129,12 +1129,13 @@ export async function runWorkflowWithQuickJS(params: {
11291129
// handed back on a write that the VM has not been given yet. Every write
11301130
// made from this view goes through `createEvent` below so it names the
11311131
// position it was decided against and its response is queued here.
1132-
// Same durability contract as the node:vm suspension handler: a forced
1133-
// hook creation that is still the last event this run wrote owes its
1134-
// victim a wake, because the invocation that created it died before
1135-
// publishing one. Repaid here, on the log as loaded, before this
1136-
// invocation writes anything.
1137-
await republishOwedForceClaimVictimWake(world, runId, events);
1132+
// Same durability contract as the node:vm suspension handler: every
1133+
// recent forced hook creation in the log may still owe its victim a wake,
1134+
// because the invocation that created it may have died before publishing
1135+
// one, so it is republished under the hook's idempotency key (see
1136+
// `forcedCreationsOwingWake`). Once per invocation, on the log as loaded;
1137+
// the forced creations this invocation makes publish their own.
1138+
await republishOwedForceClaimVictimWakes(world, runId, events);
11381139

11391140
const logView = new QuickJSLogView(events, loadedCursor);
11401141
const createEvent: EventCreator = async (data, eventParams) => {

0 commit comments

Comments
 (0)