Skip to content

replication: recover outbound subscription wedged connected:false after peer restart (#466) - #467

Merged
kriszyp merged 4 commits into
mainfrom
kris/466-reconnect-recovery
Jun 24, 2026
Merged

kriszyp merged 4 commits into
mainfrom
kris/466-reconnect-recovery

Conversation

@kriszyp

@kriszyp kriszyp commented Jun 24, 2026 •

Copy link
Copy Markdown
Member

Fixes #466.

Summary

A node reconnecting to a freshly-restarted peer could wedge connected:false permanently — no socket, no SYN, no scheduled retry timer, and findWedgedNodeUrls/reconcileWorkers never re-driving it. This fixes the two recovery gaps that were code-confirmed in v5.1.10 (the third, watchdog-rearm, is a noted follow-up — see below).

The two fixes

PRIMARY — connect() reschedules on a createWebSocket rejection.
connect() is re-driven via setTimeout(() => this.connect()) from the close handler and forceReconnect() with no .catch(). connect() awaits createWebSocket(), which can reject before any socket / open / error / close listener exists — the no-valid-cert throw, or SNICallback.initialize() rejecting while a freshly-restarted peer rebuilds its TLS secure contexts. The old finally cleared reconnectScheduled and let the rejection escape as an unhandled rejection: no socket, no timer, connected:false forever.

connect() now catches that rejection and funnels it into a shared scheduleReconnect() (extracted from the close handler; forceReconnect() uses it too). The flag stays true with a fresh backoff retry, and the connection self-heals once the peer's certs settle. reconnectScheduled is cleared only on the success path (after this.socket is assigned — the original #420 semantic). The two unguarded connection.connect() call sites in replicator.ts (initial subscribe) are covered by the same internal handling.

BACKSTOP — findWedgedNodeUrls catches the never-connected entry, and the re-drive actually reconnects it.
The predicate required entry.disconnectedAt != null, which only disconnectedFromNode stamps (close / forceReconnect). A connect() that never fired 'open' leaves the entry connected:undefined with no disconnectedAt, invisible to the 30s wedge reconcile. The entry now records createdAt at creation; the predicate uses connected !== true and downSince = disconnectedAt ?? createdAt, and the reconcile re-drive loop mirrors it (connected === true skip).

The reconcile's re-driven subscribe-to-node also now carries forceReconnect:true. subscribeToNode reuses the cached connection when isReusableConnection is true, and subscribe() alone only re-emits the listener — so a never-opened (still-reusable) connection would otherwise be re-subscribed and stay wedged. forceReconnect() drives an independent reconnect and no-ops when a retry is already pending. (The original #233/#289 wedge healed without this only because that case is intentionallyUnsubscribed, hence non-reusable.) Two guards keep this from misfiring (both from cross-model review): the force only applies to a reused connection (a freshly-created one already called connect() — forcing it mid-await could open a duplicate socket), and the reconcile loop re-checks the full findWedgedNodeUrls predicate per database (a peer URL is flagged because some db is wedged, but a sibling db may be healthy or still in its initial-connect grace period).

Where to look / what to verify

  • Backstop false-positive guard (the main risk): the widened predicate and the forceReconnect:true re-drive must not churn a connection that is legitimately mid-initial-connect or a healthy slow base copy. A healthy connection flips connected:true via connectedToNode on 'open' within seconds, and a base copy runs while connected:true, so neither is flagged. The reconcile re-drive re-applies the same per-db predicate (connected !== true + downSince = disconnectedAt ?? createdAt past threshold + still desired), and only force-reconnects a reused connection — so a fresh sibling db or an in-flight dial is left alone. The createdAt fallback only matters for an entry that never reached 'open', still threshold-gated by a real timestamp.
  • reconnectScheduled state machine: the success path clears the flag after socket assignment; the reject path leaves it true with a pending timer; neither leaves it stuck-true-without-retry nor false-without-retry. The close handler now sets the flag (it previously did not) — this is more aligned with the existing #420 anti-double-arm intent, not less.

Deferred follow-up (NOT in this PR)

The receive watchdog is stop()-ed under backpressure (addPauseReason), so a leg that dies while paused at high backpressure won't fire forceReconnect() until resume. Re-arming it (or a separate liveness check not suppressed by pauseReasons) under sustained backpressure is a separate change. The two fixes here already break the permanent wedge; that would tighten the base-copy-stall case further.

Tests

  • unitTests/replication/connectReschedulesOnRejection.test.mjs (new): a createWebSocket rejection reschedules (no unhandled rejection / permanent stuck), reconnectScheduled ends consistent with a pending retry, the armed retry actually fires, and an intentionally-unsubscribed connection does not revive. Verified failing on the pre-fix build, passing after.
  • unitTests/replication/findWedgedNodeUrls.test.mjs (extended): flags a never-connected entry via createdAt past threshold; does not flag a fresh entry or one whose recent disconnectedAt should win over a stale createdAt. The never-connected case fails pre-fix, passes after.
  • unitTests/replication/forceReconnect.test.mjs (extended): forceReconnect() drives recovery on a never-opened, still-reusable connection (no socket yet).
  • Full unit suite green (179 passing).

Note: the wiring of forceReconnect:true through the worker subscribe-to-node message into subscribeToNode is exercised at the primitive level (the forceReconnect test above) but not end-to-end in a unit test — the worker's connection cache (connections) is module-scoped with no test seam, and a full reconnect-after-restart integration test is macOS-loopback-flaky in the harness (the same secure TLS issue noted in #466). It's a one-line if (request.forceReconnect) conditional.

🤖 Generated with Claude Code (model: Claude Opus 4.8)

…er peer restart (#466)

A node reconnecting to a freshly-restarted peer could wedge connected:false
permanently — no socket, no SYN, no retry timer, and the reconcile never
re-driving it. Two layered recovery gaps:

PRIMARY: connect() is re-driven via setTimeout(() => this.connect()) (close
handler and forceReconnect) with no .catch(). connect() awaits
createWebSocket(), which can reject before any socket/listener exists (no valid
cert yet, or SNICallback.initialize() failing while a restarted peer rebuilds
its TLS secure contexts). The old finally cleared reconnectScheduled and let the
rejection escape as an unhandled rejection — the only pending retry vanished.
Now connect() catches that rejection and funnels it into a shared
scheduleReconnect() (the close handler and forceReconnect use it too), so the
flag stays true with a fresh backoff retry and the connection self-heals once
the peer's certs settle. The two unguarded connect() call sites in replicator.ts
(initial subscribe) are covered by the same internal handling.

BACKSTOP: findWedgedNodeUrls required entry.disconnectedAt != null, which only
disconnectedFromNode stamps — a never-connected entry (connect() that never
fired 'open') stays connected:undefined with no disconnectedAt and was invisible
to the 30s wedge reconcile. The entry now records createdAt at creation; the
predicate uses connected !== true and disconnectedAt ?? createdAt as the
down-since clock (and the reconcile re-drive loop mirrors it). A real timestamp
still gates the threshold, so a fresh / mid-initial-connect / healthy entry is
never flagged.

Deferred follow-up (not in this PR): the receive watchdog is stop()'d under
backpressure (addPauseReason), so a leg that dies while paused at high
backpressure won't fire forceReconnect until resume — re-arming it under
sustained backpressure is a separate change.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@gemini-code-assist

Copy link
Copy Markdown

Warning

You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again!

@kriszyp
kriszyp requested review from kylebernhardy and ldt1996 June 24, 2026 01:26
@claude

claude Bot commented Jun 24, 2026

Copy link
Copy Markdown
Contributor

Reviewed; no blockers found.

… backstop actually heals (#466)

Cross-model (Codex) review of the #466 backstop: the reconcile re-posting
subscribe-to-node does not reconnect a never-opened entry on its own. The
worker's getSubscriptionConnection() returns the cached NodeReplicationConnection
whenever isReusableConnection is true (not finished, not intentionally
unsubscribed), and subscribe() alone only re-emits 'subscriptions-updated'. A
never-connected wedge (connect() rejected, no socket, still reusable) would just
get re-subscribed and stay wedged — the original #233/#289 wedge healed only
because that case is intentionallyUnsubscribed, hence non-reusable.

The reconcile now sets forceReconnect:true on the re-drive request, and
subscribeToNode calls connection.forceReconnect() when set. forceReconnect tears
the dead socket down and arms a fresh connect independent of the cache, and
no-ops when a retry is already pending (its reconnectScheduled / isFinished /
intentionallyUnsubscribed guards) — so the backstop heals the never-connected
case without churning a healthy or already-retrying connection.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@kriszyp
kriszyp marked this pull request as ready for review June 24, 2026 01:41
@kriszyp
kriszyp requested a review from a team as a code owner June 24, 2026 01:41
@gemini-code-assist

Copy link
Copy Markdown

Warning

You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again!

@kriszyp
kriszyp requested review from heskew and removed request for a team June 24, 2026 01:41
kriszyp and others added 2 commits June 23, 2026 19:43
…ive (#466)

Cross-model (Codex) second-pass review of the backstop force-reconnect:

1. subscribeToNode only force-reconnects a REUSED connection now. When the
   cached connection was non-reusable (e.g. the empty-subscription intentional
   close), getSubscriptionConnection creates a fresh NodeReplicationConnection
   and calls connect() before returning; calling forceReconnect() on that
   in-flight connect (before reconnectScheduled is set) could open a duplicate
   socket. getSubscriptionConnection now reports whether it reused vs. created,
   and the force only fires for the reused case.

2. The reconcile re-drive loop re-checks the full findWedgedNodeUrls predicate
   per database (connected !== true, downSince = disconnectedAt ?? createdAt,
   past WEDGE_RECONCILE_THRESHOLD_MS, still desired) instead of re-driving every
   non-connected entry on a flagged URL. A peer URL lands in wedgedNodeUrls
   because some db is wedged, but a sibling db may be healthy or still inside its
   initial-connect grace period — re-checking keeps a just-created entry or a
   not-desired entry from being force-reconnected and interrupting an in-flight
   dial.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Add DESIGN.md invariant #10 capturing the three reconnect-recovery drivers
(close-handler/forceReconnect retry, receive watchdog, main-thread wedge
reconcile), the createWebSocket-rejection trap they all shared, and the two
backstop subtleties (never-opened entry uses connected!==true + createdAt; the
reconcile must forceReconnect only a reused connection on a per-db wedged entry).
Hard-won recovery-layering knowledge that wasn't previously captured.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
@kriszyp
kriszyp marked this pull request as draft June 24, 2026 01:52
@kriszyp
kriszyp removed the request for review from heskew June 24, 2026 01:52
@kriszyp
kriszyp marked this pull request as ready for review June 24, 2026 01:56
@gemini-code-assist

Copy link
Copy Markdown

Warning

You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again!

@kriszyp

kriszyp commented Jun 24, 2026

Copy link
Copy Markdown
Member Author

Dependency note: the reconnect-recovery backstop this PR arms runs behind shouldReplicateFromNode (isDesired gate in findWedgedNodeUrls/reconcileWorkers). On a repeatedly-upgraded node that predicate's self-record point-lookup can be #352-undecodable → returns falsy → the backstop never sees the wedged peer. So this backstop is effectively a no-op for the exact field wedge it targets until #474 (decode-resilient self-record gate) lands. Recommend sequencing/landing #474 with or before this. CDP-confirmed on the field cluster: stuck connection had isFinished:true, intentionallyUnsubscribed:true, retries:0, with the self-record decode warning present.

kriszyp added a commit that referenced this pull request Jun 25, 2026
)

Adds the deferred third recovery layer from #466 / PR #467. While a receive
leg is paused for back-pressure the byte-silence receiveWatchdog is stopped
(ws.pause() freezes bytesRead) and the active sendPing is exempt, so a leg
that dies mid-pause — e.g. a system base copy stalled at ~100% back-pressure
whose peer restarted — had no recovery driver and could wedge connected:false
forever, removing the #424 forceReconnect path for exactly that case.

A pause-stall watchdog (createPauseStallWatchdog, a thin wrapper over the
existing createReceiveWatchdog stall-timer) now guards the paused window,
keyed on a local consumerProgress counter that advances on signals surviving
ws.pause(): onCommit (apply loop drained a queued batch) and blob-stream
drains. Armed on pause, stopped on resume; exactly one of {receiveWatchdog,
pauseStallWatchdog} is armed at a time. It fires forceReconnect only after a
sustained window of ZERO consumer progress, so a pause that is legitimately
making progress re-arms every window and never trips.

Cross-model review (Codex + agy): also makes resetPingTimer pause-aware so a
frame handler can't re-arm the byte watchdog while paused; threshold defaults
to max(PING_TIMEOUT*2, blobTimeout*2) and the slow-single-operation residual
(benign: re-streams from the durable cursor) is documented.

Tests: unitTests/replication/pauseStallWatchdog.test.mjs. Replication unit
suite green (187 passing). Single-repo (no core change).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
kriszyp added a commit that referenced this pull request Jun 25, 2026
…tion delayed close (#471) (#475)

* fix(replication): don't finish a still-desired peer on empty-subscription delayed close (#471)

In replicateOverWS, when the wire subscription went 0-length, scheduleClose fired
close(1008, ..., intentional=true) unconditionally. That marks the NodeReplicationConnection
isFinished/intentionallyUnsubscribed, emits 'finished' (removing it from the worker connections
map), and never reconnects. But after a base-copy resync a still-desired peer (its
nodeSubscriptions still populated) can briefly map to a 0-length wire subscription, which
permanently wedged it connected:false with no reconnect (live 4-node preprod).

New pure predicate shouldFinishEmptySubscriptionClose(connection) gates intentional on whether
the peer is genuinely no longer desired: a genuine unsubscribe sets nodeSubscriptions to []
(replicator.assignReplicationSource subscribe([],false) and unsubscribeFromNode->unsubscribe()),
so empty/absent => finish (idle cleanup unchanged); populated => not finished, so the close
handler reconnects and the connection self-heals.

Sibling to harper-pro#466/#467 (outbound reconnect-recovery) and #420 (open-but-idle wedge).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>

* fix(replication): discriminate empty-subscription close on database presence, not nodeSubscriptions

Codex cross-model review: the wire subscription is a 1:1 map of connection.nodeSubscriptions, so at
close time that array is empty in BOTH the genuine-unsubscribe and spurious-empty cases — gating on
it was a no-op for the real path. The only two empty-subscribe([]) sources are database removal
(assignReplicationSource — database gone, genuinely terminal) and subscribeToNode's
nodes.filter(shouldReplicateFromNode) collapsing to [] while the database is still present (spurious,
e.g. the #470 self-gate misread for a still-desired peer). So discriminate on database presence:
finish only when getDatabases()[databaseName] is gone; an empty subscription while the database is
still present stays retryable and self-heals. Preserves the prior behavior exactly for the
database-absent case.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

1 participant