From 8b0e96c3141ae01e6985cfc661d84234b97c2646 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 03:16:35 +0800 Subject: [PATCH] test(ecstore): exercise real rebalance fences --- crates/ecstore/src/core/pools.rs | 102 ++++++++- crates/ecstore/src/core/sets.rs | 14 +- crates/ecstore/src/object_api/types.rs | 87 +++++++- .../ecstore/src/services/rebalance/control.rs | 187 +++++++++++++++- .../ecstore/src/services/rebalance/entry.rs | 202 ++++++++++++++++++ crates/ecstore/src/services/rebalance/mod.rs | 60 ++++++ 6 files changed, 635 insertions(+), 17 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 91cf38f66..854a3854b 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -19,8 +19,8 @@ use crate::bucket::{ LifecycleExpiryConfigs, bucket_lifecycle_audit::LcEventSrc, bucket_lifecycle_ops::{ - LifecycleOps, apply_expiry_rule_for_data_movement, apply_expiry_rule_in, eval_action_from_lifecycle, - lifecycle_delete_all_versions_blocked_by_replication, + LifecycleOps, apply_expiry_on_transitioned_object, apply_expiry_rule_for_data_movement, apply_expiry_rule_in, + eval_action_from_lifecycle, lifecycle_delete_all_versions_blocked_by_replication, }, get_expiry_configs, lifecycle::IlmAction, @@ -914,6 +914,79 @@ where }) } +#[cfg(test)] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) enum PoolActivationStartKind { + Rebalance, + Decommission, +} + +#[cfg(test)] +struct PoolActivationStartProbeState { + kind: PoolActivationStartKind, + attempted: std::sync::atomic::AtomicBool, + notify: tokio::sync::Notify, +} + +#[cfg(test)] +static POOL_ACTIVATION_START_PROBES: std::sync::OnceLock>>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +pub(crate) struct PoolActivationStartProbe { + state: Arc, +} + +#[cfg(test)] +impl PoolActivationStartProbe { + pub(crate) fn install(kind: PoolActivationStartKind) -> Self { + let state = Arc::new(PoolActivationStartProbeState { + kind, + attempted: std::sync::atomic::AtomicBool::new(false), + notify: tokio::sync::Notify::new(), + }); + POOL_ACTIVATION_START_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("pool activation start probe should not be poisoned") + .push(Arc::clone(&state)); + Self { state } + } + + pub(crate) async fn wait_until_attempted(&self) { + while !self.state.attempted.load(Ordering::Acquire) { + self.state.notify.notified().await; + } + } +} + +#[cfg(test)] +impl Drop for PoolActivationStartProbe { + fn drop(&mut self) { + let mut probes = POOL_ACTIVATION_START_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("pool activation start probe should not be poisoned"); + probes.retain(|state| !Arc::ptr_eq(state, &self.state)); + } +} + +#[cfg(test)] +pub(crate) fn observe_pool_activation_start_attempt(kind: PoolActivationStartKind) { + let probes = POOL_ACTIVATION_START_PROBES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("pool activation start probe should not be poisoned") + .iter() + .filter(|state| state.kind == kind) + .cloned() + .collect::>(); + for state in probes { + state.attempted.store(true, Ordering::Release); + state.notify.notify_one(); + } +} + fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) { *pool_meta = previous_pool_meta; } @@ -2462,7 +2535,15 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement( let Ok(bucket_incarnation_id) = store.bucket_incarnation_id_from_disk(bucket).await else { return Ok(false); }; - let _ = apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, lock_lost_signal).await; + let _ = match lock_lost_signal { + Some(signal) => { + apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, Some(signal)).await + } + None => { + apply_expiry_on_transitioned_object(store, &object_info, &event, event_source, bucket_incarnation_id) + .await + } + }; } Ok(false) } @@ -2471,7 +2552,12 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement( return Ok(false); } let applied = !apply_actions - || apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, lock_lost_signal).await; + || match lock_lost_signal { + Some(signal) => { + apply_expiry_rule_for_data_movement(store, &event, event_source, &object_info, Some(signal)).await + } + None => apply_expiry_rule_in(store, &event, event_source, &object_info).await, + }; resolve_data_movement_lifecycle_expiry_result(action, apply_actions, applied) } _ => Ok(false), @@ -2573,6 +2659,8 @@ impl ECStore { .first() .cloned() .ok_or_else(|| Error::other("decommission start rebalance metadata load failed: no storage pools available"))?; + #[cfg(test)] + observe_pool_activation_start_attempt(PoolActivationStartKind::Decommission); let activation_fence = acquire_pool_rebalance_activation_locks(rebalance_pool.clone()).await?; let mut rebalance_meta = RebalanceMeta::new(); @@ -4116,7 +4204,11 @@ impl ECStore { validate_start_decommission_request(&indices, self.single_pool())?; self.ensure_decommission_rebalance_idle_after_refresh().await?; - ensure_decommission_start_local_leader(&self.endpoints(), &indices)?; + #[cfg(test)] + let endpoints = self.instance_endpoints().unwrap_or_else(|| self.endpoints()); + #[cfg(not(test))] + let endpoints = self.endpoints(); + ensure_decommission_start_local_leader(&endpoints, &indices)?; for idx in indices.iter().copied() { ensure_valid_decommission_pool_index(self.pools.len(), idx)?; diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 3acf1a705..31682e2e2 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -1228,6 +1228,14 @@ pub(crate) async fn make_local_two_set_sets() -> (Vec, Arc) -> (Vec, Arc) { + make_local_two_set_sets_for_pool_with_ctx(ctx, 0).await +} + +#[cfg(test)] +pub(crate) async fn make_local_two_set_sets_for_pool_with_ctx( + ctx: Arc, + pool_idx: usize, +) -> (Vec, Arc) { use crate::layout::endpoint::Endpoint; use rustfs_lock::client::local::LocalClient; @@ -1243,7 +1251,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc) let temp_dir = tempfile::tempdir().expect("tempdir should be created"); let mut endpoint = Endpoint::try_from(temp_dir.path().to_str().expect("tempdir path should be utf8")) .expect("endpoint should parse"); - endpoint.set_pool_index(0); + endpoint.set_pool_index(pool_idx); endpoint.set_set_index(set_index); endpoint.set_disk_index(disk_index); let disk = new_disk( @@ -1279,7 +1287,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc) 2, 1, set_index, - 0, + pool_idx, endpoints, format.clone(), lockers, @@ -1292,7 +1300,7 @@ pub(crate) async fn make_local_two_set_sets_with_ctx(ctx: Arc) let sets = Arc::new(Sets { id: format.id, disk_set: disk_sets, - pool_idx: 0, + pool_idx, endpoints: PoolEndpoints { legacy: false, set_count: 2, diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 1bbff7a7f..7e004a1f0 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -24,7 +24,7 @@ use crate::storage_api_contracts::{ pub struct NamespaceLockFence { signals: Arc>>, #[cfg(test)] - forced_lost: Arc, + forced_lost: Arc>>, } impl Debug for NamespaceLockFence { @@ -40,13 +40,17 @@ impl NamespaceLockFence { Self { signals: Arc::default(), #[cfg(test)] - forced_lost: Arc::new(std::sync::atomic::AtomicBool::new(false)), + forced_lost: Arc::new(vec![Arc::new(std::sync::atomic::AtomicBool::new(false))]), } } pub(crate) fn is_lock_lost(&self) -> bool { #[cfg(test)] - if self.forced_lost.load(std::sync::atomic::Ordering::Acquire) { + if self + .forced_lost + .iter() + .any(|lost| lost.load(std::sync::atomic::Ordering::Acquire)) + { return true; } self.signals.iter().any(|signal| signal.is_lost()) @@ -62,25 +66,80 @@ impl NamespaceLockFence { } Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned()); #[cfg(test)] - if other.forced_lost.load(std::sync::atomic::Ordering::Acquire) { - self.forced_lost.store(true, std::sync::atomic::Ordering::Release); + if !Arc::ptr_eq(&self.forced_lost, &other.forced_lost) { + Arc::make_mut(&mut self.forced_lost).extend(other.forced_lost.iter().cloned()); } } #[cfg(test)] pub(crate) fn lost_for_test() -> Self { let fence = Self::new(); - fence.forced_lost.store(true, std::sync::atomic::Ordering::Release); + fence.forced_lost[0].store(true, std::sync::atomic::Ordering::Release); fence } #[cfg(test)] pub(crate) fn loss_handle_for_test() -> (Self, Arc) { let fence = Self::new(); - (fence.clone(), Arc::clone(&fence.forced_lost)) + (fence.clone(), Arc::clone(&fence.forced_lost[0])) } } +#[cfg(test)] +static NAMESPACE_LOCK_SIGNAL_TEST_FENCES: std::sync::OnceLock>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +pub(crate) struct NamespaceLockSignalTestFence { + signal_key: usize, +} + +#[cfg(test)] +impl NamespaceLockSignalTestFence { + pub(crate) fn install_with_loss_handle( + signal: &Arc, + loss_handle: Arc, + ) -> Self { + let fence = NamespaceLockFence { + signals: Arc::default(), + forced_lost: Arc::new(vec![loss_handle]), + }; + let signal_key = Arc::as_ptr(signal) as usize; + let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("namespace lock signal test fence should not be poisoned"); + assert!( + !fences.iter().any(|(key, _)| *key == signal_key), + "namespace lock signal test fence must be unique" + ); + fences.push((signal_key, fence)); + Self { signal_key } + } +} + +#[cfg(test)] +impl Drop for NamespaceLockSignalTestFence { + fn drop(&mut self) { + let mut fences = NAMESPACE_LOCK_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("namespace lock signal test fence should not be poisoned"); + fences.retain(|(key, _)| *key != self.signal_key); + } +} + +#[cfg(test)] +pub(crate) fn namespace_lock_signal_test_fence_is_lost(signal: &Arc) -> bool { + NAMESPACE_LOCK_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("namespace lock signal test fence should not be poisoned") + .iter() + .find(|(key, _)| *key == Arc::as_ptr(signal) as usize) + .is_some_and(|(_, fence)| fence.is_lock_lost()) +} + #[derive(Debug)] pub struct ObjectLockConfigSnapshot { store_id: Option, @@ -402,9 +461,23 @@ impl ObjectOptions { } pub(crate) fn add_namespace_lock_lost_signal(&mut self, signal: Arc) { + #[cfg(test)] + let test_fence = NAMESPACE_LOCK_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("namespace lock signal test fence should not be poisoned") + .iter() + .find(|(key, _)| *key == Arc::as_ptr(&signal) as usize) + .map(|(_, fence)| fence.clone()); self.namespace_lock_fence .get_or_insert_with(NamespaceLockFence::new) .add_signal(signal); + #[cfg(test)] + if let Some(test_fence) = test_fence { + self.namespace_lock_fence + .get_or_insert_with(NamespaceLockFence::new) + .extend(&test_fence); + } } pub(crate) fn ensure_namespace_lock_fence(&mut self) { diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index 7e1860c6b..530aac3f7 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -35,6 +35,10 @@ use uuid::Uuid; #[cfg(test)] static FAIL_NEXT_REBALANCE_ACTIVATION_SAVE: std::sync::Mutex> = std::sync::Mutex::new(None); +#[cfg(test)] +static REBALANCE_DISK_STATS_OVERRIDES: std::sync::OnceLock>>> = + std::sync::OnceLock::new(); + #[cfg(test)] pub(super) fn fail_next_rebalance_activation_save_for_test(rebalance_id: &str) { *FAIL_NEXT_REBALANCE_ACTIVATION_SAVE @@ -42,6 +46,24 @@ pub(super) fn fail_next_rebalance_activation_save_for_test(rebalance_id: &str) { .expect("rebalance activation save failure hook should not be poisoned") = Some(rebalance_id.to_string()); } +#[cfg(test)] +fn set_rebalance_disk_stats_override_for_test(store_id: Uuid, disk_stats: Vec) { + REBALANCE_DISK_STATS_OVERRIDES + .get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new())) + .lock() + .expect("rebalance disk stats override should not be poisoned") + .insert(store_id, disk_stats); +} + +#[cfg(test)] +fn take_rebalance_disk_stats_override_for_test(store_id: Uuid) -> Option> { + REBALANCE_DISK_STATS_OVERRIDES + .get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new())) + .lock() + .expect("rebalance disk stats override should not be poisoned") + .remove(&store_id) +} + fn ensure_rebalance_activation_pool_meta_allowed(meta: &PoolMeta) -> Result<()> { if pool_meta_has_active_decommission(meta) { return Err(Error::DecommissionAlreadyRunning); @@ -62,7 +84,14 @@ pub(super) struct RebalanceRunGuard { impl RebalanceRunGuard { pub(super) fn ensure_held(&self, stage: &str) -> Result<()> { - if self.persisted_guard.is_lock_lost() { + #[cfg(test)] + let forced_lost = self + .persisted_guard + .lock_lost_signal() + .is_some_and(|signal| crate::object_api::namespace_lock_signal_test_fence_is_lost(&signal)); + #[cfg(not(test))] + let forced_lost = false; + if self.persisted_guard.is_lock_lost() || forced_lost { return Err(Error::other(format!("rebalance distributed run fence lost during {stage}"))); } Ok(()) @@ -368,6 +397,8 @@ impl ECStore { where S: EcstoreObjectIO + StorageNamespaceLocking, { + #[cfg(test)] + crate::core::pools::observe_pool_activation_start_attempt(crate::core::pools::PoolActivationStartKind::Rebalance); let activation_fence = acquire_pool_rebalance_activation_locks(pool.clone()).await?; let mut pool_meta = PoolMeta::default(); pool_meta.load_no_lock(pool.clone()).await?; @@ -591,6 +622,13 @@ impl ECStore { disk_stats[disk.pool_index as usize].available_space += disk.available_space; } + #[cfg(test)] + if let Some(overridden) = take_rebalance_disk_stats_override_for_test(self.id) { + disk_stats = overridden; + total_cap = disk_stats.iter().map(|stat| stat.total_space).sum(); + total_free = disk_stats.iter().map(|stat| stat.available_space).sum(); + } + let percent_free_goal = percent_free_ratio(total_free, total_cap); validate_rebalance_disk_stats_coverage(&disk_stats)?; @@ -1019,8 +1057,153 @@ impl ECStore { #[cfg(test)] mod tests { use super::*; + use crate::core::pools::{PoolActivationStartKind, PoolActivationStartProbe}; use crate::object_api::NamespaceLockFence; - use crate::set_disk::hermetic_set_disks_isolated; + use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause, hermetic_set_disks_isolated}; + + 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( + rebalance_store.id, + vec![ + DiskStat { + total_space: 100, + available_space: 0, + }, + DiskStat { + total_space: 100, + available_space: 100, + }, + ], + ); + let (first_object, competing_object, competing_kind) = match paused_kind { + PoolActivationStartKind::Rebalance => { + (REBAL_META_NAME, crate::core::pools::POOL_META_NAME, PoolActivationStartKind::Decommission) + } + PoolActivationStartKind::Decommission => { + (crate::core::pools::POOL_META_NAME, REBAL_META_NAME, PoolActivationStartKind::Rebalance) + } + }; + let first_barrier = PutObjectCommitBarrier::install( + crate::disk::RUSTFS_META_BUCKET, + first_object, + PutObjectCommitPause::BeforeQuotaRename, + ); + let competing_barrier = PutObjectCommitBarrier::install( + crate::disk::RUSTFS_META_BUCKET, + competing_object, + PutObjectCommitPause::BeforeQuotaRename, + ); + + let mut rebalance_task; + let mut decommission_task; + if paused_kind == PoolActivationStartKind::Rebalance { + let store = Arc::clone(&rebalance_store); + rebalance_task = + tokio::spawn(async move { store.init_rebalance_start(vec!["bucket".to_string()]).await.map(|_| ()) }); + first_barrier.wait_until_paused().await; + let probe = PoolActivationStartProbe::install(competing_kind); + let store = Arc::clone(&decommission_store); + decommission_task = tokio::spawn(async move { store.start_decommission(vec![0]).await }); + tokio::time::timeout(std::time::Duration::from_secs(15), probe.wait_until_attempted()) + .await + .expect("real decommission start should reach activation lock acquisition"); + } else { + let store = Arc::clone(&decommission_store); + decommission_task = tokio::spawn(async move { store.start_decommission(vec![0]).await }); + first_barrier.wait_until_paused().await; + let probe = PoolActivationStartProbe::install(competing_kind); + let store = Arc::clone(&rebalance_store); + rebalance_task = + tokio::spawn(async move { store.init_rebalance_start(vec!["bucket".to_string()]).await.map(|_| ()) }); + tokio::time::timeout(std::time::Duration::from_secs(15), probe.wait_until_attempted()) + .await + .expect("real rebalance start should reach activation lock acquisition"); + } + + let mut observed_rebalance_result = None; + let mut observed_decommission_result = None; + let competing_reached_commit = tokio::time::timeout(std::time::Duration::from_secs(15), async { + if paused_kind == PoolActivationStartKind::Rebalance { + tokio::select! { + _ = competing_barrier.wait_until_paused() => true, + result = &mut decommission_task => { + observed_decommission_result = Some(result.expect("decommission activation task should not panic")); + false + } + } + } else { + tokio::select! { + _ = competing_barrier.wait_until_paused() => true, + result = &mut rebalance_task => { + observed_rebalance_result = Some(result.expect("rebalance activation task should not panic")); + false + } + } + } + }) + .await + .expect("competing activation should either block to timeout or reach its commit"); + if competing_reached_commit { + competing_barrier.release(); + } + drop(competing_barrier); + first_barrier.release(); + drop(first_barrier); + + let rebalance_result = match observed_rebalance_result { + Some(result) => result, + None => tokio::time::timeout(std::time::Duration::from_secs(15), &mut rebalance_task) + .await + .expect("real rebalance activation should finish") + .expect("rebalance activation task should not panic"), + }; + let decommission_result = match observed_decommission_result { + Some(result) => result, + None => tokio::time::timeout(std::time::Duration::from_secs(15), &mut decommission_task) + .await + .expect("real decommission activation should finish") + .expect("decommission activation task should not panic"), + }; + assert!( + !competing_reached_commit, + "the competing real start reached its commit while the other activation lock was held" + ); + assert_ne!( + rebalance_result.is_ok(), + decommission_result.is_ok(), + "exactly one real activation entry may commit" + ); + + let mut persisted_rebalance = RebalanceMeta::new(); + let rebalance_committed = persisted_rebalance + .load(rebalance_store.pools[0].clone()) + .await + .is_ok_and(|_| is_rebalance_conflicting_with_decommission(&persisted_rebalance)); + let mut persisted_pool = PoolMeta::default(); + persisted_pool + .load_no_lock(rebalance_store.pools[0].clone()) + .await + .expect("persisted pool metadata should remain readable"); + let decommission_committed = pool_meta_has_active_decommission(&persisted_pool); + assert_eq!( + usize::from(rebalance_committed) + usize::from(decommission_committed), + 1, + "at most one persisted activation state may be active" + ); + assert_eq!(rebalance_committed, paused_kind == PoolActivationStartKind::Rebalance); + assert_eq!(decommission_committed, paused_kind == PoolActivationStartKind::Decommission); + } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_and_decommission_starts_commit_at_most_one_side() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], async { + assert_real_activation_start_race(PoolActivationStartKind::Rebalance).await; + assert_real_activation_start_race(PoolActivationStartKind::Decommission).await; + }) + .await; + } #[test] fn rebalance_activation_rejects_persisted_decommission_despite_idle_local_snapshot() { diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 37370415e..271da9133 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -50,6 +50,66 @@ fn ensure_rebalance_entry_active(cancel: &CancellationToken) -> Result<()> { Ok(()) } +#[cfg(test)] +static REBALANCE_RUN_SIGNAL_TEST_FENCES: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +struct RebalanceRunSignalTestFence { + rebalance_id: String, + loss_handle: Arc, +} + +#[cfg(test)] +impl RebalanceRunSignalTestFence { + fn install(rebalance_id: &str) -> Self { + let loss_handle = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let previous = REBALANCE_RUN_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new())) + .lock() + .expect("rebalance run signal test fence should not be poisoned") + .insert(rebalance_id.to_string(), Arc::clone(&loss_handle)); + assert!(previous.is_none(), "rebalance run signal test fence must be unique"); + Self { + rebalance_id: rebalance_id.to_string(), + loss_handle, + } + } + + fn mark_lost(&self) { + self.loss_handle.store(true, std::sync::atomic::Ordering::Release); + } +} + +#[cfg(test)] +impl Drop for RebalanceRunSignalTestFence { + fn drop(&mut self) { + REBALANCE_RUN_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new())) + .lock() + .expect("rebalance run signal test fence should not be poisoned") + .remove(&self.rebalance_id); + } +} + +#[cfg(test)] +fn attach_rebalance_run_signal_test_fence( + rebalance_id: &str, + signal: &Arc, +) -> Option { + let loss_handle = REBALANCE_RUN_SIGNAL_TEST_FENCES + .get_or_init(|| std::sync::Mutex::new(std::collections::HashMap::new())) + .lock() + .expect("rebalance run signal test fence should not be poisoned") + .get(rebalance_id) + .cloned()?; + Some(crate::object_api::NamespaceLockSignalTestFence::install_with_loss_handle( + signal, + loss_handle, + )) +} + impl ECStore { async fn finish_rebalance_entry_after_cleanup( &self, @@ -179,6 +239,10 @@ impl ECStore { ensure_rebalance_entry_active(&cancel)?; let run_guard = self.rebalance_run_guard(rebalance_id.as_ref(), "rebalance entry").await?; let lock_lost_signal = run_guard.lock_lost_signal(); + #[cfg(test)] + let _run_signal_test_fence = lock_lost_signal + .as_ref() + .and_then(|signal| attach_rebalance_run_signal_test_fence(rebalance_id.as_ref(), signal)); let mut rebalanced: usize = 0; let mut expired: usize = 0; @@ -725,12 +789,16 @@ impl ECStore { #[cfg(test)] mod tests { use super::*; + use crate::object_api::PutObjReader; use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; + use crate::set_disk::{PutObjectCommitBarrier, PutObjectCommitPause}; use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; use crate::storage_api_contracts::object::ObjectOperations as _; + use http::HeaderMap; use rustfs_filemeta::{FileInfo, FileMeta}; use std::time::Duration as StdDuration; use time::OffsetDateTime; + use tokio::io::AsyncReadExt; #[tokio::test] async fn rebalance_stats_wait_for_source_cleanup_result() { @@ -983,4 +1051,138 @@ mod tests { "stop should publish the terminal state after entry cleanup drains" ); } + + #[tokio::test] + #[serial_test::serial] + async fn real_rebalance_run_fence_loss_before_target_commit_preserves_target_and_source() { + let rebalance_id = "rebalance-target-commit-fence"; + let (_temp_dirs, store, _unused_store) = crate::services::rebalance::test_two_pool_stores(Some(RebalanceMeta { + id: rebalance_id.to_string(), + percent_free_goal: 1.0, + cancel: Some(CancellationToken::new()), + pool_stats: vec![ + RebalanceStats { + participating: true, + init_capacity: 100, + buckets: vec![crate::disk::RUSTFS_META_BUCKET.to_string()], + info: RebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }, + RebalanceStats::default(), + ], + ..Default::default() + })) + .await; + let bucket = crate::disk::RUSTFS_META_BUCKET; + let object = "rebalance-commit-fence-object"; + let version_id = uuid::Uuid::new_v4(); + let payload = b"rebalance target commit must not survive its run fence".repeat(1024); + let source_set = store.pools[0].get_disks_by_key(object); + let target_set = store.pools[1].get_disks_by_key(object); + let mut writer = PutObjReader::from_vec(payload.clone()); + let source_before = source_set + .put_object( + bucket, + object, + &mut writer, + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect("source version should be written"); + let source_versions = source_set + .load_file_info_versions_exact(bucket, object) + .await + .expect("source metadata should be readable") + .expect("source version should exist"); + let mut file_meta = FileMeta::new(); + for version in &source_versions.versions { + file_meta + .add_version(version.clone()) + .expect("source version should encode into a metacache entry"); + } + let entry = MetaCacheEntry { + name: object.to_string(), + metadata: file_meta.marshal_msg().expect("source metadata should marshal"), + cached: Some(file_meta), + reusable: false, + }; + let run_signal_fence = RebalanceRunSignalTestFence::install(rebalance_id); + let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::BeforeQuotaRename); + let entry_store = Arc::clone(&store); + let entry_set = Arc::clone(&source_set); + let mut entry_task = tokio::spawn(async move { + entry_store + .rebalance_entry( + bucket.to_string(), + 0, + entry, + entry_set, + Arc::new(RebalanceBucketConfigs::default()), + Arc::from(rebalance_id), + CancellationToken::new(), + ) + .await + }); + barrier.wait_until_paused().await; + run_signal_fence.mark_lost(); + barrier.release(); + drop(barrier); + let entry_error = tokio::time::timeout(StdDuration::from_secs(5), &mut entry_task) + .await + .expect("fenced real rebalance entry should finish") + .expect("rebalance entry task should not panic") + .expect_err("lost run fence must reject the target commit"); + assert!(entry_error.to_string().contains("run fence lost")); + let target_error = target_set + .get_object_info( + bucket, + object, + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + ..Default::default() + }, + ) + .await + .expect_err("target version must remain absent after the lost commit fence"); + assert!( + crate::error::is_err_object_not_found(&target_error) || crate::error::is_err_version_not_found(&target_error), + "target version must remain absent after the lost commit fence" + ); + let mut source_body = Vec::new(); + let mut source_after = source_set + .get_object_reader( + bucket, + object, + None, + HeaderMap::new(), + &ObjectOptions { + versioned: true, + version_id: Some(version_id.to_string()), + no_lock: true, + ..Default::default() + }, + ) + .await + .expect("source version must remain readable"); + assert_eq!(source_after.object_info.version_id, source_before.version_id); + assert_eq!(source_after.object_info.data_dir, source_before.data_dir); + assert_eq!(source_after.object_info.mod_time, source_before.mod_time); + assert_eq!(source_after.object_info.size, source_before.size); + assert_eq!(source_after.object_info.etag, source_before.etag); + source_after + .stream + .read_to_end(&mut source_body) + .await + .expect("source body should drain"); + assert_eq!(source_body, payload, "source version must remain byte-identical"); + } } diff --git a/crates/ecstore/src/services/rebalance/mod.rs b/crates/ecstore/src/services/rebalance/mod.rs index f297335c6..4e44727e3 100644 --- a/crates/ecstore/src/services/rebalance/mod.rs +++ b/crates/ecstore/src/services/rebalance/mod.rs @@ -79,5 +79,65 @@ pub(crate) async fn test_store_with_persisted_rebalance_meta( (temp_dirs, store) } +#[cfg(test)] +pub(crate) async fn test_two_pool_stores( + rebalance_meta: Option, +) -> ( + Vec, + std::sync::Arc, + std::sync::Arc, +) { + use crate::core::pools::PoolMeta; + use crate::layout::endpoints::{EndpointServerPools, SetupType}; + + let ctx = std::sync::Arc::new(crate::runtime::instance::InstanceContext::new()); + ctx.update_erasure_type(SetupType::DistErasure).await; + let (mut temp_dirs, first_pool) = + crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 0).await; + let (second_temp_dirs, second_pool) = + crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(std::sync::Arc::clone(&ctx), 1).await; + temp_dirs.extend(second_temp_dirs); + let pools = vec![first_pool, second_pool]; + { + let local_disk_map = ctx.local_disk_map(); + let mut local_disk_map = local_disk_map.write().await; + for pool in &pools { + for set in &pool.disk_set { + for disk in set.disks.read().await.iter().flatten() { + local_disk_map.insert(disk.endpoint().to_string(), Some(disk.clone())); + } + } + } + } + let pool_meta = PoolMeta::new(&pools, &PoolMeta::default()); + pool_meta + .save(pools.clone()) + .await + .expect("baseline pool metadata should be persisted"); + if let Some(meta) = rebalance_meta.as_ref() { + meta.save(pools[0].clone()) + .await + .expect("active rebalance metadata should be persisted"); + } + let endpoint_pools: EndpointServerPools = pools.iter().map(|pool| pool.endpoints.clone()).collect::>().into(); + ctx.set_endpoints(endpoint_pools.clone()); + let make_store = || { + std::sync::Arc::new(crate::store::ECStore { + id: uuid::Uuid::new_v4(), + disk_map: std::collections::HashMap::new(), + pools: pools.clone(), + peer_sys: crate::cluster::rpc::S3PeerSys::new_with_instance_ctx(&endpoint_pools, std::sync::Arc::clone(&ctx)), + pool_meta: tokio::sync::RwLock::new(pool_meta.clone()), + rebalance_meta: tokio::sync::RwLock::new(rebalance_meta.clone()), + decommission_cancelers: tokio::sync::RwLock::new(vec![None, None]), + start_gate: tokio::sync::Mutex::new(()), + pool_meta_save_gate: tokio::sync::Mutex::new(()), + ctx: std::sync::Arc::clone(&ctx), + bucket_fence_registry: std::sync::Arc::default(), + }) + }; + (temp_dirs, make_store(), make_store()) +} + #[cfg(test)] mod rebalance_unit_tests;