Repository navigation
fix(replication): reconcile-level fallback for connected:true/Receiving copy-stall wedge - #463
Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces a defense-in-depth mechanism to detect and recover from stalled replication subscriptions that remain connected but fail to make progress. It adds periodic checks to identify these stalled nodes using shared status buffers and triggers a forced reconnection when necessary. Comprehensive unit tests have also been added to cover these new detection paths. The reviewer identified two potential issues regarding empty or undefined nodes arrays that could lead to runtime TypeErrors or malformed request objects, and provided actionable code suggestions to defensively guard against them.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| if ( | ||
| entry.connected === false || | ||
| !entry.worker || | ||
| !httpWorkers.includes(entry.worker) || | ||
| !isDesired(entry.nodes?.[0], database) | ||
| ) | ||
| continue; |
There was a problem hiding this comment.
If entry.nodes is empty or undefined, entry.nodes?.[0] will evaluate to undefined. Passing undefined to isDesired (which maps to shouldReplicateFromNode) will cause a runtime TypeError because shouldReplicateFromNode accesses properties like node.isLeader and node.subscriptions without optional chaining. Since reconcileWorkers runs in a setInterval on the main thread, an unhandled exception here would crash the entire Node.js process.
We should defensively guard against !entry.nodes?.[0] before calling isDesired to prevent any potential runtime crashes.
if (
entry.connected === false ||
!entry.worker ||
!httpWorkers.includes(entry.worker) ||
!entry.nodes?.[0] ||
!isDesired(entry.nodes[0], database)
)| const entry = entries?.get(databaseName); | ||
| const worker = entry?.worker; | ||
| const nodes = entry?.nodes; | ||
| if (!entry || !worker || !nodes) continue; |
There was a problem hiding this comment.
If entry.nodes is empty, nodes will be [] and nodes[0] will be undefined. Spreading undefined (i.e., ...nodes[0]) in the request object creation won't throw a runtime error in modern JavaScript, but it will result in a request object missing crucial node properties (like url and name). This would cause the force-reconnect-node message to be malformed and fail to reconnect.
We should defensively guard against !nodes[0] along with the other checks.
| const entry = entries?.get(databaseName); | |
| const worker = entry?.worker; | |
| const nodes = entry?.nodes; | |
| if (!entry || !worker || !nodes) continue; | |
| const entry = entries?.get(databaseName); | |
| const worker = entry?.worker; | |
| const nodes = entry?.nodes; | |
| if (!entry || !worker || !nodes || !nodes[0]) continue; |
|
Reviewed; no blockers found. |
|
Warning You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again! |
…ng copy-stall wedge Add findStalledReceivingNodeUrls as a defense-in-depth companion to findWedgedNodeUrls: the reconcile only re-drives connected:false entries, so a base copy parked connected:true with the received-version watermark frozen (ping-alive, harper-pro#453) is invisible to it. The worker-local copy-progress watchdog (#454) is the primary recovery; this main-thread net only matters if that watchdog itself fails. Progress is measured by the per-record RECEIVED_TIME watermark read from the process-shared status buffer (advances during a healthy copy even while version is suppressed), so a slow-but-progressing copy is never torn down. Recovery is forceReconnectToNode (a re-subscribe is a no-op for a still-connected entry). Threshold is 15min, far longer than the worker watchdogs; re-drives only fire once apply progress resumes since the last forced reconnect, so a kick that changed nothing (caught-up/cosmetic Receiving) does not churn the connection. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
…DESIGN.md Note the connected:false (findWedgedNodeUrls) vs connected:true/Receiving (findStalledReceivingNodeUrls) recovery duality, the main-thread RECEIVED_TIME progress signal, and why recovery is forceReconnect not re-subscribe. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2576fd0 to
a81697f
Compare
Summary
Adds a conservative, reconcile-level fallback that recovers a replication subscription stuck
connected:true/lastReceivedStatus: "Receiving"/ received-version frozen with no apply progress — the ping-alive base-copy wedge from harper-pro#453.findStalledReceivingNodeUrls(new, pure/unit-tested): companion tofindWedgedNodeUrls. The existing reconcile only re-drivesconnected:falseentries (the Replication wedges permanently after simultaneous cluster restart (reconciler skips open-but-idle sockets) → blocks replicated deploys #420/fix(replication): recover open-but-idle wedged subscriptions via watchdog-driven reconnect #424/Replication: ~9/11 outgoing peers show connected:false in connectionReplicationMap after restart despite being reachable #289 wedge family); a copy that parksconnected:truewith the received-version watermark suppressed is invisible to it, and the byte-level receive watchdog can't see it either because keepalive pings keep the socket'sbytesReadadvancing.forceReconnectToNode(new, worker side): drops + reconnects the cached connection. A re-subscribe is a no-op for aconnected:trueentry (getSubscriptionConnectionreturns the still-connected connection unchanged viaisReusableConnection), so recovery must force a reconnect, mirroring what the worker-local watchdog does.reconcileWorkersgains a third, independent branch that force-reconnects exactly the stalled(node, database)entries, staggered like the existing wedged path.Purpose
Defense-in-depth net deliberately deferred from #454. The worker-local copy-progress watchdog (PR #454, branch
kris/copy-progress-watchdog) is the primary recovery for this wedge and handles it functionally. A Gemini cross-model review of #454 recommended a reconcile-level fallback so that if the local watchdog itself fails (event-loop block, software bug, misconfiguration) the wedge can still never be permanent. This is that fallback. Related: #453 (incident), #454/#460 (watchdog + follow-ups), #420/#424/#289 (the connected:false family this extends).How it stays safe (where to look)
The hard constraint is never tear down a healthy slow-but-progressing copy. The signal that makes this safe:
replicationConnection.ts,if (!inCopyMode)guard), butRECEIVED_TIME(the per-record apply watermark) andRECEIVING_STATUS = Receivingare written together on every applied copy record. So a copy that is making progress — however slowly — keeps bumpingRECEIVED_TIME, and only a genuinely frozen one trips the check.status === ReceivingimpliesRECEIVED_TIME > 0(set in the same place), so a never-connected/zeroed buffer never reads as stalled.receiveStallReconnectAtthrottles re-drives to once per threshold.auditStore.getUserSharedBuffer, keyed[database, nodeName]) thatcluster_statusalready reports from.Attention for the reviewer (open items / deliberate trade-offs)
RECEIVED_TIMEstaleness, not literallyversion === 0. During a re-copy the frozen version sits at the prior watermark (not 0), so a version-equality gate would miss it;RECEIVED_TIMEstaleness covers fresh-clone, re-copy, and a stalled normal mid-batch apply uniformly, all with the same safe recovery. Theversion === 0phrasing in the incident is a symptom of being mid-copy, not the discriminator.ws.pause()) with zero applied records for the entire 15-min window would be reconnected. Consequence is bounded — a reconnect resumes from the persisted copy cursor (no data loss) — and the long threshold makes a 15-min continuous pause the only trigger. I could not find a tighter main-thread-visible signal to distinguish paused-but-healthy from wedged (the back-pressure ratio in the shared buffer is send-side); flagging the choice rather than hiding it.findStalledReceivingNodeUrls.test.mjs) andforceReconnectitself is already covered (Replication wedges permanently after simultaneous cluster restart (reconciler skips open-but-idle sockets) → blocks replicated deploys #420). A true e2e test needs a watchdog-disable hook + threshold injection + simulated long stall; that's the larger, separately-tunable surface this was scoped out of for v5.1 — reasonable as a follow-up.Cross-model review — deliberate/deferred items
Codex + Gemini both reviewed; both rated correctness/testing strongly. Two non-blocking items I chose not to action here:
subscriptionManager↔replicationConnection(constants). The cycle already exists (replicationConnection importsdisconnectedFromNode/connectedToNode/ensureNodefrom subscriptionManager) andclusterStatus.tsalready imports these same constants from replicationConnection. It's runtime-safe (constants/functions only, none touched at top-level load) and Harper's conventions explicitly tolerate cycles. Extracting the position constants to a shared module is a reasonable cleanup but would touch unrelated files — left for a follow-up rather than widening this PR.getReceiveStatus(Map lookup +Object.valuesscan + nativegetUserSharedBufferevery 5s per active connection). Gemini suggested caching theFloat64Arrayon the entry. Deferred: the cost is sub-microsecond and bounded to connected/desired entries, and caching risks reading a stale(db, node)buffer ifnodes[0]changes on failover — not worth the correctness footgun for the micro-saving.🤖 Generated by Claude (Opus 4.8). Diff and commit history are the source of truth.