Skip to content

Commit 6f1e4c2

Browse files
Merge remote-tracking branch 'origin/main' into codex/hook-retention-local
# Conflicts: # packages/world-local/src/index.ts # packages/world-local/src/storage/events-storage.ts
2 parents 28d27f5 + 679dfa9 commit 6f1e4c2

42 files changed

Lines changed: 2697 additions & 1320 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
---
2+
---
3+
4+
Log the run id of each sequential-steps benchmark iteration alongside two Datadog APM links — the trigger request's trace and a span search for the run — so a suspicious STSO distribution can be opened in APM directly instead of hunted down by deployment id and time window.
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
'@workflow/core': patch
3+
'workflow': patch
4+
---
5+
6+
Simplify hardened serialization: intrinsic captures assert at import instead of degrading per use, `URLSearchParams` now serializes natively on Node 18, and bound-function getters are reported as guest code (previously a getter defined via `fn.bind()` executed with an empty report).

‎.changeset/lazy-hook-resumption.md‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,11 @@
1+
---
2+
"workflow": minor
3+
"@workflow/core": minor
4+
"@workflow/world": minor
5+
"@workflow/world-vercel": minor
6+
"@workflow/world-local": minor
7+
"@workflow/web-shared": patch
8+
"@workflow/world-postgres": patch
9+
---
10+
11+
Lazy hook resumption: on a fast path, `resumeHook()` writes the `hook_received` event and dispatches the workflow queue message concurrently instead of sequentially, cutting a round trip off resume latency. A `(runId, resumeId)` dedup constraint keeps the two writers converging on exactly one event; the runtime falls back to the sequential path when dedup support is unavailable or when `WORKFLOW_DISABLE_LAZY_HOOK_RESUME=1`. `resumeHook()` resolves to a `ResumedHook` (a `Hook` plus an optional `resilientResume: true` flag, set only when the direct write failed transiently and the resume was recovered via the queue consumer's re-ensure) — preserving the contract introduced alongside the resilient-resume work.

‎.changeset/sidebar-request-id.md‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@workflow/web-shared': minor
3+
---
4+
5+
Surface an analytics event's `vercelId` in the run sidebar's Metadata list as a copyable "Request ID", and stop rendering the literal text `null` for a provenance id the analytics contract left empty. Also order Metadata rows missing from the panel's explicit order list after the listed ones instead of ahead of them, so a failed step's Error Code no longer renders above the step's own name.

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ showSections={["parameters"]}
5050

5151
### Returns
5252

53-
Returns a `Promise<ResumedHook>`. `ResumedHook` extends `Hook` with an optional `resilientResume` flag, which is `true` when the `hook_received` event write failed with a transient error (429/5xx) and the payload will instead be delivered through the workflow queue (see the [resilient hook resume changelog](/docs/changelog/resilient-resume)). The `Hook` it extends resolves to:
53+
Returns a `Promise<ResumedHook>` — a `Hook` extended with an optional `resilientResume` flag. Resolving means the resume was accepted and the workflow will continue, whether the `hook_received` event was written directly or, on the parallel fast path, delivered through the workflow queue for the runtime to materialize (see the [lazy hook resume changelog](/docs/changelog/resilient-resume)). `resilientResume` is `true` only when the direct event write failed transiently and the resume was recovered through the queue; on the happy path it is absent. The resolved hook:
5454

5555
<TSDoc
5656
definition={`

‎docs/content/docs/v5/changelog/resilient-resume.mdx‎

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,16 +7,16 @@ description: resumeHook() now tolerates transient event storage failures, as lon
77

88
## Motivation
99

10-
When event storage is degraded but the queue is up, `resumeHook()` previously failed entirely because the `hook_received` event write happens before the queue dispatch. This change brings `resumeHook()` to feature parity with [resilient `start()`](/docs/changelog/resilient-start): a transient event-write failure (429/5xx) no longer fails the resume when the payload can be delivered through the queue instead.
10+
`resumeHook()` used to write the `hook_received` event and dispatch the workflow queue message strictly one after the other, so every resume paid two sequential round trips and a transient event-storage failure failed the whole resume even when the queue was healthy. This change runs both writes **concurrently** — cutting a round trip off resume latency — and, on the same path, brings `resumeHook()` to parity with [resilient `start()`](/docs/changelog/resilient-start): a transient event-write failure no longer fails the resume when the payload can still be delivered through the queue.
1111

1212
## Design
1313

14-
- `resumeHook()` writes the `hook_received` event first, then dispatches to the workflow queue. The order is deliberately sequential: dispatching in parallel could let the workflow runtime process the message and materialize a duplicate `hook_received` before the direct write commits.
15-
- If the event write fails with a retryable error (429/5xx), the queue dispatch carries a `hookInput` payload — the dehydrated hook payload plus a client-minted `resumeId` idempotency key and the hook token. The workflow runtime materializes the missing `hook_received` event from `hookInput` during replay, deduplicating by `resumeId` in case the direct write actually committed or the queue message is redelivered.
16-
- Replay itself also deduplicates: `hook_received` events sharing a `resumeId` belong to the same resume attempt, and only the first in the event log is delivered to workflow code. Even if concurrent queue redeliveries materialize the event twice, the payload reaches the workflow exactly once.
17-
- If the event write fails with any other error, or the queue dispatch fails, `resumeHook()` throws as before.
18-
- `resumeHook()` now returns `ResumedHook` (exported from `workflow/api`), which extends `Hook` with an optional `resilientResume` flag. The flag is `true` when the fallback path was taken — the resume is accepted and will be delivered through the queue.
14+
- On the fast path, `resumeHook()` writes the `hook_received` event and dispatches the workflow queue message **concurrently** (`Promise.allSettled`). The queue message carries a `hookInput` payload — the dehydrated hook payload plus a client-minted `resumeId` idempotency key, the hook token, and a payload digest.
15+
- A `(runId, resumeId)` dedup constraint keeps the two writers converging on **exactly one** `hook_received` event: whichever lands first wins, and the other is resolved server-side as success rather than a duplicate. The queue consumer idempotently re-ensures the event from `hookInput` before replay, so the resume is guaranteed even if the direct write never commits.
16+
- Replay also deduplicates: `hook_received` events sharing a `resumeId` belong to the same resume attempt, and only the first in the event log is delivered to workflow code. Even if a redelivery materializes the event twice, the payload reaches the workflow exactly once.
17+
- **Queue dispatch failure is fatal** — the run was not re-triggered, so no consumer will re-ensure the event, and `resumeHook()` throws. A transient event-write failure (429/5xx, a transport error, or an expected `(runId, resumeId)` conflict with the consumer's own re-ensure) is swallowed because the queue delivery still guarantees the resume; a terminal run surfaces as `HookNotFoundError`, and any other event-write error is rethrown.
18+
- `resumeHook()` returns `ResumedHook` (exported from `workflow/api`), which extends `Hook` with an optional `resilientResume` flag. The flag is `true` only when the direct write failed transiently and the resume was recovered through the queue; on the happy path and the sequential fallback it is absent.
1919

2020
## Compatibility
2121

22-
The resilient path is only taken when the target workflow run's deployment is recent enough to understand `hookInput` on the queue message. Because runs keep executing on the deployment they were created on, a resume that targets a run from an older deployment fails fast with the original error instead — the previous behavior, allowing the caller to retry.
22+
The parallel fast path is gated per resume: it activates only when both the target run's queue consumer and the live backend independently attest dedup support (re-checked on every resume, so rollout and rollback both degrade safely). Otherwise — for oversized payloads, legacy runs, or with `WORKFLOW_DISABLE_LAZY_HOOK_RESUME=1` — `resumeHook()` falls back to the original sequential write-then-dispatch path. Because runs keep executing on the deployment they were created on, a resume targeting a run from an older deployment simply uses the sequential path.

‎docs/content/docs/v5/configuration/runtime-tuning.mdx‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,13 @@ For example, a workflow can run a 10-minute inline step even with `WORKFLOW_REPL
3838
- Default: `3`
3939
- Recovery replays before replay divergence is recorded as corruption.
4040

41+
### `WORKFLOW_DISABLE_LAZY_HOOK_RESUME`
42+
43+
- Default: enabled (lazy hook resume on)
44+
- Resuming a hook persists the `hook_received` event and publishes the workflow invocation concurrently, cutting a round trip off resume latency. On this parallel path the queue message also carries the payload, so a transient event-write failure still resumes the run — the queue consumer re-ensures the `hook_received` event before replay. A backend `(runId, resumeId)` constraint keeps the two writers converging on exactly one event.
45+
- The runtime falls back to the sequential path automatically when the consumer or backend does not attest dedup support (or the payload is too large to inline on the queue message). On the sequential path the event is written *before* dispatch and its failure fails the resume — the fallback trades that resilience away to stay safe when dedup is not enforced, it does not preserve it.
46+
- Set `1` to force the sequential path as a kill switch. The chosen strategy is reported on the resume span as `workflow.hook.resume_strategy`.
47+
4148
### `WORKFLOW_PRECONDITION_GUARD`
4249

4350
- Default: enabled

‎packages/core/e2e/benchmark.test.ts‎

Lines changed: 50 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -193,6 +193,9 @@ interface StreamIterationResult {
193193

194194
interface SequentialIterationResult {
195195
runId: string;
196+
/** Datadog trace id for the `/api/bench` request that started this run, when
197+
* the deployment's route reports one (older deployments won't). */
198+
traceId?: string;
196199
/** STSO gaps preceding an 'inline' step (same warm process as the step
197200
* before it) — the framework's pure step-to-step overhead. */
198201
stsoInlineMs: number[];
@@ -221,6 +224,8 @@ interface BenchTriggerResponse {
221224
runId: string;
222225
/** Date.now() stamped in the route immediately before start(). */
223226
clientStart: number;
227+
/** Datadog trace id for this request, when the route reports one. */
228+
traceId?: string;
224229
}
225230

226231
function withTimeout<T>(
@@ -272,7 +277,13 @@ async function triggerBenchRun(
272277
`bench trigger for ${workflowFn} returned malformed body: ${JSON.stringify(data)?.slice(0, 200)}`
273278
);
274279
}
275-
return { runId: data.runId, clientStart: data.clientStart };
280+
return {
281+
runId: data.runId,
282+
clientStart: data.clientStart,
283+
// Optional: a deployment built before the route reported it simply logs
284+
// the run id without a trace link.
285+
traceId: typeof data.traceId === 'string' ? data.traceId : undefined,
286+
};
276287
}
277288

278289
/** Poll a run's return value to completion (the handle polls internally). */
@@ -340,7 +351,7 @@ async function runStreamIteration(
340351
async function runSequentialIteration(
341352
stepCount: number
342353
): Promise<SequentialIterationResult> {
343-
const { runId, clientStart } = await triggerBenchRun(
354+
const { runId, clientStart, traceId } = await triggerBenchRun(
344355
'benchSequentialStepsWorkflow',
345356
[stepCount]
346357
);
@@ -366,6 +377,7 @@ async function runSequentialIteration(
366377

367378
return {
368379
runId,
380+
traceId,
369381
stsoInlineMs,
370382
stsoQueueHopMs,
371383
woMs: workflowOverheadMs(clientStart, steps),
@@ -636,6 +648,29 @@ const SCENARIO_DESCRIPTIONS = [
636648
},
637649
];
638650

651+
// Datadog APM permalink for a trace id. The benchmark deployment exports its
652+
// OTel spans to Datadog, and `/api/bench` returns the trace id of the request
653+
// that started each run.
654+
const DATADOG_TRACE_URL = 'https://app.datadoghq.com/apm/trace/';
655+
656+
/**
657+
* Datadog APM search for the spans tagged with a given `workflow.run.id`.
658+
*
659+
* The permalink above opens the *trigger's* trace, which under the default
660+
* `WORKFLOW_TRACE_MODE=linked` holds only `workflow.start` plus span links out
661+
* to the per-invocation trace roots — an entry point to the run rather than the
662+
* run itself. This search lands straight on the execution spans, which is where
663+
* an STSO investigation actually goes, so both are logged.
664+
*
665+
* Depends on `workflow.run.id` being an indexed span tag in the Datadog org. If
666+
* it isn't, this returns an empty search and the trace permalink stays the way
667+
* in; neither link is load-bearing for the benchmark itself.
668+
*/
669+
function datadogRunSearchUrl(runId: string): string {
670+
const query = encodeURIComponent(`@workflow.run.id:${runId}`);
671+
return `https://app.datadoghq.com/apm/traces?query=${query}`;
672+
}
673+
639674
describe('workflow benchmarks', () => {
640675
// Preflight: prove the deployment executes workflows (and the trigger route
641676
// works) before any scenario spends its attempt budget. Without this, a
@@ -777,6 +812,19 @@ describe('workflow benchmarks', () => {
777812
extraAttempts: Math.max(2, Math.ceil(SEQUENTIAL_ITERATIONS * 0.5)),
778813
}
779814
);
815+
// Name the runs behind the STSO histograms in this job's own log, right
816+
// where they were produced. When a bucket looks wrong the investigation
817+
// starts in APM, and this saves the usual hunt by deployment id + time
818+
// window. Logged rather than rendered into the PR comment so it is also
819+
// there for a local `pnpm bench` and for a run whose comment step never
820+
// gets to execute.
821+
for (const { runId, traceId } of results) {
822+
console.log(
823+
`[bench] ${SCENARIO_SEQUENTIAL} run ${runId}` +
824+
(traceId ? ` — trace ${DATADOG_TRACE_URL}${traceId}` : '') +
825+
` — spans ${datadogRunSearchUrl(runId)}`
826+
);
827+
}
780828
// Report STSO split by whether the step that ends the gap was 'inline'
781829
// (same warm process as the step before it — pure framework overhead) or
782830
// a 'queue-hop' (first step of a fresh process — dispatch + reinit cost).

‎packages/core/e2e/e2e.test.ts‎

Lines changed: 0 additions & 103 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,6 @@ import {
2222
test,
2323
} from 'vitest';
2424
import { getTrustedSourcesHeaders } from '../../../scripts/trusted-sources-headers.mjs';
25-
import { QUEUE_HOOK_INPUT_MIN_VERSION } from '../src/capabilities';
2625
import type { Run } from '../src/runtime';
2726
import {
2827
getHookByToken,
@@ -31,7 +30,6 @@ import {
3130
healthCheck,
3231
start as rawStart,
3332
resumeHook,
34-
setWorld,
3533
} from '../src/runtime';
3634
import {
3735
cliCancel,
@@ -4044,107 +4042,6 @@ describe('e2e', () => {
40444042
}
40454043
);
40464044

4047-
// ============================================================
4048-
// Resilient resume: hook payload delivered even when hook_received fails
4049-
// ============================================================
4050-
test(
4051-
'resilient resume: hookWorkflow receives payload when hook_received returns 500',
4052-
{ timeout: 60_000 },
4053-
async () => {
4054-
const token = `resilient-resume-${Math.random().toString(36).slice(2)}`;
4055-
const customData = Math.random().toString(36).slice(2);
4056-
4057-
// Start the hook-awaiting workflow normally
4058-
const run = await start(await e2e('hookWorkflow'), [token, customData]);
4059-
4060-
// Wait for the hook to be registered
4061-
await sleep(5_000);
4062-
4063-
// The resilient path is gated on the target run's recorded
4064-
// `@workflow/core` version understanding `hookInput` on the queue
4065-
// payload (see QUEUE_HOOK_INPUT_MIN_VERSION in capabilities.ts).
4066-
// This build's own not-yet-bumped version sits below that published
4067-
// cutoff — deliberately, there is no own-version escape hatch — so
4068-
// simulate a run recorded by a capable release. The target deployment
4069-
// here genuinely understands `hookInput` (it is built from this same
4070-
// source), so the simulation is honest.
4071-
const capableVersion = QUEUE_HOOK_INPUT_MIN_VERSION;
4072-
4073-
// Build a stubbed world whose events.create throws a 500 on the
4074-
// hook_received write, but passes all other events through. The queue
4075-
// dispatch should still succeed, and the workflow runtime should
4076-
// materialize the missing hook_received event from `hookInput` on the
4077-
// queue message (resilient resume).
4078-
const realWorld = await getWorld();
4079-
const stubbedWorld: World = {
4080-
...realWorld,
4081-
events: {
4082-
...realWorld.events,
4083-
create: (async (...args: Parameters<World['events']['create']>) => {
4084-
const [, event] = args;
4085-
if (event.eventType === 'hook_received') {
4086-
throw new WorkflowWorldError('Simulated storage outage', {
4087-
status: 500,
4088-
});
4089-
}
4090-
return realWorld.events.create(...args);
4091-
}) as World['events']['create'],
4092-
},
4093-
runs: {
4094-
...realWorld.runs,
4095-
// Fallback path (hooks without a stored resumeContext): rewrite
4096-
// the run's recorded core version to the capable release.
4097-
get: (async (...args: Parameters<World['runs']['get']>) => {
4098-
const run = await realWorld.runs.get(...args);
4099-
return {
4100-
...run,
4101-
executionContext: {
4102-
...run.executionContext,
4103-
workflowCoreVersion: capableVersion,
4104-
},
4105-
};
4106-
}) as World['runs']['get'],
4107-
},
4108-
};
4109-
4110-
const hook = await getHookByToken(token);
4111-
expect(hook.runId).toBe(run.runId);
4112-
// Fast path (hooks with a stored resumeContext): rewrite the recorded
4113-
// core version on the hook object we pass to resumeHook().
4114-
if (hook.resumeContext) {
4115-
hook.resumeContext.workflowCoreVersion = capableVersion;
4116-
}
4117-
4118-
// Swap in the stubbed world for the duration of the resumeHook() call.
4119-
// `resumeHook` uses `getWorld()` internally (no `world` option), so we
4120-
// use `setWorld()` to replace the cached instance and restore the real
4121-
// one afterwards.
4122-
setWorld(stubbedWorld);
4123-
let resumedHook: Awaited<ReturnType<typeof resumeHook>>;
4124-
try {
4125-
resumedHook = await resumeHook(hook, {
4126-
message: 'via-resilient-resume',
4127-
customData: (hook.metadata as any)?.customData,
4128-
done: true,
4129-
});
4130-
} finally {
4131-
setWorld(realWorld);
4132-
}
4133-
4134-
// The direct hook_received write failed with 500, so resumeHook should
4135-
// have taken the resilient path and flagged the returned hook.
4136-
expect(resumedHook.resilientResume).toBe(true);
4137-
4138-
// Despite hook_received failing, the workflow should still receive
4139-
// the payload via the runtime's queue-payload fallback.
4140-
const returnValue = await run.returnValue;
4141-
expect(returnValue).toHaveLength(1);
4142-
expect(returnValue[0].message).toBe('via-resilient-resume');
4143-
expect(returnValue[0].customData).toBe(customData);
4144-
expect(returnValue[0].done).toBe(true);
4145-
}
4146-
);
4147-
41484045
test(
41494046
'getterStepWorkflow - getter functions with "use step" directive',
41504047
{ timeout: 60_000 },

0 commit comments

Comments
 (0)