Skip to content

Fix cluster entity recovery when replayed requests defect during handler rebuild - #8637

Open
tim-smart wants to merge 3 commits into
mainfrom
agent/nelson/eff-1632-defect-replay-recovery
Open

tim-smart wants to merge 3 commits into
mainfrom
agent/nelson/eff-1632-defect-replay-recovery

Conversation

@tim-smart

Copy link
Copy Markdown
Contributor

Problem

When an entity handler defects, EntityManager rebuilds the handlers and replays unfinished requests inside the ResourceRef acquisition. That has three problems:

  • If a replayed request defects again, its defect reaches onDefect while isRestartingDueToDefect is still set, and it is dropped. The request stays in activeRequests with no handler and never completes.
  • Queuing that dropped defect isn't enough. The replay snapshot includes requests that arrived during acquisition and are still waiting for their first dispatch. Those get submitted by both the replay and their original writer, so a persisted request can run again after its successful reply has been stored.
  • If shutdown starts while replacement handlers are being acquired, the rebuild still replays application requests into the retiring entity.

Fix

  • Each handler generation now has its own defect guard, replacing the global restart flag. A retiring generation starts at most one rebuild, and a defect from the replacement generation can start the next one.
  • Replay now runs after the replacement is acquired and published. It stops if another rebuild takes over, and it is skipped once the activation has been removed.
  • A replayReady latch holds new arrivals until the backlog has been resubmitted. Shutdown opens it so waiting writers and EOF can proceed.
  • Active requests track whether they reached a handler. Only delivered requests go into a replay snapshot.
  • Application Request writes to an activation that is no longer registered are interrupted instead of being written to a retiring server. Callers already handle this outcome, since ResourceMap.get returns the same interrupt once the map closes:
    • The storage delivery loop treats an interrupt as a no-op, so the persisted message stays in storage for redelivery.
    • Local volatile calls go through RpcClient, which resumes the caller with the interrupt.
    • RunnerServer turns an interrupted sharding.send into an interrupt exit for the remote caller.
    • Before this change, the same request ended with an interrupt reply once the handler scope closed, or after entityTerminationTimeout.

The implementation and regression tests are by Adrian Gierakowski (@adrian-gierakowski), cherry-picked from the cluster-defect-replay-lifecycle and cluster-defect-replay-repro branches of rhinofi/effect with authorship kept.

Tests

packages/effect/test/cluster/DefectRecovery.test.ts covers repeated synchronous replay defects, coalescing of concurrent defects, arrivals during acquisition for persisted and volatile RPCs (including a check of the encoded reply store at handler entry), and shutdown during replacement acquisition. On unmodified main, 4 of the 5 fail and the concurrent-defect control passes.

  • DefectRecovery.test.ts: 5/5 pass on each of 3 repeated runs.
  • packages/effect/test/cluster and packages/effect/test/workflow: 207 tests pass.
  • oxlint, dprint and tsc -b for packages/effect pass.

Closes EFF-1632

adrian-gierakowski and others added 3 commits September 30, 2026 19:35
Keep the complete regression suite in this shared test-only commit so
both proposed fixes are evaluated against exactly the same scenarios.
Use upstream cluster layers and the real in-memory message-storage driver;
no Rhino service code or PostgreSQL instance is required.

On rc.118, replacement handler acquisition replays unfinished requests
before publication. A synchronous replay defect reaches recovery while
its restart guard is set and is discarded. The repeated-defect test
requires four attempts after three defects; vanilla rc.118 stops at two.
A concurrent-defect control requires only one rebuild for two defects
from the same handler generation.

Add arrival-during-acquisition cases for persisted and volatile RPCs.
After another replay defect, a request still waiting for first dispatch
must not be submitted both by recovery and by its original waiting
writer. Assert ordering, handler generations, results and one execution
of the arriving request.

For persistence, query MemoryDriver.encoded.repliesFor at each handler
entry and record the requestId plus whether its successful WithExit reply
already exists. The queued-defect proposal executes the same requestId
twice: the first entry sees no success reply, the second sees the stored
success from the first execution. This is a correctness risk for
non-idempotent side effects, beyond redundant work or duplicate replies.
The passing behavior executes once and leaves a successful reply stored.
This verifies storage-write-before-reexecution ordering, not PostgreSQL
crash durability or transaction isolation.

Set Persisted on each RPC directly, since group annotations do not
override existing RPC annotations. Assert stored success presence for
persisted RPCs and absence for volatile RPCs to verify the fixture.

Also cover shutdown starting during replacement acquisition. Once
acquisition finishes, shutdown must complete before its termination
timeout and must not replay application requests in the retiring entity.

Validation on unmodified rc.118: four regressions fail and the concurrent
control passes. pnpm check, focused oxlint and explicit-config dprint
checks pass. The lifecycle proposal passes all five shared cases; the
queued proposal's remaining failures are documented in its own commit.
The shared regression commit demonstrates lost recovery when a replayed
request defects during replacement acquisition. Queuing that defect
restores recovery but leaves duplicate execution after a successful reply
has been persisted, and application replay after shutdown has begun.
This alternative addresses all five shared scenarios.

ResourceRef completes acquisition and publication before EntityManager
replays unfinished requests. No ResourceRef API changes are required.
Replace the global restart guard with one defect guard per handler
generation: a generation reports its first defect only, while a defect
from a published replacement can immediately start another rebuild.

Bind replay to the acquired resource's identity and stop its loop when
another rebuild takes ownership. Hold incoming writes behind a replay
latch until the backlog has been submitted so new arrivals cannot overtake
it. Track whether an active request reached handlers; a request waiting
for its original dispatch must not also be included in a replay snapshot.
This prevents the queued proposal's same-requestId re-execution after a
WithExit/Success reply is already stored, as asserted by the shared tests.

Preserve retry backoff, request payloads, stream progress and shutdown
interruption suppression. When shutdown removes an activation, open the
replay latch so waiting writers can observe shutdown. Reject queued
application writes unless their exact activation is still registered,
while allowing control messages to drain. Skip application replay after
shutdown even if replacement acquisition completes successfully.

The first gated prototype left EOF waiting for the termination timeout
when shutdown began during acquisition. Independent review reproduced
that issue; opening the gate and checking activation identity fixed it.
A real-clock probe then completed in 23ms instead of the 1000ms timeout.

Compared with 52b46b7 on cluster-defect-replay-queued, this removes the
pending-defect queue and moves replay out of acquisition, at the cost
of an admission latch and per-request delivery bookkeeping. All regression
tests, including the encoded-store persistence assertion, live in the
shared parent commit rather than in either implementation commit.

Validation: all five shared regressions and all 184 cluster tests pass;
pnpm check, oxlint and formatting checks pass. Both runtime implementations
are unchanged by moving tests into their common parent. Two adversarial
review passes resolved the shutdown finding with no outstanding findings.
Co-authored-by: Adrian Gierakowski <agierakowski@gmail.com>
@changeset-bot

changeset-bot Bot commented Sep 30, 2026

Copy link
Copy Markdown

🦋 Changeset detected

Latest commit: 7d50b0b

The changes in this PR will be included in the next version bump.

This PR includes changesets to release 31 packages
Name Type
effect Patch
@effect/opentelemetry Patch
@effect/vitest Patch
@effect/ai-anthropic Patch
@effect/ai-openai-compat Patch
@effect/ai-openai Patch
@effect/ai-openrouter Patch
@effect/ai-typesafe Patch
@effect/atom-react Patch
@effect/atom-solid Patch
@effect/atom-vue Patch
@effect/platform-browser Patch
@effect/platform-bun Patch
@effect/platform-deno Patch
@effect/platform-node-shared Patch
@effect/platform-node Patch
@effect/sql-clickhouse Patch
@effect/sql-d1 Patch
@effect/sql-libsql Patch
@effect/sql-mssql Patch
@effect/sql-mysql2 Patch
@effect/sql-pg Patch
@effect/sql-pglite Patch
@effect/sql-sqlite-bun Patch
@effect/sql-sqlite-do Patch
@effect/sql-sqlite-node Patch
@effect/sql-sqlite-react-native Patch
@effect/sql-sqlite-wasm Patch
@effect/docgen Patch
@effect/doctest Patch
@effect/openapi-generator Patch

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

@github-actions

Copy link
Copy Markdown
Contributor

Bundle Size Analysis

Generated from PR build output; treat the content below as untrusted.

File Name Current Size Previous Size Difference
arbitrary-combinators.ts 38.71 KB 38.71 KB 0.00 KB (0.00%)
basic.ts 6.88 KB 6.88 KB 0.00 KB (0.00%)
batching.ts 9.97 KB 9.97 KB 0.00 KB (0.00%)
brand.ts 6.57 KB 6.57 KB 0.00 KB (0.00%)
cache.ts 10.66 KB 10.66 KB 0.00 KB (0.00%)
config.ts 21.98 KB 21.98 KB 0.00 KB (0.00%)
differ.ts 21.12 KB 21.12 KB 0.00 KB (0.00%)
http-client.ts 22.10 KB 22.10 KB 0.00 KB (0.00%)
http-router.ts 33.35 KB 33.35 KB 0.00 KB (0.00%)
logger.ts 10.92 KB 10.92 KB 0.00 KB (0.00%)
metric.ts 8.83 KB 8.83 KB 0.00 KB (0.00%)
optic.ts 6.80 KB 6.80 KB 0.00 KB (0.00%)
pubsub.ts 15.00 KB 15.00 KB 0.00 KB (0.00%)
queue.ts 11.91 KB 11.91 KB 0.00 KB (0.00%)
schedule.ts 11.03 KB 11.03 KB 0.00 KB (0.00%)
schema-bigdecimal.ts 13.58 KB 13.58 KB 0.00 KB (0.00%)
schema-binary.ts 39.74 KB 39.74 KB 0.00 KB (0.00%)
schema-class.ts 20.86 KB 20.86 KB 0.00 KB (0.00%)
schema-fromJsonSchemaDocument.ts 31.88 KB 31.88 KB 0.00 KB (0.00%)
schema-representation-roundtrip.ts 27.07 KB 27.07 KB 0.00 KB (0.00%)
schema-string-transformation.ts 14.32 KB 14.32 KB 0.00 KB (0.00%)
schema-string.ts 11.82 KB 11.82 KB 0.00 KB (0.00%)
schema-template-literal.ts 15.89 KB 15.89 KB 0.00 KB (0.00%)
schema-toArbitrary.ts 38.25 KB 38.25 KB 0.00 KB (0.00%)
schema-toCodeDocument.ts 25.09 KB 25.09 KB 0.00 KB (0.00%)
schema-toCodecJson.ts 20.05 KB 20.05 KB 0.00 KB (0.00%)
schema-toEquivalence.ts 20.25 KB 20.25 KB 0.00 KB (0.00%)
schema-toFormatter.ts 20.35 KB 20.35 KB 0.00 KB (0.00%)
schema-toJsonSchemaDocument.ts 24.81 KB 24.81 KB 0.00 KB (0.00%)
schema-toRepresentation.ts 20.31 KB 20.31 KB 0.00 KB (0.00%)
schema.ts 20.06 KB 20.06 KB 0.00 KB (0.00%)
stm.ts 12.91 KB 12.91 KB 0.00 KB (0.00%)
stream.ts 9.82 KB 9.82 KB 0.00 KB (0.00%)

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants