diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index b2b7a9dfe..09e2f32cb 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -870,9 +870,26 @@ fn activation_pool_meta_lock_error(err: rustfs_lock::LockError) -> Error { } } -pub(crate) async fn acquire_pool_rebalance_activation_locks( - pool: Arc, -) -> Result<(rustfs_lock::NamespaceLockGuard, rustfs_lock::NamespaceLockGuard)> +pub(crate) struct PoolRebalanceActivationFence { + pool_meta_guard: rustfs_lock::NamespaceLockGuard, + rebalance_meta_guard: rustfs_lock::NamespaceLockGuard, +} + +impl PoolRebalanceActivationFence { + pub(crate) fn ensure_held(&self) -> Result<()> { + ensure_activation_locks_held(self.pool_meta_guard.is_lock_lost(), self.rebalance_meta_guard.is_lock_lost()) + } +} + +fn ensure_activation_locks_held(pool_lock_lost: bool, rebalance_lock_lost: bool) -> Result<()> { + if pool_lock_lost || rebalance_lock_lost { + return Err(Error::other("activation lock lost before metadata commit or worker admission")); + } + + Ok(()) +} + +pub(crate) async fn acquire_pool_rebalance_activation_locks(pool: Arc) -> Result where S: crate::storage_api_contracts::namespace::NamespaceLocking< Error = Error, @@ -891,7 +908,10 @@ where .await .map_err(activation_rebalance_meta_lock_error)?; - Ok((pool_meta_guard, rebalance_meta_guard)) + Ok(PoolRebalanceActivationFence { + pool_meta_guard, + rebalance_meta_guard, + }) } fn rollback_decommission_pool_meta(pool_meta: &mut PoolMeta, previous_pool_meta: PoolMeta) { @@ -2525,7 +2545,7 @@ impl ECStore { .first() .cloned() .ok_or_else(|| Error::other("decommission start rebalance metadata load failed: no storage pools available"))?; - let (_pool_meta_guard, _rebalance_meta_guard) = acquire_pool_rebalance_activation_locks(rebalance_pool.clone()).await?; + let activation_fence = acquire_pool_rebalance_activation_locks(rebalance_pool.clone()).await?; let mut rebalance_meta = RebalanceMeta::new(); match rebalance_meta @@ -2570,6 +2590,7 @@ impl ECStore { latest_pool_meta.queue_buckets(idx, decom_buckets.clone()); } + activation_fence.ensure_held()?; latest_pool_meta.save_no_lock(self.pools.clone()).await?; { let mut pool_meta = self.pool_meta.write().await; @@ -5285,7 +5306,7 @@ mod pools_tests { apply_decommission_status_space_info, bind_decommission_cancelers, bind_missing_decommission_cancelers, cancel_decommission_canceler, classify_decommission_terminal_state, count_decommission_item, decommission_cancel_signal_result, decommission_item_size, decommission_meta_bucket_options, - decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, + decommission_start_pool_state, dedup_indices, default_decommission_bucket_concurrency, ensure_activation_locks_held, ensure_decommission_cancel_allowed, ensure_decommission_clear_allowed, ensure_decommission_listing_disks_available, ensure_decommission_not_rebalancing, ensure_decommission_start_allowed, ensure_decommission_start_keeps_active_pool, ensure_decommission_start_local_leader, ensure_decommission_start_pool_states, @@ -5431,6 +5452,16 @@ mod pools_tests { ); } + #[test] + fn test_activation_fence_rejects_either_lost_guard_before_commit() { + assert!(ensure_activation_locks_held(false, false).is_ok()); + for (pool_lost, rebalance_lost) in [(true, false), (false, true), (true, true)] { + let err = ensure_activation_locks_held(pool_lost, rebalance_lost) + .expect_err("either lost activation guard must fence the commit"); + assert!(err.to_string().contains("activation lock lost")); + } + } + #[test] fn test_apply_decommission_status_space_info_adds_idle_pool_usage() { let status = apply_decommission_status_space_info( diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index 0083fd61b..595980d73 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -16,7 +16,9 @@ use super::{ RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats, RebalanceStopPropagationRecord, encode_rebalance_stop_propagation_record, }; -use crate::core::pools::{PoolMeta, acquire_pool_rebalance_activation_locks, pool_meta_has_active_decommission}; +use crate::core::pools::{ + PoolMeta, PoolRebalanceActivationFence, acquire_pool_rebalance_activation_locks, pool_meta_has_active_decommission, +}; use crate::error::{Error, Result}; use crate::object_api::ObjectOptions; use crate::set_disk::get_lock_acquire_timeout; @@ -38,9 +40,20 @@ fn ensure_rebalance_activation_pool_meta_allowed(meta: &PoolMeta) -> Result<()> Ok(()) } -async fn merge_and_save_rebalance_meta_no_lock(pool: Arc, local_snapshot: &RebalanceMeta, stage: &str) -> Result<()> +pub(super) enum RebalanceWorkerActivationFence { + Ready(PoolRebalanceActivationFence), + NotStartedTerminal, +} + +async fn merge_and_save_rebalance_meta_no_lock( + pool: Arc, + local_snapshot: &RebalanceMeta, + stage: &str, + before_save: F, +) -> Result<()> where S: EcstoreObjectIO, + F: FnOnce() -> Result<()>, { let opts = ObjectOptions { no_lock: true, @@ -59,6 +72,7 @@ where Err(err) => return Err(Error::other(format!("rebalance meta load before save failed during {stage}: {err}"))), } + before_save()?; merged.save_with_opts(pool, opts).await } @@ -123,7 +137,7 @@ impl ECStore { .await .map_err(rebalance_meta_lock_error)?; - merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage).await + merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, || Ok(())).await } async fn save_rebalance_activation_meta_with_merge( @@ -135,19 +149,23 @@ impl ECStore { where S: EcstoreObjectIO + StorageNamespaceLocking, { - let (_pool_meta_guard, _rebalance_meta_guard) = acquire_pool_rebalance_activation_locks(pool.clone()).await?; + 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?; ensure_rebalance_activation_pool_meta_allowed(&pool_meta)?; - merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage).await + merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, || activation_fence.ensure_held()).await } - pub(super) async fn fence_rebalance_worker_activation(&self, pool: Arc, expected_id: &str) -> Result + pub(super) async fn fence_rebalance_worker_activation( + &self, + pool: Arc, + expected_id: &str, + ) -> Result where S: EcstoreObjectIO + StorageNamespaceLocking, { - let (_pool_meta_guard, _rebalance_meta_guard) = acquire_pool_rebalance_activation_locks(pool.clone()).await?; + 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?; ensure_rebalance_activation_pool_meta_allowed(&pool_meta)?; @@ -169,7 +187,12 @@ impl ECStore { ))); } - Ok(is_rebalance_conflicting_with_decommission(&persisted)) + activation_fence.ensure_held()?; + if !is_rebalance_conflicting_with_decommission(&persisted) { + return Ok(RebalanceWorkerActivationFence::NotStartedTerminal); + } + + Ok(RebalanceWorkerActivationFence::Ready(activation_fence)) } #[tracing::instrument(skip_all)] diff --git a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs index 38d4d7460..e52694b23 100644 --- a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs @@ -32,7 +32,10 @@ use super::migration::{ MigrationBackend, MigrationVersionResult, migrate_entry_version, migrate_entry_version_with_retry_wait, rebalance_delete_marker_opts, }; -use super::runtime::{should_fail_repeated_rebalance_bucket_defer, source_cleanup_defer_attempt}; +use super::runtime::{ + RebalanceLocalActivationOutcome, commit_local_rebalance_worker_activation, should_fail_repeated_rebalance_bucket_defer, + source_cleanup_defer_attempt, +}; use super::worker::{ RebalanceEntryCleanupResult, ensure_rebalance_listing_disks_available, is_transient_rebalance_error, parse_rebalance_max_attempts, rebalance_listing_retry_delay, rebalance_migration_retry_delay, @@ -2708,6 +2711,43 @@ async fn test_start_rebalance_for_id_rejects_stopped_metadata() { assert!(err.to_string().contains("was stopped before start")); } +#[tokio::test] +async fn test_stop_at_activation_barrier_prevents_worker_token_commit() { + let meta = Arc::new(tokio::sync::RwLock::new(RebalanceMeta { + id: "rebalance-a".to_string(), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + })); + let fence_reached = Arc::new(tokio::sync::Barrier::new(2)); + let stop_committed = Arc::new(tokio::sync::Barrier::new(2)); + + let stop_meta = Arc::clone(&meta); + let stop_fence_reached = Arc::clone(&fence_reached); + let stop_committed_signal = Arc::clone(&stop_committed); + let stop = tokio::spawn(async move { + stop_fence_reached.wait().await; + stop_meta.write().await.stopped_at = Some(OffsetDateTime::now_utc()); + stop_committed_signal.wait().await; + }); + + fence_reached.wait().await; + stop_committed.wait().await; + stop.await.expect("stop barrier task should finish"); + + let mut meta = meta.write().await; + let outcome = commit_local_rebalance_worker_activation(&mut meta, "rebalance-a", tokio_util::sync::CancellationToken::new()) + .expect("stopped metadata should produce a non-start outcome"); + assert_eq!(outcome, RebalanceLocalActivationOutcome::NotStartedTerminal); + assert!(meta.cancel.is_none(), "stopped rebalance must not receive a worker token"); +} + fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc { let endpoint_pools: crate::layout::endpoints::EndpointServerPools = Vec::new().into(); Arc::new(crate::store::ECStore { diff --git a/crates/ecstore/src/services/rebalance/runtime.rs b/crates/ecstore/src/services/rebalance/runtime.rs index e349afbe2..72aafe071 100644 --- a/crates/ecstore/src/services/rebalance/runtime.rs +++ b/crates/ecstore/src/services/rebalance/runtime.rs @@ -1,3 +1,4 @@ +use super::control::RebalanceWorkerActivationFence; use super::meta::{ apply_rebalance_save_option, apply_rebalance_terminal_event, classify_rebalance_terminal_event, clone_first_arc, complete_rebalance_pools_at_goal, complete_rebalance_pools_with_empty_queue, ensure_valid_rebalance_pool_index, @@ -39,6 +40,30 @@ pub(super) fn source_cleanup_defer_attempt(deferred_attempts: &mut HashMap Result { + if meta.id != expected_id { + return Err(Error::other(format!( + "rebalance metadata changed before local worker activation: expected {expected_id}, found {}", + meta.id + ))); + } + if meta.stopped_at.is_some() || !is_rebalance_in_progress(meta) { + return Ok(RebalanceLocalActivationOutcome::NotStartedTerminal); + } + meta.cancel = Some(cancel); + Ok(RebalanceLocalActivationOutcome::Started) +} + impl ECStore { #[tracing::instrument(skip_all)] pub async fn start_rebalance(self: &Arc) -> Result<()> { @@ -59,15 +84,17 @@ impl ECStore { rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.clone() }; let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?; - if !self.fence_rebalance_worker_activation(pool, &expected_id).await? { - return Ok(()); - } + let activation_fence = match self.fence_rebalance_worker_activation(pool, &expected_id).await? { + RebalanceWorkerActivationFence::Ready(fence) => fence, + RebalanceWorkerActivationFence::NotStartedTerminal => return Ok(()), + }; let decommission_running = self.is_decommission_running().await; let cancel_tx = CancellationToken::new(); let rx = cancel_tx.clone(); let mut meta_to_save = None; + let activation_outcome; { let mut rebalance_meta = self.rebalance_meta.write().await; @@ -94,10 +121,12 @@ impl ECStore { if complete_rebalance_pools_with_empty_queue(meta, now) { meta_to_save = Some(meta.clone()); } - meta.cancel = Some(cancel_tx); + activation_fence.ensure_held()?; + activation_outcome = commit_local_rebalance_worker_activation(meta, &expected_id, cancel_tx)?; drop(rebalance_meta); } + drop(activation_fence); if let Some(meta) = meta_to_save { let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?; @@ -108,6 +137,10 @@ impl ECStore { )?; } + if activation_outcome != RebalanceLocalActivationOutcome::Started { + return Ok(()); + } + let participants = if let Some(ref meta) = *self.rebalance_meta.read().await { resolve_rebalance_participants(meta.pool_stats.as_slice(), self.pools.len()) } else {