Skip to content

Replication W4: Per-origin transaction-log convergence (keystone) #434

Description

@kriszyp

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 (RocksTransactionLogStore nodeLogs[]/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

Design direction

  1. Promote per-origin cursors to the primary mechanism (not proxied-secondary) — a per-origin cursor vector (startByLog floor map) replaces the single scalar cursor.
  2. 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.
  3. Transitive per-origin position propagation — carry each origin's own monotonic position end-to-end so a relay never re-streams an applied tail (removes the Replication: transitive/proxied re-delivery floods peers with already-applied out-of-order writes (reduce volume; complements harper#1310) #399 cause; the current leading-dup fast-skip only absorbs the symptom).
  4. Per-origin cluster_status positions (Per-originating-node transaction log partitioning (dependent on new transaction log format) #192) — feeds the W8 lag gauges.

Scope

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 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.
  • Per-origin position/lag is reportable.

🤖 Filed by Claude on behalf of Kris.

Metadata

Metadata

Assignees

No one assigned

    Labels

    area:replicationReplication, cluster sync, peer connectionsenhancementNew feature or request

    Type

    Fields

    Priority

    P2

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions