mirror of
https://github.com/vercel/workflow.git
synced 2026-09-14 19:59:43 +08:00
perf/script-precompile-module-init
443 Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
1ffc078e8e |
perf(core): precompile workflow vm.Script at module init
Compile the workflow bundle's `vm.Script` for each known workflow source filename when `workflowEntrypoint` is constructed (module-init time), rather than lazily on the first queue delivery's replay. Builders inline the deduplicated, sorted set of workflow filenames into generated routes via the new `workflowFilenames` entrypoint option, so the first replay is a cache hit instead of paying the bundle parse/compile on the critical path. |
||
|
|
ab2e9b8d07 | [core] Send workflowName with step events (#2511) | ||
|
|
a92c16debd |
Reject empty-string hook tokens in createHook() (#2490)
createHook() used `options.token ?? ctx.generateNanoid()`, so a nullish token fell back to a generated one but an empty string `""` was accepted verbatim — a meaningless, non-deterministic token that is almost always an accidental value (e.g. an unset variable). Throw a clear error when an explicit empty-string token is passed; `undefined`/`null` still auto-generate, and non-empty strings are unchanged. Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
939890d4c2 |
perf(core): cache compiled workflow-bundle vm.Script across replays (#2471)
* perf(core): cache compiled workflow-bundle vm.Script across replays The inline replay loop calls runWorkflow on every iteration, and each call re-parsed the entire workflow bundle string via vm.runInContext. For a bundle containing many workflow definitions (the production shape: one workflow called per replay), this re-scans every definition on every replay. Cache the compiled vm.Script per process, keyed by (workflowCode, filename), and run it against the fresh context instead of recompiling. Compilation is a pure function of (code, filename), so the result is byte-identical to the previous re-parse-every-time behaviour — determinism is preserved. filename is part of the key because it drives source attribution in stack traces (consumed by remapErrorStack). Measured per-replay savings scale with bundle size (and multiply by replay count): ~34% for a 50-workflow app, ~59% for 155 workflows, ~80% for 400. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * perf(core): bound script cache with LRU; soften determinism claim; add tests Addresses review on #2471: - Bound `scriptCache` to a small LRU (cap 8 bundle versions). Production serves one bundle per process so the bound is never reached; it exists for dev/watch mode, where each edit produces a new bundle string that would otherwise be pinned forever (~0.8MB/edit, monotonic). Touch-on-access keeps the latest bundle hot; evicting a `code` entry drops its per-filename scripts together, restoring pre-cache GC behaviour. - Document precisely why keying includes `filename` (intentional: drives stack-trace attribution via `remapErrorStack`; NOT a dedupe key), and that the whole bundle is compiled once per distinct filename. - Soften the "byte-identical including thrown errors" claim to same-workflow-function + same-`filename`-attribution, noting the one caveat: a lookup-expression error's line number shifts to line 1 of the separate lookup Script. Updated in both the code comment and the PR description. - Add tests: cache-is-bounded regression (eviction past the cap), LRU recency (hot bundle survives churn), and a realistic multi-workflow collision test (distinct code/filename never returns the wrong Script, results carry their own bundle marker). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
16b36703e2 |
perf(core): drain consumable replay events synchronously (#2473)
* docs(core): document scheduleWhenIdle macrotask is load-bearing Revert the synchronous consume-loop drain optimization: it caused a replay divergence (ReplayDivergenceError on step_started → CorruptedEventLogError) in the world-testing inline-batches parallel workflow on the Windows CI runner. The per-event `process.nextTick` in the consume loop is load-bearing — it guarantees at most one event is consumed per macrotask, letting the cross-VM `resolve → workflow VM body → subscribe()` chain register the next operation's consumer before the drain advances. A synchronous drain races ahead of that registration. What remains is a documentation comment on `scheduleWhenIdle` capturing why its initial `setTimeout(0)` must not be downgraded to a microtask (empirically: queueMicrotask breaks hook/sleep Promise.race ordering → CorruptedEventLogError). No behavior change. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * docs(core): use empty changeset for comment-only macrotask doc The scheduleWhenIdle change is a pure code comment with no consumer-facing effect, so it does not warrant a patch bump / changelog entry. Replace the patch changeset with an empty one to satisfy the changeset-bot convention without claiming a release. Per the PR template's `pnpm changeset --empty` guidance for non-releasing changes. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> * perf(core): drain consumable replay events synchronously The EventsConsumer rescheduled `process.nextTick(this.consume)` after every consumed event, so replaying N already-consumable events (structural lifecycle events, step_created/step_started, completed deliveries) cost N macrotask hops — O(N) per consume wave across a sequential replay. Drain consecutively consumable events within a single synchronous pass instead. This is safe because callbacks only ever consume events with a consumer that is already registered; new consumers are registered by workflow VM body code that runs asynchronously off ctx.promiseQueue after a delivery resolve(). When the next event's consumer is not yet registered, no callback consumes it and we fall through to the existing cross-VM-safe deferred unconsumed-event check, exactly as before. A null end-of-events sentinel never continues the drain, so it cannot spin past end-of-log. scheduleWhenIdle is intentionally left unchanged: its initial setTimeout(0) is load-bearing for cross-VM propagation (pendingDeliveries is already 0 between a delivery resolve() and the VM body registering its next subscriber). Replacing it with queueMicrotask empirically breaks hook/sleep Promise.race ordering (CorruptedEventLogError); a comment now records this. Re-validated after a premature revert: the windows-unit flake that prompted the revert reproduces on unmodified main at the same rate (local 8-way harness: opt 4/80 vs main 7/80; main historical windows-unit ~13%), so it is a pre-existing flake, not a regression from this change. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
e7ef9d823b |
perf(core): lazy inline step start (save one world round-trip per step) (#2478)
* perf(core): lazy inline step start to save a world round-trip per step The owned-inline runtime path used to write step_created (suspension handler) and then step_started (executeStep) as two separate world round-trips for a step it already owns and is about to run inline. This defers the step_created write: executeStep sends a single step_started carrying the step input, and the world creates the step on the fly (materializing the step entity plus a synthetic step_created event so replay still observes it). Mirrors the existing resilient run_started -> run_created pattern. Exactly-one ownership is preserved by the world's atomic create-claim: the loser of a concurrent lazy step_started gets EntityConflictError, which executeStep maps to `skipped`, so it never runs the body. A lazy step_started is only ever sent for a brand-new step (the suspension handler defers only steps with no prior step_created), so crash recovery still re-runs a `running` step via the normal non-lazy step_started. Worlds updated: world-local, world-postgres (implicit create + synthetic step_created event), world-vercel (routes the input as the v4 frame payload and threads the server's stepCreated flag). @workflow/world adds optional `input` to step_started and a `stepCreated` EventResult signal. Rollout: server-first. The matching workflow-server change must deploy before this ships; the Vercel world targets a single Vercel-operated backend (server always >= SDK). For local/postgres the world ships in the same package as the runtime, so there is no version skew. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * fix(core): materialize deferred step before failing unregistered step on lazy inline path The lazy inline step-start optimization defers a step's step_created write, expecting executeStep to materialize the step via a lazy step_started carrying its input. For an UNREGISTERED step, executeStep bails out before sending that step_started and writes step_failed directly — but the step entity was never created, so the world's "step must exist" ordering guard rejects the step_failed and the run wedges (times out). This regressed the StepNotRegisteredError e2e tests uniformly across every framework/world (the ghost step never reached `failed`). Fix: on the lazy path, send the lazy step_started first to materialize the step (entity + synthetic step_created, keeping replay correct), then write step_failed. The lazy step_started's atomic create-claim preserves exactly-one-owner: a concurrent winner makes ours reject with EntityConflictError → skipped, so the failure is never written twice. Adds world-level regression tests (world-local, world-postgres) asserting a lazy step_started followed by step_failed marks the step failed. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
2074f91b86 |
perf(core): skip per-step events.list via inline event-log delta (#2475)
* perf(core): skip per-step events.list via inline event-log delta In the inline sequential loop, the runtime re-read its own just-written step events with an incremental events.list every iteration — pure latency on the Vercel world. Add an opt-in CreateEventParams.sinceCursor so a step-terminal write can return the event-log delta since that cursor (EventResult.events/cursor/hasMore), and have the inline loop consume it in place of the fetch. The delta is computed identically to events.list against the same log, so the consumed prefix is byte-for-byte what a fetch would return. The fast path is gated conservatively to the single-step sequential case with no open hooks/waits (so no out-of-band hook_received/wait_completed can land in the snapshot→replay window), and falls back to the normal fetch on any World that does not return a delta. world-local implements the delta; world-vercel/world-postgres are unchanged and fall back. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * perf(world-vercel): forward sinceCursor over the v4 wire for inline delta Adds `sinceCursor` to the v4 POST frame meta so a step-terminal write can ask the server for the authoritative event-log delta on the response (events/cursor/hasMore), letting the inline loop skip a follow-up events.list. The server-side computation ships in vercel/workflow-server#538; older servers ignore the field and the runtime falls back to events.list (no behavior change). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * test(inline-delta): cover truncated multi-page delta -> hasMore fallback The inline-delta query in world-local intentionally omits a `limit`, so a delta larger than one page is truncated and reports `hasMore: true`. The runtime consume gate only stashes a delta when `!hasMore` and otherwise falls back to the exhaustive `events.list` loop, so a partial page can never be consumed as the complete delta. Make that contract explicit with a comment at the query site, and add tests pinning it: a world-local test proving the delta truncates and surfaces `hasMore: true` byte-identically to `events.list(sinceCursor)`, and an executeStep test proving `hasMore: true` is threaded verbatim. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * docs(world): clarify sinceCursor returns the first delta page, not the full set The CreateEventParams.sinceCursor docstring said the result is "exactly the delta an events.list(...) call would return," which read as the full set. It is the first page of that delta; hasMore signals more. Spell out the single-page-or-fallback contract so other World adapters implement sinceCursor consistently, and note that an in-band burst larger than one page bypasses the fast path (correct, but forgoes the saved round-trip). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
fe333088b7 |
Version Packages (beta) (#2491)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
f193d6e8ef |
Version Packages (beta) (#2451)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
3c79c56af2 |
fix(core): bump payload-compression cutoff to 5.0.0-beta.18 (#2470)
The gzip/zstd FORMAT_VERSION_TABLE entries were gated on 5.0.0-beta.16, but beta.16 and beta.17 were both published (2026-06-15) before the compression PR (#2394) merged (2026-06-16) — neither contains the compression read path. The next published version is beta.18 (pending Version Packages #2451), which is the first that can decode these payloads. With the cutoff at beta.16, getRunCapabilities() reported beta.16/.17 targets as compression-capable, so a cross-deployment start()/resumeHook() (or a resilient-start probe resolving to such a target) would write zstd/gzip payloads the target cannot decode — silent replay corruption, exactly the TODO(release) hazard noted on those lines. Bump both entries (and the doc comments) to beta.18 and extend the capability test to assert beta.16/beta.17 are treated as incapable. |
||
|
|
cb181392b9 |
feat(cli): print run deep links with --url, fix dashboard route (#2467)
Add a `--url` flag to `inspect`/`web` that prints a run's observability dashboard deep link to stdout and exits — no browser, no local server — so scripts and agents can share a link instead of opening a UI. Fix the Vercel dashboard URL to the current `…/workflows/runs/<id>?environment=<env>` route (drop the legacy `/observability` segment) and respect `--env`. Apply the same route fix to the e2e helpers, CI aggregation scripts, and the nextjs-turbopack workbench. Document deep-linking in the workflow skill and observability docs. Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
5f0b845211 |
RFC: compress serialized payload refs — zstd (gzip fallback), specVersion 5 (#2394)
* feat(core,world): gzip-compress serialized payloads behind specVersion 5 Add a composable 'gzip' format prefix layer to the serialization pipeline (compress before encrypt: encr(gzip(devl))), cutting stored payload bytes by ~70-87% on real-world-style workloads. Compression is gated on run specVersion 5 (new SPEC_VERSION_SUPPORTS_COMPRESSION) and on target-deployment capabilities for cross-deployment writes; payloads under 1KB or that don't compress meaningfully are stored unchanged. Reads dispatch on the format prefix so both compressed and uncompressed data are always readable. WORKFLOW_DISABLE_COMPRESSION=1 disables writes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test(core): add CPU/perf compression benchmark + shared workloads Split the compression benchmark into reproducible size and CPU scripts sharing deterministic workloads (lib/workloads.mjs). The CPU benchmark measures serialize/deserialize overhead per payload, total CPU across thousands of events, and compares gzip levels/brotli/deflate. Documents how to run the size, CPU, and end-to-end (bench.bench.ts) benchmarks against local and Vercel in scripts/README.md. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(world-vercel): advertise specVersion 5 to enable compression on Vercel Now that workflow-server declares spec-5 support (vercel/workflow-server#520), bump the Vercel world's advertised specVersion from 4 to 5 so new Vercel runs are stamped spec 5 and become eligible for gzip payload compression. Payloads stay opaque to the server (compression is client-side); spec 5 is a superset of spec 4, so initial run attributes still work. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(core): emit OTel span attributes for compression impact Track gzip payload compression on both the serialize (write) and deserialize (read) paths via span attributes: workflow.serialization.{operation,compressed,uncompressed_bytes, stored_bytes,compression_ratio}. Sizes are measured at the compression boundary (pre-encryption), so they reflect compression's effect rather than the at-rest size. The compression codec stays pure — compress/decompress optionally populate a CompressionStats sink, threaded through CodecOptions to the mode serializers and read by the dehydrate/hydrate wrappers, which set attributes on the active span. Telemetry failures are swallowed so they can never break the serialize/deserialize data path. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * feat(core,web-shared): prefer zstd compression codec (gzip fallback) Switch the payload compression codec to zstd, which benchmarks 3–7× faster than gzip at an equal-or-better ratio on representative workloads (compression runs at every step boundary, so the write CPU is a per-step tax). zstd uses node:zlib (>= 22.15); gzip via the portable CompressionStream remains the fallback when zstd is unavailable, and WORKFLOW_COMPRESSION_CODEC=gzip forces it. Reads dispatch on the format prefix, so 'zstd' and 'gzip' payloads are both always decodable. zstd is Node-only (Web CompressionStream has no zstd), so the browser o11y read path registers a WASM-backed decoder (@tootallnate/zstd-wasm) via a new registerZstdDecoder hook; node:zlib handles Node-side reads (runtime replay, CLI, server o11y). A new workflow.serialization.codec span attribute reports which codec applied. gzip and zstd read support co-ship, so the existing specVersion-5 capability gate is unchanged. Verified end-to-end: spec-5 runs store zstd-prefixed payloads on disk and replay/complete correctly; the WASM decoder round-trips node:zlib zstd output. Benchmarks updated to compare zstd vs gzip. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
4b7a7203bf |
fix(core): make deploymentId 'latest' a no-op in non-Vercel worlds (#2397)
* fix(core): make deploymentId 'latest' a no-op in non-Vercel worlds
Previously, start({ deploymentId: 'latest' }) threw a WorkflowRuntimeError
in any World that doesn't implement resolveLatestDeploymentId() (local dev,
Postgres). That meant a workflow which opts into 'latest' on Vercel would
fail outright in local development.
Resolving 'latest' only means something in worlds with atomic, immutable
deployments. In other worlds there is nothing to resolve between, so instead
of throwing we now log a warning and fall back to the current deployment,
making 'latest' an effective no-op there.
- start.ts: warn + fall back to currentDeploymentId instead of throwing
- start.test.ts: replace the "should throw" test with a warn + fallback test
- e2e.test.ts: assert 'latest' completes (no-op) on non-Vercel worlds
- docs: note the no-op behavior in v4 + v5 start.mdx
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
* fix(core): warn once for deploymentId 'latest' no-op; harden test cleanup
Address PR review:
- Gate the 'latest'-has-no-effect warning behind a once-per-process guard
(mirrors the warnOnce pattern in constants.ts) so a workflow that hardcodes
'latest' for Vercel doesn't flood local/Postgres dev logs on every run.
Exposes _resetLatestNoOpWarnForTests() (@internal) for unit tests.
- start.test.ts: reset the guard in beforeEach and restore spies in afterEach
via vi.restoreAllMocks() so a throwing assertion can't leak the
runtimeLogger.warn spy into later tests; drop the manual mockRestore().
- Add a test asserting the warning fires exactly once across repeated
'latest' starts while every run still falls back to the current deployment.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
|
||
|
|
d4dd6f9c01 | Fix lazy Next workflow HMR (#2438) | ||
|
|
df402c416b |
Version Packages (beta) (#2428)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
926a5e7c6a |
otel: explicit traceparent injection + linked-trace mode for bounded per-invocation traces (#2363)
* otel: explicit traceparent injection + linked-trace mode for bounded per-invocation traces
- Add WORKFLOW_TRACE_MODE ('linked' default, 'continuous' legacy) to the
workflow and step queue handlers. In linked mode, WORKFLOW_V2/STEP spans
start a new trace root with span links to the incoming delivery context
and the run-origin context, and re-enqueued messages forward the
ORIGINAL run-origin trace carrier unchanged.
- world-vercel now explicitly injects W3C traceparent/tracestate/baggage
headers on outgoing workflow-server HTTP requests from inside the
client span (no-op without an OTEL SDK registered).
- New workflow.trace.mode span attribute; unit tests for both modes and
for header injection.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* changeset: call out behavioral telemetry changes of the linked default
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* docs: add v5 observability tracing page
Documents OTEL spans/attributes, linked trace mode and WORKFLOW_TRACE_MODE,
span links, context propagation, and the v4 behavior-change callout.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* otel: human-friendly span names for workflow and step spans
WORKFLOW_V2/STEP prefixes with full machine names (workflow//./src/...//fn)
become workflow.execute / step.execute / workflow.start with the short
function name. New workflowDisplayName/stepDisplayName helpers in
@workflow/utils handle both raw and queue-sanitized name forms; full names
remain in the workflow.name/step.name attributes.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* changeset: merge span-name and linked-trace notes into one changeset
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* docs: update trace-shape prose to renamed span names
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* docs: replace ascii trace diagram with mermaid
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* address review: empty carriers, shared trace helpers, mode warning, name edge cases, consumer span kind
- Treat an empty ({}) trace carrier as absent everywhere the trace-mode
logic branches, so linked mode falls back to a fresh origin instead of
forwarding a useless {} forever; workflow.trace.propagated now reports
whether a usable carrier arrived.
- Extract the duplicated linked-mode logic into shared telemetry helpers
getNextTraceCarrier() and buildInvocationSpanLinks(), used by both the
workflow and step queue handlers; resume-hook now uses
linkToTraceCarrier (gaining the isSpanContextValid guard).
- Warn once per distinct unrecognized WORKFLOW_TRACE_MODE value instead
of silently selecting linked.
- shortNameFromSanitized: map default/__default to the module short name
(mirroring parseName) and document the `$`-sanitization limitation.
- Queue-delivered workflow.execute spans now use the CONSUMER span kind,
matching queue-delivered step.execute spans; docs span table and
changeset updated accordingly.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
|
||
|
|
5711c1e9d6 | Version Packages (beta) (#2390) | ||
|
|
b3cc513220 | [ci] Increase dev.test.ts cleanup hook timeout (#2416) | ||
|
|
4763a760bf |
test: e2e coverage for run-idempotency conflict-handling strategies (#2387)
* test: e2e coverage for run-idempotency conflict-handling strategies Covers the patterns documented in foundations/idempotency: - claim-only hook mutex: token claimed and held with no payload data, duplicate identifies the owner, token released after completion - adopt the owner's result via conflict.returnValue - signal the owner: duplicate forwards its payload via resumeHook - supersede: duplicate cancels the owner and reclaims the token - route-side resume-or-start retry pattern reaching the started run Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test: fix adopt-owner-result race — gate owner completion on observed conflict On slow runtimes the duplicate's first invocation could land after the owner completed and released the token, making the duplicate a fresh owner that waits forever for a payload (90s timeout across CI matrices). Poll the duplicate's event log for hook_conflict before resuming the owner, and widen the test timeout for the added gate budget. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * review: assert superseded owner's returnValue rejection; empty changeset - Await run1.returnValue and assert WorkflowRunCancelledError so the cancellation is verified end-to-end and no rejection leaks from the supersede test. - Test-only PR: use an empty changeset. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * ci: retrigger preview deployments (turbopack deployment for |
||
|
|
628795aa87 |
Add allowReservedAttributes option to start() (#2385)
* Add allowReservedAttributes option to start() experimental_setAttributes already exposes allowReservedAttributes for framework-level callers that own a $-prefixed sub-namespace, and the run_created / run_started event schemas plus the local and Postgres worlds already accept and validate the flag. start() was the one gap: it always validated initial attributes with the reserved prefix disallowed and had no way to opt out, so framework code could not seed reserved attributes at run creation. Thread the option through start(): - StartOptions.allowReservedAttributes, passed to client-side validation and forwarded on the run_created eventData - carried in the queue runInput (new RunInputSchema field) and forwarded to run_started so the resilient/lazy run creation path validates identically Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Add e2e coverage for reserved initial attributes via allowReservedAttributes Verified locally against the nextjs-turbopack dev server: the reserved key passes client and server validation, lands on the run at creation, and survives the workflow's own attr_set writes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
58ddc62d02 |
Version Packages (beta) (#2364)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
303b6da28a |
[core] Add wire-level framing for byte streams (#1853)
Co-authored-by: Peter Wielander <peter.wielander@vercel.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
01c8c0878a |
Replace hook.hasConflict with hook.getConflict() returning the conflicting Run (#2373)
* feat: replace hook.hasConflict with hook.getConflict (Promise<Run | null>)
hasConflict's boolean didn't expose WHICH run owns the token, so the
duplicate run couldn't act on the conflict. getConflict resolves with
null once registration commits, or with a Run handle for the conflicting
run — letting the workflow return/log the owner's runId, inspect its
status, await its result, or cancel it and continue, all in code.
The workflow-mode create-hook module exposes the bundle's compiled Run
class (durable step-proxy methods) on a well-known symbol so the host-
side hook consumer can construct the conflicting run inside the VM.
Contexts without the class (plain unit tests) fall back to a { runId }
object, which is also the documented v4 shape (no native Run
serialization in v4).
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* fix: never resolve getConflict with a non-Run fallback shape
getConflict's contract is Promise<Run | null>. In the degenerate cases
where a real Run cannot be constructed — a hook_conflict event persisted
by an old world without conflictingRunId, or a context that never loaded
the workflow-mode create-hook module — reject with HookConflictError
instead of resolving with a { runId }-shaped impostor.
Test harnesses now register the Run class on the (VM) globalThis like
real bundles do.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* refactor: make getConflict a method — hook.getConflict()
A property getter that triggers registration/suspension reads as passive
state; a method makes the side effect explicit. Update implementation,
types, tests, e2e workflows, docs, and changeset.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* review: guard Run class registration, fix anchors, clarify changeset
- Only register WORKFLOW_RUN_CLASS when the workflow runtime is present
(WORKFLOW_CREATE_HOOK installed on globalThis), so host imports of the
workflow-mode module neither mutate the host global nor expose the
non-step-proxy host Run.
- Drop #run-idempotency link fragments — that section lands in the
stacked docs PR (#2011), which restores the anchored links.
- Note in docs that getConflict() rejects with HookConflictError for
legacy hook_conflict events lacking the owner's run ID.
- Changeset now calls out the hasConflict -> getConflict() replacement.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* refactor: resolve the conflicting Run through the serialization class registry
Replace the bespoke WORKFLOW_RUN_CLASS global with the registry the
serialization pipeline already uses to revive Run instances:
- The SWC plugin already auto-registers the workflow bundle's compiled
Run in globalThis[workflow-class-registry], but under a path-derived
classId the host cannot know statically. The workflow-mode create-hook
module now aliases it under a stable id (class//workflow//Run) via a
new aliasSerializationClass() helper (a plain registry entry —
registerSerializationClass cannot be reused since the plugin's IIFE
already defined the non-configurable classId property).
- createConflictingRun() looks the class up with
getSerializationClass(RUN_CLASS_ID, ctx.globalThis) and constructs
through its WORKFLOW_DESERIALIZE hook, exactly as the Instance reviver
would for a serialized Run crossing from a step into the workflow.
- Because the registry is keyed per-global, no environment guard is
needed: a stray host-side import registers the host Run on the host
registry, which is the correct class for that context. The
WORKFLOW_CREATE_HOOK guard, the ??=, and the WORKFLOW_RUN_CLASS symbol
are all deleted.
Verified: 1156 core unit tests; compiled workbench bundle contains the
stable alias alongside the plugin's path-derived registration with zero
WORKFLOW_RUN_CLASS references; all 5 hookGetConflict e2e tests pass
against a local nextjs-turbopack dev server, including conflict
resolution reading conflict.status through a durable step.
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
Co-authored-by: Nathan Rajlich <n@n8.io>
|
||
|
|
e163422551 |
Add hook.hasConflict for early hook conflict detection (#2015)
* feat: add hook ready promise * test: cover hook ready continuation scheduling * feat: replace hook.ready with hook.hasConflict (Promise<boolean>) - hook.hasConflict resolves true when the token is owned by another active hook, false once registration is committed — no throw, so workflows can branch on conflicts early. Awaiting it suspends the workflow to commit the hook registration (createHook alone does not). - Chain the already-created fast-path through promiseQueue so resolution order matches event-log order (review feedback). - Skip inline step execution when a suspension has an awaited hook creation so the hasConflict continuation can advance independently of step execution (review feedback). - Update unit tests, e2e tests, workbench workflows, and v4/v5 docs. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * docs: fix inconsistent hasConflict bullet in create-webhook reference State both resolution values explicitly (true = token already owned, false = registered) instead of a parenthetical that only described the false case. * docs: require docs preview links in PR descriptions for docs changes Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * docs: restore SWC Plugin heading in AGENTS.md Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: Nathan Rajlich <n@n8.io> |
||
|
|
b3279f8b17 |
[core] V2: unify wait+step queue dispatch in suspension processing (#1925)
* [core] V2: pre-schedule the wait timer before inline-executing a step Fix `Promise.race(step, sleep)` semantics in V2 mixed suspensions without losing inline step execution. Inline `await executeStep(...)` blocks the V2 handler for the full step duration, but `wait_completed` events are only created on the *next* loop iteration's "complete elapsed waits" pass. So if the sleep is shorter than the step, replay always picked the step because the wait_completed event hadn't been written yet — `sleepWinsRaceWorkflow` returned `'step'` instead of `'sleep'`. Fix: when a suspension contains both an owned inline step and at least one pending wait, queue a delayed self-message with `delaySeconds = suspensionResult.timeoutSeconds` *before* starting inline execution. The queued continuation fires in a separate function invocation while the step is still running. That parallel invocation's "complete elapsed waits" pass writes wait_completed, replay observes the elapsed wait, and `Promise.race` resolves with the sleep correctly. The original (still-running) inline invocation finishes its step, sees `run_completed` on the next loop iteration, and exits. This preserves inline-step execution speed for the step-wins case: the step finishes inline and the workflow returns directly. The eagerly-queued wait continuation fires after the step has won and just no-ops on the terminal run. Test plan: - New e2e tests `sleepWinsRaceWorkflow` and `stepWinsRaceWorkflow` exercising `Promise.race` between a step function and `sleep()`, in both directions. - Verified locally against `nextjs-turbopack` workbench: both pass. Event log confirms `wait_completed` is created at t≈1s after `wait_created` (1s sleep) instead of at t≈11s after the inline step finishes. Eager-processing changelog updated with a "Mixed Suspensions" section describing the pre-scheduled wait approach and its rationale. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * [world-local] Honor delaySeconds before message delivery The local queue's `queue()` enqueue path ignored the `delaySeconds` option entirely — every message was delivered immediately, regardless of the requested delay. VQS-side queues (used by world-vercel and world-postgres) honor delaySeconds at the broker, so this brings world-local in line with production semantics. The runtime needs this to land before the wait-as-continuation unification in the next commit: that change starts queueing wait timers as fresh delayed continuations instead of returning `{ timeoutSeconds }`. Without delaySeconds support, those wait continuations would fire instantly in dev and trigger spurious replays. Sleep happens outside the queue's worker semaphore so a delayed message doesn't tie up a worker slot during its delay window — other immediate messages are free to dispatch in parallel. New tests in queue.test.ts cover: - delaySeconds > 0 → setTimeout called with the right ms value - delaySeconds === 0 → no setTimeout (immediate dispatch) - delaySeconds omitted → no setTimeout (immediate dispatch) * [core] V2: unify wait+step queue dispatch in suspension processing Replace the asymmetric "steps go to the queue, waits become a { timeoutSeconds } return value" pattern with a single Promise.all batch that queues every pending operation we are not running inline. Before this change, suspension processing had three branches that all needed to keep the wait/step asymmetry consistent: - pendingSteps.length === 0 returned { timeoutSeconds } - inlineStep + waits eagerly queued a delayed self-message AND set inlineStep to undefined (Option A) AND returned { timeoutSeconds } - inlineStep retry path returned { timeoutSeconds } if there were waits After this change, every suspension goes through one path: for non-inline pendingSteps: queue stepId message if timeoutSeconds defined: queue delayed continuation await Promise.all(dispatches) if !inlineStep: return await executeStep(inlineStep) Behaviorally, this restores inline step execution even when the suspension also has a wait (Option A's carve-out is no longer necessary): the wait timer fires in a separate function invocation on the queue, in parallel with the inline step. If the sleep wins the race, that parallel invocation observes wait_completed via the "complete elapsed waits" pass and finishes the run; if the step wins, the wait continuation fires later and no-ops on the terminal run via the existing terminal-event check. Other cleanups: - The inline-step retry path no longer needs to forward suspensionResult.timeoutSeconds — the wait timer was already enqueued as part of the unified dispatch above. - A dead post-step `if (timeoutSeconds && pendingSteps.length === 1)` block (just a comment, no body) is removed; the loop's "complete elapsed waits" pass handles the same case correctly. - Step queueing now uses a shared `traceCarrier` rather than re-serializing per step. Retry/throttle and hook-conflict paths still return { timeoutSeconds } since their semantics are "redeliver THIS message after a delay" rather than "schedule a fresh wait timer." Those can be unified in a follow-up. Test plan: - New e2e tests `sleepWinsRaceWorkflow` and `stepWinsRaceWorkflow` pass against the `nextjs-turbopack` workbench. - Event log inspection confirms wait_completed fires at t≈1s (after wait_created at t≈0s) for the sleep-wins case, and that the inline step runs only once (no duplicate step_started events that the earlier eager-queue approach produced in dev). - All 842 @workflow/core unit tests pass. - All 346 @workflow/world-local unit tests pass (with the delaySeconds support added in the previous commit). Requires the world-local delaySeconds fix in the prior commit; without it, wait continuations would fire instantly in dev and the parallel replay would re-enter handleSuspension before the wait elapsed (recoverable via existing redelivery, but inefficient). * [docs] V2 unified suspension dispatch + changeset Update the "Mixed Suspensions" section in eager-processing.mdx to describe the unified parallel-dispatch model: - All non-inline pendingSteps are queued with stepId - The wait timer (if any) is queued as a delayed continuation - All dispatched in one Promise.all batch - One owned step is then inline-executed (if any) The doc previously described Option A (the carve-out where waits forced all steps to be queued); the unified model removes that carve-out and explains why the wait continuation works in parallel with the inline step. Also notes the dependency on world-local's new delaySeconds support (landed earlier in the same PR series). Changeset bumps both @workflow/core and @workflow/world-local since both packages have user-observable behavior changes. * [core] Dedupe wait continuations on the wait's correlationId While a wait is pending, every replay pass over the run re-observes it (once per step completion in Promise.all([steps..., sleep()]), etc.) and would enqueue another delayed continuation — each a spurious replay when the wait elapses, and each a fresh message that resets the delivery-attempt runaway guard. Key the continuation on the wait's correlationId so the worlds' idempotency dedupe collapses them. Near-elapsed waits (<= 2s) are enqueued without the key: a continuation delivered marginally early (clock skew; the ceil() on the delay can leave a ~0 margin) re-observes its wait as pending and must be able to enqueue a fresh short-delay retry. VQS idempotency records persist until message-retention TTL — reusing the key there would drop the retry and stall the run permanently. Also adapts wait-completion-replay tests (from #2038) to the unified dispatch model: the hook-branch step now executes inline (registered in the test world, which now returns a step entity from step_started), so each scenario performs one extra loop-iteration event fetch. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * [core] Always key wait continuations; bucket the key for near-elapsed waits CI caught sleepWinsRaceWorkflow failing across the world-postgres lanes: world-postgres serializes KEY-LESS workflow messages per run (inflightWorkflowRuns), so a key-less wait continuation parks behind the flow message that is inline-executing the racing step — wait_completed lands after step_completed and the race resolves to the step. Keyed messages take the concurrent dedupe path, so the continuation must always carry an idempotency key. The near-elapsed exception (<= 2s) now uses a second-bucketed suffix instead of omitting the key: an early-delivered continuation re-observes its wait as pending and re-enqueues with >= 1s delay, which guarantees a later bucket — a fresh key that dedupe windows cannot drop — while same-instant duplicates still collapse. Verified against a local world-postgres setup (express workbench, Graphile worker): sleepWins/stepWins pass 3/3 with wait_completed at t+1s; the event log confirms the continuation fires in parallel with the in-flight inline step. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * [core] Clamp wait-continuation delays; chain long waits with hop-keyed dedupe Addresses PR review: the unified dispatch passed delaySeconds to the queue unclamped while keying the continuation on the bare wait correlationId. On world-vercel (23h max delay, 24h VQS message retention) a sleep() longer than the max either failed the dispatch or was delivered early with its re-enqueue silently dropped by the still-live idempotency record - stalling the run permanently. - New runtime/wait-continuation.ts owns delay + idempotency-key selection: delays clamp to 23h and longer waits chain across hops, with the hop index suffixed to the key so re-observations within a hop window dedupe while each hop delivery gets a fresh key. Near- elapsed threshold and max delay are named constants; full rationale moved out of the runtime.ts comment block. Unit tests pin the key selection including chain advancement. - SuspensionHandlerResult: timeoutSeconds/timeoutWaitCorrelationId collapsed into waitTimeout?: { seconds, correlationId } so the pairing can't drift (review nit). - runtime.test.ts ack-ordering harness adapted to the unified model: step_created now answers EntityConflictError so the handler observes the step without owning it and must queue it (the carve-out the tests relied on - "pending wait disables inline execution" - is exactly what this branch removes). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * [world-local] Abort pending queue sleeps on close() Addresses PR review: a pending delayed message kept the dev process's event loop alive for its full delay, and close() only closed the HTTP agent - a sleep that fired afterwards attempted delivery against the closed agent and logged a spurious "[local world] Queue operation failed" error during test/CLI shutdown. One AbortController owned by the queue now cancels the delaySeconds sleep, the timeoutSeconds re-delivery sleep, and the retry backoff on close(); the resulting AbortError is already swallowed by the existing isAbortError check. close() is idempotent since shutdown paths may invoke it twice. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * [docs] Wait-continuation clamping + hop chaining; changeset eager-processing.mdx pseudocode now shows the continuation's idempotency key and clamped delay (PR review nit); the dedupe prose covers the two key variations (hop suffix for chained long waits, second bucket for near-elapsed waits). Changeset mentions long-sleep chaining and world-local's abort-on-close. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com> Co-authored-by: Pranay Prakash <pranay.gp@gmail.com> |
||
|
|
f2a7bdeb0a |
fix(world-local,world-postgres): make duplicate hook_created idempotent (#2295)
* fix(world-local): make duplicate hook_created idempotent Duplicate processing of the same hook_created — same runId, hookId, and token, e.g. cross-process replay or queue redelivery — was being recorded as a hook_conflict in the event log, which then replayed as a self- conflict HookConflictError. The fix mirrors the existing step_created duplicate-correlation path: when the exclusive token claim fails and the existing claim has the same (runId, hookId), throw EntityConflictError so the runtime's existing concurrent-replay catch path swallows it. Different runId or hookId reusing the same token still produces a real hook_conflict. The persisted token claim already carried hookId; only the read schema was dropping it. The schema now preserves hookId (marked optional for backward compatibility with older claim files). Fixes #2283 * fix(world-postgres): make duplicate hook_created idempotent world-postgres has the same gap as world-local was just fixed for: the duplicate-token check in events.create unconditionally writes a hook_conflict event when an existing hook with the same token is found, even when the existing hook has the same (runId, hookId) as the incoming event. The unique partial index on workflow_events does not catch this because the duplicate path inserts hook_conflict, not hook_created. Mirror the world-local fix: when the existing hook's (runId, hookId) matches the incoming event, throw EntityConflictError so the runtime's existing concurrent-replay catch path swallows it. Different runId or hookId reusing the same token still produces a real hook_conflict. Refs #2283 * test(e2e): add regression test for hook_conflict from same-tick replay race Regression test for #1665 / #2283. A parent workflow awaits 6 child workflows with Promise.all; each child does a tiny step and creates one webhook. Awaited children flatten into the parent run, so all webhook creations land on the same workflow body. When their step resolutions align in the same tick the workflow body is re-walked and each pass submits hook_created with the same deterministic (correlationId, token). Before the world-side idempotency fix, the world wrote hook_conflict events for the duplicates and the workflow failed with HookConflictError. With the fix, duplicates throw EntityConflictError (swallowed by the suspension handler), no hook_conflict events appear in the log, and the webhooks resolve normally. Verified locally against world-local: the test fails reliably (3/3) on the unfixed code and passes reliably (5/5) on the fixed code. * test(e2e): rewrite parallelStepsThenWebhookWorkflow to match the actual #1665 repro The earlier version invoked another 'use workflow' function directly from inside the parent workflow, which is not a valid child-workflow invocation (child workflows must be spawned via start()) and didn't mirror the bug shape on #1665 anyway. Rewrite the workflow as a single 'use workflow' function that exactly mirrors Paolo's minimal repro: await Promise.all([stepA(), stepB()]); using webhook = createWebhook(); await webhook; The for-loop runs N independent iterations of that sequence in series, each disposing its webhook via 'using' before the next, to give the timing-sensitive race multiple chances to fire. The race is hard to force deterministically on fast local dev — but the same (runId, hookId) idempotency invariant is covered deterministically by the new unit tests in world-local and world-postgres. This e2e test serves as a higher-level regression net: its assertions (no hook_conflict event in the log, no HookConflictError-failed run) are correct whether the race fires or not, and will catch any future regression on a run that does hit it. * fix(world-local,world-postgres): recover crash-orphaned hook claims/rows instead of suppressing the retry Addresses review feedback on PR #2295. The original idempotency fix made duplicate same-(runId, hookId) hook_created submissions throw EntityConflictError so the suspension handler's concurrent-replay catch path swallows them. But the claim file (world-local) and hook row (world-postgres) are written before the durable hook_created event, and the writes are not atomic. A process / DB interruption between the claim/hook write and the event write leaves an orphaned claim/hook row; the retry then matched the same (runId, hookId), threw EntityConflictError, got swallowed, and the run was permanently left with no hook_created event in the log. world-local: - Add a per-(runId, hookId) in-process mutex (withHookLock) mirroring the existing withStepLock, so two same-tick concurrent calls serialize on the entity write and the dedup branch never observes an in-flight winner mid-write. - In the dedup branch, when the existing claim is for the same (runId, hookId) we are trying to create, check whether the durable hook entity actually exists on disk: - exists → real duplicate: throw EntityConflictError as before. - missing → orphaned claim from a prior crash: fall through and complete the partial write (write the hook entity with overwrite, then emit hook_created via the outer code path). world-postgres: - In the dedup branch, when the existing hook row matches the incoming (runId, hookId), check whether a hook_created event for this (runId, correlationId) already exists in the event log: - exists → real duplicate: throw EntityConflictError as before. - missing → orphaned hook row from a prior crash between hook INSERT and events INSERT: skip the hook insert (the row is already there) and let the outer code path emit hook_created, completing the partial write. Tests: - world-local: pre-seed an orphaned token claim with no matching hook entity, retry hook_created, assert hook entity and hook_created event both land (no hook_conflict, no EntityConflictError). - world-postgres: pre-seed an orphaned hook row with no matching hook_created event, retry, assert hook_created event lands (no hook_conflict, no EntityConflictError). Both tests fail on the prior implementation (EntityConflictError thrown on retry, exact symptom from the review). * fix(world-local): probe the event log (not the hook entity) to detect duplicate hook_created Addresses follow-up review on PR #2295. The previous dedup branch checked whether the durable hook entity existed on disk. But the hook entity is written before the `hook_created` event, and the two writes are not atomic, so a crash between them leaves both the claim file and the hook entity on disk with no event in the log. The dedup branch then matched on `(runId, hookId)`, found the hook entity, threw EntityConflictError, and the suspension handler swallowed the retry — permanently losing `hook_created` from the event log. The fix mirrors what the world-postgres branch already does: probe the run's event log for an existing `hook_created` event for the same `(runId, correlationId)`. The event is the durable record of a successful hook creation; the claim file and hook entity are partial- write artifacts that may exist without the event. - exists → real duplicate: throw EntityConflictError so the runtime's concurrent-replay catch path swallows it. - missing → orphaned partial write (crash at any point before the event landed): re-write the hook entity (with overwrite: true, in case a stale partial copy exists) and let the outer code path emit the hook_created event. Added a new helper findHookCreatedEvent that runs a filtered paginatedFileSystemQuery with limit:1 over the run's events. Regression test "should recover an orphaned hook entity with no matching hook_created event" added — pre-creates a hook, deletes just the hook_created event from disk to simulate a crash between the entity write and the event write, asserts the retry emits a fresh hook_created event (no hook_conflict, no swallowed EntityConflictError). I verified this test fails on the prior fix (throws `EntityConflictError: Hook "hook_orphan_entity_1" already created`, exactly as pranaygp reported) and passes on this commit. The previous test ("should recover an orphaned hook token claim with no matching hook entity") continues to pass — the event-log probe is a strict superset of the entity probe, since a missing entity always also implies a missing event. * fix(world-local): converge same-hook creation across workers via canonical eventId Addresses follow-up review on PR #2295. The previous fix made the dedup branch probe the event log to decide real-duplicate vs orphan-recovery, but the probe and the recovery write are not a single atomic operation. Two workers sharing a data directory (or two retries that lose `writeExclusive(constraintPath)` back to back) could both pass the probe (each observing no hook_created event yet), both fall through to the recovery write, and both append a hook_created event with a different eventId — producing two events in the log for the same (runId, hookId). The in-process `withHookLock` mutex does not help here because it is process-local and tag-specific. The fix persists `eventId` in the durable token claim file (written by the original `writeExclusive(constraintPath)`). On a same-(runId, hookId) dedup match, retries adopt that canonical eventId and rebuild the event with a deterministic createdAt derived from the eventId (a ULID). The outer event write switches from `writeJSON` (check-then-write, TOCTOU) to `writeExclusive` (O_CREAT|O_EXCL via temp-file + hard-link, atomic across processes). Either worker may win the publish; the other throws EntityConflictError which the runtime's existing concurrent-replay catch path swallows. Net result: exactly one hook_created event per logical creation. Backward compatibility: a claim file written before this commit lacks `eventId`. Retries that read such a claim fall back to the event-log probe + fresh-eventId recovery — the legacy behavior that does not converge across workers but cannot regress for freshly- written claims after upgrade. world-postgres already converges across workers via the partial unique index on workflow_events_entity_creation_unique (runId+correlationId+eventType for hook/step/wait_created): the loser's INSERT raises 23505 which is already translated to EntityConflictError. Regression tests: - world-local: `converges same-hook creation across workers to one event` uses two tagged storage instances sharing one data directory and fires 25 paired Promise.allSettled hook_created calls. Expected 25 hook_created events total; before this fix yielded 50. - world-postgres: `converges same-hook creation across concurrent calls to one event` exercises the same shape against the real Postgres unique index. Already converges; the test is a guard against future regressions to the catch path. Verified the world-local test fails on c7b23e1b5 with exactly the shape pranaygp reported (50 events for 25 logical creations) and passes on this commit. The earlier orphaned-claim and orphaned- entity recovery tests also continue to pass. * fix(world-local): converge legacy hook claims via recovery-marker sidecar; replace tag-proxy test with real subprocess workers Addresses follow-up review on PR #2295. Two distinct issues, both flagged by pranaygp as P1: 1. The fallback path for token claims written by versions before eventId was persisted inline (legacy claims after upgrade) still permitted the same cross-process corruption the inline fast path was fixed to prevent. Two processes both reading a legacy claim each generated their own eventId, landed their writeExclusive(eventPath) calls at different paths, and appended two hook_created events for the same (runId, hookId). Existing persisted claims after a real upgrade are exactly the state the crash-recovery branch needs to repair, so leaving the legacy path non-convergent is silent corruption, not backward compatibility. 2. The committed cross-worker convergence test used two tagged storage instances sharing one directory as a proxy for separate processes. But tags change the destination filename (events/wrun_X-evnt_Y.worker-a.json vs ...worker-b.json), so two tagged workers can each writeExclusive their own event at different paths and both fulfill. The Map-by-eventId deduplication in the assertion then masked the duplicate publication, so the test passed for the wrong reason. Implementation: - New HookRecoveryMarkerSchema (`{ eventId, hookId, runId }`) and HookRecoveryMarkerPath helper. The marker is a sidecar at hooks/tokens/<hash>.recovery.json, written via writeExclusive so the first cross-process retry pins its candidate eventId as canonical; subsequent retries read the marker and adopt that eventId. Together with the existing writeExclusive(eventPath) in the outer publish, this gives the legacy-fallback path the same single-event convergence guarantee as the inline-eventId fast path. - pinCanonicalEventIdForLegacyClaim() encapsulates the marker write-or-read. A stale marker for a different (runId, hookId) (token-reuse with leaked state) is overwritten best-effort — the common cross-worker race for the same hook still converges; only the narrow stale-token-reuse case loses convergence. - hook_disposed now also deletes the recovery marker when it deletes the token constraint file, preventing a future legacy recovery for a recycled token from latching onto a stale eventId. - The dedup branch unified: existingClaim.eventId for new claims, pinCanonicalEventIdForLegacyClaim() for legacy ones. Removed the now-redundant findHookCreatedEvent helper — the writeExclusive(eventPath) in the outer publish is the authoritative duplicate-vs-orphan detector. Tests: - New test fixture test-fixtures/hook-race-worker.ts (TypeScript, run via child_process.fork with tsx as execPath — tsx is a transitive dev dep via vitest). Each subprocess gets its own createStorage(testDir) so the in-process hookLocks Map cannot serialize across workers. - Replaced the tag-proxy test with "converges same-hook creation across separate OS processes to one event". Spawns workerCount subprocesses, releases them from a barrier into the same hook_created, asserts exactly one fulfilled + (N-1) rejected with EntityConflictError, and asserts directly on the raw events.list() result (no Map dedup) that the number of hook_created entries equals the number of logical creations. - Added "converges same-hook creation across processes when only a legacy token claim exists". Same shape, but pre-seeds the legacy claim format (`{ token, hookId, runId }` with no eventId) before each race. Verified to FAIL on 7ce66551b (both subprocesses fulfill, no convergence) and pass on this commit. - Also verified the new-eventId subprocess test FAILS when the event write is reverted to writeJSON (TOCTOU), confirming it exercises the writeExclusive-based cross-process arbitration. Both prior orphaned-claim / orphaned-entity recovery tests also continue to pass. * fix(world-local): per-lifetime recovery markers, restore event-log probe, fix CI tsx resolution Addresses three P1 review comments on PR #2295. 1. Stale recovery marker leaking across token-reuse lifetimes (pranaygp): The previous marker path used `hashToken(token)` so a stale marker for run A could leak into run B's recovery when the same token was reused after run A terminated through normal lifecycle. `deleteAllHooksForRun()` and tagged `world.clear()` deleted the token constraint and hook entity but NOT the marker sidecar, so the next legacy claim on the same token entered the stale-marker overwrite branch and the workers overwrote it non-atomically, yielding divergent publication. Fix: - Marker path now hashes `(token, runId, hookId)` together (`hookRecoveryMarkerPath` in storage/helpers.ts). Different lifetimes can never share a marker, so the stale-marker overwrite branch is removed entirely. - `hookRecoveryMarkerPath` is moved to helpers.ts and shared across events-storage.ts, hooks-storage.ts, and index.ts. - `deleteAllHooksForRun()` and tagged `world.clear()` now also delete the recovery marker for each hook (disk hygiene; per- lifetime identity makes leaks no longer corrupting). - `hook_disposed` now uses the new per-lifetime marker path too. 2. Duplicate `hook_created` event when a legacy claim's event was already published (VADE bot, also implied by pranaygp's analysis): Removing the event-log probe from the legacy fallback let a post- upgrade retry pin a new canonical eventId via the marker and publish a duplicate event at that path, even when the original pre-upgrade writer had already successfully published the event with its own (different) eventId. Fix: - Restore `findExistingHookCreatedEventId()` (renamed and made to return the eventId for clearer semantics). - Legacy fallback now probes the event log BEFORE pinning the marker; if a matching `hook_created` event already exists, throw `EntityConflictError` so the runtime's concurrent-replay catch path swallows the retry. - Inline-`eventId` fast path does NOT need the probe — the claim itself is the durable convergence key. 3. CI failure: tsx not resolvable under pnpm isolated linking (pranaygp; confirmed by ubuntu/windows unit test 60s timeouts): The previous test hard-coded `node_modules/.bin/tsx` assuming tsx would be hoisted there. But tsx was only a transitive peer dep via vitest, and pnpm's isolated linking does NOT link transitive peer deps into the workspace bin after a fresh install — so neither root nor package-local `.bin/tsx` existed in CI, the subprocess fork never started, and the barrier hung until vitest killed the test. Fix: - Add `tsx` as a direct `devDependency` of `@workflow/world- local` (pinned to 4.20.6 to match the existing transitive resolution). - Resolve via `import.meta.resolve('tsx/package.json')` and read the `bin` field dynamically, so we adapt to wherever pnpm links tsx for this package — not a hard-coded layout. - Lazy-init the resolver (no module-load IIFE) so an absent tsx fails only the convergence tests, not all 376 tests in the file. - Surface a clear error message if resolution fails, calling out the cause (transitive vs direct deps) for future readers. Also: harden the barrier helper so `error` events and pre-ready exits resolve BOTH `readyPromises` and `donePromises`, then `SIGKILL` siblings. Previously a broken child only resolved `donePromises`, leaving `Promise.all(readyPromises)` pending until the per-test timeout (60s in CI). Regression tests added: - `legacy claim whose hook_created event was already published does not append a duplicate event` — pre-seeds a legacy claim AND a pre-existing `hook_created` event with a different eventId, asserts the retry throws EntityConflictError and the log still has exactly the original event. - `converges legacy claim recovery across run lifetimes after token reuse` — runs pranaygp's full lifecycle path: race subprocess workers on run A's legacy claim, terminate run A via `run_completed` (triggers `deleteAllHooksForRun`), reuse the token in a legacy claim for run B, race subprocess workers again, asserts exactly one fulfillment + one `EntityConflictError` per race and exactly one `hook_created` event per run. Both new tests verified to fail on 2c673e436 (after rebuilding): the published-event test throws via duplicate publish instead of EntityConflictError, the token-reuse test sees both run B workers fulfill (2 events instead of 1). The existing orphaned-claim and orphaned-entity recovery tests also continue to pass. CI loop confirmed to be repaired locally by spawning subprocesses via the new resolver and intentionally breaking the worker fixture to verify the helper fails fast (~500ms) instead of hanging at the barrier. * fix(world-local): defer hook entity write until event publish commits Addresses karthikscale3's P1 review comment on PR #2295. The dedup-recovery path used to write the hook entity BEFORE the outer event publish proved whether the attempt was repairing a missing event or just colliding with an already-published `hook_created`. For already-committed duplicates, the event write then throws `EntityConflictError`, but the hook entity had already been overwritten with the retry's payload — leaving the durable hook entity and the event log inconsistent (e.g. the entity reflects the retry's metadata while the event still carries the original). karthikscale3 reproduced this on the prior head by creating `hook_created` with metadata `{ v: "a" }`, then retrying the same `(runId, hookId, token)` with metadata `{ v: "b" }` and `isWebhook: false`: the retry threw `EntityConflictError` but `hooks.get()` returned the retry's payload. Fix: defer the hook entity write until AFTER the outer `writeExclusive(eventPath)` commits. The branch now only captures the entity-to-write and its overwrite options; the actual write happens immediately after the event publish in the shared trailing block. A retry that ends in `EntityConflictError` (the event was already published) now leaves the entity untouched. The first-writer happy path and all recovery paths (orphaned- claim, orphaned-entity, cross-worker convergence, legacy claim, token-reuse across lifetimes) are unaffected — they all reach the event publish successfully, then the entity write runs as before. Regression test `does not mutate an already-committed hook entity when a duplicate hook_created retry collides` added to world-local: runs karthikscale3's exact scenario and asserts the persisted entity still carries the original metadata and isWebhook. Verified to fail on the prior commit (persisted metadata = 0xbb instead of 0xaa) and pass on this commit after rebuilding. Parallel guard test `does not mutate an already-committed hook entity when a duplicate hook_created retry collides` added to world-postgres. Postgres already protected this via `onConflictDoNothing()` on the hook INSERT, but the test guards against a future regression that adds an UPDATE/UPSERT to the dedup path. * refactor(world-local): per-instance in-process locks; drop tsx subprocess test plumbing You were right that the tsx subprocess machinery was overkill for a storage-level convergence test. Replaced with a simple two-instance in-process test that exercises the same cross-process semantics without spawning anything. The trick: `stepLocks` and `hookLocks` were module-level Maps shared by all `createEventsStorage` calls in the same process. Move them inside the function so each `createStorage(dir)` call gets its own lock map. Two storage instances sharing one data directory then behave exactly like two separate OS processes: - independent in-process `hookLocks` Maps (no in-process serialization between them), and - a shared filesystem (so the on-disk `writeExclusive` claim / marker / event publish primitives are the only thing arbitrating convergence). This is also a real architectural improvement — the global lock map was always a leaky abstraction that made unit-test simulation of the cross-process path awkward. Changes: - `stepLocks` and `hookLocks` moved from module scope into `createEventsStorage`. `withStepLock` and `withHookLock` wrappers collapsed into direct `withInProcessLock(map, key, fn)` calls at the two call sites that need them. - The three convergence regression tests in `storage.test.ts` now use `const workerA = createStorage(testDir); const workerB = createStorage(testDir);` and race `Promise.allSettled` of `events.create` from both — no subprocess, no IPC, no barrier helper, no `raceHookCreatedAcrossProcesses`. Same assertions (exactly one fulfillment + N-1 `EntityConflictError` per race, raw `events.list()` shows exactly one `hook_created` per logical creation — no Map dedup) so the regression catches are identical. - Removed: `tsx` devDep, `test-fixtures/hook-race-worker.ts`, `HOOK_RACE_WORKER` / `resolveTsxLoaderUrl` / `TSX_BIN` / `raceHookCreatedAcrossProcesses` and the `fork`/`fileURLToPath` imports they pulled in. Verified (after rebuilding world-local): - All 379 tests pass on macOS in ~1s (was ~6.7s with subprocesses). - Convergence tests confirmed to still catch the bugs: temporarily reverted the `eventId = canonicalEventId` adoption → both workers fulfilled (2 events instead of 1). Temporarily reverted the legacy-claim marker pin → same: both workers fulfilled. - No subprocess machinery means no Windows-specific quirks (cli.mjs shebang, .cmd wrappers, .bin hoisting under pnpm isolated linking, etc.) that produced the Windows CI 60s timeouts. - World-postgres still has its own parallel guard test for the karthikscale3 "no-mutate-on-duplicate" regression; that one exercises real DB concurrency and is unaffected by this change. Full repo `pnpm test` (43 packages) and the `parallelStepsThenWebhookWorkflow` e2e test against world-local both green. * fix(world-local): repair event-first hook orphans from the persisted event; skip #1665 e2e on world-postgres - A crash between the hook_created event publish and the deferred hook entity write left the event committed with the entity missing and unrepairable (retries threw EntityConflictError without materializing the entity). Retries now rebuild the entity from the PERSISTED event's payload — never the retry's eventData — via a race-safe writeExclusive, on both the canonical-eventId collision path and the legacy-claim probe path. - Skip parallelStepsThenWebhookWorkflow e2e on world-postgres: the same-tick replay pattern surfaces a separate pre-existing step_started ordering bug there (#2331). --------- Co-authored-by: Peter Wielander <peter.wielander@vercel.com> |
||
|
|
564a47c504 |
fix: settle aborted parallel steps before completing abortParallelWorkflow (#2244)
* fix: wait for aborted parallel steps to settle * test: assert aborted results for parallel abort workflow |
||
|
|
ae8d6feeda |
Add native v4 workflow attribute events (#2226)
* Add native workflow attribute events * Fix abbreviated attributes docs sample * Document attribute replay ordering for step races * Address native attribute review feedback * Validate before claiming attr_set dedup lock; clearer start() attribute errors - world-local: claim the attr_set correlation lock only after validation, so a validation failure does not permanently mark the correlationId as written and wedge the run in a re-invoke loop on retry - world-postgres: distinguish a concurrently-deleted run from a cap violation when the guarded attributes update matches no rows - core: reject non-string initial attribute values in start() with a clear error instead of a downstream schema failure Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Add attribute edge-case tests across all layers - core: normalizeAttributeChanges unit tests (non-object inputs, FatalError wrapping, key/value/batch limits, boundary lengths, UTF-8 byte counting) - core: start() rejects reserved keys, oversized keys/values, and over-cap initial attribute batches before any write - world-local + world-postgres: per-run cap enforced against existing attributes (upsert-at-cap allowed, removal frees room), oversized values rejected on attr_set, invalid initial attributes rejected on run_created - e2e: validation DX workflow asserting every invalid write throws a catchable FatalError naming the violated rule and limit, with the run staying healthy; start() rejects invalid initial attributes client-side Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Remove accidentally committed local e2e diagnostics artifact Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Bump world-vercel to spec version 4 for native attributes The deployed workflow-server (vercel/workflow-server#469) materializes native attr_set events and accepts initial run attributes, but world-vercel still advertised spec v3 — so start(..., { attributes }) rejected itself client-side ('requires spec version 4') on every Vercel deployment, failing the new e2e seeding test across the prod matrix. New runs are now stamped v4. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Reject duplicate correlated attr_set before materializing in Postgres A redelivered duplicate — including one carrying different changes for the same correlationId — previously re-applied the run attributes update and only then failed the event insert, leaving the snapshot out of sync with the event log. Pre-check the event log for the correlationId before mutating; the unique index still guards the truly-concurrent race, which is idempotent (deterministic replay carries identical changes). Also apply the suggested docs wording for initial attributes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Apply suggestion from @VaguelySerious Signed-off-by: Peter Wielander <mittgfu@gmail.com> * Fail the run on World-rejected attribute writes; un-nest runtime test Two fixes from review: - runtime.test.ts: the pre-existing test "propagates transient step_created failures..." was accidentally nested inside the new attribute-race test, failing the new test ("Calling the test function inside another test function is not allowed") and preventing the old test from running. Restored it verbatim at describe level. - A workflow-body attr_set the World rejects as invalid (e.g. the cumulative per-run attribute cap, which only the World can check) is deterministic: redelivering the orchestrator message replays the same write into the same rejection, wedging the run in redelivery with no terminal event. handleSuspension now wraps such rejections in FatalError, and workflowEntrypoint fails the run with the validation error instead of rejecting the delivery. Transient storage errors still propagate and retry via redelivery. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Signed-off-by: Peter Wielander <mittgfu@gmail.com> Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> Co-authored-by: Peter Wielander <peter.wielander@vercel.com> |
||
|
|
05e46fa3f6 |
Version Packages (beta) (#2326)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
884bd8db76 | [ci] Fix flaky windows unit tests (#2359) | ||
|
|
4e8a9657c9 | Fix e2e failure reporting under vitest 4 and preserve fetch error causes (#2355) | ||
|
|
a813382216 |
[core] Fix process crash from rejected waitUntil promises (#2336)
Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
bf44d4dd0a |
[core] Remove duplicate waitUntil for suspension handler async operations (#2345)
Co-authored-by: Peter Wielander <mittgfu@gmail.com> |
||
|
|
4670c4b92d |
feat(core): add optional namespace for queue topic prefix (#2305)
* feat(core): add optional namespace for queue prefix * fix(world-postgres): job queue prefix validation * fix: changeset description Co-authored-by: Peter Wielander <mittgfu@gmail.com> Signed-off-by: Will Sather <56037657+willsather@users.noreply.github.com> * fix: add world-postgres to changeset * fix: world-postgres handle namespaced job queue names * fix: resolve namespace via env var in core runtime * fix: world-postgres job queue name task handler * fix(world-postgres): honor namespace on consumer side * Fix namespaced queue routing reliability (#2340) * Fix namespaced queue routing reliability * Inline queue namespace in generated routes * Avoid loading Vercel functions during runtime import --------- Signed-off-by: Will Sather <56037657+willsather@users.noreply.github.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> Co-authored-by: JJ Kasper <jj@jjsweb.site> |
||
|
|
eb976db35b | [core] Forward-port stream reconnect to getReadable level (#2318) | ||
|
|
73e64bba03 |
Version Packages (beta) (#2254)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
bb6ff9ac99 |
Patch vulnerable package dependencies (#2301)
* chore: patch package dependency vulnerabilities Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> * Prefer direct dependency upgrades for security fixes --------- Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> |
||
|
|
aa628b7a8f | fix: bump devalue to 5.8.1 (#2292) | ||
|
|
ccd37e9a59 |
Handle lazy stream key request failures (#2257)
Signed-off-by: Pranay Prakash <pranay.gp@gmail.com> |
||
|
|
0fd0891cc4 |
[core] Preserve event-log order in hook-vs-sleep replay races (#2171) (#2185)
Co-authored-by: Nathan Rajlich <n@n8.io> |
||
|
|
8d75491a07 |
[core] Surface workflowCoreVersion from healthCheck() result (#1854)
Co-authored-by: Peter Wielander <peter.wielander@vercel.com> |
||
|
|
0b8b077345 | Fix Next lazy step route imports (#2263) | ||
|
|
ff66ee9f2b |
Version Packages (beta) (#2216)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
12c35b54eb | fix(core): skip terminal event replay on main (#2215) | ||
|
|
52d63d1b61 |
fix(core): advance workflow clock only for consumed events (#2206) (#2211)
* fix(core): advance workflow clock only for consumed events
* test(core): hammer replay branch clock determinism
* fix(core): address consumed event replay review feedback
* test(core): bound replay regression concurrency
(cherry picked from commit
|
||
|
|
2a3b11bcb4 |
Retry replay divergence before failing event logs (#2212)
(cherry picked from commit
|
||
|
|
275316fac4 |
Version Packages (beta) (#2183)
Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com> |
||
|
|
8f68d3525c | fix(core): resolve forwarded stream keys across deployments (#2191) | ||
|
|
ae3c833acd | [e2e] Improve error labeling in event-log-race-repro CI job (#2190) | ||
|
|
1ee63b870a |
[core] Harden event pagination response parsing (#2180) (#2179)
Co-authored-by: Peter Wielander <mittgfu@gmail.com> |