Conversation
…ent invocation A concurrent replay that already created a step under the same correlation id used to be tolerated with an info log regardless of what it created. When the two replays bound the id to different step calls, the step body ran with the winner's arguments and the loser's branch received its result: silent cross-wiring when both calls were the same step function (a fan-out mapping one step over a list). Read the persisted step on the 409 and compare its name and decrypted, decompressed input with ours on every duplicate-create path (sequential, resilient dispatch, batch of one, batched fan-out, folded lazy-inline pair); a mismatch is a non-deterministic replay and now fails the run as CORRUPTED_EVENT_LOG instead of continuing. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
🦋 Changeset detectedLatest commit: a99361d The changes in this PR will be included in the next version bump. This PR includes changesets to release 16 packages
Not sure what this means? Click here to learn what changesets are. Click here if you're a maintainer who wants to add another changeset to this PR |
🧪 E2E Test Results✅ All tests passed
|
| Passed | Failed | Skipped | Total | |
|---|---|---|---|---|
| ✅ ▲ Vercel Production | 3662 | 0 | 685 | 4347 |
| ✅ 💻 Local Development | 3922 | 0 | 586 | 4508 |
| ✅ 📦 Local Production | 3922 | 0 | 586 | 4508 |
| ✅ 🐘 Local Postgres | 3922 | 0 | 586 | 4508 |
| ✅ 🪟 Windows | 320 | 0 | 2 | 322 |
| ✅ 🌐 Cross-language Conformance | 68 | 0 | 74 | 142 |
| ✅ vercel-http-transport | 823 | 0 | 143 | 966 |
| ✅ vercel-multi-region | 27 | 0 | 0 | 27 |
| ✅ vercel-ws-transport | 557 | 0 | 87 | 644 |
| Total | 17223 | 0 | 2749 | 19972 |
Details by Category
✅ ▲ Vercel Production
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-node | 133 | 0 | 28 |
| ✅ astro-quickjs | 133 | 0 | 28 |
| ✅ example-node | 133 | 0 | 28 |
| ✅ example-quickjs | 133 | 0 | 28 |
| ✅ express-node | 133 | 0 | 28 |
| ✅ express-quickjs | 133 | 0 | 28 |
| ✅ fastify-node | 133 | 0 | 28 |
| ✅ fastify-quickjs | 133 | 0 | 28 |
| ✅ hono-node | 133 | 0 | 28 |
| ✅ hono-quickjs | 133 | 0 | 28 |
| ✅ nest-node | 133 | 0 | 28 |
| ✅ nest-quickjs | 133 | 0 | 28 |
| ✅ nextjs-turbopack-node | 158 | 0 | 3 |
| ✅ nextjs-turbopack-quickjs | 158 | 0 | 3 |
| ✅ nextjs-webpack-node | 158 | 0 | 3 |
| ✅ nextjs-webpack-quickjs | 158 | 0 | 3 |
| ✅ nitro-node | 133 | 0 | 28 |
| ✅ nitro-quickjs | 133 | 0 | 28 |
| ✅ nuxt-node | 133 | 0 | 28 |
| ✅ nuxt-quickjs | 133 | 0 | 28 |
| ✅ python-node | 66 | 0 | 95 |
| ✅ sveltekit-node | 152 | 0 | 9 |
| ✅ sveltekit-quickjs | 152 | 0 | 9 |
| ✅ tanstack-start-node | 133 | 0 | 28 |
| ✅ tanstack-start-quickjs | 133 | 0 | 28 |
| ✅ vite-node | 133 | 0 | 28 |
| ✅ vite-quickjs | 133 | 0 | 28 |
✅ 💻 Local Development
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 27 |
| ✅ astro-stable-quickjs | 134 | 0 | 27 |
| ✅ express-stable-node | 134 | 0 | 27 |
| ✅ express-stable-quickjs | 134 | 0 | 27 |
| ✅ fastify-stable-node | 134 | 0 | 27 |
| ✅ fastify-stable-quickjs | 134 | 0 | 27 |
| ✅ hono-stable-node | 134 | 0 | 27 |
| ✅ hono-stable-quickjs | 134 | 0 | 27 |
| ✅ nest-stable-node | 134 | 0 | 27 |
| ✅ nest-stable-quickjs | 134 | 0 | 27 |
| ✅ nextjs-turbopack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 1 |
| ✅ nitro-stable-node | 134 | 0 | 27 |
| ✅ nitro-stable-quickjs | 134 | 0 | 27 |
| ✅ nuxt-stable-node | 134 | 0 | 27 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 27 |
| ✅ sveltekit-stable-node | 153 | 0 | 8 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 8 |
| ✅ tanstack-start-node | 134 | 0 | 27 |
| ✅ tanstack-start-quickjs | 134 | 0 | 27 |
| ✅ vite-stable-node | 134 | 0 | 27 |
| ✅ vite-stable-quickjs | 134 | 0 | 27 |
✅ 📦 Local Production
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 27 |
| ✅ astro-stable-quickjs | 134 | 0 | 27 |
| ✅ express-stable-node | 134 | 0 | 27 |
| ✅ express-stable-quickjs | 134 | 0 | 27 |
| ✅ fastify-stable-node | 134 | 0 | 27 |
| ✅ fastify-stable-quickjs | 134 | 0 | 27 |
| ✅ hono-stable-node | 134 | 0 | 27 |
| ✅ hono-stable-quickjs | 134 | 0 | 27 |
| ✅ nest-stable-node | 134 | 0 | 27 |
| ✅ nest-stable-quickjs | 134 | 0 | 27 |
| ✅ nextjs-turbopack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 1 |
| ✅ nitro-stable-node | 134 | 0 | 27 |
| ✅ nitro-stable-quickjs | 134 | 0 | 27 |
| ✅ nuxt-stable-node | 134 | 0 | 27 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 27 |
| ✅ sveltekit-stable-node | 153 | 0 | 8 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 8 |
| ✅ tanstack-start-node | 134 | 0 | 27 |
| ✅ tanstack-start-quickjs | 134 | 0 | 27 |
| ✅ vite-stable-node | 134 | 0 | 27 |
| ✅ vite-stable-quickjs | 134 | 0 | 27 |
✅ 🐘 Local Postgres
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ astro-stable-node | 134 | 0 | 27 |
| ✅ astro-stable-quickjs | 134 | 0 | 27 |
| ✅ express-stable-node | 134 | 0 | 27 |
| ✅ express-stable-quickjs | 134 | 0 | 27 |
| ✅ fastify-stable-node | 134 | 0 | 27 |
| ✅ fastify-stable-quickjs | 134 | 0 | 27 |
| ✅ hono-stable-node | 134 | 0 | 27 |
| ✅ hono-stable-quickjs | 134 | 0 | 27 |
| ✅ nest-stable-node | 134 | 0 | 27 |
| ✅ nest-stable-quickjs | 134 | 0 | 27 |
| ✅ nextjs-turbopack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-turbopack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-turbopack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-stable-quickjs | 160 | 0 | 1 |
| ✅ nextjs-webpack-canary-node | 141 | 0 | 20 |
| ✅ nextjs-webpack-canary-quickjs | 141 | 0 | 20 |
| ✅ nextjs-webpack-stable-node | 160 | 0 | 1 |
| ✅ nextjs-webpack-stable-quickjs | 160 | 0 | 1 |
| ✅ nitro-stable-node | 134 | 0 | 27 |
| ✅ nitro-stable-quickjs | 134 | 0 | 27 |
| ✅ nuxt-stable-node | 134 | 0 | 27 |
| ✅ nuxt-stable-quickjs | 134 | 0 | 27 |
| ✅ sveltekit-stable-node | 153 | 0 | 8 |
| ✅ sveltekit-stable-quickjs | 153 | 0 | 8 |
| ✅ tanstack-start-node | 134 | 0 | 27 |
| ✅ tanstack-start-quickjs | 134 | 0 | 27 |
| ✅ vite-stable-node | 134 | 0 | 27 |
| ✅ vite-stable-quickjs | 134 | 0 | 27 |
✅ 🪟 Windows
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ nextjs-turbopack-node | 160 | 0 | 1 |
| ✅ nextjs-turbopack-quickjs | 160 | 0 | 1 |
✅ 🌐 Cross-language Conformance
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ python | 68 | 0 | 74 |
✅ vercel-http-transport
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ example | 133 | 0 | 28 |
| ✅ express | 133 | 0 | 28 |
| ✅ hono | 133 | 0 | 28 |
| ✅ nextjs-turbopack | 158 | 0 | 3 |
| ✅ nitro | 133 | 0 | 28 |
| ✅ vite | 133 | 0 | 28 |
✅ vercel-multi-region
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ nextjs-turbopack | 27 | 0 | 0 |
✅ vercel-ws-transport
| App | Passed | Failed | Skipped |
|---|---|---|---|
| ✅ example | 133 | 0 | 28 |
| ✅ express | 133 | 0 | 28 |
| ✅ nextjs-turbopack | 158 | 0 | 3 |
| ✅ vite | 133 | 0 | 28 |
📊 Workflow Benchmarkscommit Backend:
Streams
📈 STSO distribution vs main (inline / queue-hop histograms)1020 steps (inline) Cumulative STSO time: main 137118ms → this run 126145ms (Δ -10973ms, -8%) 📈 CRTT drill-down vs main (RTT distributions & profiles)RTT over stream progress (avg per tenth of stream, bars scaled min→max): RTT by chunk size (avg per log size bin, ~160B → ~12KB serialized, bars scaled min→max): Delivery jitter over stream progress (avg positive CDV per tenth of stream, bars scaled min→max): ℹ️ Metric definitions & methodologyStreams: first-chunk RTT (the stream-open path, before any buffering/backpressure), CRTT percentiles, and worst delivery stall (CDV max). Cells are medians across iterations; per-run values in the artifacts. No 🔴/🟢 marks until targets attach. The collapsed STSO distribution section above buckets every step gap, split inline (same warm process — pure framework overhead) vs queue-hop (fresh process — dispatch, reinit, replay). The collapsed CRTT drill-down: per-variant RTT histograms (fixed log bins, Best/P75/P90/P99 deltas compare against the most recent benchmark run on Metrics — TTFS: time to first step body (in-deployment start() → first step body) · Fan-out TTFS: fan-out time to first step (in-deployment start() → first of the parallel step bodies to complete) · Fan-out TTLS: fan-out time to last step (in-deployment start() → last of the parallel step bodies to complete, i.e. when the Promise.all resolves) · STSO: step-to-step overhead (gap between consecutive step bodies) · WO: workflow overhead (whole-run time outside step bodies, in-deployment anchored) · CRTT: chunk round-trip time (per-chunk write → read latency, one clock domain: deployment → stream backend → same deployment) · CDV: chunk delay variation / delivery jitter (inter-arrival gap minus inter-write gap per seq-adjacent pair; skew-free; the row is each run's MAX positive value, so one stall moves it) Scenarios — step: one trivial no-op step, no stream; no hooks, so the run stays in turbo mode (in-process fast path) · stream: one streaming step; no hooks, so the run stays in turbo mode (in-process fast path) · hook + stream: registers a hook before one step, which exits turbo mode (dispatch path) · 1020 steps: 1020 trivial sequential steps; STSO is measured between consecutive steps in the given step ranges, and WO is the whole-run overhead outside step bodies · Promise.all(100 steps): 100 trivial no-op steps started together in a single Promise.all; Fan-out TTFS is the first of them to complete and Fan-out TTLS the last, both from the in-deployment clientStart, so their gap is the spread the runtime adds across the fan-out · paced control (100/s, 60B): the control: 300 tiny (~60B) deltas metronome-paced at 100/s — zero workload structure, so it reads the transport floor and flush cadence, and disambiguates transport-wide vs workload-specific when a replay row moves · size sweep (100/s, 160B-12KB): same pacing as the control with deltas padded in rotation across seven log-spaced sizes (~160B–12KB) — rotation decouples size from stream position, so it isolates whether chunk size causes latency · replay gateway-gpt-5.4-nano-2000t (1x): raw provider SSE cadence captured at the AI gateway boundary (gpt-5.4-nano, the most popular gateway model; per-token deltas p50 208B = the modal production chunk size), replayed exactly as measured — the typical customer's workload; its CDV is the typical customer's real delivery jitter · replay eve-gpt-5.6-sol-2000t (1x): a captured eve turn (gpt-5.6-sol, the most-used demanding eve model; ~2000 output tokens = production p50 turn length) replayed exactly as measured — eve's envelope protocol re-ships the cumulative message so sizes ramp 142B→13KB; the demanding outlier tenant's reality · replay eve-gpt-5.6-sol-2000t (2x): the same eve capture at 2x — the headroom/stress row; real fast-tier models emit the same chunk sizes at proportionally higher rate, so time compression is a faithful speed model · first chunk (pooled): every run's seq-0 RTT pooled across all stream scenarios — the first chunk precedes any workload differentiation, so pooling samples one shared stream-open path with exact percentiles Replay cadences (semantic sha256) — eve-gpt-5.6-sol-2000t 🔴 marks a percentile over its target (within target is left unmarked). Targets (p75/p90/p99, ms) — TTFS 200/300/600 All timestamps are deployment-side; runs are triggered in-deployment, so the CI runner and api.vercel.com sit outside every measured window. TTFS = Cold starts stay in the numbers (real bursty-workload latency, inflates P75+); Best is the warm floor. |
Sim WorldSimulated world deterministic testing for races. Traces 🟠 world-sim scenario book — 1 fail of 41 total
Full trace: |
About these numbersSizes are gzip; parentheses show the change against
|
| // won still runs the body; create won + claim lost skips). | ||
| inlineClaims.set(entry.correlationId, { owned: false }); | ||
| if (entry.kind === 'inline-created') { | ||
| // The created row lost: the entity already existed. Make |
There was a problem hiding this comment.
🟡 Changes recommended
Several correctness, coverage, and test-fixture issues remain unresolved.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
Adds duplicate step_created verification to detect replay collisions and fail corrupted runs as CORRUPTED_EVENT_LOG.
Changes:
- Compares persisted step names and decrypted/decompressed inputs.
- Integrates verification across suspension-handler creation paths.
- Adds regression tests and a core patch changeset.
File summaries
| File | Summary |
|---|---|
packages/core/src/runtime/suspension-handler.ts |
Integrates duplicate verification into creation paths. |
packages/core/src/runtime/suspension-handler.test.ts |
Adds conflict-verification tests. |
packages/core/src/runtime/step-create-conflict.ts |
Implements persisted step comparison. |
packages/core/src/runtime.ts |
Treats detected corruption as terminal. |
.changeset/step-created-conflict-guard.md |
Records the core patch release. |
Review details
Suppressed comments (2)
packages/core/src/runtime.ts:3429
- This terminal handling only covers a
CorruptedEventLogErrorthrown while awaitinghandleSuspension. With deferred fan-out enabled,verifyDuplicateBatchedStepCreateruns insidedeferredBatchWorkfor trailing chunks; those rejections are joined later at the dispatch/step joins and rethrown outside this branch, causing redelivery instead of writingrun_failed. Route deferred-work errors through the same non-retryable handling before retrying or acking.
if (
!FatalError.is(suspensionError) &&
!CorruptedEventLogError.is(suspensionError)
) {
packages/core/src/runtime/step-create-conflict.ts:145
- The stated best-effort behavior for legacy v1 inputs is not implemented here: when both values are non-binary, this branch compares JSON strings and returns
different. Thus a duplicate from a legacy run with different arguments throwsCorruptedEventLogError, even though non-binary inputs are supposed to be logged and treated as incomparable. Returnincomparablefor any non-Uint8Arraypair instead.
// Legacy specVersion 1 runs store step input as plain JSON.
try {
return JSON.stringify(a) === JSON.stringify(b) ? 'equal' : 'different';
} catch {
- Files reviewed: 5/5 changed files
- Comments generated: 5
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| /** Same check as the single-write paths, for a `step_created` that the | ||
| * batch commit reported as a duplicate. */ | ||
| const verifyDuplicateBatchedStepCreate = ( | ||
| entry: (typeof batchQueue)[number], | ||
| conflictMessage: string | ||
| ): Promise<void> => | ||
| verifyDuplicateStepCreate({ |
| world, | ||
| runId, | ||
| correlationId: entry.correlationId, | ||
| stepName: entry.stepName ?? '', |
| : 'Wait already exists, continuing', | ||
| { | ||
| if (entry.kind === 'step') { | ||
| await verifyDuplicateBatchedStepCreate(entry, err.message); |
| function isUnserializablePlaceholder(plaintext: Uint8Array): boolean { | ||
| try { | ||
| const { payload } = decodeFormatPrefix(plaintext); | ||
| return new TextDecoder() | ||
| .decode(payload) | ||
| .includes(UNSERIALIZABLE_STEP_INPUT_MARKER); |
| stepDispatchIdempotencyKey, | ||
| } from './helpers.js'; | ||
| import { ReplayRecoveryReporter } from './replay-recovery-reporter.js'; | ||
| import { verifyDuplicateStepCreate } from './step-create-conflict.js'; |
VaguelySerious
left a comment
There was a problem hiding this comment.
AI review: blocking issues found
| world, | ||
| runId, | ||
| correlationId: entry.correlationId, | ||
| stepName: entry.stepName ?? '', |
There was a problem hiding this comment.
AI Review: Blocking
entry.stepName ?? '' is '' for every inline-created entry, so a benign lost inline pre-claim fails the run.
The inline-created push (the pair fold, ~L1238) sets event.eventData.stepName but never the top-level entry.stepName — only kind: 'step' (~L1386) does. So the call at L1697 passes stepName: '', and persisted.stepName !== '' is unconditionally true for any real step. Every lost pre-claim on the pair-folded path throws CorruptedEventLogError.
That path is not exotic: the fold is default-on (WORKFLOW_BATCH_TRANSITIONS defaults to true), and the comment right above L1697 describes the 409 as the ordinary outcome — "an earlier delivery's create, or a racing handler's claim" — whose expected result is owned: false and a skipped body. It now fails the run instead.
suspension-handler.test.ts "records a lost claim when the pair's create-claim is refused" doesn't catch this because createBatchWorld has no steps property at all, so world.steps.get throws TypeError: Cannot read properties of undefined and the guard's own catch swallows it as "unreadable persisted step, continuing". Against a World that implements steps.get it throws.
Reproduced locally: pair-folded s1 + eager s2/s3, World returns 409 on both s1 rows, steps.get returns exactly the same invocation (same name, same args):
CorruptedEventLogError: Step s1 was already created as "s1", but this replay
invoked "" under the same correlation id.
at verifyDuplicateStepCreate src/runtime/step-create-conflict.ts:80:11
at commitChunk src/runtime/suspension-handler.ts:1697:19
Adding stepName: queueItem.stepName, to the inline-created push makes it pass. Worth making the type stop allowing it too — stepName is optional on the entry, so nothing forces the two step-carrying kinds to set it; either narrow it per-kind or make verifyDuplicateBatchedStepCreate return early when stepName is missing rather than comparing against ''.
| }); | ||
| const world = { | ||
| events: { create: eventsCreate }, | ||
| steps: { get: stepsGet }, |
There was a problem hiding this comment.
AI Review: Note
The new conflictingWorld is the only mock in this file with a steps property. createBatchWorld (used by the whole handleSuspension batched fan-out block, including the two 409 tests the new code now runs through) has none, so every batch-path duplicate silently takes the "could not be read" branch and the added verification is never exercised there. That is what hides the inline-created bug above.
Adding steps: { get: ... } to createBatchWorld would make the batch tests exercise the same code production does, and would have turned the existing lost-pre-claim test red.
| return; | ||
| } | ||
|
|
||
| if (persisted.stepName !== stepName) { |
There was a problem hiding this comment.
AI Review: Note
The name half of this check is already covered elsewhere, and now reports under two different codes depending on timing.
step.ts (the step consumer) compares eventStepName !== stepName for any incoming event with a stepName, step_created included, and raises ReplayDivergenceError. That fires on every replay where the winner's step_created is in the log. This new check fires only on the one pass that loses the create race, and classifies the same condition as CORRUPTED_EVENT_LOG.
So the genuinely new coverage in this PR is the input comparison below, and it is the narrowest of the two: it only ever runs on the losing-create pass. On any later replay hasCreatedEvent is true, no step_created is attempted, no 409 comes back, and the arguments are never compared — the cross-wired result is delivered exactly as before. Worth saying so in the doc comment, since the current wording reads as if the guard covers the collision generally.
If you want the args check to cover every replay rather than just the race window, the natural place is next to the existing name check in the step consumer, where both sides are already in hand and it costs no network.
|
|
||
| let persisted: { stepName: string; input?: unknown }; | ||
| try { | ||
| persisted = await world.steps.get(runId, correlationId); |
There was a problem hiding this comment.
AI Review: Note
This turns what was a log line into a blocking network round trip on the suspension write path, and it over-fetches.
No params, so resolveData defaults to 'all' → remoteRefBehavior=resolve. The Vercel World resolves and returns input, output and error, with remote refs materialized. The check reads stepName and input. On a conflict where the winner has already completed with a large output, this downloads a payload it never looks at, then decompresses both inputs and holds four buffers (two stored, two plaintext) live at once, in a function serving many concurrent runs. There is no resolveData value that fetches input without output today, but the fetch is at least worth a comment acknowledging the cost, and the pair could be compared without materializing both plaintexts when the stored bytes are already equal (the common benign case) — bytesEqual on the raw values first, decrypt/decompress only on mismatch.
Second, on the batch path this is awaited inside the sequential for loop in commitChunk, so N conflicting rows in one chunk cost N serialized round trips. The single-write and resilient-dispatch paths keep their per-op concurrency; the batch path doesn't. In practice the fan-out tail is thin — per-run step_created bursts over the last hour are p50 1, p90 2, p99 5, only 0.29% ≥ 8 wide — but the max observed is 103, so the worst case is ~100 sequential GETs added to one suspension. Collecting the verifications and settling them together would remove that without changing semantics.
| if (!FatalError.is(suspensionError)) { | ||
| if ( | ||
| !FatalError.is(suspensionError) && | ||
| !CorruptedEventLogError.is(suspensionError) |
There was a problem hiding this comment.
AI Review: Note
Two ways a CorruptedEventLogError raised in a suspension op never reaches this branch:
settlePhasethrowsreasons.find(PreconditionFailedError.is) ?? reasons[0]. A concurrent 412 in the same phase wins by design, and the caller then restarts the replay from a corrected log — wherehasCreatedEventis now true, so no create is attempted and the collision is never re-detected.- Absent a 412 it throws
reasons[0], positionally. An unrelated transient rejection earlier inopsmasks the corruption, the message is redelivered, and the next pass sees the create already in the log.
Both are pre-existing behaviours of settlePhase rather than something this PR introduces, and neither is worse than today's silent continue. But they mean the guard is best-effort in a second sense beyond the one documented in step-create-conflict.ts: it can miss a collision it would have caught, purely on the timing of a sibling write. Worth a line in the doc comment so the production signal isn't read as a complete count.
| } | ||
| // Legacy specVersion 1 runs store step input as plain JSON. | ||
| try { | ||
| return JSON.stringify(a) === JSON.stringify(b) ? 'equal' : 'different'; |
There was a problem hiding this comment.
AI Review: Nit
JSON.stringify renders Map and Set as {}, so two structurally different collections on the legacy v1 path compare equal and the guard reports a benign duplicate. Harmless (it degrades to today's behaviour) but the doc comment above describes this as a comparison, so a // legacy v1 collections compare equal here note would keep the next reader from trusting it.
Detection half of the correlation-id draw-order bug: when a concurrent replay already created a step under the same correlation id, verify it is the same step invocation before continuing, and fail the run if it is not.
Why
Correlation ids are ordinals of one per-run draw sequence. When two replays reach a
useStepcall in different orders, one id gets bound to two different invocations; the World keeps whicheverstep_createdlanded first, and the loser's branch goes on to await an id whose body runs with someone else's arguments. #3554 and #3700 make the schedule deterministic so this should no longer happen, but the failure mode when it does is the worst kind: the step consumer only checks the step name on an incoming event, and every duplicate-create path in the suspension handler swallows the 409 withStep already exists, continuing. Two branches mapping the same step over a list swap results silently and the run completes. A production run on4.8.5did exactly that (fix for that channel: thestablebackport of #3700 is in a separate PR).What this does
New
runtime/step-create-conflict.ts,verifyDuplicateStepCreate: on a duplicatestep_created, read the persisted step (world.steps.get) and comparestepNamewith ours, andmaybeDecryptanddecompress, so the AES-GCM nonce and the compression layer can't produce a false difference.Same invocation: continue exactly as before. Different name or arguments: throw
CorruptedEventLogError.runtime.tstreats that like aFatalErrorfrom the suspension handler (fails the run with the classified code, hereCORRUPTED_EVENT_LOG) rather than redelivering into the same collision.Best-effort by design: an unreadable persisted step, a missing input, a non-binary (legacy v1) input, or a serialization-failure placeholder logs and continues, so the guard can never turn a transient read failure into a failed run.
Wired into every path that tolerates a duplicate
step_created: sequential write, resilient dispatch, batch of one, batched fan-out, and the folded lazy-inline pair (verified when the created row loses; a lost started row alone means we created it).Not covered here: the unfolded lazy-inline claim in
step-executor.ts, where a loststep_startedclaim returnsskippedbefore the workflow continues to await the id. It has the same hazard and the same inputs available; it needs its own look at how an error thrown fromexecuteSteppropagates, so it's a follow-up.Testing
suspension-handler.test.ts(same input, different input, different name, encrypted compare by plaintext, unreadable step, missing input, placeholder, resilient-dispatch path).@workflow/coresuite green;tsc --noEmitclean.🤖 Generated with Claude Code