[core] Detect wedged waits instead of wake-looping on them forever - #3541
[core] Detect wedged waits instead of wake-looping on them forever#3541pranaygp wants to merge 1 commit into
Conversation
On worlds that persist the wait entity and its event-log row in separate writes (world-vercel), a request can commit the entity and then fail before the row insert. Every retry of the event write then conflicts (409) against the committed entity while the log stays permanently short one row. sleep() resolves only from a wait_completed row and the elapsed-wait pass can only complete waits whose wait_created row it can read, so the run replays into the same conflict forever: a ~1s wake loop that never errors and never completes. Make the contradiction loud, and terminal past a generous threshold: - wait_completed (elapsed-wait pass): when the create conflicts AND the follow-up reload still cannot produce the row, warn and report workflow.wait.wedge_suspected on the invocation span; once the clock is more than WORKFLOW_WAIT_WEDGE_FAIL_AFTER_SECONDS (default 600) past the wait's resumeAt, fail the run as CORRUPTED_EVENT_LOG. - wait_created (suspension handler): resumeAt cannot anchor this site (an uncreated wait recomputes it from the live clock every replay), so the anchor is the scheduling instant embedded in the wait's replay-stable correlation id. Past the threshold the contradiction is verified against a fresh event-log read before failing, so a concurrent writer's row landing after this replay's snapshot is never mistaken for a wedge. Escalation is stateless on purpose: every wake is a fresh queue message, so there is no attempt counter to persist — but "how long has this contradiction persisted against a replay-stable anchor" is derivable on every observation. Benign concurrent-handler races (the conflicting row is readable) stay silent exactly as before. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Signed-off-by: Pranay Prakash <pranay.gp@gmail.com>
🦋 Changeset detectedLatest commit: f81f331 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 E2E Test SummarySummary
Details by Category✅ ▲ Vercel Production
✅ 💻 Local Development
✅ 📦 Local Production
✅ 🐘 Local Postgres
✅ 🪟 Windows
✅ vercel-multi-region
|
📊 Workflow Benchmarkscommit Backend:
📈 STSO distribution vs main (inline / queue-hop histograms)1020 steps (inline) Cumulative STSO time: main 194368ms → this run 171408ms (Δ -22960ms, -12%) ℹ️ Metric definitions & methodologyThe collapsed STSO distribution section above buckets every step gap of the sequential-steps run (not a sampled window), split by whether the step ending the gap ran inline — in the same warm process as the step before it, so the gap is pure framework overhead — or after a queue-hop — the first step of a fresh process, which pays queue dispatch, client reinit and event-log replay. Bars overlay the two runs: 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, deployment clocks) · 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) · SL: stream latency (in-deployment write → read propagation, readAt - writtenAt) · SO: stream overhead (end-to-end write+consume time beyond the modelled generation window) 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 · stream latency: parallel reader/writer steps on a dedicated stream; SL is the in-deployment write->read propagation (readAt - writtenAt) · stream overhead (text): writer streams 300 variable-length text token deltas paced at 100/s for 3s (a haiku-size LLM's token throughput) while a parallel reader drains the whole stream; SO is the end-to-end write+consume time beyond the 3s generation window (overhead/backpressure) · stream overhead (structured): same workload as stream overhead (text), but each delta is an AI-SDK-style structured object ({ type: 'text-delta', id, text }) instead of a raw string, so the SO gap vs the text scenario is the added serialization cost 🔴 marks a percentile over its target (within target is left unmarked). Targets (p75/p90/p99, ms) — TTFS 200/300/600 · SL 50/60/125 · SO 250/500/1000 All metrics are measured from deployment-side timestamps only. Runs are triggered by an in-deployment route that stamps the anchor ( Cold starts are kept in the numbers on purpose — they are part of real bursty-workload latency. The workbench deployment cold-starts the |
There was a problem hiding this comment.
Pull request overview
Adds SDK-side detection for “wedged waits” (409 conflict on wait event writes where the corresponding event-log row is never readable), so runs stop silently wake-looping forever and instead warn for a configurable window before failing as CORRUPTED_EVENT_LOG. This fits into @workflow/core runtime durability/corruption detection, complementing server-side recovery for related wedge classes.
Changes:
- Introduces stateless, time-anchored wait-wedge classification and error messaging (
runtime/wait-wedge.ts) with a tunable threshold (WORKFLOW_WAIT_WEDGE_FAIL_AFTER_SECONDS). - Adds runtime integration at both wedge sites (
wait_completedinruntime.ts,wait_createdinsuspension-handler.ts), including telemetry reporting (workflow.wait.wedge_suspected). - Adds unit + queue-handler integration tests and documents the new environment variable.
Reviewed changes
Copilot reviewed 8 out of 8 changed files in this pull request and generated 1 comment.
Show a summary per file
| File | Description |
|---|---|
| packages/core/src/telemetry/semantic-conventions.ts | Adds the workflow.wait.wedge_suspected semantic convention for span reporting. |
| packages/core/src/runtime/wait-wedge.ts | New wedge detection utilities: thresholding, ULID anchor decoding, fresh-read verification, shared error message. |
| packages/core/src/runtime/wait-wedge.test.ts | Unit tests for classification, env override behavior, ULID decoding, and verification-read behavior. |
| packages/core/src/runtime/wait-wedge-detection.test.ts | End-to-end-ish handler tests covering both wedge sites and benign concurrent-winner races. |
| packages/core/src/runtime/suspension-handler.ts | Adds wedge detection/escalation on wait_created conflict path (suspension handler). |
| packages/core/src/runtime.ts | Adds wedge detection/escalation on wait_completed conflict path (elapsed-wait pass) and routes CorruptedEventLogError to terminal handling. |
| docs/content/docs/v5/configuration/runtime-tuning.mdx | Documents WORKFLOW_WAIT_WEDGE_FAIL_AFTER_SECONDS behavior and default. |
| .changeset/wait-wedge-detection.md | Changeset for the new runtime behavior (patch). |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| runId, | ||
| queueItem.correlationId | ||
| )); | ||
| if (suspectWedge) { |
There was a problem hiding this comment.
Restructured in ee40e67. For the record, the aliased-condition form did typecheck (TS 4.4+ narrows const booleans built from narrowing conjunctions), but the fragility point was fair — and the rework for the other review thread rebuilt this branch anyway: the !== undefined check is now part of the if condition directly, so the narrowing is structural rather than aliased, and the timestamp is no longer used in the body beyond the guard.
Sim WorldSimulated world deterministic testing for races. Traces 🟠 Mint-ordered log — 3 fail of 41 total
Full trace: 🟢 Append-only log — 0 fail of 41 total
Full trace: |
|
|
||
| /** Effective threshold. Override: `WORKFLOW_WAIT_WEDGE_FAIL_AFTER_SECONDS`. */ | ||
| export const getWaitWedgeFailAfterSeconds = (): number => | ||
| envNumber( |
There was a problem hiding this comment.
There was a problem hiding this comment.
Confirmed and fixed in ee40e67 — this finding was correct. generateUlid is ulid(fixedTimestamp) (workflow.ts), where fixedTimestamp is the run's creation epoch held constant for replay stability, so the correlation-id ULID is a run-wide constant and cannot anchor a per-wait escalation: any run older than the threshold would fail as CORRUPTED_EVENT_LOG on one benign concurrent-suspension race whenever the winner's row wasn't yet visible to the verification read.
The wait_created site is now warn-only: the run-epoch timestamp is kept solely as a cheap pre-filter (conflicts in runs younger than the threshold skip detection entirely, so the hot path is unchanged), the fresh-log-read verification still gates the warning, and escalation for genuinely stale wait_created wedges is left to workflow-server#782's recovery (backfill fresh, cancel stale), which classifies against entity state the SDK cannot see. The wait_completed site keeps warn→fail — its resumeAt anchor is durable and genuinely per-wait. Docs, tests, and the PR description are updated to match.
There was a problem hiding this comment.
Thanks for confirming and fixing. Removing the ULID-anchored escalation for the wait_created site is the right call — since generateUlid is ulid(fixedTimestamp) with a run-wide constant, there's no sound per-wait anchor available there, and the correctly-anchored wait_completed path (on the durable resumeAt) is unaffected. The tradeoff — a genuine wait_created wedge reverts to the pre-PR silent wake-loop rather than being detected — is strictly safer than failing healthy long-running runs; sound detection would need a durable per-wait anchor (e.g. a World read of the wait entity).
Root cause
On worlds that persist the wait entity and its event-log row in separate, non-transactional writes (world-vercel / DynamoDB), a request can commit the wait entity and then fail before the event-row insert (crash, dropped connection). Every retry of the event write then conflicts (409) against the committed entity, while the log stays permanently short one row. The SDK swallows the conflict as "my write already landed" — which is half-true: the entity landed, the row didn't.
The consequences are an invisible infinite loop, not an error:
sleep()resolves only from await_completedrow (workflow/sleep.ts), and the elapsed-wait pass can only complete waits whosewait_createdrow it can read.resumeAt, every pass arms a fresh ~1s wake (the near-elapsed continuation key is second-bucketed, so dedup never collapses them —runtime/wait-continuation.ts).runningforever, burning an invocation per second, with nothing but aninfo-level "already exists, skipping" log line.Steps and runs had the same wedge class and got server-side recovery (workflow-server #704, #707); waits are the remaining unhealed sibling. A companion workflow-server PR (#782) heals fresh wait wedges by backfill and cancels stale ones; this PR is the SDK-side detection so the contradiction is loud while it persists — and, where a safe anchor exists, terminal once it is provable.
What this does
wait_completed(elapsed-wait pass,runtime.ts): warn, then fail. When the create conflicts AND the follow-up reload still cannot produce the row — the server says "completed", the log says "pending" — log a warning and reportworkflow.wait.wedge_suspectedon the invocation span. Once the clock is more than the threshold past the wait'sresumeAt(durable, adopted from thewait_createdrow, identical on every wake), fail the run asCORRUPTED_EVENT_LOG(same terminal path as the slot-gap check). The benign race (conflicting row IS readable after reload) stays silent exactly as before.wait_created(suspension handler): warn only. No per-wait replay-stable time anchor exists for an uncreated wait, so this site never fails the run:resumeAtis recomputed from the live clock on every replay, so it always sits in the future.ulid(fixedTimestamp)(a run-wide constant, held fixed precisely so ids are replay-stable). An earlier revision of this PR escalated on that ULID; as review pointed out, that would fail any sufficiently old run withCORRUPTED_EVENT_LOGon a single benign concurrent-suspension race.The run epoch still works as a cheap pre-filter (conflicts in runs younger than the threshold skip detection entirely, so the hot path is untouched), and a fresh event-log read gates the warning so a concurrent winner's late-landing row is not reported as a wedge. Terminating a genuinely stale
wait_createdwedge is owned by workflow-server #782's stale-cancel tier, which classifies against entity state the SDK cannot see.Why stateless, time-based escalation (where it applies): every wake of the loop is a fresh queue message (fresh delivery attempt = 1), so there is no attempt counter to persist across invocations. "How long has this contradiction persisted against a replay-stable time anchor" is derivable on every observation, and a healthy wait completes within seconds of its target.
Threshold
WORKFLOW_WAIT_WEDGE_FAIL_AFTER_SECONDS, default 600 (10 minutes), documented indocs/content/docs/v5/configuration/runtime-tuning.mdxnext to the other wait tunables. For wedged completions it is the warn→fail boundary; for wedged creations it is the detection pre-filter. The generous default means eventually-consistent read staleness cannot plausibly trigger a failure; the wedge, once real, is permanent — 10 minutes only bounds how long the loop burns invocations.Failure shape
Reuses
CorruptedEventLogError→run_failedwitherrorCode: CORRUPTED_EVENT_LOG(no new error code; the log genuinely cannot produce a row the World attests exists, which is this code's meaning, and it flows through existing classification, dashboards, and error docs). The throw happens in the elapsed-wait pass of the replay loop, which already routes it to the terminal path (same as the slot-gap check) — the suspension handler no longer throws, so no error-routing changes remain in this PR.Tests
runtime/wait-wedge.test.ts— unit: threshold classification + env override, run-epoch ULID decoding, fresh-read verification (found / missing / fail-open on read errors).runtime/wait-wedge-detection.test.ts— drives the real queue handler with a fake World (same harness pattern aswait-completion-replay.test.ts) through both wedges: benign concurrent-winner races stay silent and the run completes; wedged completions warn inside the threshold and fail withCORRUPTED_EVENT_LOGpast it; wedged creations warn (span attribute + log) but never fail and keep the run's normal suspension behavior.🤖 Generated with Claude Code