diff --git a/hyperbytedb/src/adapters/cluster/sync_client.rs b/hyperbytedb/src/adapters/cluster/sync_client.rs index 2f92af1..8280d40 100644 --- a/hyperbytedb/src/adapters/cluster/sync_client.rs +++ b/hyperbytedb/src/adapters/cluster/sync_client.rs @@ -7,6 +7,7 @@ use crate::domain::cluster::membership::{NodeState, SharedMembership}; use crate::domain::cluster::sync::{ JoinRequest, JoinResponse, MetadataSnapshot, SyncManifest, WalSyncResponse, }; +use crate::domain::sharding::{ShardMap, ShardMapJson}; use crate::error::HyperbytedbError; use crate::ports::metadata::MetadataPort; use crate::ports::points_sink::PointsSinkPort; @@ -133,6 +134,7 @@ impl SyncClient { ); self.sync_metadata(&peer_addr).await?; + self.sync_shard_map(&peer_addr).await?; let updated_wal_seq = self.wal.last_sequence().await?; let applied = self.wal_catchup(&peer_addr, updated_wal_seq).await?; @@ -178,6 +180,23 @@ impl SyncClient { ); self.sync_metadata(&peer_addr).await?; + // Install the committed map before any region data movement so the + // joiner's map_version matches the cluster at Active. + // + // Best-effort: `/internal/shard/map` only exists on peers that have + // sharding enabled and run a build that serves it, so a 404 here says + // nothing about whether this node can catch up its WAL. Failing hard + // used to abort the join before `wal_catchup`, and after the retries + // were exhausted the node went Active having synced nothing at all. + // Skipping is safe: without a map there are no regions to move below, + // and the leader's own catch-up check gates placement onto this node. + if let Err(e) = self.sync_shard_map(&peer_addr).await { + tracing::warn!( + error = %e, + peer = %peer_addr, + "shard map install failed; continuing with WAL catch-up" + ); + } let mut applied = 0u64; if let Some(ref shard_map) = self.shard_map { @@ -328,6 +347,42 @@ impl SyncClient { Ok(()) } + /// Copy the peer's committed shard map so `map_version` matches before + /// region rows are transferred onto this joiner. + pub async fn sync_shard_map(&self, peer_addr: &str) -> Result<(), HyperbytedbError> { + let Some(shard_map) = self.shard_map.as_ref() else { + return Ok(()); + }; + let remote = fetch_shard_map(self.client.clone(), peer_addr).await?; + let version = remote.map_version; + + // Never move the map backwards. `map_version` counts local applies and + // the sync peer is whichever Active node came first out of a HashMap — + // not the leader, and not necessarily the most-applied node. Installing + // an older snapshot would drop ops Raft will never redeliver, because + // `replace_map` bypasses the state machine and leaves `last_applied` + // untouched; every later op would then fail StaleEpoch/UnknownRegion + // and the node would silently stop owning regions it holds data for. + let local = shard_map.snapshot().await?.map_version; + if !remote_map_is_newer(version, local) { + tracing::info!( + peer = %peer_addr, + remote_map_version = version, + local_map_version = local, + "peer shard map is not newer; keeping local map" + ); + return Ok(()); + } + + shard_map.replace_map(remote).await?; + tracing::info!( + peer = %peer_addr, + map_version = version, + "installed peer shard map before region movement" + ); + Ok(()) + } + async fn import_metadata_entry( &self, entry: &crate::domain::cluster::sync::MetadataEntry, @@ -487,6 +542,39 @@ impl SyncClient { } } +/// GET `/internal/shard/map` from `peer_addr` and return the committed snapshot. +/// Whether a peer's shard map may replace the local one. +/// +/// Strictly-newer only. Equal is a no-op, and older must be refused: a sync +/// peer is whichever Active node came first out of a `HashMap`, so it is not +/// necessarily the leader or the most-applied node. +#[must_use] +pub fn remote_map_is_newer(remote_version: u64, local_version: u64) -> bool { + remote_version > local_version +} + +pub async fn fetch_shard_map( + client: reqwest::Client, + peer_addr: &str, +) -> Result { + let url = format!("http://{peer_addr}/internal/shard/map"); + let resp = + client.get(&url).send().await.map_err(|e| { + HyperbytedbError::PeerUnreachable(format!("shard map request failed: {e}")) + })?; + if !resp.status().is_success() { + return Err(HyperbytedbError::SyncFailed(format!( + "shard map request failed: {}", + resp.status() + ))); + } + let json: ShardMapJson = resp + .json() + .await + .map_err(|e| HyperbytedbError::SyncFailed(format!("parse shard map: {e}")))?; + Ok(ShardMap::from(json)) +} + /// Fail closed when the peer advertises a higher WAL watermark but returns no /// readable entries — usually because the leader truncated past the gap. fn verify_catchup_progress( @@ -509,9 +597,19 @@ fn verify_catchup_progress( #[cfg(test)] mod tests { - use super::verify_catchup_progress; + use super::{remote_map_is_newer, verify_catchup_progress}; use crate::error::HyperbytedbError; + #[test] + fn remote_map_installs_only_when_strictly_newer() { + assert!(remote_map_is_newer(10, 0), "fresh joiner takes the map"); + assert!(remote_map_is_newer(10, 9)); + assert!(!remote_map_is_newer(10, 10), "equal is a no-op"); + // The case that silently strands a node: syncing from a less-applied + // peer would drop ops Raft never redelivers. + assert!(!remote_map_is_newer(8, 10)); + } + #[test] fn verify_catchup_progress_ok_when_caught_up() { verify_catchup_progress(5, 5, 0, 5).unwrap(); diff --git a/hyperbytedb/src/adapters/http/shard_handlers.rs b/hyperbytedb/src/adapters/http/shard_handlers.rs index 3f11735..05217cb 100644 --- a/hyperbytedb/src/adapters/http/shard_handlers.rs +++ b/hyperbytedb/src/adapters/http/shard_handlers.rs @@ -540,11 +540,15 @@ pub async fn handle_shard_transfer( Json(req): Json, ) -> impl IntoResponse { if req.stage { - // Pre-commit staging (split optimization): the range is not owned by - // the destination in the committed map yet, so region/epoch checks are - // skipped. Guard rails: sender must be a known member and the local - // measurement must exist; rows outside `[start, end)` are filtered - // during apply. + // Pre-commit staging (split optimization, and replica placement onto a + // joiner): the range is not owned by the destination in the committed + // map yet, so region/epoch checks are skipped. The only guard rail is + // that the sender must be a known member; rows outside `[start, end)` + // are filtered during apply. + // + // There is deliberately no local-measurement check. A joiner receiving + // a measurement it has never seen has no catalog row yet, and + // `apply_transfer_push` creates one via `prepare_batch_metadata`. let Some(membership) = state.membership.as_ref() else { return ( StatusCode::SERVICE_UNAVAILABLE, @@ -561,20 +565,6 @@ pub async fn handle_shard_transfer( .into_response(); } drop(m); - if state - .metadata - .get_measurement(&req.db, &req.rp, &req.measurement) - .await - .ok() - .flatten() - .is_none() - { - return ( - StatusCode::NOT_FOUND, - Json(serde_json::json!({"error": "measurement not found for staging"})), - ) - .into_response(); - } } else { let Some(ctx) = state.shard_routing.as_ref() else { return sharding_disabled(); @@ -710,19 +700,33 @@ pub async fn handle_shard_rehome( .into_response(); } - let outcome = match push_region_transfer_data( - peer_client, - &state.metadata, - &state.wal, - Some(&state.query_port), - state.node_id, - &ctx_key(&req), - region, - req.dest_primary, - state.max_points_per_request, - ) - .await - { + let outcome = match if req.stage { + crate::application::shard_transfer::stage_region_transfer_data( + peer_client, + &state.metadata, + &state.wal, + Some(&state.query_port), + state.node_id, + &ctx_key(&req), + region, + req.dest_primary, + state.max_points_per_request, + ) + .await + } else { + push_region_transfer_data( + peer_client, + &state.metadata, + &state.wal, + Some(&state.query_port), + state.node_id, + &ctx_key(&req), + region, + req.dest_primary, + state.max_points_per_request, + ) + .await + } { Ok(o) => o, Err(e) => { counter!("hyperbytedb_shard_transfer_failures_total").increment(1); @@ -734,18 +738,19 @@ pub async fn handle_shard_rehome( } }; - if let Err(e) = complete_region_transfer( - peer_client, - &state.metadata, - Some(&state.points_sink), - state.node_id, - &ctx_key(&req), - region, - req.dest_primary, - outcome.transfer_id, - req.drop_source, - ) - .await + if !req.stage + && let Err(e) = complete_region_transfer( + peer_client, + &state.metadata, + Some(&state.points_sink), + state.node_id, + &ctx_key(&req), + region, + req.dest_primary, + outcome.transfer_id, + req.drop_source, + ) + .await { counter!("hyperbytedb_shard_transfer_failures_total").increment(1); return ( diff --git a/hyperbytedb/src/adapters/sharding/rocksdb_shard_map.rs b/hyperbytedb/src/adapters/sharding/rocksdb_shard_map.rs index 1b7d91d..5867444 100644 --- a/hyperbytedb/src/adapters/sharding/rocksdb_shard_map.rs +++ b/hyperbytedb/src/adapters/sharding/rocksdb_shard_map.rs @@ -207,6 +207,47 @@ fn assemble_map(db: &DB) -> Result { Ok(map) } +/// Atomically persist the full map (join catch-up / snapshot install). +fn persist_full_map(db: &DB, map: &ShardMap) -> Result<(), HyperbytedbError> { + let existing = load_spaces(db)?; + let incoming: std::collections::HashSet<&MeasurementKey> = map.spaces.keys().collect(); + let mut batch = WriteBatch::default(); + batch.put( + META_KEY, + serde_json::to_vec(&PersistedShardMapMeta { + map_version: map.map_version, + next_region_id: map.next_region_id, + }) + .map_err(|e| { + HyperbytedbError::ShardMap(crate::error::ChainedError::with_context( + "shard map meta serialize", + e, + )) + })?, + ); + for space in &existing { + if !incoming.contains(&space.key) { + batch.delete(space_key(&space.key)); + } + } + for space in map.spaces.values() { + batch.put( + space_key(&space.key), + serde_json::to_vec(&PersistedShardSpace { + space: space.clone(), + }) + .map_err(|e| { + HyperbytedbError::ShardMap(crate::error::ChainedError::with_context( + "shard space serialize", + e, + )) + })?, + ); + } + db.write(batch) + .map_err(|e| HyperbytedbError::Storage(e.to_string().into())) +} + /// Atomically persist the global counters plus the one space an op touched. fn persist_space( db: &DB, @@ -326,6 +367,20 @@ impl ShardMapPort for RocksDbShardMap { Ok(map) } + async fn replace_map(&self, map: ShardMap) -> Result<(), HyperbytedbError> { + let _guard = self.apply_lock.lock().await; + for space in map.spaces.values() { + space + .validate() + .map_err(|e| HyperbytedbError::ShardMap(e.into()))?; + } + map.validate_global_region_ids() + .map_err(|e| HyperbytedbError::ShardMap(e.into()))?; + persist_full_map(&self.db, &map)?; + *self.cache.write() = Arc::new(map); + Ok(()) + } + async fn node_owns_measurement( &self, node_id: u64, @@ -392,6 +447,38 @@ mod tests { assert_eq!(map.next_region_id, 2); } + #[tokio::test] + async fn replace_map_installs_peer_snapshot() { + let dir = tempfile::tempdir().unwrap(); + let local = RocksDbShardMap::open(dir.path(), true).unwrap(); + local + .apply_op(ShardMapOp::BootstrapMeasurement { + key: MeasurementKey::new("db", "rp", "old"), + region: region(1, 0, u64::MAX), + }) + .await + .unwrap(); + + let incoming = crate::domain::sharding::ShardMap { + map_version: 4, + next_region_id: 3, + spaces: [( + MeasurementKey::new("db", "rp", "cpu"), + crate::domain::sharding::MeasurementShardSpace { + key: MeasurementKey::new("db", "rp", "cpu"), + regions: vec![region(2, 0, u64::MAX)], + }, + )] + .into_iter() + .collect(), + }; + local.replace_map(incoming).await.unwrap(); + let snap = local.snapshot().await.unwrap(); + assert_eq!(snap.map_version, 4); + assert!(snap.space("db", "rp", "old").is_none()); + assert!(snap.space("db", "rp", "cpu").is_some()); + } + #[tokio::test] async fn legacy_single_key_format_migrates_on_open() { let dir = tempfile::tempdir().unwrap(); diff --git a/hyperbytedb/src/application/cluster/bootstrap.rs b/hyperbytedb/src/application/cluster/bootstrap.rs index 9f90593..60214d8 100644 --- a/hyperbytedb/src/application/cluster/bootstrap.rs +++ b/hyperbytedb/src/application/cluster/bootstrap.rs @@ -98,6 +98,7 @@ impl ClusterBootstrap { wal: &Arc, points_sink: Option>, max_points_per_request: usize, + shard_map: Option>, ) -> anyhow::Result<()> { { let mut m = self.membership.write().await; @@ -107,7 +108,7 @@ impl ClusterBootstrap { tracing::info!("startup phase: syncing with cluster before accepting traffic"); - let sync_client = SyncClient::with_points_sink( + let mut sync_client = SyncClient::with_points_sink( config.node_id, config.cluster_addr.clone(), self.membership.clone(), @@ -117,6 +118,9 @@ impl ClusterBootstrap { max_points_per_request, self.peer_addrs.clone(), ); + if let Some(map) = shard_map { + sync_client = sync_client.with_shard_map(map); + } let dbs = metadata.list_databases().await?; let has_data = !dbs.is_empty(); diff --git a/hyperbytedb/src/application/runtime/mod.rs b/hyperbytedb/src/application/runtime/mod.rs index e2b90e2..5d77157 100644 --- a/hyperbytedb/src/application/runtime/mod.rs +++ b/hyperbytedb/src/application/runtime/mod.rs @@ -111,6 +111,9 @@ pub async fn serve(config: HyperbytedbConfig) -> anyhow::Result<()> { &wal_port, Some(sink_port), config.server.max_points_per_request, + rocks_shard_map + .as_ref() + .map(|m| m.clone() as Arc), ) .await?; Some( diff --git a/hyperbytedb/src/application/shard_routing.rs b/hyperbytedb/src/application/shard_routing.rs index b31f40a..f224e8d 100644 --- a/hyperbytedb/src/application/shard_routing.rs +++ b/hyperbytedb/src/application/shard_routing.rs @@ -274,6 +274,19 @@ pub async fn bootstrap_measurement_local( Ok(()) } +/// Effective replica count for a region: never more than live membership. +/// +/// A 1-member cluster therefore has RF=1 even when `configured` is 3. +/// RF is also never raised above the configured target. +#[must_use] +pub fn effective_replication_factor(configured: usize, active_members: usize) -> usize { + let target = configured.max(1); + if active_members == 0 { + return 0; + } + target.min(active_members) +} + /// Pick the peer set for a brand-new measurement's first region. /// /// Candidates are ordered by `hash(measurement, node_id)` so different @@ -282,6 +295,9 @@ pub async fn bootstrap_measurement_local( /// region to the same RF nodes — a deterministic cluster-wide hotspot. Pure /// function of the key + membership, so every coordinator computing the op /// agrees on the same peer set and primary. +/// +/// Peer count is [`effective_replication_factor`]: n=1 yields `[self]`, +/// not an empty set truncated against a configured RF of 2 or 3. async fn select_bootstrap_peers( ctx: &ShardRoutingContext, db: &str, @@ -298,6 +314,7 @@ async fn select_bootstrap_peers( drop(membership); peers.sort_unstable(); peers.dedup(); + let member_count = peers.len(); // Spread placement: deterministic per-(measurement, node) hash order. use std::hash::{Hash, Hasher}; peers.sort_by_key(|node_id| { @@ -308,7 +325,10 @@ async fn select_bootstrap_peers( node_id.hash(&mut h); h.finish() }); - peers.truncate(ctx.config.replication_factor.max(1)); + peers.truncate(effective_replication_factor( + ctx.config.replication_factor, + member_count, + )); Ok(peers) } @@ -996,4 +1016,32 @@ mod scatter_tests { "primaries must spread across nodes, got {primaries:?}" ); } + + #[test] + fn effective_rf_never_exceeds_membership() { + assert_eq!(effective_replication_factor(3, 1), 1); + assert_eq!(effective_replication_factor(3, 2), 2); + assert_eq!(effective_replication_factor(2, 4), 2); + assert_eq!(effective_replication_factor(0, 3), 1); + assert_eq!(effective_replication_factor(3, 0), 0); + } + + #[tokio::test] + async fn bootstrap_n1_peers_are_self() { + let membership = membership_with(&[(1, "127.0.0.1:1")]); + let config = ShardingConfig { + replication_factor: 3, + ..Default::default() + }; + let ctx = test_ctx(membership, config, 1); + + let op = build_bootstrap_op(&ctx, "db", "autogen", "cpu") + .await + .unwrap(); + let ShardMapOp::BootstrapMeasurement { region, .. } = op else { + panic!("expected bootstrap op"); + }; + assert_eq!(region.peers, vec![1], "n=1 must own the region"); + assert_eq!(region.primary, 1); + } } diff --git a/hyperbytedb/src/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index 4c92d7e..633da8f 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -9,13 +9,16 @@ use crate::adapters::cluster::peer_client::PeerClient; use crate::adapters::cluster::raft::HyperbytedbRaft; use crate::adapters::sharding::rocksdb_shard_map::RocksDbShardMap; use crate::application::shard_peer_resolution::is_active_peer; +use crate::application::shard_routing::effective_replication_factor; use crate::application::shard_transfer::{ complete_region_transfer, push_region_transfer_data, run_region_transfer, stage_region_transfer_data, }; use crate::config::ShardingConfig; use crate::domain::cluster::membership::{NodeState, SharedMembership}; -use crate::domain::sharding::{MeasurementKey, ShardMapOp, ShardRegion, ShardRehomeRequest}; +use crate::domain::sharding::{ + MeasurementKey, ShardMap, ShardMapOp, ShardRegion, ShardRehomeRequest, +}; use crate::error::HyperbytedbError; use crate::ports::metadata::MetadataPort; use crate::ports::points_sink::PointsSinkPort; @@ -30,6 +33,10 @@ type RegionHeartbeatRow = (u64, u64, u64, u64, u64, u64); /// drive splits/rebalances forever. const HEARTBEAT_TTL_INTERVALS: u64 = 3; +/// Warn once a region has gone this many consecutive ticks with candidates +/// available but none caught up to the committed map. +const PLACEMENT_STALL_WARN_TICKS: u32 = 5; + /// Split only after the region's cooldown has elapsed. `last_split_at == 0` /// means "never split" (bootstrap must stamp wall-clock time); treating 0 as /// "long ago" let a single hot region binary-split every tick (1→2→4→8). @@ -116,6 +123,15 @@ pub struct ShardScheduler { /// enqueued — their primary changes are provisioned by the MV backfill /// path. rollup_dests: tokio::sync::RwLock>, + /// Peer `map_version` probes for the tick in progress, keyed by node id, + /// cleared when each tick starts. `/internal/shard/map` serializes every + /// space and region in the cluster to answer one `u64`, so probing per + /// region cost O(regions x map size) bytes on every tick of a join. + peer_map_versions: tokio::sync::RwLock>, + /// Consecutive ticks a region skipped placement with no caught-up + /// candidate, keyed by region id. A cluster-wide placement stall is + /// otherwise only visible at `debug!`. + placement_stalls: tokio::sync::RwLock>, #[cfg(test)] test_force_leader: bool, #[cfg(test)] @@ -154,6 +170,8 @@ impl ShardScheduler { unhealthy_primaries: tokio::sync::RwLock::new(HashMap::new()), reconciliation_queue: std::sync::Mutex::new(Vec::new()), rollup_dests: tokio::sync::RwLock::new(HashMap::new()), + peer_map_versions: tokio::sync::RwLock::new(HashMap::new()), + placement_stalls: tokio::sync::RwLock::new(HashMap::new()), #[cfg(test)] test_force_leader: false, #[cfg(test)] @@ -173,6 +191,58 @@ impl ShardScheduler { self } + /// `peer`'s committed `map_version`, probed at most once per tick. + async fn cached_peer_map_version( + &self, + peer: u64, + addr: &str, + client: &reqwest::Client, + ) -> Result { + if let Some(v) = self.peer_map_versions.read().await.get(&peer) { + return Ok(*v); + } + let version = fetch_peer_map_version(client, addr).await?; + self.peer_map_versions.write().await.insert(peer, version); + Ok(version) + } + + /// Record that `region_id` could not place a replica this tick, and warn + /// once the stall has persisted. A placement blocked forever by a lagging + /// candidate is otherwise silent above `debug!`. + async fn note_placement_stall(&self, region_id: u64, candidates: usize) { + let mut stalls = self.placement_stalls.write().await; + let count = stalls.entry(region_id).or_insert(0); + *count = count.saturating_add(1); + if *count == PLACEMENT_STALL_WARN_TICKS { + counter!("hyperbytedb_shard_placement_stalled_total").increment(1); + tracing::warn!( + region_id, + candidates, + ticks = *count, + "region has no caught-up replica candidate; placement stalled" + ); + } + } + + async fn clear_placement_stall(&self, region_id: u64) { + self.placement_stalls.write().await.remove(®ion_id); + } + + /// True when `key` is a rollup (SummingMergeTree) destination. + /// + /// `apply_transfer_push` has no idempotence guard, so a staged region that + /// fails to commit its map op is staged again on the next tick. A + /// `ReplacingMergeTree` measurement collapses the redelivery; an additive + /// rollup destination sums it twice and is corrupted permanently. Contended + /// lock counts as "yes" — skipping a placement costs a tick, guessing wrong + /// costs the data. + async fn is_rollup_dest(&self, key: &MeasurementKey) -> bool { + match self.rollup_dests.try_read() { + Ok(set) => set.contains_key(key), + Err(_) => true, + } + } + #[cfg(test)] async fn seed_rollup_dest(&self, key: &MeasurementKey) { self.rollup_dests.write().await.insert(key.clone(), ()); @@ -267,6 +337,8 @@ impl ShardScheduler { } async fn tick(&self) -> Result<(), HyperbytedbError> { + // Peer map versions are only valid for the tick that probed them. + self.peer_map_versions.write().await.clear(); let map = self.shard_map.snapshot().await?; self.drain_reconciliation(&map).await; let now = SystemTime::now() @@ -362,6 +434,35 @@ impl ShardScheduler { } } + // A placement that commits changes the region's epoch, peers or + // primary, which makes the `region` borrowed from this tick's + // snapshot stale. Everything below reads that borrow, and + // `try_rebalance` would spend a full region copy before its + // proposal failed the epoch CAS — so yield the region and pick + // it up fresh on the next tick. + let mut placed = false; + for step in ["live", "idle", "primary"] { + let outcome = match step { + "live" => self.try_place_live_member(&space.key, region).await, + "idle" => self.try_place_idle_member(&space.key, region).await, + _ => self.try_place_primary(&space.key, region).await, + }; + match outcome { + Ok(true) => { + placed = true; + break; + } + Ok(false) => {} + Err(e) => { + tracing::debug!(error = %e, region_id = region.region_id, step, "placement skipped"); + } + } + } + if placed { + self.release_operator(region.region_id).await; + continue; + } + if let Err(e) = self.try_rebalance(&space.key, region, &hb).await { tracing::debug!(error = %e, region_id = region.region_id, "rebalance skipped"); } @@ -939,22 +1040,358 @@ impl ShardScheduler { Ok(()) } + /// Place a live Active member that is not yet a region peer, when + /// effective RF has room. Stages rows onto the joiner, then commits + /// `AddPeer` — never commit-then-stage (that's heal `MovePeer`). + async fn try_place_live_member( + &self, + key: &MeasurementKey, + region: &ShardRegion, + ) -> Result { + if !self.config.add_peer_proposals_enabled || self.is_rollup_dest(key).await { + return Ok(false); + } + + // Resolve the peer client before probing the joiner: without one there + // is no way to stage rows, so the round trip below would be wasted. + let Some(pc) = self.peer_client.as_ref() else { + return Ok(false); + }; + + // Re-read the region: `drain_reconciliation` runs before this loop and + // can bump epochs, so the caller's snapshot may already be stale. + // Acting on it stages a whole region copy that the epoch CAS then + // rejects. + let map = self.shard_map.snapshot().await?; + let Some(region) = current_region(&map, key, region.region_id) else { + return Ok(false); + }; + + let candidates: Vec<(u64, String)> = { + let m = self.membership.read().await; + let ids: Vec = m.active_peers(0).into_iter().map(|n| n.node_id).collect(); + live_replica_candidates(region, &ids, self.config.replication_factor) + .into_iter() + .filter_map(|id| m.get_node(id).map(|n| (id, n.addr.clone()))) + .collect() + }; + if candidates.is_empty() { + self.clear_placement_stall(region.region_id).await; + return Ok(false); + } + + // Take the first candidate that has caught up. Stopping at the first + // candidate outright let one lagging node block every region forever, + // because the choice was made before the probe and never advanced. + let cluster_ver = map.map_version; + let mut chosen = None; + for (id, addr) in &candidates { + let ver = match self + .cached_peer_map_version(*id, addr, pc.http_client()) + .await + { + Ok(v) => v, + Err(e) => { + tracing::debug!(error = %e, candidate = id, "map version probe failed"); + continue; + } + }; + if joiner_map_caught_up(cluster_ver, ver) { + chosen = Some(*id); + break; + } + tracing::debug!( + region_id = region.region_id, + candidate = id, + cluster_ver, + candidate_ver = ver, + "candidate map not caught up" + ); + } + let Some(joiner) = chosen else { + self.note_placement_stall(region.region_id, candidates.len()) + .await; + return Ok(false); + }; + self.clear_placement_stall(region.region_id).await; + + if region.primary == self.node_id { + let outcome = stage_region_transfer_data( + pc, + &self.metadata, + &self.wal, + self.query_port.as_ref(), + self.node_id, + key, + region, + joiner, + self.max_points_per_request, + ) + .await?; + if !outcome.verified() { + return Err(HyperbytedbError::ShardMap( + format!( + "live placement stage unverified: exported={} applied={}", + outcome.exported, outcome.applied + ) + .into(), + )); + } + } else { + request_region_stage( + pc.as_ref(), + &self.membership, + region.primary, + key, + region, + joiner, + ) + .await?; + } + + self.propose(ShardMapOp::AddPeer { + key: key.clone(), + region_id: region.region_id, + to_peer: joiner, + epoch: region.epoch, + }) + .await?; + counter!("hyperbytedb_shard_live_placements_total").increment(1); + Ok(true) + } + + /// Rebalance one replica slot onto the latest-joined member. + /// + /// `try_place_live_member` only fires while a region is below effective RF. + /// At 3 nodes and RF 3 every region is already complete, so a 4th process + /// needs an existing peer to step aside instead. Rows are staged onto the + /// target before `MovePeer` commits — same stage-then-commit ordering, so + /// no peer is ever published holding nothing. `MovePeer` appends the + /// target, which makes it the region's newest peer and therefore the next + /// primary-placement candidate. + /// + /// # This is not join-triggered + /// + /// The target is the active member with the greatest `joined_at`, which is + /// whoever joined last — there is no recency window, so "last" may mean + /// months ago. The rule fires whenever that member is under-loaded relative + /// to some region's non-primary peer, which means a cluster that has been + /// stable but unbalanced will start moving replicas on the first tick after + /// an upgrade, not in response to any join. + /// + /// Movement is bounded and convergent — it stops as soon as memberships + /// even out — and each copy is staged and verified before the map changes, + /// so this is a scheduling surprise rather than a correctness risk. It is + /// accepted while sharding is beta and stays deliberately ungated; revisit + /// before sharding graduates, when an unannounced rebalance on upgrade + /// stops being acceptable. + async fn try_place_idle_member( + &self, + key: &MeasurementKey, + region: &ShardRegion, + ) -> Result { + if self.is_rollup_dest(key).await { + return Ok(false); + } + + let Some(pc) = self.peer_client.as_ref() else { + return Ok(false); + }; + + let latest_member = { + let m = self.membership.read().await; + m.active_peers(0) + .into_iter() + .max_by_key(|n| (n.joined_at, n.node_id)) + .map(|n| (n.node_id, n.addr.clone())) + }; + let Some((latest_member, addr)) = latest_member else { + return Ok(false); + }; + + // Re-read the region: an `AddPeer` earlier in this same tick leaves the + // caller's snapshot stale, and acting on it would stage a full copy + // only for the proposal to bounce off the epoch CAS. + let map = self.shard_map.snapshot().await?; + let Some(region) = current_region(&map, key, region.region_id) else { + return Ok(false); + }; + let memberships = region_memberships(&map); + let Some(displaced) = idle_member_replica_swap(region, latest_member, &memberships) else { + return Ok(false); + }; + + let joiner_ver = self + .cached_peer_map_version(latest_member, &addr, pc.http_client()) + .await?; + if !joiner_map_caught_up(map.map_version, joiner_ver) { + tracing::debug!( + region_id = region.region_id, + latest_member, + cluster_ver = map.map_version, + joiner_ver, + "skip idle placement: rebalance target map not caught up" + ); + return Ok(false); + } + + if region.primary == self.node_id { + let outcome = stage_region_transfer_data( + pc, + &self.metadata, + &self.wal, + self.query_port.as_ref(), + self.node_id, + key, + region, + latest_member, + self.max_points_per_request, + ) + .await?; + if !outcome.verified() { + return Err(HyperbytedbError::ShardMap( + format!( + "idle placement stage unverified: exported={} applied={}", + outcome.exported, outcome.applied + ) + .into(), + )); + } + } else { + request_region_stage( + pc.as_ref(), + &self.membership, + region.primary, + key, + region, + latest_member, + ) + .await?; + } + + tracing::info!( + region_id = region.region_id, + displaced, + latest_member, + "rebalancing region replica onto latest-joined member" + ); + self.propose(ShardMapOp::MovePeer { + key: key.clone(), + region_id: region.region_id, + from_peer: displaced, + to_peer: latest_member, + epoch: region.epoch, + }) + .await?; + counter!("hyperbytedb_shard_idle_placements_total").increment(1); + Ok(true) + } + + /// Hand a region's primary to its newest peer so a joiner starts taking + /// writes instead of serving only as a replica. + /// + /// Rows move and are verified *before* `TransferPrimary` commits, the same + /// ordering drain and rebalance use: an empty primary must never own a + /// range. `try_rebalance` cannot do this job — it reacts to a 2x byte-load + /// imbalance between peers, and a joiner that has just been staged holds a + /// copy of the same rows, so the imbalance it looks for never appears. + async fn try_place_primary( + &self, + key: &MeasurementKey, + region: &ShardRegion, + ) -> Result { + if self.is_rollup_dest(key).await { + return Ok(false); + } + + let Some(pc) = self.peer_client.as_ref() else { + return Ok(false); + }; + + // Re-read the region: a peer placement earlier in this same tick leaves + // the caller's snapshot stale, and a stale epoch would fail the + // `TransferPrimary` proposal only after a full region copy had run. + let map = self.shard_map.snapshot().await?; + let Some(region) = current_region(&map, key, region.region_id) else { + return Ok(false); + }; + let counts = primary_counts(&map); + let active: Vec = { + let m = self.membership.read().await; + m.active_peers(0).into_iter().map(|n| n.node_id).collect() + }; + let Some(new_primary) = primary_placement_candidate(region, &counts, &active) else { + return Ok(false); + }; + + // Source from whoever holds the authoritative rows. Exporting from the + // leader when leadership and primary diverge finds nothing, and a + // 0-exported/0-applied transfer "verifies" vacuously onto an empty node. + if region.primary == self.node_id { + run_region_transfer( + pc, + &self.metadata, + &self.wal, + self.query_port.as_ref(), + self.points_sink.as_ref(), + self.node_id, + key, + region, + new_primary, + self.max_points_per_request, + // Placement keeps the same range and only changes who leads it; + // keep the source copy as a replica rather than stranding data + // if the ownership change later fails. + false, + ) + .await?; + } else { + request_region_rehome( + pc, + &self.membership, + region.primary, + key, + region, + new_primary, + false, + ) + .await?; + } + + tracing::info!( + region_id = region.region_id, + from = region.primary, + to = new_primary, + "placing region primary on newest peer" + ); + self.propose(ShardMapOp::TransferPrimary { + key: key.clone(), + region_id: region.region_id, + new_primary, + epoch: region.epoch, + }) + .await?; + counter!("hyperbytedb_shard_primary_placements_total").increment(1); + Ok(true) + } + async fn try_rebalance( &self, key: &MeasurementKey, region: &ShardRegion, hb: &[(u64, u64, u64, u64, u64, u64)], ) -> Result<(), HyperbytedbError> { + // Only peers that have actually reported can be scored. A peer with no + // live heartbeat row is unmeasured, not empty: scoring it as 0 bytes + // made any freshly added replica look infinitely lighter than the + // primary and handed it ownership on the strength of missing telemetry. let mut loads: Vec<(u64, u64)> = region .peers .iter() - .map(|node| { - let bytes = hb - .iter() + .filter_map(|node| { + hb.iter() .find(|(r, n, _, _, _, _)| *r == region.region_id && *n == *node) - .map(|(_, _, _, b, _, _)| *b) - .unwrap_or(0); - (*node, bytes) + .map(|(_, _, _, b, _, _)| (*node, *b)) }) .collect(); if loads.len() < 2 { @@ -1169,6 +1606,189 @@ fn failover_watermark_safe(candidate_watermark: u64, max_peer_watermark: u64) -> candidate_watermark > 0 || max_peer_watermark == 0 } +/// Region data movement onto a joiner starts only after it has caught up to +/// the cluster's committed `map_version`. Staging rows against a lagging map +/// would apply under the wrong epoch / peer set. +/// +/// Ahead is fine, and has to be: the cluster version is read from a snapshot +/// taken a moment before the probe, so a joiner that applied an op in between +/// legitimately reports a higher number. Requiring equality rejected that node +/// and — because the candidate was chosen before the probe and never advanced +/// — stalled replica growth for every region indefinitely. +#[must_use] +pub fn joiner_map_caught_up(cluster_map_version: u64, joiner_map_version: u64) -> bool { + joiner_map_version >= cluster_map_version +} + +/// Members eligible to become a replica of `region`, lowest id first. +/// +/// Returns every candidate rather than just the best one so the caller can +/// advance past a node that fails its catch-up probe. Returning a single +/// candidate meant one permanently-lagging node was re-chosen every tick and +/// blocked replica growth cluster-wide. +#[must_use] +pub fn live_replica_candidates( + region: &ShardRegion, + active_member_ids: &[u64], + configured_rf: usize, +) -> Vec { + if region.transfer_outstanding() { + return Vec::new(); + } + let target = effective_replication_factor(configured_rf, active_member_ids.len()); + if region.peers.len() >= target { + return Vec::new(); + } + let mut candidates: Vec = active_member_ids + .iter() + .copied() + .filter(|id| !region.peers.contains(id)) + .collect(); + candidates.sort_unstable(); + candidates +} + +/// Next live member to add as a replica, or `None` when RF is full, the +/// region carries transfer debt, or every Active member is already a peer. +#[must_use] +pub fn live_replica_candidate( + region: &ShardRegion, + active_member_ids: &[u64], + configured_rf: usize, +) -> Option { + if region.transfer_outstanding() { + return None; + } + let target = effective_replication_factor(configured_rf, active_member_ids.len()); + if region.peers.len() >= target { + return None; + } + active_member_ids + .iter() + .copied() + .filter(|id| !region.peers.contains(id)) + .min() +} + +/// Count the regions each node is primary for, across every space in the map. +#[must_use] +pub fn primary_counts(map: &ShardMap) -> HashMap { + let mut counts: HashMap = HashMap::new(); + for space in map.spaces.values() { + for region in &space.regions { + *counts.entry(region.primary).or_insert(0) += 1; + } + } + counts +} + +/// The committed state of `region_id` in `map`, or `None` if it is gone (split, +/// merged, or its measurement dropped since the caller's snapshot). +#[must_use] +pub fn current_region<'a>( + map: &'a ShardMap, + key: &MeasurementKey, + region_id: u64, +) -> Option<&'a ShardRegion> { + map.spaces + .get(key)? + .regions + .iter() + .find(|r| r.region_id == region_id) +} + +/// Count the regions each node is a peer of, across every space in the map. +#[must_use] +pub fn region_memberships(map: &ShardMap) -> HashMap { + let mut counts: HashMap = HashMap::new(); + for space in map.spaces.values() { + for region in &space.regions { + for peer in ®ion.peers { + *counts.entry(*peer).or_insert(0) += 1; + } + } + } + counts +} + +/// Non-primary peer that should yield its replica slot to `latest_member`, or +/// `None` when this region is already well placed. +/// +/// An RF-complete region never triggers `AddPeer`, so a member that is not +/// already a peer would own nothing — every region legitimately has its full +/// replica count, and pure load balancing has no reason to disturb them. This +/// rule breaks that tie in one direction only: `latest_member` may displace a +/// peer, and only one carrying strictly more region memberships than it. +/// +/// `latest_member` must be the active member with the greatest `joined_at`. +/// The caller establishes that, and it is what makes the swap converge: the +/// identity is stable until membership itself changes, so the node just +/// displaced cannot reclaim the slot — it is not the latest member. Passing an +/// arbitrary node, or the least-loaded one, forfeits that and lets a slot trade +/// back and forth indefinitely. +/// +/// Note this is *ordering*, not recency: the greatest `joined_at` is simply +/// whichever member joined last, whether that was seconds or months ago. See +/// [`ShardScheduler::try_place_idle_member`] for what that means in practice. +#[must_use] +pub fn idle_member_replica_swap( + region: &ShardRegion, + latest_member: u64, + memberships: &HashMap, +) -> Option { + if region.transfer_outstanding() || region.peers.contains(&latest_member) { + return None; + } + let latest_member_load = memberships.get(&latest_member).copied().unwrap_or(0); + region + .peers + .iter() + .copied() + .filter(|p| *p != region.primary) + .map(|p| (memberships.get(&p).copied().unwrap_or(0), p)) + .filter(|(load, _)| *load > latest_member_load) + .max() + .map(|(_, peer)| peer) +} + +/// Node that should take this region's primary so a newly added replica starts +/// serving writes, or `None` when the primary is already well placed. +/// +/// The candidate is the region's newest peer — `AddPeer` appends, and region +/// sorting only orders by range, so `peers.last()` is the most recently placed +/// replica. It takes the primary only while it carries strictly fewer primaries +/// than the current one, which is what stops a join from stampeding every +/// region's primary onto the joiner and what makes the rule terminate: once the +/// newest peer *is* the primary the rule no longer applies, so there is no +/// oscillation back to the previous owner. +#[must_use] +pub fn primary_placement_candidate( + region: &ShardRegion, + primary_counts: &HashMap, + active_member_ids: &[u64], +) -> Option { + if region.transfer_outstanding() { + return None; + } + let newest = *region.peers.last()?; + if newest == region.primary || !active_member_ids.contains(&newest) { + return None; + } + let current = primary_counts.get(®ion.primary).copied().unwrap_or(0); + let candidate = primary_counts.get(&newest).copied().unwrap_or(0); + (candidate < current).then_some(newest) +} + +/// Read a peer's committed `map_version` via `/internal/shard/map`. +pub async fn fetch_peer_map_version( + client: &reqwest::Client, + peer_addr: &str, +) -> Result { + let map = + crate::adapters::cluster::sync_client::fetch_shard_map(client.clone(), peer_addr).await?; + Ok(map.map_version) +} + #[allow(clippy::too_many_arguments)] async fn push_and_drop_range( peer_client: &Arc, @@ -1245,6 +1865,7 @@ async fn request_region_rehome( epoch: range.epoch, dest_primary, drop_source, + stage: false, }; let url = format!("http://{addr}/internal/shard/rehome"); let resp = peer_client @@ -1263,10 +1884,55 @@ async fn request_region_rehome( Ok(()) } -/// Reconciliation disposition for a failed post-commit transfer movement. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub(crate) enum TransferDisposition { - /// Transient (transport, apply lag 404/409, 5xx): retry next drain. +/// Ask the current primary to stage `[start, end)` onto a joiner that is +/// not yet a committed peer. +async fn request_region_stage( + peer_client: &PeerClient, + membership: &SharedMembership, + target_node: u64, + key: &MeasurementKey, + range: &ShardRegion, + dest: u64, +) -> Result<(), HyperbytedbError> { + let addr = { + let m = membership.read().await; + m.get_node(target_node).map(|n| n.addr.clone()) + } + .ok_or_else(|| { + HyperbytedbError::PeerUnreachable(format!("stage target node {target_node} unknown")) + })?; + let req = ShardRehomeRequest { + db: key.db.clone(), + rp: key.rp.clone(), + measurement: key.measurement.clone(), + start: range.start, + end: range.end, + epoch: range.epoch, + dest_primary: dest, + drop_source: false, + stage: true, + }; + let url = format!("http://{addr}/internal/shard/rehome"); + let resp = peer_client + .http_client() + .post(&url) + .json(&req) + .timeout(Duration::from_secs(REHOME_TIMEOUT_SECS)) + .send() + .await + .map_err(|e| HyperbytedbError::PeerUnreachable(e.to_string()))?; + if !resp.status().is_success() { + return Err(HyperbytedbError::TransferRejected { + status: resp.status().as_u16(), + }); + } + Ok(()) +} + +/// Reconciliation disposition for a failed post-commit transfer movement. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum TransferDisposition { + /// Transient (transport, apply lag 404/409, 5xx): retry next drain. Retryable, /// Source/destination identity collapsed or moved (400 dest-is-self, /// collision): re-resolve roles and retry. @@ -1888,6 +2554,249 @@ mod tests { assert!(failover_watermark_safe(0, 0)); } + #[test] + fn joiner_map_catchup_blocks_a_lagging_joiner() { + assert!(joiner_map_caught_up(3, 3)); + assert!(!joiner_map_caught_up(3, 0)); + assert!(!joiner_map_caught_up(3, 2)); + // Ahead is allowed. The gate exists to stop staging against a *lagging* + // map; a joiner that applied an op after the cluster snapshot was read + // is not lagging, and rejecting it stalled placement permanently. + assert!(joiner_map_caught_up(3, 4)); + } + + fn sample_region_peers(peers: Vec, primary: u64) -> ShardRegion { + ShardRegion { + region_id: 1, + start: 0, + end: u64::MAX, + epoch: ShardEpoch::default(), + peers, + primary, + last_split_at: 0, + transfer_verified: None, + transfer_first_seen: None, + } + } + + #[test] + fn live_replica_candidate_picks_lowest_id_joiner() { + let region = sample_region_peers(vec![1], 1); + assert_eq!(live_replica_candidate(®ion, &[1, 2], 3), Some(2)); + } + + #[test] + fn live_replica_candidate_none_when_rf_full() { + let region = sample_region_peers(vec![1, 2, 3], 1); + assert_eq!(live_replica_candidate(®ion, &[1, 2, 3, 4], 3), None); + } + + #[test] + fn live_replica_candidate_none_when_already_peer() { + let region = sample_region_peers(vec![1, 2], 1); + assert_eq!(live_replica_candidate(®ion, &[1, 2], 3), None); + } + + #[test] + fn catch_up_accepts_a_candidate_that_is_ahead() { + assert!(joiner_map_caught_up(5, 5), "equal is caught up"); + // The cluster version comes from a snapshot taken before the probe, so + // a candidate that applied an op in between is legitimately ahead. + // Rejecting it stalled placement permanently. + assert!(joiner_map_caught_up(5, 6), "ahead is caught up"); + assert!(!joiner_map_caught_up(5, 4), "behind is not"); + } + + #[test] + fn live_replica_candidates_returns_every_option_in_order() { + let region = sample_region_peers(vec![1], 1); + assert_eq!( + live_replica_candidates(®ion, &[1, 3, 2], 3), + vec![2, 3], + "all non-peers, lowest first, so a lagging one can be skipped" + ); + } + + #[test] + fn live_replica_candidates_empty_when_rf_full_or_indebted() { + let full = sample_region_peers(vec![1, 2, 3], 1); + assert!(live_replica_candidates(&full, &[1, 2, 3, 4], 3).is_empty()); + + let mut indebted = sample_region_peers(vec![1], 1); + indebted.transfer_verified = Some(false); + assert!(live_replica_candidates(&indebted, &[1, 2], 3).is_empty()); + } + + #[test] + fn live_replica_candidate_none_with_transfer_debt() { + let mut region = sample_region_peers(vec![1], 1); + region.transfer_verified = Some(false); + assert_eq!(live_replica_candidate(®ion, &[1, 2], 3), None); + } + + /// Build a map whose single space holds `regions`, for primary counting. + fn map_of(regions: Vec) -> ShardMap { + let key = MeasurementKey::new("db", "autogen", "cpu"); + let mut map = ShardMap::default(); + map.spaces.insert( + key.clone(), + crate::domain::sharding::MeasurementShardSpace { key, regions }, + ); + map + } + + #[test] + fn primary_counts_tallies_every_space() { + let map = map_of(vec![ + sample_region_peers(vec![1, 2], 1), + sample_region_peers(vec![1, 2], 2), + sample_region_peers(vec![1, 2], 1), + ]); + let counts = primary_counts(&map); + assert_eq!(counts.get(&1), Some(&2)); + assert_eq!(counts.get(&2), Some(&1)); + } + + #[test] + fn primary_placement_hands_only_region_to_the_joiner() { + let region = sample_region_peers(vec![1, 2], 1); + let counts = primary_counts(&map_of(vec![region.clone()])); + assert_eq!( + primary_placement_candidate(®ion, &counts, &[1, 2]), + Some(2), + "a joiner with no primaries must take the region so writes land on it" + ); + } + + /// The rule must terminate: after the joiner takes the primary, evaluating + /// the same region again must not hand it straight back. + #[test] + fn primary_placement_does_not_oscillate() { + let placed = sample_region_peers(vec![1, 2], 2); + let counts = primary_counts(&map_of(vec![placed.clone()])); + assert_eq!( + primary_placement_candidate(&placed, &counts, &[1, 2]), + None, + "newest peer already holds the primary; nothing left to place" + ); + } + + #[test] + fn primary_placement_stops_at_an_even_split() { + // Four regions, primaries already 2/2 across the pair. + let regions = vec![ + sample_region_peers(vec![1, 2], 1), + sample_region_peers(vec![1, 2], 1), + sample_region_peers(vec![1, 2], 2), + sample_region_peers(vec![1, 2], 2), + ]; + let counts = primary_counts(&map_of(regions.clone())); + assert_eq!( + primary_placement_candidate(®ions[0], &counts, &[1, 2]), + None, + "moving another primary would only invert the imbalance" + ); + } + + #[test] + fn primary_placement_drains_a_lopsided_owner() { + let regions = vec![ + sample_region_peers(vec![1, 2], 1), + sample_region_peers(vec![1, 2], 1), + sample_region_peers(vec![1, 2], 1), + ]; + let counts = primary_counts(&map_of(regions.clone())); + assert_eq!( + primary_placement_candidate(®ions[0], &counts, &[1, 2]), + Some(2) + ); + } + + #[test] + fn primary_placement_skips_inactive_and_indebted() { + let region = sample_region_peers(vec![1, 2], 1); + let counts = primary_counts(&map_of(vec![region.clone()])); + assert_eq!( + primary_placement_candidate(®ion, &counts, &[1]), + None, + "joiner is not an active member" + ); + + let mut indebted = region.clone(); + indebted.transfer_verified = Some(false); + assert_eq!( + primary_placement_candidate(&indebted, &counts, &[1, 2]), + None, + "no ownership mutation while transfer debt is outstanding" + ); + } + + #[test] + fn idle_member_displaces_a_loaded_replica_not_the_primary() { + let region = sample_region_peers(vec![1, 2, 3], 1); + let memberships = region_memberships(&map_of(vec![region.clone()])); + let displaced = + idle_member_replica_swap(®ion, 4, &memberships).expect("latest member takes a slot"); + assert_ne!(displaced, 1, "the primary must never be displaced"); + assert!(displaced == 2 || displaced == 3); + } + + /// The swap must converge. `latest_member` stays the same node for as long as + /// membership is unchanged, so once it holds the slot the rule stops firing + /// and the slot cannot trade back and forth. + #[test] + fn idle_member_swap_converges() { + let before = sample_region_peers(vec![1, 2, 3], 1); + let displaced = idle_member_replica_swap( + &before, + 4, + ®ion_memberships(&map_of(vec![before.clone()])), + ) + .expect("first pass moves the latest member in"); + + let mut after = before.clone(); + after.peers.retain(|p| *p != displaced); + after.peers.push(4); + assert_eq!( + idle_member_replica_swap(&after, 4, ®ion_memberships(&map_of(vec![after.clone()]))), + None, + "second pass with the same latest member must be a no-op" + ); + } + + #[test] + fn idle_member_swap_skips_debt_and_balanced_regions() { + let mut indebted = sample_region_peers(vec![1, 2, 3], 1); + indebted.transfer_verified = Some(false); + let memberships = region_memberships(&map_of(vec![indebted.clone()])); + assert_eq!( + idle_member_replica_swap(&indebted, 4, &memberships), + None, + "no peer mutation while transfer debt is outstanding" + ); + + // Newcomer already carries as many memberships as every replica here. + let region = sample_region_peers(vec![1, 2, 3], 1); + let balanced = region_memberships(&map_of(vec![ + region.clone(), + sample_region_peers(vec![4, 5, 6], 4), + ])); + assert_eq!(idle_member_replica_swap(®ion, 4, &balanced), None); + } + + #[test] + fn primary_placement_never_moves_off_the_newest_peer() { + // Node 3 is both newest and primary while node 1 carries none. Evening + // that out is byte-load rebalance's job, not placement's — placement + // exists to give a joiner work, and node 3 already has it. + let region = sample_region_peers(vec![1, 2, 3], 3); + let counts = primary_counts(&map_of(vec![region.clone()])); + assert_eq!( + primary_placement_candidate(®ion, &counts, &[1, 2, 3]), + None + ); + } + #[tokio::test] #[serial_test::serial(chdb)] async fn try_failover_proposes_transfer_primary() { @@ -2303,6 +3212,27 @@ mod tests { assert!(q.is_empty(), "SummingMergeTree spaces must never enqueue"); } + /// Placement stages rows and `apply_transfer_push` is not idempotent, so a + /// placement that stages and then fails its proposal re-delivers on the + /// next tick. `SummingMergeTree` destinations would sum the redelivery. + #[tokio::test] + #[serial_test::serial(chdb)] + async fn placement_refuses_rollup_destinations() { + let harness = ReconcileTestHarness::new().await; + let rollup = MeasurementKey::new("db", "autogen", "rollup_dest"); + let raw = MeasurementKey::new("db", "autogen", "cpu"); + harness.scheduler.seed_rollup_dest(&rollup).await; + + assert!( + harness.scheduler.is_rollup_dest(&rollup).await, + "seeded rollup destination must be refused by every placement path" + ); + assert!( + !harness.scheduler.is_rollup_dest(&raw).await, + "raw ReplacingMergeTree measurements stay eligible" + ); + } + #[tokio::test] #[serial_test::serial(chdb)] async fn rebuild_skips_flagged_rollup_destinations() { @@ -2366,6 +3296,58 @@ mod tests { /// Minimal standalone harness for reconciliation tests: real RocksDB shard /// map + scheduler internals without a full raft bootstrap where possible. + /// Minimal stand-in for a joiner: answers the map-version probe with + /// `map_version`, and accepts a staging push by echoing the line count it + /// received. Returns the address it bound to. + async fn spawn_mock_peer(map_version: u64) -> (String, Arc) { + use axum::routing::{get, post}; + let probes = Arc::new(std::sync::atomic::AtomicU64::new(0)); + let probe_counter = probes.clone(); + let app = axum::Router::new() + .route( + "/internal/shard/map", + get(move || { + let probes = probe_counter.clone(); + async move { + probes.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + axum::Json(serde_json::json!({ + "map_version": map_version, + "next_region_id": 2, + "spaces": [], + })) + } + }), + ) + .route( + "/internal/shard/transfer", + post(|body: axum::body::Bytes| async move { + let payload: serde_json::Value = + serde_json::from_slice(&body).unwrap_or_default(); + let applied = payload + .get("body") + .and_then(|b| b.as_array()) + .map(|bytes| { + let raw: Vec = bytes + .iter() + .filter_map(|v| v.as_u64().map(|n| n as u8)) + .collect(); + String::from_utf8_lossy(&raw) + .lines() + .filter(|l| !l.trim().is_empty()) + .count() as u64 + }) + .unwrap_or(0); + axum::Json(serde_json::json!({ "applied": applied })) + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap().to_string(); + tokio::spawn(async move { + let _ = axum::serve(listener, app).await; + }); + (addr, probes) + } + struct ReconcileTestHarness { _dir: tempfile::TempDir, shard_map: Arc, @@ -2387,6 +3369,235 @@ mod tests { } } + /// Harness whose scheduler owns the region and can actually reach a peer, + /// so `tick()` runs the placement steps end to end. + struct PlacementTestHarness { + _dir: tempfile::TempDir, + scheduler: ShardScheduler, + proposals: Arc>>, + /// `/internal/shard/map` hit counts, keyed by mock peer node id. + probes: HashMap>, + } + + impl PlacementTestHarness { + /// `peers` is `(node_id, map_version_delta)`; a negative delta makes + /// that peer report a map behind the cluster's. `regions` is how many + /// measurements to bootstrap, each one region owned by this node. + async fn new( + add_peer_proposals_enabled: bool, + peers: &[(u64, i64)], + regions: usize, + ) -> Self { + use crate::adapters::chdb::native_adapter::ChdbNativeAdapter; + use crate::adapters::chdb::query_adapter::ChdbQueryAdapter; + use crate::adapters::chdb::session::SharedSession; + use crate::adapters::cluster::peer_client::PeerClient; + use crate::adapters::cluster::replication_log::ReplicationLog; + use crate::adapters::metadata::rocksdb_meta::RocksDbMetadata; + use crate::adapters::wal::rocksdb_wal::RocksDbWal; + use crate::application::cluster::bootstrap::ClusterBootstrap; + use crate::application::materialized_view_service::MaterializedViewService; + use crate::ports::points_sink::PointsSinkPort; + + let dir = tempfile::tempdir().unwrap(); + let meta_dir = dir.path().join("meta"); + let wal_dir = dir.path().join("wal"); + let chdb_dir = dir.path().join("chdb"); + for p in [&meta_dir, &wal_dir, &chdb_dir] { + std::fs::create_dir_all(p).unwrap(); + } + + let chdb = SharedSession::new_eager(chdb_dir.to_str().unwrap(), 1).unwrap(); + let chdb_adapter = Arc::new(ChdbQueryAdapter::from_shared(chdb.clone(), 0)); + let sink: Arc = Arc::new(ChdbNativeAdapter::new(chdb)); + let wal = Arc::new(RocksDbWal::open(&wal_dir).unwrap()); + let metadata = Arc::new(RocksDbMetadata::open(&meta_dir).unwrap()); + let mv_service = Arc::new(MaterializedViewService::new( + metadata.clone(), + chdb_adapter, + sink.clone(), + )); + + let mut cluster_cfg = crate::config::HyperbytedbConfig::load(None) + .unwrap() + .cluster; + cluster_cfg.enabled = true; + cluster_cfg.node_id = 1; + cluster_cfg.cluster_addr = "127.0.0.1:18100".into(); + cluster_cfg.replication_log_dir = dir.path().join("repl").to_string_lossy().into(); + cluster_cfg.raft_dir = dir.path().join("raft").to_string_lossy().into(); + cluster_cfg.raft_heartbeat_interval_ms = Some(200); + cluster_cfg.raft_election_timeout_ms = Some(500); + + let bootstrap = ClusterBootstrap::init(&cluster_cfg, 1000).unwrap(); + let shard_map = Arc::new(RocksDbShardMap::open(&meta_dir, true).unwrap()); + let location_cache = Arc::new(ShardLocationCache::new()); + let raft = bootstrap + .start_raft( + &cluster_cfg, + metadata.clone(), + mv_service, + sink.clone(), + wal.clone(), + Some((shard_map.clone(), location_cache)), + None, + ) + .await + .unwrap(); + tokio::time::sleep(std::time::Duration::from_millis(150)).await; + + // Regions owned by this node, each one peer short of effective RF. + for i in 0..regions.max(1) { + let mut region = sample_region_peers(vec![1], 1); + region.region_id = (i + 1) as u64; + shard_map + .apply_op(ShardMapOp::BootstrapMeasurement { + key: MeasurementKey::new("db", "autogen", format!("cpu{i}")), + region, + }) + .await + .unwrap(); + } + let map_version = shard_map.snapshot().await.unwrap().map_version; + + let mut probes = HashMap::new(); + { + let mut m = bootstrap.membership.write().await; + m.add_node(NodeInfo { + node_id: 1, + addr: "127.0.0.1:18100".into(), + state: NodeState::Active, + joined_at: 0, + last_heartbeat: 0, + needs_sync: false, + }); + for (id, delta) in peers { + let reported = map_version.saturating_add_signed(*delta); + let (addr, counter) = spawn_mock_peer(reported).await; + probes.insert(*id, counter); + m.add_node(NodeInfo { + node_id: *id, + addr, + state: NodeState::Active, + joined_at: *id as i64, + last_heartbeat: 0, + needs_sync: false, + }); + } + } + + let peer_client = Arc::new(PeerClient::new( + 1, + "127.0.0.1:18100".into(), + bootstrap.membership.clone(), + Arc::new(ReplicationLog::open(dir.path().join("repl-peer")).unwrap()), + 2, + 8192, + 8, + 8 * 1024 * 1024, + )); + + let sharding = crate::config::ShardingConfig { + replication_factor: 3, + add_peer_proposals_enabled, + ..Default::default() + }; + let proposals = Arc::new(std::sync::Mutex::new(Vec::::new())); + let scheduler = ShardScheduler::new( + shard_map, + bootstrap.membership.clone(), + raft, + Some(peer_client), + metadata, + wal, + None, + Some(sink), + 1, + sharding, + 10_000, + ) + .with_test_force_leader(true) + .with_test_propose_sink(proposals.clone()); + + Self { + _dir: dir, + scheduler, + proposals, + probes, + } + } + + fn probe_count(&self, peer: u64) -> u64 { + self.probes[&peer].load(std::sync::atomic::Ordering::Relaxed) + } + } + + /// Guards the tick wiring. Without this, deleting the three placement calls + /// from `tick()` leaves every other test in the suite green. + #[tokio::test] + #[serial_test::serial(chdb)] + async fn tick_places_a_live_joiner_as_a_region_peer() { + let harness = PlacementTestHarness::new(true, &[(2, 0)], 1).await; + harness.scheduler.tick_once_for_test().await.unwrap(); + let ops = harness.proposals.lock().unwrap(); + assert!( + ops.iter() + .any(|op| matches!(op, ShardMapOp::AddPeer { to_peer: 2, .. })), + "tick must propose AddPeer to place the live joiner, got {ops:?}" + ); + } + + /// H.2: one lagging candidate must not block the others. Node 2 reports a + /// map behind the cluster and is chosen first by id; placement has to move + /// on to node 3 instead of stalling here every tick forever. + #[tokio::test] + #[serial_test::serial(chdb)] + async fn placement_skips_a_lagging_candidate_for_a_caught_up_one() { + let harness = PlacementTestHarness::new(true, &[(2, -1), (3, 0)], 1).await; + harness.scheduler.tick_once_for_test().await.unwrap(); + let ops = harness.proposals.lock().unwrap(); + assert!( + ops.iter() + .any(|op| matches!(op, ShardMapOp::AddPeer { to_peer: 3, .. })), + "must place onto the caught-up candidate, got {ops:?}" + ); + assert!( + !ops.iter() + .any(|op| matches!(op, ShardMapOp::AddPeer { to_peer: 2, .. })), + "must not place onto the lagging candidate, got {ops:?}" + ); + } + + /// H.4: `/internal/shard/map` serializes the whole cluster map to answer + /// one `u64`. Probing per region made a join cost O(regions x map size) + /// bytes per tick; one probe per peer per tick is the contract. + #[tokio::test] + #[serial_test::serial(chdb)] + async fn map_version_is_probed_once_per_peer_per_tick() { + let harness = PlacementTestHarness::new(true, &[(2, -1)], 4).await; + harness.scheduler.tick_once_for_test().await.unwrap(); + assert_eq!( + harness.probe_count(2), + 1, + "four regions must share one probe, not one probe each" + ); + } + + /// The upgrade gate has to hold at the tick, not just in config: a leader + /// that proposes `AddPeer` mid-rolling-restart wedges un-upgraded voters. + #[tokio::test] + #[serial_test::serial(chdb)] + async fn tick_withholds_add_peer_while_the_upgrade_gate_is_closed() { + let harness = PlacementTestHarness::new(false, &[(2, 0)], 1).await; + harness.scheduler.tick_once_for_test().await.unwrap(); + let ops = harness.proposals.lock().unwrap(); + assert!( + !ops.iter() + .any(|op| matches!(op, ShardMapOp::AddPeer { .. })), + "AddPeer must not be proposed while the gate is closed, got {ops:?}" + ); + } + impl ReconcileTestHarness { async fn new() -> Self { use crate::adapters::chdb::native_adapter::ChdbNativeAdapter; diff --git a/hyperbytedb/src/config.rs b/hyperbytedb/src/config.rs index 83a5f2d..392a6e1 100644 --- a/hyperbytedb/src/config.rs +++ b/hyperbytedb/src/config.rs @@ -79,6 +79,12 @@ pub struct ShardingConfig { /// older than the `ClearVerified` op cannot decode it from the Raft log. #[serde(default = "default_transfer_clear_proposals_enabled")] pub transfer_clear_proposals_enabled: bool, + /// Propose `AddPeer` shard-map ops to place a joiner as a region replica. + /// Disable for the duration of a rolling upgrade: a node on a build older + /// than the `AddPeer` op cannot decode it, rejects the append RPC carrying + /// it, and stops replicating until every node is upgraded. + #[serde(default = "default_add_peer_proposals_enabled")] + pub add_peer_proposals_enabled: bool, } impl Default for ShardingConfig { @@ -100,6 +106,7 @@ impl Default for ShardingConfig { scatter_max_peer_attempts: default_scatter_max_peer_attempts(), peer_heal_enabled: default_peer_heal_enabled(), transfer_clear_proposals_enabled: default_transfer_clear_proposals_enabled(), + add_peer_proposals_enabled: default_add_peer_proposals_enabled(), } } } @@ -160,6 +167,21 @@ fn default_transfer_clear_proposals_enabled() -> bool { true } +/// On, matching `transfer_clear_proposals_enabled`: the flag is an opt-out for +/// a mixed-version rollout, not an opt-in for the feature. +/// +/// The hazard is real but bounded. A follower on a build that does not know +/// `AddPeer` fails to deserialize the whole append RPC — axum rejects the body +/// before the handler runs — so it stops replicating and the leader can lose +/// commit quorum for the duration. It is self-healing: replication resumes as +/// soon as every node runs a build that knows the op, and ingest does not go +/// through Raft. Defaulting off instead would leave region placement inert for +/// everyone who enables sharding, which is the worse trade while sharding is +/// beta and clusters are rebuilt more often than rolling-upgraded. +fn default_add_peer_proposals_enabled() -> bool { + true +} + impl HyperbytedbConfig { /// Validate cross-field constraints after Figment merge. pub fn validate(&self) -> Result<(), String> { @@ -1243,4 +1265,14 @@ mod replicate_body_limit_tests { 64 * 1024 * 1024 ); } + + /// Shard-map op gates are opt-*outs* for a rolling upgrade, not opt-ins for + /// the feature. Defaulting one off leaves the behaviour it guards inert for + /// everyone who never finds the flag. + #[test] + fn shard_map_op_gates_default_on() { + let s = super::ShardingConfig::default(); + assert!(s.add_peer_proposals_enabled); + assert!(s.transfer_clear_proposals_enabled); + } } diff --git a/hyperbytedb/src/domain/sharding/ops.rs b/hyperbytedb/src/domain/sharding/ops.rs index 0a8256e..467d2eb 100644 --- a/hyperbytedb/src/domain/sharding/ops.rs +++ b/hyperbytedb/src/domain/sharding/ops.rs @@ -58,6 +58,17 @@ pub enum ShardMapOp { #[serde(default)] epoch: ShardEpoch, }, + /// Grow a region's replica set. Does not change primary. + /// + /// Used when effective RF has room (n grew). Transfer data onto `to_peer` + /// *before* proposing this op — commit-then-stage is the heal path only. + AddPeer { + key: MeasurementKey, + region_id: u64, + to_peer: u64, + #[serde(default)] + epoch: ShardEpoch, + }, TransferPrimary { key: MeasurementKey, region_id: u64, @@ -84,6 +95,7 @@ impl ShardMapOp { | ShardMapOp::Split { key, .. } | ShardMapOp::Merge { key, .. } | ShardMapOp::MovePeer { key, .. } + | ShardMapOp::AddPeer { key, .. } | ShardMapOp::TransferPrimary { key, .. } | ShardMapOp::ClearVerified { key, .. } => key, } @@ -267,6 +279,38 @@ pub fn apply_shard_map_op( } sort_space_regions(space); } + ShardMapOp::AddPeer { + key, + region_id, + to_peer, + epoch, + } => { + let space = map + .spaces + .get_mut(&key) + .ok_or(ShardMapApplyError::UnknownSpace)?; + let region = space + .regions + .iter_mut() + .find(|r| r.region_id == region_id) + .ok_or(ShardMapApplyError::UnknownRegion(region_id))?; + if region.epoch != epoch { + return Err(ShardMapApplyError::StaleEpoch("AddPeer")); + } + if region.transfer_outstanding() { + return Err(ShardMapApplyError::Invalid( + "cannot add peer while transfer debt is outstanding".into(), + )); + } + if region.peers.contains(&to_peer) { + return Err(ShardMapApplyError::Invalid(format!( + "peer {to_peer} already in region" + ))); + } + region.peers.push(to_peer); + region.epoch = epoch.bump_conf_ver(); + sort_space_regions(space); + } ShardMapOp::TransferPrimary { key, region_id, @@ -509,6 +553,86 @@ mod tests { assert_eq!(starts, vec![0, split_key]); } + #[test] + fn add_peer_appends_replica_without_moving_primary() { + let key = MeasurementKey::new("db", "autogen", "cpu"); + let region = sample_region(1, 0, u64::MAX, 1); + let mut map = ShardMap::default(); + apply_shard_map_op( + &mut map, + ShardMapOp::BootstrapMeasurement { + key: key.clone(), + region: region.clone(), + }, + ) + .unwrap(); + apply_shard_map_op( + &mut map, + ShardMapOp::AddPeer { + key, + region_id: 1, + to_peer: 3, + epoch: region.epoch, + }, + ) + .unwrap(); + let r = &map.spaces.values().next().unwrap().regions[0]; + assert_eq!(r.primary, 1); + assert!(r.peers.contains(&3)); + assert_eq!(r.epoch.conf_ver, region.epoch.conf_ver + 1); + } + + #[test] + fn add_peer_rejects_duplicate_and_stale_epoch() { + let key = MeasurementKey::new("db", "autogen", "cpu"); + let region = sample_region(1, 0, u64::MAX, 1); + let mut map = ShardMap::default(); + apply_shard_map_op( + &mut map, + ShardMapOp::BootstrapMeasurement { + key: key.clone(), + region: region.clone(), + }, + ) + .unwrap(); + let dup = apply_shard_map_op( + &mut map, + ShardMapOp::AddPeer { + key: key.clone(), + region_id: 1, + to_peer: 1, + epoch: region.epoch, + }, + ) + .unwrap_err(); + assert!(dup.to_string().contains("already in region"), "{dup}"); + + apply_shard_map_op( + &mut map, + ShardMapOp::AddPeer { + key: key.clone(), + region_id: 1, + to_peer: 3, + epoch: region.epoch, + }, + ) + .unwrap(); + let stale = apply_shard_map_op( + &mut map, + ShardMapOp::AddPeer { + key, + region_id: 1, + to_peer: 4, + epoch: region.epoch, + }, + ) + .unwrap_err(); + assert!( + stale.to_string().contains("stale epoch on AddPeer"), + "{stale}" + ); + } + #[test] fn transfer_primary_changes_owner() { let key = MeasurementKey::new("db", "autogen", "cpu"); @@ -639,6 +763,32 @@ mod tests { } } + #[test] + fn add_peer_rejects_outstanding_transfer_debt() { + let key = MeasurementKey::new("db", "autogen", "cpu"); + let region = flagged_region(1, 0, u64::MAX, 1, 1000); + let mut map = ShardMap::default(); + apply_shard_map_op( + &mut map, + ShardMapOp::BootstrapMeasurement { + key: key.clone(), + region: region.clone(), + }, + ) + .unwrap(); + let err = apply_shard_map_op( + &mut map, + ShardMapOp::AddPeer { + key, + region_id: 1, + to_peer: 3, + epoch: region.epoch, + }, + ) + .unwrap_err(); + assert!(err.to_string().contains("transfer debt"), "{err}"); + } + #[test] fn split_flags_only_changed_primary_children() { let key = MeasurementKey::new("db", "autogen", "cpu"); diff --git a/hyperbytedb/src/domain/sharding/transfer.rs b/hyperbytedb/src/domain/sharding/transfer.rs index c02555b..3827706 100644 --- a/hyperbytedb/src/domain/sharding/transfer.rs +++ b/hyperbytedb/src/domain/sharding/transfer.rs @@ -165,6 +165,10 @@ pub struct ShardRehomeRequest { pub dest_primary: u64, #[serde(default)] pub drop_source: bool, + /// When true, the destination is not yet a committed peer (`AddPeer` + /// staging). Uses the stage transfer that skips dest ownership checks. + #[serde(default)] + pub stage: bool, } #[cfg(test)] diff --git a/hyperbytedb/src/domain/sharding/types.rs b/hyperbytedb/src/domain/sharding/types.rs index 4fd903f..8be463a 100644 --- a/hyperbytedb/src/domain/sharding/types.rs +++ b/hyperbytedb/src/domain/sharding/types.rs @@ -217,6 +217,20 @@ impl From<&ShardMap> for ShardMapJson { } } +impl From for ShardMap { + fn from(json: ShardMapJson) -> Self { + let mut spaces = HashMap::new(); + for space in json.spaces { + spaces.insert(space.key.clone(), space); + } + Self { + map_version: json.map_version, + next_region_id: json.next_region_id, + spaces, + } + } +} + #[cfg(test)] mod locate_tests { use super::*; diff --git a/hyperbytedb/src/ports/sharding.rs b/hyperbytedb/src/ports/sharding.rs index 0a75808..dc11bd0 100644 --- a/hyperbytedb/src/ports/sharding.rs +++ b/hyperbytedb/src/ports/sharding.rs @@ -27,6 +27,12 @@ pub trait ShardMapPort: Send + Sync { async fn apply_op(&self, op: ShardMapOp) -> Result; + /// Replace the local map with a peer snapshot (join catch-up). + /// + /// Used so a joiner's `map_version` matches the cluster before region + /// data movement starts. Does not change membership or region peers. + async fn replace_map(&self, map: ShardMap) -> Result<(), HyperbytedbError>; + /// Returns true when this node holds a replica of any series for the measurement. async fn node_owns_measurement( &self, @@ -73,6 +79,10 @@ impl ShardMapPort for DisabledShardMap { Err(HyperbytedbError::Internal("sharding is disabled".into())) } + async fn replace_map(&self, _map: ShardMap) -> Result<(), HyperbytedbError> { + Err(HyperbytedbError::Internal("sharding is disabled".into())) + } + async fn node_owns_measurement( &self, _node_id: u64, diff --git a/hyperbytedb/tests/common/sharding_cluster.rs b/hyperbytedb/tests/common/sharding_cluster.rs index b9d56bb..808109d 100644 --- a/hyperbytedb/tests/common/sharding_cluster.rs +++ b/hyperbytedb/tests/common/sharding_cluster.rs @@ -47,6 +47,8 @@ pub struct ShardedTestNode { pub location_cache: Arc, pub membership: SharedMembership, pub query_port: Arc, + /// Shared with other in-process peers — libchdb allows one session per process. + pub chdb: SharedSession, flush: Arc, handle: tokio::task::JoinHandle<()>, shutdown: Option>, @@ -123,7 +125,7 @@ pub async fn start_sharded_node( let wal = Arc::new(RocksDbWal::open(&wal_dir).unwrap()); let metadata = Arc::new(RocksDbMetadata::open(&meta_dir).unwrap()); let chdb_adapter = Arc::new(ChdbQueryAdapter::from_shared(chdb.clone(), 0)); - let sink: Arc = Arc::new(ChdbNativeAdapter::new(chdb)); + let sink: Arc = Arc::new(ChdbNativeAdapter::new(chdb.clone())); let flush: Arc = Arc::new(FlushServiceImpl::new(wal.clone(), 0, sink.clone())); @@ -284,12 +286,140 @@ pub async fn start_sharded_node( location_cache, membership: shared_membership, query_port: chdb_adapter, + chdb, flush, handle, shutdown: Some(shutdown_tx), } } +pub async fn fetch_shard_map( + client: &reqwest::Client, + url: &str, +) -> hyperbytedb::domain::sharding::ShardMap { + let resp = client + .get(format!("{url}/internal/shard/map")) + .send() + .await + .unwrap(); + assert!( + resp.status().is_success(), + "shard map fetch failed: {}", + resp.status() + ); + let json: hyperbytedb::domain::sharding::ShardMapJson = resp.json().await.unwrap(); + json.into() +} + +pub async fn install_shard_map_from_peer(from: &ShardedTestNode, onto: &ShardedTestNode) { + let client = reqwest::Client::new(); + let map = fetch_shard_map(&client, &from.url).await; + onto.shard_map.replace_map(map).await.unwrap(); + let snap = onto.shard_map.snapshot().await.unwrap(); + onto.location_cache.refresh_from_map(&snap); +} + +/// Start `node_id` sharing `existing`'s membership and chDB session (libchdb +/// is process-global). Marks the joiner Active before serving. +pub async fn start_sharded_joiner( + dir: &Path, + existing: &ShardedTestNode, + node_id: u64, + opts: &ShardedClusterOptions, +) -> ShardedTestNode { + let listener = bind_ephemeral().await; + let addr = listener.local_addr().unwrap().to_string(); + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_millis() as i64; + { + let mut m = existing.membership.write().await; + m.add_node(NodeInfo { + node_id, + addr, + state: NodeState::Active, + joined_at: now, + last_heartbeat: now, + needs_sync: false, + }); + } + start_sharded_node( + dir, + node_id, + listener, + existing.membership.clone(), + opts, + existing.chdb.clone(), + ) + .await +} + +pub async fn apply_add_peer_on_nodes( + nodes: &[&ShardedTestNode], + db: &str, + rp: &str, + measurement: &str, + to_peer: u64, +) { + for node in nodes { + let map = node.shard_map.snapshot().await.unwrap(); + let region = map + .space(db, rp, measurement) + .and_then(|s| s.regions.first()) + .cloned() + .expect("region for AddPeer"); + let op = ShardMapOp::AddPeer { + key: MeasurementKey::new(db, rp, measurement), + region_id: region.region_id, + to_peer, + epoch: region.epoch, + }; + node.shard_map.apply_op(op).await.unwrap(); + let snap = node.shard_map.snapshot().await.unwrap(); + node.location_cache.refresh_from_map(&snap); + } +} + +pub async fn apply_move_peer_on_nodes( + nodes: &[&ShardedTestNode], + db: &str, + rp: &str, + measurement: &str, + from_peer: u64, + to_peer: u64, +) { + for node in nodes { + let map = node.shard_map.snapshot().await.unwrap(); + let region = map + .space(db, rp, measurement) + .and_then(|s| s.regions.first()) + .cloned() + .expect("region for MovePeer"); + let op = ShardMapOp::MovePeer { + key: MeasurementKey::new(db, rp, measurement), + region_id: region.region_id, + from_peer, + to_peer, + epoch: region.epoch, + }; + node.shard_map.apply_op(op).await.unwrap(); + let snap = node.shard_map.snapshot().await.unwrap(); + node.location_cache.refresh_from_map(&snap); + } +} + +pub async fn start_sharded_single_node(dir: &Path, opts: ShardedClusterOptions) -> ShardedTestNode { + let chdb_dir = dir.join("chdb-shared"); + std::fs::create_dir_all(&chdb_dir).unwrap(); + let chdb = SharedSession::new_eager(chdb_dir.to_str().unwrap(), 1).unwrap(); + + let l1 = bind_ephemeral().await; + let a1 = l1.local_addr().unwrap().to_string(); + let membership = build_shared_membership(&[(1, a1)]); + start_sharded_node(dir, 1, l1, membership, &opts, chdb).await +} + pub async fn start_sharded_pair_cluster( dir: &Path, opts: ShardedClusterOptions, @@ -375,6 +505,17 @@ pub async fn promote_region_primary_on_all_nodes( rp: &str, measurement: &str, new_primary: u64, +) { + let refs: Vec<&ShardedTestNode> = nodes.iter().collect(); + promote_region_primary_on_nodes(&refs, db, rp, measurement, new_primary).await; +} + +pub async fn promote_region_primary_on_nodes( + nodes: &[&ShardedTestNode], + db: &str, + rp: &str, + measurement: &str, + new_primary: u64, ) { for node in nodes { let map = node.shard_map.snapshot().await.unwrap(); diff --git a/hyperbytedb/tests/sharding_cluster_integration.rs b/hyperbytedb/tests/sharding_cluster_integration.rs index 577ab0c..0493aa1 100644 --- a/hyperbytedb/tests/sharding_cluster_integration.rs +++ b/hyperbytedb/tests/sharding_cluster_integration.rs @@ -750,6 +750,710 @@ async fn aggregate_read_from_primary_when_replica_lags() { ); } +/// P1.1: a 1-member cluster with sharding on accepts /write and owns the +/// first region (peers=[self], primary=self) even when configured RF is 3. +#[tokio::test] +#[serial(chdb)] +async fn one_member_cluster_owns_region_after_first_write() { + let dir = tempfile::tempdir().unwrap(); + let opts = ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }; + let node = start_sharded_single_node(dir.path(), opts).await; + + let client = reqwest::Client::new(); + create_db(&client, &node.url, "p1db").await; + let resp = write_line( + &client, + &node.url, + "p1db", + "cpu,host=solo value=1 1000000000", + ) + .await; + assert_eq!( + resp.status(), + reqwest::StatusCode::NO_CONTENT, + "n=1 sharded write must succeed: {}", + resp.status() + ); + + let map = node.shard_map.snapshot().await.unwrap(); + let space = map + .space("p1db", "autogen", "cpu") + .expect("first write must bootstrap a region"); + assert!( + !space.regions.is_empty(), + "shard map must have at least one region" + ); + let region = &space.regions[0]; + assert_eq!(region.peers, vec![1], "n=1 peers must be [self]"); + assert_eq!(region.primary, 1, "n=1 primary must be self"); +} + +/// P1.2: after a joiner is Active, its map_version equals the cluster's +/// before any region data movement (peer-set change) starts. +#[tokio::test] +#[serial(chdb)] +async fn joiner_map_version_matches_before_region_movement() { + let dir = tempfile::tempdir().unwrap(); + let opts = ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }; + let node1 = start_sharded_single_node(dir.path(), opts).await; + let client = reqwest::Client::new(); + create_db(&client, &node1.url, "p1db").await; + let resp = write_line( + &client, + &node1.url, + "p1db", + "cpu,host=solo value=1 1000000000", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + + let cluster_before = node1.shard_map.snapshot().await.unwrap(); + assert!( + cluster_before.map_version >= 1, + "first write must bump map_version" + ); + let region = cluster_before + .space("p1db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .expect("region after write"); + assert_eq!(region.peers, vec![1]); + + let l2 = bind_ephemeral().await; + let a2 = l2.local_addr().unwrap().to_string(); + { + let mut m = node1.membership.write().await; + m.add_node(hyperbytedb::domain::cluster::membership::NodeInfo { + node_id: 2, + addr: a2, + state: hyperbytedb::domain::cluster::membership::NodeState::Active, + joined_at: 0, + last_heartbeat: 0, + needs_sync: false, + }); + } + // libchdb is process-global; share node 1's session (pair-cluster harness). + let node2 = start_sharded_node( + dir.path(), + 2, + l2, + node1.membership.clone(), + &ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }, + node1.chdb.clone(), + ) + .await; + + let empty = node2.shard_map.snapshot().await.unwrap(); + assert_eq!(empty.map_version, 0, "joiner starts with an empty map"); + + install_shard_map_from_peer(&node1, &node2).await; + + let after = node2.shard_map.snapshot().await.unwrap(); + assert_eq!( + after.map_version, cluster_before.map_version, + "joiner map_version must equal the cluster's before movement" + ); + let joined_region = after + .space("p1db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .expect("installed region"); + assert_eq!( + joined_region.peers, + vec![1], + "catch-up must not add the joiner as a peer" + ); + assert!( + hyperbytedb::application::shard_scheduler::joiner_map_caught_up( + cluster_before.map_version, + after.map_version + ) + ); +} + +/// P1.3: after join, each pre-existing region has the joiner as a committed +/// peer, series rows for that range exist on the joiner, RF = min(config, 2). +#[tokio::test] +#[serial(chdb)] +async fn joiner_receives_existing_region_as_replica() { + use hyperbytedb::application::shard_scheduler::live_replica_candidate; + use hyperbytedb::application::shard_transfer::stage_region_transfer_data; + use hyperbytedb::ports::metadata::MetadataPort; + use hyperbytedb::ports::query::QueryPort; + use hyperbytedb::ports::wal::WalPort; + use std::sync::Arc; + + let dir = tempfile::tempdir().unwrap(); + let opts = ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }; + let node1 = start_sharded_single_node(dir.path(), opts).await; + let client = reqwest::Client::new(); + create_db(&client, &node1.url, "p1db").await; + let resp = write_line( + &client, + &node1.url, + "p1db", + "cpu,host=solo value=1 1000000000", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + flush_node(&node1).await; + + let before = node1.shard_map.snapshot().await.unwrap(); + let region = before + .space("p1db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after write"); + assert_eq!(region.peers, vec![1]); + let expected_sid = series_id("cpu", &BTreeMap::from([("host".into(), "solo".into())])); + assert!( + region.contains(expected_sid), + "written series must fall in the bootstrapped region" + ); + + let node2 = start_sharded_joiner( + dir.path(), + &node1, + 2, + &ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }, + ) + .await; + create_db(&client, &node2.url, "p1db").await; + install_shard_map_from_peer(&node1, &node2).await; + + let active = [1u64, 2]; + assert_eq!( + live_replica_candidate(®ion, &active, 3), + Some(2), + "live joiner must be the placement candidate" + ); + + let peer_client = Arc::new( + hyperbytedb::adapters::cluster::peer_client::PeerClient::new( + 1, + node1.addr.clone(), + node1.membership.clone(), + Arc::new( + hyperbytedb::adapters::cluster::replication_log::ReplicationLog::open( + dir.path().join("repl-place"), + ) + .unwrap(), + ), + 2, + 8192, + 8, + 8 * 1024 * 1024, + ), + ); + let metadata: Arc = node1.metadata.clone(); + let wal: Arc = node1.wal.clone(); + let query_port: Arc = node1.query_port.clone(); + let key = MeasurementKey::new("p1db", "autogen", "cpu"); + let outcome = stage_region_transfer_data( + &peer_client, + &metadata, + &wal, + Some(&query_port), + 1, + &key, + ®ion, + 2, + 10_000, + ) + .await + .expect("stage onto joiner"); + assert!( + outcome.verified(), + "stage must confirm every exported point: exported={} applied={}", + outcome.exported, + outcome.applied + ); + assert!( + outcome.exported >= 1, + "pre-existing measurement must copy at least one point" + ); + + apply_add_peer_on_nodes(&[&node1, &node2], "p1db", "autogen", "cpu", 2).await; + + for node in [&node1, &node2] { + let map = node.shard_map.snapshot().await.unwrap(); + let placed = map + .space("p1db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .expect("region after AddPeer"); + assert!( + placed.peers.contains(&2), + "node {} map missing joiner peer: {:?}", + node.node_id, + placed.peers + ); + assert_eq!(placed.primary, 1, "AddPeer must not move primary"); + assert_eq!( + placed.peers.len(), + 2, + "effective RF = min(configured=3, members=2)" + ); + } + + let dest_series = node2 + .metadata + .list_series_ids("p1db", "autogen", "cpu") + .await + .unwrap(); + assert!( + dest_series.contains(&expected_sid), + "joiner metadata must hold the series for the copied range: {dest_series:?}" + ); +} + +/// Points written to `node`'s own WAL at or after `from_seq`. +/// +/// A forwarded write leaves nothing here — the coordinator relays it to the +/// primary — so this is what separates "node 2 accepted the write" from +/// "node 2 proxied it". +async fn local_wal_points( + node: &ShardedTestNode, + from_seq: u64, +) -> Vec { + use hyperbytedb::ports::wal::WalPort; + let wal: std::sync::Arc = node.wal.clone(); + wal.read_from(from_seq) + .await + .expect("read wal") + .into_iter() + .flat_map(|(_, entry)| entry.points) + .collect() +} + +/// P1.4/P1.5: a joiner that is only a replica must not be credited as a pass — +/// the primary has to move onto it and writes have to apply locally there. +/// +/// Replica-only state is asserted first (write to node 2 forwards, node 2's WAL +/// stays empty), then the primary is placed and the same write lands locally. +#[tokio::test] +#[serial(chdb)] +async fn joiner_takes_region_primary_and_applies_writes_locally() { + use hyperbytedb::application::shard_scheduler::{primary_counts, primary_placement_candidate}; + use hyperbytedb::application::shard_transfer::run_region_transfer; + use hyperbytedb::ports::metadata::MetadataPort; + use hyperbytedb::ports::points_sink::PointsSinkPort; + use hyperbytedb::ports::query::QueryPort; + use hyperbytedb::ports::wal::WalPort; + use std::sync::Arc; + + let dir = tempfile::tempdir().unwrap(); + let opts = || ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }; + let node1 = start_sharded_single_node(dir.path(), opts()).await; + let client = reqwest::Client::new(); + create_db(&client, &node1.url, "p14db").await; + let resp = write_line( + &client, + &node1.url, + "p14db", + "cpu,host=solo value=1 1000000000", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + flush_node(&node1).await; + + let region = node1 + .shard_map + .snapshot() + .await + .unwrap() + .space("p14db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after write"); + assert_eq!(region.primary, 1); + + let node2 = start_sharded_joiner(dir.path(), &node1, 2, &opts()).await; + create_db(&client, &node2.url, "p14db").await; + install_shard_map_from_peer(&node1, &node2).await; + + // --- P1.3 midpoint: node 2 is a committed replica, primary is still 1. --- + let peer_client = Arc::new( + hyperbytedb::adapters::cluster::peer_client::PeerClient::new( + 1, + node1.addr.clone(), + node1.membership.clone(), + Arc::new( + hyperbytedb::adapters::cluster::replication_log::ReplicationLog::open( + dir.path().join("repl-p14"), + ) + .unwrap(), + ), + 2, + 8192, + 8, + 8 * 1024 * 1024, + ), + ); + let metadata: Arc = node1.metadata.clone(); + let wal: Arc = node1.wal.clone(); + let query_port: Arc = node1.query_port.clone(); + let sink: Arc = node1.points_sink.clone(); + let key = MeasurementKey::new("p14db", "autogen", "cpu"); + + hyperbytedb::application::shard_transfer::stage_region_transfer_data( + &peer_client, + &metadata, + &wal, + Some(&query_port), + 1, + &key, + ®ion, + 2, + 10_000, + ) + .await + .expect("stage onto joiner"); + apply_add_peer_on_nodes(&[&node1, &node2], "p14db", "autogen", "cpu", 2).await; + + // P1.5: replica-only is not a pass. A write aimed at node 2 forwards to the + // primary and leaves node 2's own WAL untouched. + let replica_seq = node2.wal.last_sequence().await.unwrap(); + let resp = write_line( + &client, + &node2.url, + "p14db", + "cpu,host=solo value=2 1000000001", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + assert!( + local_wal_points(&node2, replica_seq + 1).await.is_empty(), + "replica-only joiner must forward the write, not apply it locally" + ); + + // --- P1.4: place the primary on the joiner. --- + let staged = node1 + .shard_map + .snapshot() + .await + .unwrap() + .space("p14db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after AddPeer"); + let counts = primary_counts(&node1.shard_map.snapshot().await.unwrap()); + assert_eq!( + primary_placement_candidate(&staged, &counts, &[1, 2]), + Some(2), + "the joiner holds no primaries and must be chosen to take this one" + ); + + // Transfer-then-commit: rows move and verify before ownership changes. + run_region_transfer( + &peer_client, + &metadata, + &wal, + Some(&query_port), + Some(&sink), + 1, + &key, + &staged, + 2, + 10_000, + false, + ) + .await + .expect("verified handoff to the joiner"); + let nodes = [node1, node2]; + promote_region_primary_on_all_nodes(&nodes, "p14db", "autogen", "cpu", 2).await; + + for node in &nodes { + let placed = node + .shard_map + .snapshot() + .await + .unwrap() + .space("p14db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after TransferPrimary"); + assert_eq!( + placed.primary, 2, + "node {} map must show the joiner as primary", + node.node_id + ); + assert!( + placed.peers.contains(&1), + "the previous owner stays a replica: {:?}", + placed.peers + ); + } + + // P1.4: the write now applies on node 2 itself. + let primary_seq = nodes[1].wal.last_sequence().await.unwrap(); + let resp = write_line( + &client, + &nodes[1].url, + "p14db", + "cpu,host=solo value=3 1000000002", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + let landed = local_wal_points(&nodes[1], primary_seq + 1).await; + assert!( + landed.iter().any(|p| p.measurement == "cpu"), + "joiner is primary but the write did not reach its own WAL: {landed:?}" + ); +} + +/// P1.6: the 1->2 case must not be a special case. A 4th process joining a +/// 3-node RF-complete cluster gets a replica slot (no region is under RF, so an +/// existing peer steps aside) and then a primary, exactly like the joiner in +/// `joiner_takes_region_primary_and_applies_writes_locally`. +#[tokio::test] +#[serial(chdb)] +async fn fourth_node_joining_rf_complete_cluster_takes_a_region() { + use hyperbytedb::application::shard_scheduler::{ + idle_member_replica_swap, primary_counts, primary_placement_candidate, region_memberships, + }; + use hyperbytedb::application::shard_transfer::{ + run_region_transfer, stage_region_transfer_data, + }; + use hyperbytedb::ports::metadata::MetadataPort; + use hyperbytedb::ports::points_sink::PointsSinkPort; + use hyperbytedb::ports::query::QueryPort; + use hyperbytedb::ports::wal::WalPort; + use std::sync::Arc; + + let dir = tempfile::tempdir().unwrap(); + let opts = || ShardedClusterOptions { + sharding: hyperbytedb::config::ShardingConfig { + enabled: true, + replication_factor: 3, + scatter_peer_timeout_ms: 500, + scatter_max_peer_attempts: 3, + ..Default::default() + }, + ..Default::default() + }; + let nodes = start_sharded_three_node_cluster(dir.path(), opts()).await; + bootstrap_region_on_all_nodes(&nodes, "p16db", "autogen", "cpu", vec![1, 2, 3], 1).await; + + let client = reqwest::Client::new(); + for node in &nodes { + create_db(&client, &node.url, "p16db").await; + } + let resp = write_line( + &client, + &nodes[0].url, + "p16db", + "cpu,host=rfc value=7 1000000000", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + flush_node(&nodes[0]).await; + + let region = nodes[0] + .shard_map + .snapshot() + .await + .unwrap() + .space("p16db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after write"); + assert_eq!(region.peers.len(), 3, "cluster starts RF-complete"); + + let node4 = start_sharded_joiner(dir.path(), &nodes[0], 4, &opts()).await; + create_db(&client, &node4.url, "p16db").await; + install_shard_map_from_peer(&nodes[0], &node4).await; + + // No region is under RF, so replica *addition* correctly declines; the + // newcomer has to displace a replica instead. + assert_eq!( + hyperbytedb::application::shard_scheduler::live_replica_candidate( + ®ion, + &[1, 2, 3, 4], + 3 + ), + None, + "RF is already satisfied; AddPeer must not fire" + ); + let memberships = region_memberships(&nodes[0].shard_map.snapshot().await.unwrap()); + let displaced = + idle_member_replica_swap(®ion, 4, &memberships).expect("newcomer takes a replica slot"); + assert_ne!(displaced, region.primary, "the primary is never displaced"); + + let peer_client = Arc::new( + hyperbytedb::adapters::cluster::peer_client::PeerClient::new( + 1, + nodes[0].addr.clone(), + nodes[0].membership.clone(), + Arc::new( + hyperbytedb::adapters::cluster::replication_log::ReplicationLog::open( + dir.path().join("repl-p16"), + ) + .unwrap(), + ), + 2, + 8192, + 8, + 8 * 1024 * 1024, + ), + ); + let metadata: Arc = nodes[0].metadata.clone(); + let wal: Arc = nodes[0].wal.clone(); + let query_port: Arc = nodes[0].query_port.clone(); + let sink: Arc = nodes[0].points_sink.clone(); + let key = MeasurementKey::new("p16db", "autogen", "cpu"); + + // P1.3 shape: stage before the map commits the newcomer as a peer. + let staged = stage_region_transfer_data( + &peer_client, + &metadata, + &wal, + Some(&query_port), + 1, + &key, + ®ion, + 4, + 10_000, + ) + .await + .expect("stage onto the newcomer"); + assert!( + staged.verified() && staged.exported >= 1, + "pre-existing rows must reach the newcomer: exported={} applied={}", + staged.exported, + staged.applied + ); + + let all = [&nodes[0], &nodes[1], &nodes[2], &node4]; + apply_move_peer_on_nodes(&all, "p16db", "autogen", "cpu", displaced, 4).await; + + let swapped = nodes[0] + .shard_map + .snapshot() + .await + .unwrap() + .space("p16db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after MovePeer"); + assert!(swapped.peers.contains(&4), "newcomer is a committed peer"); + assert!( + !swapped.peers.contains(&displaced), + "displaced peer released its slot" + ); + assert_eq!(swapped.peers.len(), 3, "RF is unchanged by the swap"); + + // P1.4 shape: MovePeer appended the newcomer, so it is now the newest peer + // and the primary-placement rule hands it the region. + let counts = primary_counts(&nodes[0].shard_map.snapshot().await.unwrap()); + assert_eq!( + primary_placement_candidate(&swapped, &counts, &[1, 2, 3, 4]), + Some(4), + "the newcomer holds no primaries and must take this one" + ); + + run_region_transfer( + &peer_client, + &metadata, + &wal, + Some(&query_port), + Some(&sink), + 1, + &key, + &swapped, + 4, + 10_000, + false, + ) + .await + .expect("verified handoff to the newcomer"); + + promote_region_primary_on_nodes(&all, "p16db", "autogen", "cpu", 4).await; + + for node in all { + let placed = node + .shard_map + .snapshot() + .await + .unwrap() + .space("p16db", "autogen", "cpu") + .and_then(|s| s.regions.first()) + .cloned() + .expect("region after TransferPrimary"); + assert_eq!( + placed.primary, 4, + "node {} must see the newcomer as primary", + node.node_id + ); + } + + let seq = node4.wal.last_sequence().await.unwrap(); + let resp = write_line( + &client, + &node4.url, + "p16db", + "cpu,host=rfc value=8 1000000001", + ) + .await; + assert_eq!(resp.status(), reqwest::StatusCode::NO_CONTENT); + let landed = local_wal_points(&node4, seq + 1).await; + assert!( + landed.iter().any(|p| p.measurement == "cpu"), + "newcomer is primary but the write did not reach its own WAL: {landed:?}" + ); +} + async fn wait_for_database(client: &reqwest::Client, url: &str, db: &str) { for _ in 0..100 { let resp = query_sql(client, url, db, "SHOW DATABASES").await;