You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Workstream W4 of #430 · Foundation 2 / keystone · backbone for transitive correctness & topology scale
Summary
Make per-origin transaction-log positions the primary resume-cursor mechanism and carry them transitively across hops. This removes the root cause of transitive-replication flooding, enables independent per-origin resume, and makes per-origin lag observable. Per-origin is a property of the cursor, not the socket — see Connection scaling; we explicitly do not open a socket per origin. RocksDB already partitions its transaction log per origin (RocksTransactionLogStorenodeLogs[]/logById), so this is a cursor-layer + protocol convergence — not an on-disk format redesign. The advanced per-origin/transitive behavior is RocksDB-only; LMDB keeps its existing basic replication and degrades gracefully (see Migration / mixed-engine).
The per-source cursors already exist (nodes[] in the Symbol.for('seq') entry) but are a proxied-only secondary signal today.
A subscription resumes from a single scalar cursor per peer; a multi-origin failover forces a full re-scan (Use separate WS connections per originating-node subscription (avoid full audit scan resets on fail-over) #193). The substrate for the fix is already present — RocksDB stores per-origin logs (nodeLogs[]/logById, keyed by the origin nodeId, not the delivering peer) and getRange already accepts a per-log startByLog floor map (replication passes it today with a single-origin map, replicationConnection.ts:2443). The gap is that the per-origin cursor isn't yet the primary mechanism and isn't carried transitively.
Design direction
Promote per-origin cursors to the primary mechanism (not proxied-secondary) — a per-origin cursor vector (startByLog floor map) replaces the single scalar cursor.
Multiplex per-origin streams over the one-connection-per-peer (refines Use separate WS connections per originating-node subscription (avoid full audit scan resets on fail-over) #193) — independent per-origin resume and no cross-origin reset on failover, carried as the startByLog vector on a single merged stream (the merge + per-log-start mechanism already exists; today it's used with a one-entry map). No socket per origin. Optional: a small fixed socket pool per peer for apply parallelism on a hot peer — never one-per-origin.
LMDB is deprecated but must keep basic replication through the migration window (4.7 → 5.2-on-LMDB → 5.2-on-RocksDB), so a cluster can be converted node-by-node and an LMDB node can still participate while it's being flipped. Scope rules:
LMDB stays on the legacy shared-cursor path (full-mesh / simple topologies, incremental audit streaming + base copy). We do not advance it — no LMDB per-origin sub-log, no LMDB log-format redesign.
The connection model stays O(peers), not O(peers×origins). Per-origin independence lives in the cursor, not the socket: one connection per peer multiplexes many per-origin streams, each with its own startByLog floor. Socket-per-origin (#193's literal wording) is rejected — it explodes in proxied/hub topologies (a hub relaying M origins to S spokes → O(S×M) sockets), which is exactly where adaptive fan-out (#218) and transitive replication drive the topology. In a direct mesh, origins == peers, so there's no difference there; the explosion is specific to proxying. Residual costs are bounded metadata/storage, not connections:
Cursor-vector size = O(origins a given peer relays to you); sparse (a direct subscription is 1:1 → zero overhead). Send the full vector on resume, deltas on incremental updates.
Per-origin log count on a node = O(origins whose data the node carries) — bounded by residency/subscription scope in a sharded cluster; a pre-existing RocksDB property, not introduced here. (Only a fully-replicated very-large cluster would strain it — there you shard.)
Tradeoff vs separate sockets — head-of-line blocking and per-origin thread assignment — handled with per-channel credit-based flow control and, for a hot peer, a small fixed socket pool (not per-origin).
Out of scope: large fan-in tree aggregation (thousands of leaves → a central server, where even O(branches) per-origin tracking on the relay is bounded by re-origination, not by per-leaf logs) is handled by the separate re-origin relay mode — W12 / #444. W4 covers mesh/transitive only and preserves end-to-end origin identity; the two modes share the delivery-vs-conflict identity split but are otherwise independent.
Dependencies
Best sequenced after W1/W2 land (clean connection registry + cursor correctness). This is the keystone — schedule deliberately.
Effort / risk
XL / high. The backbone of "scale to different topologies."
Acceptance criteria
A 3-hop A→B→C topology does not re-stream already-applied tails on failover.
One origin's failover does not reset other origins' in-flight streams.
Workstream W4 of #430 · Foundation 2 / keystone · backbone for transitive correctness & topology scale
Summary
Make per-origin transaction-log positions the primary resume-cursor mechanism and carry them transitively across hops. This removes the root cause of transitive-replication flooding, enables independent per-origin resume, and makes per-origin lag observable. Per-origin is a property of the cursor, not the socket — see Connection scaling; we explicitly do not open a socket per origin. RocksDB already partitions its transaction log per origin (
RocksTransactionLogStorenodeLogs[]/logById), so this is a cursor-layer + protocol convergence — not an on-disk format redesign. The advanced per-origin/transitive behavior is RocksDB-only; LMDB keeps its existing basic replication and degrades gracefully (see Migration / mixed-engine).Root cause / current state
startTime, because the relayed cursor is the proxy B's local-time position, not origin A's applied position (replicationConnection.ts~3235). So B re-streams a wide already-applied tail to C on every resume; the duplicates arrive buried below the head where the head-tie fast-skip can't catch them, and historically drove the out-of-order resequencing walk to its depth cap → event-loop starvation → ping-timeouts → more failovers → more re-streams (a self-amplifying loop, also the [sharding] FATAL ERROR: NewSpace:: EnsureCurrentCapacity Allocation failed - JavaScript heap out of memory (verify on current build) #197 OOM vector).nodes[]in theSymbol.for('seq')entry) but are a proxied-only secondary signal today.nodeLogs[]/logById, keyed by the origin nodeId, not the delivering peer) andgetRangealready accepts a per-logstartByLogfloor map (replication passes it today with a single-origin map,replicationConnection.ts:2443). The gap is that the per-origin cursor isn't yet the primary mechanism and isn't carried transitively.Design direction
startByLogfloor map) replaces the single scalar cursor.startByLogvector on a single merged stream (the merge + per-log-start mechanism already exists; today it's used with a one-entry map). No socket per origin. Optional: a small fixed socket pool per peer for apply parallelism on a hot peer — never one-per-origin.cluster_statuspositions (Per-originating-node transaction log partitioning (dependent on new transaction log format) #192) — feeds the W8 lag gauges.Scope
nodes[]per-source cursors the primary resume mechanism (per-origin cursor vector viastartByLog)Retires / advances
Migration / mixed-engine (LMDB → RocksDB)
LMDB is deprecated but must keep basic replication through the migration window (4.7 → 5.2-on-LMDB → 5.2-on-RocksDB), so a cluster can be converted node-by-node and an LMDB node can still participate while it's being flipped. Scope rules:
Connection scaling (explicit non-goal: socket-per-origin)
The connection model stays O(peers), not O(peers×origins). Per-origin independence lives in the cursor, not the socket: one connection per peer multiplexes many per-origin streams, each with its own
startByLogfloor. Socket-per-origin (#193's literal wording) is rejected — it explodes in proxied/hub topologies (a hub relaying M origins to S spokes → O(S×M) sockets), which is exactly where adaptive fan-out (#218) and transitive replication drive the topology. In a direct mesh, origins == peers, so there's no difference there; the explosion is specific to proxying. Residual costs are bounded metadata/storage, not connections:Out of scope: large fan-in tree aggregation (thousands of leaves → a central server, where even O(branches) per-origin tracking on the relay is bounded by re-origination, not by per-leaf logs) is handled by the separate re-origin relay mode — W12 / #444. W4 covers mesh/transitive only and preserves end-to-end origin identity; the two modes share the delivery-vs-conflict identity split but are otherwise independent.
Dependencies
Best sequenced after W1/W2 land (clean connection registry + cursor correctness). This is the keystone — schedule deliberately.
Effort / risk
XL / high. The backbone of "scale to different topologies."
Acceptance criteria
🤖 Filed by Claude on behalf of Kris.