diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 20b4b1ebf..e7d274cad 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -68,6 +68,8 @@ use std::fmt::Display; use std::io::Cursor; use std::io::Write; use std::path::PathBuf; +#[cfg(test)] +use std::sync::atomic::AtomicBool; use std::sync::{ Arc, atomic::{AtomicUsize, Ordering}, @@ -872,11 +874,17 @@ fn activation_pool_meta_lock_error(err: rustfs_lock::LockError) -> Error { pub(crate) struct PoolRebalanceActivationFence { pool_meta_guard: rustfs_lock::NamespaceLockGuard, rebalance_meta_guard: rustfs_lock::NamespaceLockGuard, + #[cfg(test)] + forced_lost: Arc, } impl PoolRebalanceActivationFence { pub(crate) fn ensure_held(&self) -> Result<()> { - if self.pool_meta_guard.is_lock_lost() || self.rebalance_meta_guard.is_lock_lost() { + #[cfg(test)] + let forced_lost = self.forced_lost.load(Ordering::Acquire); + #[cfg(not(test))] + let forced_lost = false; + if forced_lost || self.pool_meta_guard.is_lock_lost() || self.rebalance_meta_guard.is_lock_lost() { return Err(Error::other("activation lock lost before metadata commit or worker admission")); } @@ -887,6 +895,11 @@ impl PoolRebalanceActivationFence { opts.add_namespace_lock_guard(&self.pool_meta_guard); opts.add_namespace_lock_guard(&self.rebalance_meta_guard); } + + #[cfg(test)] + fn force_lost_for_test(&self) { + self.forced_lost.store(true, Ordering::Release); + } } pub(crate) async fn acquire_pool_rebalance_activation_locks(pool: Arc) -> Result @@ -911,6 +924,8 @@ where Ok(PoolRebalanceActivationFence { pool_meta_guard, rebalance_meta_guard, + #[cfg(test)] + forced_lost: Arc::new(AtomicBool::new(false)), }) } @@ -921,6 +936,91 @@ pub(crate) enum PoolActivationStartKind { Decommission, } +#[cfg(test)] +struct PoolActivationDurableSaveBarrierState { + pool_key: usize, + arrived: tokio::sync::Notify, + release: tokio::sync::Notify, +} + +#[cfg(test)] +static POOL_ACTIVATION_DURABLE_SAVE_BARRIER: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +pub(crate) struct PoolActivationDurableSaveBarrier { + state: Arc, +} + +#[cfg(test)] +fn pool_activation_test_pool_key(pool: &Arc) -> usize { + Arc::as_ptr(pool).cast::<()>() as usize +} + +#[cfg(test)] +impl PoolActivationDurableSaveBarrier { + pub(crate) fn install(pool: &Arc) -> Self { + let state = Arc::new(PoolActivationDurableSaveBarrierState { + pool_key: pool_activation_test_pool_key(pool), + arrived: tokio::sync::Notify::new(), + release: tokio::sync::Notify::new(), + }); + let mut barrier = POOL_ACTIVATION_DURABLE_SAVE_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("pool activation durable save barrier should not be poisoned"); + assert!(barrier.is_none(), "pool activation durable save barrier must be unique"); + *barrier = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) async fn wait_until_paused(&self) { + tokio::time::timeout(std::time::Duration::from_secs(30), self.state.arrived.notified()) + .await + .expect("activation should reach the post-durable-save barrier"); + } + + pub(crate) fn release_after_fence_loss(&self) { + self.state.release.notify_one(); + } +} + +#[cfg(test)] +impl Drop for PoolActivationDurableSaveBarrier { + fn drop(&mut self) { + self.state.release.notify_one(); + let mut barrier = POOL_ACTIVATION_DURABLE_SAVE_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("pool activation durable save barrier should not be poisoned"); + if barrier.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *barrier = None; + } + } +} + +#[cfg(test)] +pub(crate) async fn pause_pool_activation_after_durable_save(pool: &Arc, fence: &PoolRebalanceActivationFence) { + let pool_key = pool_activation_test_pool_key(pool); + let barrier = { + let mut barrier = POOL_ACTIVATION_DURABLE_SAVE_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("pool activation durable save barrier should not be poisoned"); + if barrier.as_ref().is_some_and(|state| state.pool_key == pool_key) { + barrier.take() + } else { + None + } + }; + if let Some(barrier) = barrier { + barrier.arrived.notify_one(); + barrier.release.notified().await; + fence.force_lost_for_test(); + } +} + #[cfg(test)] struct PoolActivationStartProbeState { kind: PoolActivationStartKind, @@ -1869,7 +1969,7 @@ impl PoolMeta { async fn save_no_lock_with_activation_fence( &self, pools: Vec>, - activation_fence: &PoolRebalanceActivationFence, + activation_fence: PoolRebalanceActivationFence, ) -> Result<()> where S: EcstoreObjectIO, @@ -1878,16 +1978,38 @@ impl PoolMeta { if data.is_empty() { return Ok(()); } + let mut pools = pools.into_iter(); + let Some(canonical_pool) = pools.next() else { + return Ok(()); + }; + let mut opts = ObjectOptions { + max_parity: true, + no_lock: true, + ..Default::default() + }; + activation_fence.add_namespace_lock_fence(&mut opts); + activation_fence.ensure_held()?; + #[cfg(test)] + let barrier_pool = canonical_pool.clone(); + save_config_with_opts(canonical_pool, POOL_META_NAME, data.clone(), &opts).await?; + #[cfg(test)] + pause_pool_activation_after_durable_save(&barrier_pool, &activation_fence).await; + + // Pool zero is canonical. Once its save succeeds, later writes only + // replicate committed state and must not reuse the admission fence. + drop(activation_fence); for pool in pools { - let mut opts = ObjectOptions { - max_parity: true, - no_lock: true, - ..Default::default() - }; - activation_fence.add_namespace_lock_fence(&mut opts); - activation_fence.ensure_held()?; - save_config_with_opts(pool, POOL_META_NAME, data.clone(), &opts).await?; - activation_fence.ensure_held()?; + save_config_with_opts( + pool, + POOL_META_NAME, + data.clone(), + &ObjectOptions { + max_parity: true, + no_lock: true, + ..Default::default() + }, + ) + .await?; } Ok(()) @@ -2789,7 +2911,7 @@ impl ECStore { activation_fence.ensure_held()?; latest_pool_meta - .save_no_lock_with_activation_fence(self.pools.clone(), &activation_fence) + .save_no_lock_with_activation_fence(self.pools.clone(), activation_fence) .await?; { let mut pool_meta = self.pool_meta.write().await; @@ -4585,6 +4707,54 @@ mod tests { use crate::bucket::replication::{ReplicationState, ReplicationStatusType}; use serde::Serialize; + #[tokio::test] + #[serial_test::serial] + async fn decommission_activation_replicates_commit_after_post_save_fence_loss() { + let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await; + let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]); + let start_store = Arc::clone(&store); + let start_task = tokio::spawn(async move { + start_store + .save_current_pool_meta_for_decommission_start( + &[0], + vec![( + 0, + PoolSpaceInfo { + free: 50, + total: 100, + used: 50, + }, + )], + Vec::new(), + ) + .await + }); + + barrier.wait_until_paused().await; + barrier.release_after_fence_loss(); + tokio::time::timeout(std::time::Duration::from_secs(30), start_task) + .await + .expect("decommission activation should finish after its canonical commit") + .expect("decommission activation task should not panic") + .expect("post-commit fence loss must not report the committed activation as failed"); + + let local = store.pool_meta.read().await; + assert!(pool_meta_has_active_decommission(&local)); + drop(local); + + for pool in &store.pools { + let mut persisted = PoolMeta::default(); + persisted + .load_no_lock(pool.clone()) + .await + .expect("every pool should retain readable committed decommission metadata"); + assert!( + pool_meta_has_active_decommission(&persisted), + "every pool must adopt the canonical committed activation" + ); + } + } + #[test] fn ensure_pool_not_left_in_cmdline_after_decommission_allows_active_pool() { assert!(ensure_pool_not_left_in_cmdline_after_decommission(0, "http://node{1...4}/disk{1...4}", false).is_ok()); diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index d788ef4dd..3b9304617 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -232,10 +232,15 @@ where fence.add_namespace_lock_fence(&mut opts); fence.ensure_held()?; } + #[cfg(test)] + let barrier_pool = pool.clone(); merged.save_with_opts(pool, opts).await?; + #[cfg(test)] if let Some(fence) = activation_fence { - fence.ensure_held()?; + crate::core::pools::pause_pool_activation_after_durable_save(&barrier_pool, fence).await; } + // With an activation fence, a successful save is the commit point. Lease + // loss after this point cannot make the durable activation uncommitted. Ok(()) } @@ -1148,7 +1153,7 @@ impl ECStore { #[cfg(test)] mod tests { use super::*; - use crate::core::pools::{PoolActivationStartKind, PoolActivationStartProbe}; + use crate::core::pools::{PoolActivationDurableSaveBarrier, PoolActivationStartKind, PoolActivationStartProbe}; use crate::object_api::NamespaceLockFence; use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated}; @@ -1192,6 +1197,121 @@ mod tests { assert!(cancel.is_cancelled()); } + #[tokio::test] + #[serial_test::serial] + async fn rebalance_activation_adopts_commit_after_post_save_fence_loss() { + let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(None).await; + set_rebalance_disk_stats_override_for_test( + store.id, + vec![ + DiskStat { + total_space: 100, + available_space: 0, + }, + DiskStat { + total_space: 100, + available_space: 100, + }, + ], + ); + let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]); + let start_store = Arc::clone(&store); + let start_task = tokio::spawn(async move { + start_store + .init_rebalance_start(vec!["post-commit-fence-loss".to_string()]) + .await + }); + + barrier.wait_until_paused().await; + barrier.release_after_fence_loss(); + let rebalance_id = tokio::time::timeout(std::time::Duration::from_secs(30), start_task) + .await + .expect("rebalance activation should finish after its durable commit") + .expect("rebalance activation task should not panic") + .expect("post-commit fence loss must not report the committed activation as failed"); + + let local_meta = store.rebalance_meta.read().await; + let local = local_meta + .as_ref() + .expect("the committed rebalance metadata should be installed locally"); + assert_eq!(local.id, rebalance_id); + assert!(is_rebalance_conflicting_with_decommission(local)); + drop(local_meta); + + let mut persisted = RebalanceMeta::new(); + persisted + .load(store.pools[0].clone()) + .await + .expect("the committed rebalance metadata should remain readable"); + assert_eq!(persisted.id, rebalance_id); + assert!(is_rebalance_conflicting_with_decommission(&persisted)); + } + + #[tokio::test] + #[serial_test::serial] + async fn rebalance_worker_admission_installs_committed_candidate_after_post_save_fence_loss() { + let rebalance_id = "post-save-worker-admission"; + let active = RebalanceMeta { + id: rebalance_id.to_string(), + percent_free_goal: 0.5, + pool_stats: vec![ + RebalanceStats { + participating: true, + init_free_space: 50, + init_capacity: 100, + buckets: vec!["completed-pool-bucket".to_string()], + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + RebalanceStats { + participating: true, + init_capacity: 100, + buckets: vec!["active-pool-bucket".to_string()], + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + ], + ..Default::default() + }; + let (_temp_dirs, store, _other_store) = crate::services::rebalance::test_two_pool_stores(Some(active)).await; + let barrier = PoolActivationDurableSaveBarrier::install(&store.pools[0]); + let start_store = Arc::clone(&store); + let start_task = tokio::spawn(async move { start_store.start_rebalance().await }); + + barrier.wait_until_paused().await; + barrier.release_after_fence_loss(); + tokio::time::timeout(std::time::Duration::from_secs(30), start_task) + .await + .expect("rebalance worker admission should finish after its durable commit") + .expect("rebalance worker admission task should not panic") + .expect("post-commit fence loss must not reject the committed worker candidate"); + + let local_meta = store.rebalance_meta.read().await; + let local = local_meta + .as_ref() + .expect("the committed worker candidate should remain installed locally"); + assert_eq!(local.id, rebalance_id); + assert_eq!(local.pool_stats[0].info.status, RebalStatus::Completed); + if let Some(cancel) = local.cancel.as_ref() { + cancel.cancel(); + } + drop(local_meta); + + let mut persisted = RebalanceMeta::new(); + persisted + .load(store.pools[0].clone()) + .await + .expect("the committed worker candidate should remain readable"); + assert_eq!(persisted.id, rebalance_id); + assert_eq!(persisted.pool_stats[0].info.status, RebalStatus::Completed); + } + async fn assert_real_activation_start_race(paused_kind: PoolActivationStartKind) { let (_temp_dirs, rebalance_store, decommission_store) = crate::services::rebalance::test_two_pool_stores(None).await; set_rebalance_disk_stats_override_for_test( diff --git a/crates/ecstore/src/services/rebalance/runtime.rs b/crates/ecstore/src/services/rebalance/runtime.rs index 5f03b46c2..ed8da595f 100644 --- a/crates/ecstore/src/services/rebalance/runtime.rs +++ b/crates/ecstore/src/services/rebalance/runtime.rs @@ -218,10 +218,6 @@ impl ECStore { cancel_tx.cancel(); return Err(err); } - if let Err(err) = activation_fence.ensure_held() { - cancel_tx.cancel(); - return Err(err); - } let mut rebalance_meta = self.rebalance_meta.write().await; let Some(meta) = rebalance_meta.as_mut() else { cancel_tx.cancel(); @@ -237,7 +233,7 @@ impl ECStore { return Err(err); } } - if let Err(err) = activation_fence.ensure_held() { + if !must_persist && let Err(err) = activation_fence.ensure_held() { let mut rebalance_meta = self.rebalance_meta.write().await; rollback_local_rebalance_worker_activation(rebalance_meta.as_mut(), expected_id.as_ref(), &rx); return Err(err);