tidaldb/docs/planning/milestone-11/phase-6.md
jx12n 3bfde53b90 feat(m11): data-plane sharding × replication (m11p6 L0-L2)
ClusterNode hosts a BTreeMap<ShardId, Arc<ShardReplica>>: writes hash-route
to the owning shard leader, reads scatter over shard groups. In-group
shard==region preserved so the engine and tidal-net are untouched; S=1 stays
byte-for-byte (today's cluster is a 1-shard × RF=N group). Topology grows
shard-group awareness; membership, election, forward, reseed, and join_boot
thread ShardId through.

Proven by an in-process 2×2 RF=2 gRPC test plus S=1 parity, incl. tier-3
real-OS-process failover. clippy/fmt clean.
2026-06-12 23:06:41 -06:00

15 KiB
Raw Blame History

m11p6 — Sharding × Replication + Rebalancing (IN PROGRESS)

Phase spec and exit gate: docs/roadmap-to-cluster.md §4/m11p6. Closes the ROADMAP write-scaling gap ("EITHER replicated (1 leader for everything) OR sharded (no replication)"). Predecessors: p1 (ship queue/perf), p2 (one log), p3 (quorum/commit-index), p4 (election/fencing), p5 (membership/discovery/elasticity).

Goal: writes scale horizontally without giving up replication. The "replicated XOR sharded" split ends: S shard groups, each a replication group at RF (default 3) with its own WAL/relay/commit-index and its own elected leader, leaders balanced across nodes; regions demote to placement labels; any gateway hash-routes to the shard leader (write) or nearest replica (read); rebalancing = snapshot + stream + fenced cutover.

The load-bearing realization (why this is tractable, not a rewrite)

A single TidalDb is already a complete, self-contained shard (scale-arch spec §1; single-process mode already runs N real TidalDb in one process). And today's cluster is exactly "1 shard × RF=N" — every region replicates one log under one leader. So m11p6 = run today's whole single-shard cluster machinery N times, once per data-shard group, where each group spans only its replica nodes, in its own data subdir + gRPC port. The data-shard number S is purely a gateway-routing + directory + port concern; inside a group, "shard id == region id" is preserved unchanged.

Consequences:

  • Engine replication (relay/commit/election/ship/receiver/membership): NO semantic change — each runs per-group exactly as today's whole cluster.
  • tidal-net: NO change — each group owns its own GrpcTransport bound to a distinct port, peering only with that group's replica nodes.
  • S=1 is byte-for-byte today's behavior — every existing tier-3 suite (quorum, membership, election, region, runbook, chaos) runs in S=1 mode and stays green by construction.
  • The work concentrates in tidal-server: decompose the per-shard machinery out of RegionClusterState into ShardReplica; the node hosts BTreeMap<ShardId, Arc<ShardReplica>>; add hash-routing to the right group's leader. Plus topology schema, route unification, rebalancing, stress/harness, the exit gate, docs.

Design (as adopted)

1. Two types: process node + shard replica

As-built note (decided during L1): the §1-era naming below was superseded. The process node is a NEW type ClusterNode (not "keep the name RegionClusterState"); the former RegionClusterState was renamed to ShardReplica (it already WAS one group's machinery). The responsibilities are exactly as described — only the node's type name changed from the planned RegionClusterState to ClusterNode.

  • ClusterNode (NEW process node) = the process node and the axum State. Owns the node identity (region/region_name), the gateway ShardRouter + node_http map, the shared forward client, and groups: BTreeMap<ShardId, Arc<ShardReplica>> (the groups this node hosts) plus a placement map of every group (to forward writes for groups it does NOT host). The HTTP handlers stay free functions over State<Arc<ClusterNode>>; each hash-routes the target group and either calls a local ShardReplica method or forwards to that group's leader.
  • ShardReplica (the renamed RegionClusterState) = one shard group's full replication machinery, essentially the per-shard fields + methods of the old RegionClusterState: db: Arc<TidalDb>, ship_feed, ship_queue, commit + commit_watch bridge, election_runtime/election_store/election_boot, membership: MembershipView + membership_store, stream_baseline, snapshot source + reseed_marker_store, leader: RwLock<Option<RegionId>>, activation_prev, partitioned, deferred_retires, ack_default, quorum_timeout, and its own GrpcTransport bound to this group's port. Construction = today's RegionClusterState::new body, parameterized by (shard, replica nodes, bootstrap leader, group data subdir, group grpc port).

2. Shard identity, ports, and data dirs

  • In-group identity is region-id based (preserved from today). ShardReplica for data-shard S on node N opens TidalDb with NodeConfig{ shard_id: ShardId(N.region_id), peer_shards: <other group-S nodes' region ids>, .. }. Two groups on one node use the same shard_id value but live in separate data dirs + separate transports, so nothing collides. The data-shard S never enters the replication wire.
  • gRPC ports: each replica entry may set an explicit grpc_addr/grpc_bind per (node, shard); otherwise derive node.base_port + S. Legacy S=0 → offset 0 → the region's declared address verbatim.
  • Data dirs: explicit shards:<data_dir>/shard-{S:05}/ per hosted group. Legacy (synthesized 1-shard) → <data_dir> verbatim (existing on-disk clusters restart unchanged).
  • /metrics: the node owns ONE listener on metrics_addr; the engine's per-TidalDb metrics HTTP server is suppressed in cluster mode. The node renders each group's cluster_metrics with a shard="S" label and the per-shard engine metrics likewise.

3. Topology schema (backward compatible)

TopologySpec gains optional shards: Option<Vec<ShardSpec>> and RegionSpec gains optional zone: Option<String> (placement/read-affinity label), both #[serde(default)].

nodes:                       # `regions:` still accepted as the node list
  - { name: us-east, grpc_addr: ..., http_addr: ..., zone: az-a }
shards:                      # NEW — optional
  - id: 0
    leader: us-east          # term-0 / preferred leader (balances placement)
    replicas:                # nodes hosting this group's RF replicas
      - { node: us-east }    # grpc_addr/grpc_bind optional → derive base+S
      - { node: eu-west }
      - { node: ap-south }
  - { id: 1, leader: eu-west, replicas: [...] }
  - { id: 2, leader: ap-south, replicas: [...] }

Absent shards: ⇒ legacy synthesis: one ShardSpec{ id: 0, leader: topology.leader, replicas: <every region, grpc verbatim> }. This is literally today's "1 shard × RF=all-regions," so the model unifies cleanly.

Validation: shard ids dense & unique; each leader/replica.node names a declared region; RF ≥ 1; a node's shards do not collide on derived ports.

4. Routing (the new gateway layer)

Any node, any write to entity E:

  1. S = ShardRouter.route(E) (the engine's FNV-1a hash; the same router the /sharded/* path already uses).
  2. If this node hosts group S and leads it → apply locally on ShardReplica[S] (the existing stage→complete→ship→await-quorum path).
  3. If this node hosts group S but follows → forward to group S's leader (resolved from ShardReplica[S].leader), to the leader node's http_addr.
  4. If this node does not host group S → forward to any replica node of S (from topology); that node routes to the leader. Term fencing + the NotLeader retry bound the forward chain.

Reads route to a local replica of S if present, else nearest (zone) / any replica. Scatter-gather merges over shard GROUPS: one replica per group, totals summed (entity-sharded groups are disjoint — no dedup-by-max), the existing degraded/unavailable_shards/deadline semantics retained.

Route unification: /items, /embeddings, /signals become the shard- routed surface (they hash-route instead of forwarding to a single global leader); /sharded/* is retained as an alias of the same path. x-tidal-ack, x-tidal-seq, quorum await, and NotLeader/QuorumTimeout all become per-shard; NotLeader names the shard and its leader.

5. Leadership balance

ShardSpec.leader is each group's term-0 / preferred leader; the operator (and the test harness) set shard i's leader to node i for balanced placement. Each group boot-classifies its own leader exactly like today's topology.leader, but per group. On failover, the group's election picks the surviving max-applied voter (m11p4 protocol, unchanged); only the dead node's groups elect. An optional leadership-transfer-back-to-preferred is the rebalance path, not the availability mechanism.

6. Rebalancing (operator-triggered)

Shard move / replica change reuses m11p5 verbatim, per group: a conf-change on the group's own log adds a Learner replica (a new node, or an existing node gaining a replica of group S), which catches up via FetchSnapshot + StreamSegments (per-group, already per-instance once §1 lands), auto-promotes Learner→Voter, then the operator may transfer leadership; removing a replica is the m11p5 fenced removal on that group's log. Endpoints: POST /cluster/shards/{id}/replicas (add/remove) and POST /cluster/shards/{id}/transfer. Per the roadmap, automatic rate-limited rebalancing is explicitly "later" — m11p6 ships the operator verbs.

Exit gate (from the roadmap)

  • 3 shards × RF=3: ≥5,000 quorum signals/s (≥2.5× single-shard p3).
  • Kill any node → only its shard-leaderships move (<10 s), reads never stop.
  • tidal-stress sharded-vs-replicated comparison collapses into one path.

Layer plan (each keeps the workspace green)

  • L0 — this doc + topology schema (shards:/zone: + validation + legacy synthesis). S=1 unchanged.
  • L1 — extract ShardReplica; node hosts a one-entry map; suppress engine metrics server in cluster mode + node aggregates. S=1 byte-for-byte.
  • L2 — multi-shard construction + gateway hash-routing + per-shard forward + route unification + scatter-gather over groups.
  • L3 — operator rebalancing verbs.
  • L4 — tidal-stress path collapse + nodes×shards tier-3 harness + cluster_sharding.rs exit gate (run for real).
  • L5 — docs (runbook/monitoring/roadmap/CHANGELOG/k8s/spec) + memory.

Status

  • L0 topology schema (shards:/zone: + resolver + validation; legacy synthesis) + this doc — 20 topology tests green incl. 5 new shard tests.

  • L1 RegionClusterStateShardReplica rename + group-parameterized ShardReplica::new(.., group, enable_metrics) (leader/peers/voters/data- dir/metrics from the resolved group; in-group identity stays region-based). S=1 byte-for-byte green (116 lib + cluster_region 11 + cluster_routes 6 + cluster_grpc 2). NB: the ClusterNode process wrapper + multi-shard hosting moved to L2 (inseparable from the routing/handler rewiring).

  • L2 DONE + green (the sharding × replication data plane). ClusterNode (the axum State + process handle) hosts BTreeMap<ShardId, Arc<ShardReplica>>, built by ClusterNode::new(topology, region, schema, profiles, data_dir, hlc) (resolves groups; opens one ShardReplica per hosted group — own data subdir shard-<id>/ + derived/explicit gRPC port for S>1; node data dir verbatim + metrics owner for S=1). ServeState for ClusterNode starts every group's election driver and shuts each down via Arc::try_unwrap. Entity writes (/items,/embeddings,/signals,/hardnegs) hash-route via route_entity → local replica (existing leader/forward path) or forward_to_group_node (remote group). Reads (/feed,/search) scatter in-process over hosted groups and merge (sum totals, score-sort — S=1 unchanged). /cluster/status/local gains a per-shard shards[] array; region_health aggregates all groups. Verified: in-process 2 nodes × 2 shards × RF=2 test (region_sharded_writes_route_per_shard_and_reads_scatter) — writes route per-shard (A forwards shard-1 writes to B), each shard replicates to its follower, reads scatter over both groups and return all items from BOTH shards, over real gRPC. S=1 byte-for-byte: 116 lib + cluster_region 12 + cluster_routes 6 + cluster_grpc 2 + tier-3 cluster_multiproc 5 (real OS processes: replication <2s, leader-crash failover <10s). clippy -D + fmt clean. Deferred to L2-followups (tracked): per-shard admin verbs (?shard= on promote/heal/partition/etc.) and /sharded/*→unified-path aliasing — currently the admin verbs + /sharded/* target the first hosted group (correct for S=1; multi-shard admin needs the selector).

    Known L2 limitations — S>1 only, FULL-placement assumed (tracked for L4, NOT silently dropped): these are correct for S=1 and for the headline exit-gate shape (3×3 RF=3, where every node hosts every group), and the writes proven by the 2×2 test are unaffected — but a genuinely PARTIAL placement (a node hosting a strict subset of groups) hits them:

    1. Unified reads are local-only. /feed///search scatter over the groups THIS node hosts (hosted_dbs); a node that does not host every group returns a strict-subset corpus with a 200 and no degraded flag. Complete only under full placement; cross-node read fan-out (reuse the /sharded/* HttpShardContext) is the L4 item. Use /sharded/* for partial placements.
    2. No cross-shard diversity re-rank. The scatter merge keeps each group's own diversity pass then score-merges; it does not re-run diversity ACROSS groups, so an S>1 /feed can be less diverse than S=1. Summing totals over disjoint groups is exact for cardinality; a cross-shard MMR re-rank (engine entry point) is the L4 follow-up.
    3. S>1 observability is one-group-deep. Only the metrics-owner group binds the engine /metrics listener; the other hosted groups' per-shard cluster_metrics (ship/relay/commit/quorum) are not rendered. The kill-node gate uses /cluster/status/local (and now /cluster/status) shards[], which IS per-group; node-level per-shard /metrics rendering (design §2) is the L4/p8 item.
  • L3 rebalancing verbs (operator shard move; reuse m11p5 per-group)

  • L4 tidal-stress path collapse + nodes×shards tier-3 harness + cluster_sharding.rs exit gate (3×3 kill-node + per-shard ledger; run for real). NB: the ≥5,000/s throughput sub-gate is a Ref-A line item — k3s access has blocked the Ref-A runs since p1 (same standing caveat).

  • L5 docs (runbook/monitoring/roadmap/CHANGELOG/k8s/spec) + memory

Exit-gate evidence

Data-plane correctness (local, in-process, real gRPC) — DONE. The 2×2 test proves the replicated-XOR-sharded split is gone: a single gateway routes writes to the owning shard's leader (forwarding cross-node), each shard group replicates independently, and corpus reads merge across groups. The per-shard election + quorum + commit machinery is the unchanged m11p4/m11p3 code instantiated per group, and S=1 failover is re-proven over real OS processes (cluster_multiproc).

Remaining for the headline exit gate: the 3 shards × RF=3 tier-3 run over real OS processes (kill any node → only its shard-leaderships move <10 s, reads never stop; per-shard zero-acked-loss ledger) and the ≥5,000/s tidal-stress figure (Ref-A-pending). These need the L4 harness extension (per-(node,shard) ports + shards: emission) + the rebalance verbs (L3).