From f7b550a79a0af8383f9ae624a3511499bfd5f919 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Fri, 4 Sep 2026 22:53:07 +0100 Subject: [PATCH 1/9] Allow a one-member sharded cluster to own regions. Bootstrap now takes RF = min(configured, live members), so n=1 gets peers=[self] even when configured RF is 3. --- hyperbytedb/src/application/shard_routing.rs | 50 ++++++++++++++++++- hyperbytedb/tests/common/sharding_cluster.rs | 11 ++++ .../tests/sharding_cluster_integration.rs | 47 +++++++++++++++++ 3 files changed, 107 insertions(+), 1 deletion(-) 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/tests/common/sharding_cluster.rs b/hyperbytedb/tests/common/sharding_cluster.rs index b9d56bb..02f537c 100644 --- a/hyperbytedb/tests/common/sharding_cluster.rs +++ b/hyperbytedb/tests/common/sharding_cluster.rs @@ -290,6 +290,17 @@ pub async fn start_sharded_node( } } +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, diff --git a/hyperbytedb/tests/sharding_cluster_integration.rs b/hyperbytedb/tests/sharding_cluster_integration.rs index 577ab0c..93a37f3 100644 --- a/hyperbytedb/tests/sharding_cluster_integration.rs +++ b/hyperbytedb/tests/sharding_cluster_integration.rs @@ -750,6 +750,53 @@ 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"); +} + 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; From 8bdbd9937686035dce530567657d58d364f87c31 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Fri, 4 Sep 2026 23:01:20 +0100 Subject: [PATCH 2/9] Install the committed shard map on join before region movement. A joiner's map_version must match the cluster's before any peer-set change, so staging cannot run against a lagging epoch. --- .../src/adapters/cluster/sync_client.rs | 45 +++++++++ .../adapters/sharding/rocksdb_shard_map.rs | 87 ++++++++++++++++ .../src/application/cluster/bootstrap.rs | 6 +- hyperbytedb/src/application/runtime/mod.rs | 3 + .../src/application/shard_scheduler.rs | 26 +++++ hyperbytedb/src/domain/sharding/types.rs | 14 +++ hyperbytedb/src/ports/sharding.rs | 10 ++ hyperbytedb/tests/common/sharding_cluster.rs | 31 +++++- .../tests/sharding_cluster_integration.rs | 99 +++++++++++++++++++ 9 files changed, 319 insertions(+), 2 deletions(-) diff --git a/hyperbytedb/src/adapters/cluster/sync_client.rs b/hyperbytedb/src/adapters/cluster/sync_client.rs index 2f92af1..48797dd 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,9 @@ 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. + self.sync_shard_map(&peer_addr).await?; let mut applied = 0u64; if let Some(ref shard_map) = self.shard_map { @@ -328,6 +333,23 @@ 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; + 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 +509,29 @@ impl SyncClient { } } +/// GET `/internal/shard/map` from `peer_addr` and return the committed snapshot. +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( 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_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index 4c92d7e..22432e8 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -1169,6 +1169,24 @@ 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 its committed +/// `map_version` matches the cluster's. Staging rows against a lagging map +/// would apply under the wrong epoch / peer set. +#[must_use] +pub fn joiner_map_caught_up(cluster_map_version: u64, joiner_map_version: u64) -> bool { + joiner_map_version == cluster_map_version +} + +/// 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, @@ -1888,6 +1906,14 @@ mod tests { assert!(failover_watermark_safe(0, 0)); } + #[test] + fn joiner_map_catchup_blocks_movement_until_versions_match() { + assert!(joiner_map_caught_up(3, 3)); + assert!(!joiner_map_caught_up(3, 0)); + assert!(!joiner_map_caught_up(3, 2)); + assert!(!joiner_map_caught_up(3, 4)); + } + #[tokio::test] #[serial_test::serial(chdb)] async fn try_failover_proposes_transfer_primary() { 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 02f537c..b5e5ca1 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,39 @@ 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); +} + 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(); diff --git a/hyperbytedb/tests/sharding_cluster_integration.rs b/hyperbytedb/tests/sharding_cluster_integration.rs index 93a37f3..8be489b 100644 --- a/hyperbytedb/tests/sharding_cluster_integration.rs +++ b/hyperbytedb/tests/sharding_cluster_integration.rs @@ -797,6 +797,105 @@ async fn one_member_cluster_owns_region_after_first_write() { 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 + ) + ); +} + 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; From 25ead68f8b66e1b8b67c8a7d1f0f60e344501caf Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Sun, 6 Sep 2026 08:02:15 +0100 Subject: [PATCH 3/9] Stage existing region data onto a live joiner (P1.3) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A joining node became a member but never received data for measurements that already existed: `try_split` clones the parent peer set and `try_rebalance` only moves a primary among peers already in the region, so an existing region's peer set never grew to include a new process. Add `ShardMapOp::AddPeer` and a scheduler placement step that stages the region onto a live Active member before committing it as a peer. Staging first is the contract here — the heal path's commit-then-stage ordering would publish a peer that holds no rows. Apply refuses `AddPeer` on a region carrying outstanding transfer debt, on a duplicate peer, and on a stale epoch. `ShardRehomeRequest.stage` lets the scheduler ask a remote primary to run the same staging push when the leader is not the region's primary; the rehome handler skips `complete_region_transfer` for a staging push since ownership does not change. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- .../src/adapters/http/shard_handlers.rs | 81 ++++---- .../src/application/shard_scheduler.rs | 194 ++++++++++++++++++ hyperbytedb/src/domain/sharding/ops.rs | 150 ++++++++++++++ hyperbytedb/src/domain/sharding/transfer.rs | 4 + hyperbytedb/tests/common/sharding_cluster.rs | 62 ++++++ .../tests/sharding_cluster_integration.rs | 153 ++++++++++++++ 6 files changed, 605 insertions(+), 39 deletions(-) diff --git a/hyperbytedb/src/adapters/http/shard_handlers.rs b/hyperbytedb/src/adapters/http/shard_handlers.rs index 3f11735..2f57282 100644 --- a/hyperbytedb/src/adapters/http/shard_handlers.rs +++ b/hyperbytedb/src/adapters/http/shard_handlers.rs @@ -561,20 +561,8 @@ 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(); - } + // Joiner staging: dest may not have the measurement catalog row yet. + // apply_transfer_push registers it via prepare_batch_metadata. } else { let Some(ctx) = state.shard_routing.as_ref() else { return sharding_disabled(); @@ -710,19 +698,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 +736,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/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index 22432e8..ed08002 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -9,6 +9,7 @@ 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, @@ -362,6 +363,10 @@ impl ShardScheduler { } } + if let Err(e) = self.try_place_live_member(&space.key, region).await { + tracing::debug!(error = %e, region_id = region.region_id, "live placement skipped"); + } + if let Err(e) = self.try_rebalance(&space.key, region, &hb).await { tracing::debug!(error = %e, region_id = region.region_id, "rebalance skipped"); } @@ -939,6 +944,88 @@ 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<(), HyperbytedbError> { + // 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(()); + }; + + let joiner_addr = { + let m = self.membership.read().await; + let ids: Vec = m.active_peers(0).into_iter().map(|n| n.node_id).collect(); + live_replica_candidate(region, &ids, self.config.replication_factor) + .and_then(|id| m.get_node(id).map(|n| (id, n.addr.clone()))) + }; + let Some((joiner, addr)) = joiner_addr else { + return Ok(()); + }; + + let cluster_ver = self.shard_map.snapshot().await?.map_version; + let joiner_ver = fetch_peer_map_version(pc.http_client(), &addr).await?; + if !joiner_map_caught_up(cluster_ver, joiner_ver) { + tracing::debug!( + region_id = region.region_id, + joiner, + cluster_ver, + joiner_ver, + "skip live placement: joiner map not caught up" + ); + return Ok(()); + } + + 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(()) + } + async fn try_rebalance( &self, key: &MeasurementKey, @@ -1177,6 +1264,28 @@ pub fn joiner_map_caught_up(cluster_map_version: u64, joiner_map_version: u64) - joiner_map_version == cluster_map_version } +/// 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() +} + /// Read a peer's committed `map_version` via `/internal/shard/map`. pub async fn fetch_peer_map_version( client: &reqwest::Client, @@ -1263,6 +1372,52 @@ 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 + .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(()) +} + +/// 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 @@ -1914,6 +2069,45 @@ mod tests { 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 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); + } + #[tokio::test] #[serial_test::serial(chdb)] async fn try_failover_proposes_transfer_primary() { 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/tests/common/sharding_cluster.rs b/hyperbytedb/tests/common/sharding_cluster.rs index b5e5ca1..dd02da3 100644 --- a/hyperbytedb/tests/common/sharding_cluster.rs +++ b/hyperbytedb/tests/common/sharding_cluster.rs @@ -319,6 +319,68 @@ pub async fn install_shard_map_from_peer(from: &ShardedTestNode, onto: &ShardedT 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 start_sharded_single_node(dir: &Path, opts: ShardedClusterOptions) -> ShardedTestNode { let chdb_dir = dir.join("chdb-shared"); std::fs::create_dir_all(&chdb_dir).unwrap(); diff --git a/hyperbytedb/tests/sharding_cluster_integration.rs b/hyperbytedb/tests/sharding_cluster_integration.rs index 8be489b..4e90963 100644 --- a/hyperbytedb/tests/sharding_cluster_integration.rs +++ b/hyperbytedb/tests/sharding_cluster_integration.rs @@ -896,6 +896,159 @@ async fn joiner_map_version_matches_before_region_movement() { ); } +/// 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:?}" + ); +} + 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; From ece4034438558a9e6d82ca5f9a50408761a993e6 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Sun, 6 Sep 2026 08:19:58 +0100 Subject: [PATCH 4/9] Place region primaries on joined members (P1.4-P1.6) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A joiner that received region data was still only a replica: nothing moved a primary onto it, so writes kept landing on the original owner and the new process added HA but no write capacity. Add two placement steps to the scheduler tick, both stage-then-commit so an empty node never owns a range: - `try_place_primary` hands a region's primary to its newest peer while that peer carries strictly fewer primaries than the current one. The "newest peer" bias is what makes it terminate — once the joiner is primary the rule no longer applies, so ownership cannot oscillate back, and the count guard stops every region stampeding onto one node. - `try_place_idle_member` covers 3->4, where no region is under RF and `AddPeer` correctly declines. An existing replica steps aside for the newest member via `MovePeer`, which appends it and so makes it the next primary-placement candidate. `try_rebalance` scored a peer with no heartbeat row as 0 bytes, which made any freshly staged replica look infinitely lighter than the primary and handed it ownership on the strength of missing telemetry. Unmeasured peers are now skipped rather than treated as empty. Integration coverage asserts the contract rather than the midpoint: a write aimed at the replica-only joiner forwards and leaves its WAL empty, and the same write applies locally once the primary has moved. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- .../src/application/shard_scheduler.rs | 455 +++++++++++++++++- hyperbytedb/tests/common/sharding_cluster.rs | 39 ++ .../tests/sharding_cluster_integration.rs | 405 ++++++++++++++++ 3 files changed, 892 insertions(+), 7 deletions(-) diff --git a/hyperbytedb/src/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index ed08002..a616b1c 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -16,7 +16,9 @@ use crate::application::shard_transfer::{ }; 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; @@ -367,6 +369,14 @@ impl ShardScheduler { tracing::debug!(error = %e, region_id = region.region_id, "live placement skipped"); } + if let Err(e) = self.try_place_idle_member(&space.key, region).await { + tracing::debug!(error = %e, region_id = region.region_id, "idle placement skipped"); + } + + if let Err(e) = self.try_place_primary(&space.key, region).await { + tracing::debug!(error = %e, region_id = region.region_id, "primary placement skipped"); + } + if let Err(e) = self.try_rebalance(&space.key, region, &hb).await { tracing::debug!(error = %e, region_id = region.region_id, "rebalance skipped"); } @@ -1026,22 +1036,200 @@ impl ShardScheduler { Ok(()) } + /// Give a member that joined an already-replicated cluster a replica slot. + /// + /// `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 + /// newcomer before `MovePeer` commits — same stage-then-commit ordering, so + /// no peer is ever published holding nothing. `MovePeer` appends the + /// newcomer, which makes it the region's newest peer and therefore the next + /// primary-placement candidate. + async fn try_place_idle_member( + &self, + key: &MeasurementKey, + region: &ShardRegion, + ) -> Result<(), HyperbytedbError> { + let Some(pc) = self.peer_client.as_ref() else { + return Ok(()); + }; + + let newcomer = { + 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((newcomer, addr)) = newcomer else { + return Ok(()); + }; + + let map = self.shard_map.snapshot().await?; + let memberships = region_memberships(&map); + let Some(displaced) = idle_member_replica_swap(region, newcomer, &memberships) else { + return Ok(()); + }; + + let joiner_ver = fetch_peer_map_version(pc.http_client(), &addr).await?; + if !joiner_map_caught_up(map.map_version, joiner_ver) { + tracing::debug!( + region_id = region.region_id, + newcomer, + cluster_ver = map.map_version, + joiner_ver, + "skip idle placement: newcomer map not caught up" + ); + return Ok(()); + } + + 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, + newcomer, + 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, + newcomer, + ) + .await?; + } + + tracing::info!( + region_id = region.region_id, + displaced, + newcomer, + "moving region replica onto newly joined member" + ); + self.propose(ShardMapOp::MovePeer { + key: key.clone(), + region_id: region.region_id, + from_peer: displaced, + to_peer: newcomer, + epoch: region.epoch, + }) + .await?; + counter!("hyperbytedb_shard_idle_placements_total").increment(1); + Ok(()) + } + + /// 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<(), HyperbytedbError> { + let Some(pc) = self.peer_client.as_ref() else { + return Ok(()); + }; + + let map = self.shard_map.snapshot().await?; + 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(()); + }; + + // 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(()) + } + 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 { @@ -1286,6 +1474,96 @@ pub fn live_replica_candidate( .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 +} + +/// 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 `newcomer`, or `None` +/// when this region is already well placed. +/// +/// An RF-complete region never triggers `AddPeer`, so a process joining an +/// already-replicated cluster (3 nodes at RF 3, add a 4th) would own nothing +/// forever — every region legitimately has its full replica count, and pure +/// load balancing has no reason to disturb them. The bias is deliberately +/// one-way: only the newcomer displaces anyone, and only a peer carrying +/// strictly more region memberships than it. +/// +/// `newcomer` must be the newest active member — the caller establishes that, +/// and it is what makes the swap converge. Because the identity is stable until +/// membership itself changes, the node just displaced cannot turn around and +/// reclaim the slot: it is not the newcomer. Passing an arbitrary node here +/// forfeits that guarantee and can trade a slot back and forth. +#[must_use] +pub fn idle_member_replica_swap( + region: &ShardRegion, + newcomer: u64, + memberships: &HashMap, +) -> Option { + if region.transfer_outstanding() || region.peers.contains(&newcomer) { + return None; + } + let newcomer_load = memberships.get(&newcomer).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 > newcomer_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, @@ -2108,6 +2386,169 @@ mod tests { 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("newcomer takes a slot"); + assert_ne!(displaced, 1, "the primary must never be displaced"); + assert!(displaced == 2 || displaced == 3); + } + + /// The swap must converge. `newcomer` 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 newcomer 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 newcomer 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() { diff --git a/hyperbytedb/tests/common/sharding_cluster.rs b/hyperbytedb/tests/common/sharding_cluster.rs index dd02da3..808109d 100644 --- a/hyperbytedb/tests/common/sharding_cluster.rs +++ b/hyperbytedb/tests/common/sharding_cluster.rs @@ -381,6 +381,34 @@ pub async fn apply_add_peer_on_nodes( } } +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(); @@ -477,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 4e90963..0493aa1 100644 --- a/hyperbytedb/tests/sharding_cluster_integration.rs +++ b/hyperbytedb/tests/sharding_cluster_integration.rs @@ -1049,6 +1049,411 @@ async fn joiner_receives_existing_region_as_replica() { ); } +/// 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; From 0f786b746309be9d2775deaefcb1c23ff4045dc3 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Sun, 6 Sep 2026 08:24:42 +0100 Subject: [PATCH 5/9] Re-read the region before each placement transfer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both placement steps ran against the region as it looked when the tick snapshotted the map. A placement earlier in the same tick invalidates that: `AddPeer` bumps the epoch and appends a peer, so the next step would stage an entire region copy and only then have its proposal rejected by the epoch CAS — wasted transfer, repeated every tick until the snapshots happened to line up. Re-read the region from a fresh snapshot inside each step, and skip it if the region has since been split, merged, or dropped. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- .../src/application/shard_scheduler.rs | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/hyperbytedb/src/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index a616b1c..7099796 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -1065,7 +1065,13 @@ impl ShardScheduler { return Ok(()); }; + // 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(()); + }; let memberships = region_memberships(&map); let Some(displaced) = idle_member_replica_swap(region, newcomer, &memberships) else { return Ok(()); @@ -1152,7 +1158,13 @@ impl ShardScheduler { return Ok(()); }; + // 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(()); + }; let counts = primary_counts(&map); let active: Vec = { let m = self.membership.read().await; @@ -1486,6 +1498,21 @@ pub fn primary_counts(map: &ShardMap) -> HashMap { 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 { From cb1463227ffde2771a5b0666e70c53f238801ab3 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Sun, 6 Sep 2026 09:08:06 +0100 Subject: [PATCH 6/9] Fix review blockers in region placement (#1-#7) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Seven findings from reviewing the branch, four of which could corrupt data or wedge a cluster. **Rollup destinations were eligible for placement.** `apply_transfer_push` has no idempotence guard, so a region staged onto a joiner whose map op then fails is staged again next tick. `ReplacingMergeTree` measurements collapse the redelivery; `SummingMergeTree` rollup destinations sum it and are corrupted permanently. `enqueue_reconciliation` already excluded these spaces for exactly this reason — all three placement paths now share that guard, failing closed when the lock is contended. **`AddPeer` had no upgrade gate.** `ClusterRequest` is serialised into the Raft log, so a committed op an un-upgraded voter cannot deserialise wedges it. `ClearVerified` set the precedent; `AddPeer` now has the same gate and defaults **off** for the release that introduces it. **`sync_shard_map` could move the map backwards.** The sync peer is whichever Active node came first out of a HashMap — not the leader, not necessarily the most-applied node. Installing an older snapshot dropped ops Raft never redelivers, because `replace_map` bypasses the state machine and leaves `last_applied` untouched; every later op then failed StaleEpoch and the node silently stopped owning regions it held data for. Install only a strictly newer map. **A missing shard map aborted the whole join.** `/internal/shard/map` only exists on peers with sharding enabled, so a 404 from an older or unsharded peer failed `join_and_sync` *before* WAL catch-up; after the retries the node went Active having synced nothing. It is best-effort now: with no map there are no regions to move, and the leader's own catch-up check still gates placement onto this node. Also: `try_place_live_member` never re-read its region, so it staged a full copy before failing the epoch CAS (the sibling fix missed it); and the steps after a placement ran on the region as it looked *before* that placement committed — `try_rebalance` would spend an entire region copy on a stale peer list. A placement that commits now yields the region. Placement was reachable only from `tick()` and no test called it, so deleting those calls left the suite green. Adds a harness with a mock peer that drives a real tick, covering both the placement and the upgrade gate. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- .../src/adapters/cluster/sync_client.rs | 57 ++- .../src/application/shard_scheduler.rs | 345 ++++++++++++++++-- hyperbytedb/src/config.rs | 32 ++ 3 files changed, 405 insertions(+), 29 deletions(-) diff --git a/hyperbytedb/src/adapters/cluster/sync_client.rs b/hyperbytedb/src/adapters/cluster/sync_client.rs index 48797dd..8280d40 100644 --- a/hyperbytedb/src/adapters/cluster/sync_client.rs +++ b/hyperbytedb/src/adapters/cluster/sync_client.rs @@ -182,7 +182,21 @@ 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. - self.sync_shard_map(&peer_addr).await?; + // + // 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 { @@ -341,6 +355,25 @@ impl SyncClient { }; 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, @@ -510,6 +543,16 @@ 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, @@ -554,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/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index 7099796..22abbd2 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -176,6 +176,21 @@ impl ShardScheduler { self } + /// 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(), ()); @@ -365,16 +380,33 @@ impl ShardScheduler { } } - if let Err(e) = self.try_place_live_member(&space.key, region).await { - tracing::debug!(error = %e, region_id = region.region_id, "live placement skipped"); - } - - if let Err(e) = self.try_place_idle_member(&space.key, region).await { - tracing::debug!(error = %e, region_id = region.region_id, "idle placement skipped"); + // 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 let Err(e) = self.try_place_primary(&space.key, region).await { - tracing::debug!(error = %e, region_id = region.region_id, "primary placement skipped"); + if placed { + self.release_operator(region.region_id).await; + continue; } if let Err(e) = self.try_rebalance(&space.key, region, &hb).await { @@ -961,11 +993,24 @@ impl ShardScheduler { &self, key: &MeasurementKey, region: &ShardRegion, - ) -> Result<(), HyperbytedbError> { + ) -> 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(()); + 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 joiner_addr = { @@ -975,10 +1020,10 @@ impl ShardScheduler { .and_then(|id| m.get_node(id).map(|n| (id, n.addr.clone()))) }; let Some((joiner, addr)) = joiner_addr else { - return Ok(()); + return Ok(false); }; - let cluster_ver = self.shard_map.snapshot().await?.map_version; + let cluster_ver = map.map_version; let joiner_ver = fetch_peer_map_version(pc.http_client(), &addr).await?; if !joiner_map_caught_up(cluster_ver, joiner_ver) { tracing::debug!( @@ -988,7 +1033,7 @@ impl ShardScheduler { joiner_ver, "skip live placement: joiner map not caught up" ); - return Ok(()); + return Ok(false); } if region.primary == self.node_id { @@ -1033,7 +1078,7 @@ impl ShardScheduler { }) .await?; counter!("hyperbytedb_shard_live_placements_total").increment(1); - Ok(()) + Ok(true) } /// Give a member that joined an already-replicated cluster a replica slot. @@ -1049,9 +1094,13 @@ impl ShardScheduler { &self, key: &MeasurementKey, region: &ShardRegion, - ) -> Result<(), HyperbytedbError> { + ) -> Result { + if self.is_rollup_dest(key).await { + return Ok(false); + } + let Some(pc) = self.peer_client.as_ref() else { - return Ok(()); + return Ok(false); }; let newcomer = { @@ -1062,7 +1111,7 @@ impl ShardScheduler { .map(|n| (n.node_id, n.addr.clone())) }; let Some((newcomer, addr)) = newcomer else { - return Ok(()); + return Ok(false); }; // Re-read the region: an `AddPeer` earlier in this same tick leaves the @@ -1070,11 +1119,11 @@ impl ShardScheduler { // 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(()); + return Ok(false); }; let memberships = region_memberships(&map); let Some(displaced) = idle_member_replica_swap(region, newcomer, &memberships) else { - return Ok(()); + return Ok(false); }; let joiner_ver = fetch_peer_map_version(pc.http_client(), &addr).await?; @@ -1086,7 +1135,7 @@ impl ShardScheduler { joiner_ver, "skip idle placement: newcomer map not caught up" ); - return Ok(()); + return Ok(false); } if region.primary == self.node_id { @@ -1138,7 +1187,7 @@ impl ShardScheduler { }) .await?; counter!("hyperbytedb_shard_idle_placements_total").increment(1); - Ok(()) + Ok(true) } /// Hand a region's primary to its newest peer so a joiner starts taking @@ -1153,9 +1202,13 @@ impl ShardScheduler { &self, key: &MeasurementKey, region: &ShardRegion, - ) -> Result<(), HyperbytedbError> { + ) -> Result { + if self.is_rollup_dest(key).await { + return Ok(false); + } + let Some(pc) = self.peer_client.as_ref() else { - return Ok(()); + return Ok(false); }; // Re-read the region: a peer placement earlier in this same tick leaves @@ -1163,7 +1216,7 @@ impl ShardScheduler { // `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(()); + return Ok(false); }; let counts = primary_counts(&map); let active: Vec = { @@ -1171,7 +1224,7 @@ impl ShardScheduler { 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(()); + return Ok(false); }; // Source from whoever holds the authoritative rows. Exporting from the @@ -1222,7 +1275,7 @@ impl ShardScheduler { }) .await?; counter!("hyperbytedb_shard_primary_placements_total").increment(1); - Ok(()) + Ok(true) } async fn try_rebalance( @@ -2991,6 +3044,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() { @@ -3054,6 +3128,52 @@ 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 { + use axum::routing::{get, post}; + let app = axum::Router::new() + .route( + "/internal/shard/map", + get(move || async move { + 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 + } + struct ReconcileTestHarness { _dir: tempfile::TempDir, shard_map: Arc, @@ -3075,6 +3195,177 @@ 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>>, + } + + impl PlacementTestHarness { + async fn new(add_peer_proposals_enabled: bool) -> 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; + + // Region owned by this node, one peer short of effective RF. + let key = MeasurementKey::new("db", "autogen", "cpu"); + shard_map + .apply_op(ShardMapOp::BootstrapMeasurement { + key, + region: sample_region_peers(vec![1], 1), + }) + .await + .unwrap(); + let map_version = shard_map.snapshot().await.unwrap().map_version; + + let joiner_addr = spawn_mock_peer(map_version).await; + { + 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, + }); + m.add_node(NodeInfo { + node_id: 2, + addr: joiner_addr, + state: NodeState::Active, + joined_at: 1, + 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, + } + } + } + + /// 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).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:?}" + ); + } + + /// 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).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..e3493dc 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 on mixed-version clusters: nodes running builds older than the + /// `AddPeer` op cannot decode it from the Raft log, and a committed entry + /// they cannot deserialize wedges them. + #[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,15 @@ fn default_transfer_clear_proposals_enabled() -> bool { true } +/// Off for the release that introduces `AddPeer`. A leader that emits a new +/// shard-map op the moment it upgrades commits a Raft entry the rest of a +/// half-upgraded fleet cannot deserialize. Operators turn this on once every +/// node runs a build that knows the op; the default flips in a later release, +/// the same way `transfer_clear_proposals_enabled` did. +fn default_add_peer_proposals_enabled() -> bool { + false +} + impl HyperbytedbConfig { /// Validate cross-field constraints after Figment merge. pub fn validate(&self) -> Result<(), String> { @@ -1243,4 +1259,20 @@ mod replicate_body_limit_tests { 64 * 1024 * 1024 ); } + + /// A new Raft-log op must stay off by default for the release that adds it. + /// A leader that proposes one mid-rolling-restart commits an entry the + /// un-upgraded followers cannot deserialize. + #[test] + fn new_shard_map_ops_are_off_by_default() { + let s = super::ShardingConfig::default(); + assert!( + !s.add_peer_proposals_enabled, + "AddPeer is new in this release and must not be proposed until the fleet can decode it" + ); + assert!( + s.transfer_clear_proposals_enabled, + "ClearVerified shipped in an earlier release and is past its upgrade window" + ); + } } From 6987d7ef47125efdcaee6933f2d05b14ae044293 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Sun, 6 Sep 2026 11:17:20 +0100 Subject: [PATCH 7/9] Stop placement stalling on one lagging candidate (#9-#12) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Placement could stop cluster-wide, permanently and quietly. `joiner_map_caught_up` required exact `map_version` equality and `live_replica_candidate` always returned the lowest-id non-peer, so a single node that never matched was chosen first, failed the probe, and was chosen again the next tick. No other candidate was ever tried, for any region, and the only trace was a `debug!` line. - Catch-up accepts a candidate that is *ahead*. The cluster version comes from a snapshot read just before the probe, so a node that applied an op in between is legitimately ahead and was being rejected for it. Behind is still refused — that is what the gate is for. - `live_replica_candidates` returns every eligible member in id order and placement takes the first that passes its probe, instead of stopping at a candidate chosen before any probe ran. - A region that goes `PLACEMENT_STALL_WARN_TICKS` ticks with candidates but none caught up now warns and increments a counter. A stall that only shows up at debug level is a stall nobody finds. - `/internal/shard/map` serialises every space and region to answer one `u64`, and it was probed once per region per tick — O(regions x map size) bytes on every tick of a join. Probes are cached per peer for the tick. Also corrects the `handle_shard_transfer` staging comment, which still claimed a local-measurement guard that P1.3 deliberately removed. H.6 (whether the gate should compare the Raft applied index instead) is answered in the plan, no code change: the index is already exposed at /internal/raft/metrics, but `replace_map` installs a map without advancing it, so an applied-index gate would reject exactly the joiner it should admit. Fix belongs with Phase 2's snapshot-based map recovery. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- .../src/adapters/http/shard_handlers.rs | 16 +- .../src/application/shard_scheduler.rs | 304 +++++++++++++++--- 2 files changed, 267 insertions(+), 53 deletions(-) diff --git a/hyperbytedb/src/adapters/http/shard_handlers.rs b/hyperbytedb/src/adapters/http/shard_handlers.rs index 2f57282..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,8 +565,6 @@ pub async fn handle_shard_transfer( .into_response(); } drop(m); - // Joiner staging: dest may not have the measurement catalog row yet. - // apply_transfer_push registers it via prepare_batch_metadata. } else { let Some(ctx) = state.shard_routing.as_ref() else { return sharding_disabled(); diff --git a/hyperbytedb/src/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index 22abbd2..0d1eec6 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -33,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). @@ -119,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)] @@ -157,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)] @@ -176,6 +191,43 @@ 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 @@ -285,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() @@ -1013,28 +1067,53 @@ impl ShardScheduler { return Ok(false); }; - let joiner_addr = { + 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_candidate(region, &ids, self.config.replication_factor) - .and_then(|id| m.get_node(id).map(|n| (id, n.addr.clone()))) + 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() }; - let Some((joiner, addr)) = joiner_addr else { + 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 joiner_ver = fetch_peer_map_version(pc.http_client(), &addr).await?; - if !joiner_map_caught_up(cluster_ver, joiner_ver) { + 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, - joiner, + candidate = id, cluster_ver, - joiner_ver, - "skip live placement: joiner map not caught up" + candidate_ver = ver, + "candidate map not caught up" ); - return Ok(false); } + 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( @@ -1126,7 +1205,9 @@ impl ShardScheduler { return Ok(false); }; - let joiner_ver = fetch_peer_map_version(pc.http_client(), &addr).await?; + let joiner_ver = self + .cached_peer_map_version(newcomer, &addr, pc.http_client()) + .await?; if !joiner_map_caught_up(map.map_version, joiner_ver) { tracing::debug!( region_id = region.region_id, @@ -1509,12 +1590,46 @@ 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 its committed -/// `map_version` matches the cluster's. Staging rows against a lagging map +/// 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 + 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 @@ -2420,11 +2535,14 @@ mod tests { } #[test] - fn joiner_map_catchup_blocks_movement_until_versions_match() { + 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)); - assert!(!joiner_map_caught_up(3, 4)); + // 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 { @@ -2459,6 +2577,36 @@ mod tests { 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); @@ -3131,17 +3279,23 @@ mod tests { /// 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 { + 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 || async move { - axum::Json(serde_json::json!({ - "map_version": map_version, - "next_region_id": 2, - "spaces": [], - })) + 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( @@ -3171,7 +3325,7 @@ mod tests { tokio::spawn(async move { let _ = axum::serve(listener, app).await; }); - addr + (addr, probes) } struct ReconcileTestHarness { @@ -3201,10 +3355,19 @@ mod tests { _dir: tempfile::TempDir, scheduler: ShardScheduler, proposals: Arc>>, + /// `/internal/shard/map` hit counts, keyed by mock peer node id. + probes: HashMap>, } impl PlacementTestHarness { - async fn new(add_peer_proposals_enabled: bool) -> Self { + /// `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; @@ -3263,18 +3426,21 @@ mod tests { .unwrap(); tokio::time::sleep(std::time::Duration::from_millis(150)).await; - // Region owned by this node, one peer short of effective RF. - let key = MeasurementKey::new("db", "autogen", "cpu"); - shard_map - .apply_op(ShardMapOp::BootstrapMeasurement { - key, - region: sample_region_peers(vec![1], 1), - }) - .await - .unwrap(); + // 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 joiner_addr = spawn_mock_peer(map_version).await; + let mut probes = HashMap::new(); { let mut m = bootstrap.membership.write().await; m.add_node(NodeInfo { @@ -3285,14 +3451,19 @@ mod tests { last_heartbeat: 0, needs_sync: false, }); - m.add_node(NodeInfo { - node_id: 2, - addr: joiner_addr, - state: NodeState::Active, - joined_at: 1, - 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( @@ -3332,8 +3503,13 @@ mod tests { _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 @@ -3341,7 +3517,7 @@ mod tests { #[tokio::test] #[serial_test::serial(chdb)] async fn tick_places_a_live_joiner_as_a_region_peer() { - let harness = PlacementTestHarness::new(true).await; + let harness = PlacementTestHarness::new(true, &[(2, 0)], 1).await; harness.scheduler.tick_once_for_test().await.unwrap(); let ops = harness.proposals.lock().unwrap(); assert!( @@ -3351,12 +3527,48 @@ mod tests { ); } + /// 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).await; + let harness = PlacementTestHarness::new(false, &[(2, 0)], 1).await; harness.scheduler.tick_once_for_test().await.unwrap(); let ops = harness.proposals.lock().unwrap(); assert!( From e96f7a078b6c77d8801f754a03ec2e7d215f21cd Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Sun, 6 Sep 2026 22:11:07 +0100 Subject: [PATCH 8/9] Say what idle-member placement actually triggers on (#8) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `try_place_idle_member` was documented as giving "a member that joined an already-replicated cluster" a replica slot, which reads as join-triggered. It is not. The target is the active member with the greatest `joined_at` — whoever joined last, with no recency window, so "last" may mean months ago. The rule fires whenever that member is under-loaded against some region's non-primary peer, which means a stable-but-unbalanced cluster starts moving replicas on the first tick after an upgrade rather than on any join. The behaviour is deliberately unchanged. It is bounded and convergent, and every copy is staged and verified before the map changes, so this is a scheduling surprise rather than a correctness risk — acceptable while sharding is beta. Gating it was considered and rejected for now; the doc says so, and says to revisit before sharding graduates. Renames `newcomer` to `latest_member` throughout, so the identifier stops implying recency too. No logic change. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- .../src/application/shard_scheduler.rs | 90 +++++++++++-------- 1 file changed, 55 insertions(+), 35 deletions(-) diff --git a/hyperbytedb/src/application/shard_scheduler.rs b/hyperbytedb/src/application/shard_scheduler.rs index 0d1eec6..633da8f 100644 --- a/hyperbytedb/src/application/shard_scheduler.rs +++ b/hyperbytedb/src/application/shard_scheduler.rs @@ -1160,15 +1160,31 @@ impl ShardScheduler { Ok(true) } - /// Give a member that joined an already-replicated cluster a replica slot. + /// 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 - /// newcomer before `MovePeer` commits — same stage-then-commit ordering, so + /// target before `MovePeer` commits — same stage-then-commit ordering, so /// no peer is ever published holding nothing. `MovePeer` appends the - /// newcomer, which makes it the region's newest peer and therefore the next + /// 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, @@ -1182,14 +1198,14 @@ impl ShardScheduler { return Ok(false); }; - let newcomer = { + 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((newcomer, addr)) = newcomer else { + let Some((latest_member, addr)) = latest_member else { return Ok(false); }; @@ -1201,20 +1217,20 @@ impl ShardScheduler { return Ok(false); }; let memberships = region_memberships(&map); - let Some(displaced) = idle_member_replica_swap(region, newcomer, &memberships) else { + let Some(displaced) = idle_member_replica_swap(region, latest_member, &memberships) else { return Ok(false); }; let joiner_ver = self - .cached_peer_map_version(newcomer, &addr, pc.http_client()) + .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, - newcomer, + latest_member, cluster_ver = map.map_version, joiner_ver, - "skip idle placement: newcomer map not caught up" + "skip idle placement: rebalance target map not caught up" ); return Ok(false); } @@ -1228,7 +1244,7 @@ impl ShardScheduler { self.node_id, key, region, - newcomer, + latest_member, self.max_points_per_request, ) .await?; @@ -1248,7 +1264,7 @@ impl ShardScheduler { region.primary, key, region, - newcomer, + latest_member, ) .await?; } @@ -1256,14 +1272,14 @@ impl ShardScheduler { tracing::info!( region_id = region.region_id, displaced, - newcomer, - "moving region replica onto newly joined member" + 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: newcomer, + to_peer: latest_member, epoch: region.epoch, }) .await?; @@ -1695,38 +1711,42 @@ pub fn region_memberships(map: &ShardMap) -> HashMap { counts } -/// Non-primary peer that should yield its replica slot to `newcomer`, or `None` -/// when this region is already well placed. +/// 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. /// -/// An RF-complete region never triggers `AddPeer`, so a process joining an -/// already-replicated cluster (3 nodes at RF 3, add a 4th) would own nothing -/// forever — every region legitimately has its full replica count, and pure -/// load balancing has no reason to disturb them. The bias is deliberately -/// one-way: only the newcomer displaces anyone, and only a peer 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. /// -/// `newcomer` must be the newest active member — the caller establishes that, -/// and it is what makes the swap converge. Because the identity is stable until -/// membership itself changes, the node just displaced cannot turn around and -/// reclaim the slot: it is not the newcomer. Passing an arbitrary node here -/// forfeits that guarantee and can trade a slot back and forth. +/// 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, - newcomer: u64, + latest_member: u64, memberships: &HashMap, ) -> Option { - if region.transfer_outstanding() || region.peers.contains(&newcomer) { + if region.transfer_outstanding() || region.peers.contains(&latest_member) { return None; } - let newcomer_load = memberships.get(&newcomer).copied().unwrap_or(0); + 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 > newcomer_load) + .filter(|(load, _)| *load > latest_member_load) .max() .map(|(_, peer)| peer) } @@ -2716,12 +2736,12 @@ mod tests { 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("newcomer takes a slot"); + 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. `newcomer` stays the same node for as long as + /// 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] @@ -2732,7 +2752,7 @@ mod tests { 4, ®ion_memberships(&map_of(vec![before.clone()])), ) - .expect("first pass moves the newcomer in"); + .expect("first pass moves the latest member in"); let mut after = before.clone(); after.peers.retain(|p| *p != displaced); @@ -2740,7 +2760,7 @@ mod tests { assert_eq!( idle_member_replica_swap(&after, 4, ®ion_memberships(&map_of(vec![after.clone()]))), None, - "second pass with the same newcomer must be a no-op" + "second pass with the same latest member must be a no-op" ); } From 425f80a9593f26c44080d42eee47787099463786 Mon Sep 17 00:00:00 2001 From: Austin Barrington Date: Mon, 7 Sep 2026 08:21:19 +0100 Subject: [PATCH 9/9] Default the AddPeer gate on MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Shipped it off, which was the wrong call. The gate is an opt-out for a rolling upgrade, not an opt-in for the feature — `transfer_clear_proposals_ enabled` guards the identical hazard and defaults on, and defaulting off leaves region placement inert for everyone who enables sharding and never finds the flag. The hazard is bounded: a follower on a build that does not know `AddPeer` rejects the whole append RPC (axum fails the body before the handler runs), so it stops replicating and the leader can lose commit quorum for the duration. It resolves as soon as every node is upgraded, and ingest does not go through Raft. Operators disable the gate for the rollout window. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019Te8YUXjjLssk3hxwoE5Db --- hyperbytedb/src/config.rs | 42 +++++++++++++++++++-------------------- 1 file changed, 21 insertions(+), 21 deletions(-) diff --git a/hyperbytedb/src/config.rs b/hyperbytedb/src/config.rs index e3493dc..392a6e1 100644 --- a/hyperbytedb/src/config.rs +++ b/hyperbytedb/src/config.rs @@ -80,9 +80,9 @@ pub struct ShardingConfig { #[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 on mixed-version clusters: nodes running builds older than the - /// `AddPeer` op cannot decode it from the Raft log, and a committed entry - /// they cannot deserialize wedges them. + /// 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, } @@ -167,13 +167,19 @@ fn default_transfer_clear_proposals_enabled() -> bool { true } -/// Off for the release that introduces `AddPeer`. A leader that emits a new -/// shard-map op the moment it upgrades commits a Raft entry the rest of a -/// half-upgraded fleet cannot deserialize. Operators turn this on once every -/// node runs a build that knows the op; the default flips in a later release, -/// the same way `transfer_clear_proposals_enabled` did. +/// 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 { - false + true } impl HyperbytedbConfig { @@ -1260,19 +1266,13 @@ mod replicate_body_limit_tests { ); } - /// A new Raft-log op must stay off by default for the release that adds it. - /// A leader that proposes one mid-rolling-restart commits an entry the - /// un-upgraded followers cannot deserialize. + /// 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 new_shard_map_ops_are_off_by_default() { + fn shard_map_op_gates_default_on() { let s = super::ShardingConfig::default(); - assert!( - !s.add_peer_proposals_enabled, - "AddPeer is new in this release and must not be proposed until the fleet can decode it" - ); - assert!( - s.transfer_clear_proposals_enabled, - "ClearVerified shipped in an earlier release and is past its upgrade window" - ); + assert!(s.add_peer_proposals_enabled); + assert!(s.transfer_clear_proposals_enabled); } }