[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) { |
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.
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 makes wait writes transactional and backfills existing wedges; this PR is the SDK-side detection so the contradiction is loud while it persists and terminal once it is provable.
What this does
wait_completed(elapsed-wait pass,runtime.ts): 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, 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):resumeAtcannot anchor this site — an uncreated wait recomputes it from the live clock on every replay, so it always sits in the future. The anchor is the scheduling instant embedded in the wait's replay-stable correlation id (seeded RNG + replay clock ⇒ same ULID every replay). Within the threshold, behavior is unchanged (silent info — creation conflicts are the ordinary concurrent-suspension race). Past it, the contradiction is verified against a fresh event-log read before failing, so a concurrent writer's row landing after this replay's snapshot can never be mistaken for a wedge.Why stateless, time-based escalation: 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. 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 suspension-handler throw required one gate change inruntime.ts: the suspension-error catch now routesCorruptedEventLogErrorto its terminal fail-the-run path alongsideFatalError, instead of rethrowing for a redelivery that would replay into the same conflict.Tests
runtime/wait-wedge.test.ts— unit: threshold classification + env override, ULID anchor 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; contradictions inside the threshold warn and keep retrying (wake continuation still armed); contradictions past the threshold fail the run withCORRUPTED_EVENT_LOG.🤖 Generated with Claude Code