From 814e17b02a0f1152a0e2d42ff9586cfee556e24d Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 00:01:11 +0800 Subject: [PATCH] fix(ecstore): bind rebalance workers to activation id --- crates/ecstore/src/core/pools.rs | 36 ++- .../ecstore/src/services/notification_sys.rs | 10 +- .../ecstore/src/services/rebalance/control.rs | 242 ++++++++++++++++-- .../ecstore/src/services/rebalance/entry.rs | 82 ++++-- .../rebalance/rebalance_unit_tests.rs | 130 ++++++++-- .../ecstore/src/services/rebalance/runtime.rs | 65 +++-- .../ecstore/src/services/rebalance/types.rs | 2 + 7 files changed, 481 insertions(+), 86 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 09e2f32cb..28fedd49a 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -879,6 +879,11 @@ 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()) } + + pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) { + opts.add_namespace_lock_guard(&self.pool_meta_guard); + opts.add_namespace_lock_guard(&self.rebalance_meta_guard); + } } fn ensure_activation_locks_held(pool_lock_lost: bool, rebalance_lock_lost: bool) -> Result<()> { @@ -1793,6 +1798,33 @@ impl PoolMeta { Ok(()) } + async fn save_no_lock_with_activation_fence( + &self, + pools: Vec>, + activation_fence: &PoolRebalanceActivationFence, + ) -> Result<()> + where + S: EcstoreObjectIO, + { + let data = self.encode_config_data()?; + if data.is_empty() { + return Ok(()); + } + 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()?; + } + + Ok(()) + } + pub fn decommission_cancel(&mut self, idx: usize) -> bool { if let Some(stats) = self.pools.get_mut(idx) { if let Some(d) = &stats.decommission { @@ -2591,7 +2623,9 @@ impl ECStore { } activation_fence.ensure_held()?; - latest_pool_meta.save_no_lock(self.pools.clone()).await?; + latest_pool_meta + .save_no_lock_with_activation_fence(self.pools.clone(), &activation_fence) + .await?; { let mut pool_meta = self.pool_meta.write().await; *pool_meta = latest_pool_meta; diff --git a/crates/ecstore/src/services/notification_sys.rs b/crates/ecstore/src/services/notification_sys.rs index f3dd0a416..72cc11d45 100644 --- a/crates/ecstore/src/services/notification_sys.rs +++ b/crates/ecstore/src/services/notification_sys.rs @@ -1123,7 +1123,15 @@ impl NotificationSys { match store.stop_rebalance_for_id(expected_rebalance_id).await { Ok(_) => { - if let Err(err) = store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await { + let save_result = match expected_rebalance_id { + Some(expected_id) => { + store + .save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_id) + .await + } + None => store.save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt).await, + }; + if let Err(err) = save_result { error!( event = EVENT_NOTIFICATION_PEER_PROPAGATION, component = LOG_COMPONENT_ECSTORE, diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index 595980d73..7106188ac 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -45,35 +45,73 @@ pub(super) enum RebalanceWorkerActivationFence { NotStartedTerminal, } -async fn merge_and_save_rebalance_meta_no_lock( +pub(super) struct RebalanceRunGuard { + _guard: tokio::sync::OwnedRwLockReadGuard<()>, +} + +async fn merge_and_save_rebalance_meta_no_lock( pool: Arc, local_snapshot: &RebalanceMeta, stage: &str, - before_save: F, + mut opts: ObjectOptions, + activation_fence: Option<&PoolRebalanceActivationFence>, + expected_id: Option<&str>, ) -> Result<()> where S: EcstoreObjectIO, - F: FnOnce() -> Result<()>, { - let opts = ObjectOptions { - no_lock: true, - ..Default::default() - }; let mut merged = RebalanceMeta::new(); match merged.load_with_opts(pool.clone(), opts.clone()).await { Ok(()) => { + if let Some(expected_id) = expected_id { + ensure_rebalance_run_id(Some(&merged), expected_id, stage)?; + } if merge_rebalance_meta(&mut merged, local_snapshot) == RebalanceMetaMergeOutcome::RejectedActiveConflict { return Err(Error::RebalanceAlreadyRunning); } } Err(Error::ConfigNotFound) => { + if expected_id.is_some() { + return Err(rebalance_metadata_not_initialized_error(stage)); + } merged = local_snapshot.clone(); } 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 + if let Some(fence) = activation_fence { + fence.add_namespace_lock_fence(&mut opts); + fence.ensure_held()?; + } + merged.save_with_opts(pool, opts).await?; + if let Some(fence) = activation_fence { + fence.ensure_held()?; + } + Ok(()) +} + +pub(super) fn ensure_rebalance_run_id(meta: Option<&RebalanceMeta>, expected_id: &str, stage: &str) -> Result<()> { + let Some(meta) = meta else { + return Err(rebalance_metadata_not_initialized_error(stage)); + }; + if meta.id != expected_id { + return Err(Error::other(format!( + "stale rebalance worker rejected during {stage}: expected {expected_id}, found {}", + meta.id + ))); + } + Ok(()) +} + +pub(super) fn ensure_rebalance_worker_active(meta: Option<&RebalanceMeta>, expected_id: &str, stage: &str) -> Result<()> { + ensure_rebalance_run_id(meta, expected_id, stage)?; + let Some(meta) = meta else { + return Err(rebalance_metadata_not_initialized_error(stage)); + }; + if meta.stopped_at.is_some() || !is_rebalance_conflicting_with_decommission(meta) { + return Err(Error::other(format!("inactive rebalance worker rejected during {stage}: {expected_id}"))); + } + Ok(()) } pub(super) fn validate_rebalance_disk_stats_coverage(disk_stats: &[DiskStat]) -> Result<()> { @@ -122,22 +160,105 @@ fn clear_rebalance_status_refresh(current: &mut Option) { } impl ECStore { + // Transition order is start_gate -> activation_gate -> rebalance_meta; the probe read below is released before the gate. + async fn rebalance_activation_write_guard( + &self, + expected_id: Option<&str>, + stage: &str, + ) -> Result>> { + let activation_gate = { + let meta = self.rebalance_meta.read().await; + if let Some(expected_id) = expected_id { + ensure_rebalance_run_id(meta.as_ref(), expected_id, stage)?; + } + meta.as_ref().map(|meta| Arc::clone(&meta.activation_gate)) + }; + Ok(match activation_gate { + Some(gate) => Some(gate.write_owned().await), + None => None, + }) + } + + pub(super) async fn rebalance_run_guard(&self, expected_id: &str, stage: &str) -> Result { + let activation_gate = { + let meta = self.rebalance_meta.read().await; + ensure_rebalance_worker_active(meta.as_ref(), expected_id, stage)?; + Arc::clone( + &meta + .as_ref() + .ok_or_else(|| rebalance_metadata_not_initialized_error(stage))? + .activation_gate, + ) + }; + let guard = Arc::clone(&activation_gate).read_owned().await; + let meta = self.rebalance_meta.read().await; + ensure_rebalance_worker_active(meta.as_ref(), expected_id, stage)?; + let current_gate = &meta + .as_ref() + .ok_or_else(|| rebalance_metadata_not_initialized_error(stage))? + .activation_gate; + if !Arc::ptr_eq(&activation_gate, current_gate) { + return Err(Error::other(format!( + "stale rebalance activation gate rejected during {stage}: {expected_id}" + ))); + } + drop(meta); + Ok(RebalanceRunGuard { _guard: guard }) + } + pub(super) async fn save_rebalance_meta_with_merge( &self, pool: Arc, local_snapshot: &RebalanceMeta, stage: &str, ) -> Result<()> + where + S: EcstoreObjectIO + StorageNamespaceLocking, + { + self.save_rebalance_meta_with_merge_for_id(pool, local_snapshot, stage, None) + .await + } + + pub(super) async fn save_rebalance_meta_for_id_with_merge( + &self, + pool: Arc, + local_snapshot: &RebalanceMeta, + stage: &str, + expected_id: &str, + ) -> Result<()> + where + S: EcstoreObjectIO + StorageNamespaceLocking, + { + self.save_rebalance_meta_with_merge_for_id(pool, local_snapshot, stage, Some(expected_id)) + .await + } + + async fn save_rebalance_meta_with_merge_for_id( + &self, + pool: Arc, + local_snapshot: &RebalanceMeta, + stage: &str, + expected_id: Option<&str>, + ) -> Result<()> where S: EcstoreObjectIO + StorageNamespaceLocking, { let ns_lock = pool.new_ns_lock(crate::disk::RUSTFS_META_BUCKET, REBAL_META_NAME).await?; - let _guard = ns_lock + let guard = ns_lock .get_write_lock(get_lock_acquire_timeout()) .await .map_err(rebalance_meta_lock_error)?; + let mut opts = ObjectOptions { + no_lock: true, + ..Default::default() + }; + opts.add_namespace_lock_guard(&guard); - merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, || Ok(())).await + merge_and_save_rebalance_meta_no_lock(pool, local_snapshot, stage, opts, None, expected_id).await?; + if guard.is_lock_lost() { + return Err(Error::other("rebalance metadata lock lost during metadata commit")); + } + Ok(()) } async fn save_rebalance_activation_meta_with_merge( @@ -154,7 +275,18 @@ impl ECStore { 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, || activation_fence.ensure_held()).await + merge_and_save_rebalance_meta_no_lock( + pool, + local_snapshot, + stage, + ObjectOptions { + no_lock: true, + ..Default::default() + }, + Some(&activation_fence), + None, + ) + .await } pub(super) async fn fence_rebalance_worker_activation( @@ -197,6 +329,7 @@ impl ECStore { #[tracing::instrument(skip_all)] pub async fn load_rebalance_meta(&self) -> Result<()> { + let _start_guard = self.start_gate.lock().await; let mut meta = RebalanceMeta::new(); debug!( event = EVENT_REBALANCE_STATE, @@ -206,7 +339,9 @@ impl ECStore { "Loading rebalance metadata" ); let pool = clone_first_arc(&self.pools, "rebalanceMeta: no pools available")?; - if resolve_rebalance_meta_load_result(meta.load(pool).await)? { + let loaded = resolve_rebalance_meta_load_result(meta.load(pool).await)?; + let _activation_guard = self.rebalance_activation_write_guard(None, "load rebalance metadata").await?; + if loaded { { let mut rebalance_meta = self.rebalance_meta.write().await; @@ -243,9 +378,14 @@ impl ECStore { #[tracing::instrument(skip_all)] pub async fn refresh_rebalance_status_meta(&self) -> Result<()> { + let _start_guard = self.start_gate.lock().await; let pool = clone_first_arc(&self.pools, "refresh_rebalance_status_meta: no pools available")?; let mut persisted = RebalanceMeta::new(); - match persisted.load(pool).await { + let loaded = persisted.load(pool).await; + let _activation_guard = self + .rebalance_activation_write_guard(None, "refresh rebalance metadata") + .await?; + match loaded { Ok(()) => { let mut rebalance_meta = self.rebalance_meta.write().await; merge_rebalance_status_refresh(&mut rebalance_meta, persisted); @@ -311,7 +451,7 @@ impl ECStore { if let Some(meta) = rebalance_meta.as_ref() { let pool = clone_first_arc(&self.pools, "update_rebalance_stats: no pools available")?; resolve_rebalance_meta_save_result( - self.save_rebalance_meta_with_merge(pool, meta, "update_rebalance_stats") + self.save_rebalance_meta_for_id_with_merge(pool, meta, "update_rebalance_stats", meta.id.as_str()) .await, "update_rebalance_stats", )?; @@ -436,11 +576,14 @@ impl ECStore { pub async fn init_rebalance_start(self: &Arc, bucktes: Vec) -> Result { let _start_guard = self.start_gate.lock().await; - let decommission_running = self.is_decommission_running().await; { + let decommission_running = self.is_decommission_running().await; let rebalance_meta = self.rebalance_meta.read().await; validate_init_rebalance_state(decommission_running, rebalance_meta.as_ref())?; } + let _activation_guard = self + .rebalance_activation_write_guard(None, "initialize replacement rebalance") + .await?; self.init_rebalance_meta(bucktes).await } @@ -480,11 +623,35 @@ impl ECStore { #[tracing::instrument(skip(self, versions))] pub async fn update_pool_stats_batch(&self, pool_index: usize, bucket: String, versions: &[&FileInfo]) -> Result<()> { + self.update_pool_stats_batch_for_id(pool_index, bucket, versions, None).await + } + + pub(super) async fn update_pool_stats_batch_for_rebalance( + &self, + pool_index: usize, + bucket: String, + versions: &[&FileInfo], + expected_id: &str, + ) -> Result<()> { + self.update_pool_stats_batch_for_id(pool_index, bucket, versions, Some(expected_id)) + .await + } + + async fn update_pool_stats_batch_for_id( + &self, + pool_index: usize, + bucket: String, + versions: &[&FileInfo], + expected_id: Option<&str>, + ) -> Result<()> { if versions.is_empty() { return Ok(()); } let mut rebalance_meta = self.rebalance_meta.write().await; + if let Some(expected_id) = expected_id { + ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "update rebalance pool stats")?; + } if let Some(meta) = rebalance_meta.as_mut() { if !should_accept_rebalance_stats_update(meta, pool_index) { return Ok(()); @@ -499,8 +666,9 @@ impl ECStore { } #[tracing::instrument(skip(self))] - pub async fn next_rebal_bucket(&self, pool_index: usize) -> Result> { + pub async fn next_rebal_bucket(&self, pool_index: usize, expected_id: &str) -> Result> { let rebalance_meta = self.rebalance_meta.read().await; + ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "next rebalance bucket")?; debug!( event = EVENT_REBALANCE_BUCKET, component = LOG_COMPONENT_ECSTORE, @@ -514,8 +682,9 @@ impl ECStore { } #[tracing::instrument(skip(self))] - pub async fn bucket_rebalance_done(&self, pool_index: usize, bucket: String) -> Result<()> { + pub async fn bucket_rebalance_done(&self, pool_index: usize, bucket: String, expected_id: &str) -> Result<()> { let mut rebalance_meta = self.rebalance_meta.write().await; + ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "mark rebalance bucket done")?; mark_rebalance_bucket_done(rebalance_meta.as_mut(), pool_index, &bucket) } @@ -525,8 +694,10 @@ impl ECStore { bucket: &str, object: &str, message: String, + expected_id: &str, ) -> Result<()> { let mut rebalance_meta = self.rebalance_meta.write().await; + ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "record rebalance cleanup warning")?; record_rebalance_cleanup_warning_in_meta( rebalance_meta.as_mut(), pool_index, @@ -537,8 +708,15 @@ impl ECStore { ) } - pub(super) async fn defer_rebalance_bucket(&self, pool_index: usize, bucket: String, last_error: String) -> Result<()> { + pub(super) async fn defer_rebalance_bucket( + &self, + pool_index: usize, + bucket: String, + last_error: String, + expected_id: &str, + ) -> Result<()> { let mut rebalance_meta = self.rebalance_meta.write().await; + ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "defer rebalance bucket")?; let Some(meta) = rebalance_meta.as_mut() else { return Err(rebalance_metadata_not_initialized_error("defer rebalance bucket")); }; @@ -634,6 +812,8 @@ impl ECStore { #[tracing::instrument(skip(self))] pub async fn stop_rebalance_for_id(self: &Arc, expected_id: Option<&str>) -> Result<()> { + let _start_guard = self.start_gate.lock().await; + let _activation_guard = self.rebalance_activation_write_guard(expected_id, "stop rebalance").await?; let meta_to_save = { let mut rebalance_meta = self.rebalance_meta.write().await; stop_rebalance_meta_snapshot_for_id(rebalance_meta.as_mut(), OffsetDateTime::now_utc(), expected_id) @@ -642,7 +822,7 @@ impl ECStore { if let Some(meta_to_save) = meta_to_save { let pool = clone_first_arc(self.pools.as_slice(), "stop_rebalance: no pools available")?; resolve_rebalance_meta_save_result( - self.save_rebalance_meta_with_merge(pool, &meta_to_save, "stop_rebalance") + self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, "stop_rebalance", meta_to_save.id.as_str()) .await, "stop_rebalance", )?; @@ -656,6 +836,10 @@ impl ECStore { expected_id: Option<&str>, start_error: String, ) -> Result<()> { + let _start_guard = self.start_gate.lock().await; + let _activation_guard = self + .rebalance_activation_write_guard(expected_id, "rollback rebalance start") + .await?; let meta_to_save = { let mut rebalance_meta = self.rebalance_meta.write().await; rollback_rebalance_start_meta_snapshot_for_id( @@ -669,8 +853,13 @@ impl ECStore { if let Some(meta_to_save) = meta_to_save { let pool = clone_first_arc(self.pools.as_slice(), "rollback_rebalance_start: no pools available")?; resolve_rebalance_meta_save_result( - self.save_rebalance_meta_with_merge(pool, &meta_to_save, "rollback_rebalance_start") - .await, + self.save_rebalance_meta_for_id_with_merge( + pool, + &meta_to_save, + "rollback_rebalance_start", + meta_to_save.id.as_str(), + ) + .await, "rollback_rebalance_start", )?; } @@ -692,8 +881,13 @@ impl ECStore { if let Some(meta_to_save) = meta_to_save { let pool = clone_first_arc(self.pools.as_slice(), "record_rebalance_stop_propagation: no pools available")?; resolve_rebalance_meta_save_result( - self.save_rebalance_meta_with_merge(pool, &meta_to_save, "record_rebalance_stop_propagation") - .await, + self.save_rebalance_meta_for_id_with_merge( + pool, + &meta_to_save, + "record_rebalance_stop_propagation", + meta_to_save.id.as_str(), + ) + .await, "record_rebalance_stop_propagation", )?; } diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 764a68500..37a2ded82 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -50,16 +50,20 @@ impl ECStore { bucket: &str, object: &str, stats_updates: &[&FileInfo], + expected_id: &str, cleanup: impl std::future::Future>, ) -> Result { // Persisted stats can complete a pool on restart, so source cleanup must resolve first. - let cleanup_result = resolve_rebalance_entry_cleanup_delete_result(cleanup.await, bucket, object); + let run_guard = self.rebalance_run_guard(expected_id, "rebalance source cleanup").await?; + let cleanup_result = cleanup.await; + drop(run_guard); + let cleanup_result = resolve_rebalance_entry_cleanup_delete_result(cleanup_result, bucket, object); let RebalanceEntryCleanupResult::Completed { warning } = cleanup_result else { return Ok(cleanup_result); }; if let Some(message) = warning.as_ref() && let Err(err) = self - .record_rebalance_cleanup_warning(pool_index, bucket, object, message.clone()) + .record_rebalance_cleanup_warning(pool_index, bucket, object, message.clone(), expected_id) .await { error!( @@ -76,7 +80,7 @@ impl ECStore { } resolve_rebalance_stats_update_result( - self.update_pool_stats_batch(pool_index, bucket.to_string(), stats_updates) + self.update_pool_stats_batch_for_rebalance(pool_index, bucket.to_string(), stats_updates, expected_id) .await, pool_index, bucket, @@ -95,6 +99,7 @@ impl ECStore { entry: MetaCacheEntry, set: Arc, bucket_configs: Arc, + rebalance_id: Arc, // wk: Arc, ) -> Result { debug!( @@ -129,7 +134,7 @@ impl ECStore { return Ok(RebalanceEntryOutcome::Completed); } - if self.check_if_rebalance_done(pool_index).await { + if self.check_if_rebalance_done(pool_index, rebalance_id.as_ref()).await? { debug!( event = EVENT_REBALANCE_ENTRY, component = LOG_COMPONENT_ECSTORE, @@ -160,7 +165,10 @@ impl ECStore { let mut cleanup_preflight_allowed_missing = Vec::new(); let mut stats_updates = Vec::with_capacity(fivs.versions.len()); for version in fivs.versions.iter() { - if crate::core::pools::should_skip_lifecycle_for_data_movement( + let run_guard = self + .rebalance_run_guard(rebalance_id.as_ref(), "rebalance lifecycle mutation") + .await?; + let expired_by_lifecycle = crate::core::pools::should_skip_lifecycle_for_data_movement( self.clone(), &bucket, version, @@ -169,8 +177,9 @@ impl ECStore { true, &crate::bucket::lifecycle::bucket_lifecycle_audit::LcEventSrc::Rebal, ) - .await? - { + .await?; + drop(run_guard); + if expired_by_lifecycle { expired += 1; // The lifecycle expiry above physically deleted this version from the source set. // Record its identity so the source-cleanup preflight tolerates its absence, @@ -211,17 +220,32 @@ impl ECStore { let expected_bucket_incarnation_id = bucket_configs.bucket_incarnation_id; let mut transfer = |src_pool_idx: usize, bucket: String, rd: GetObjectReader| { let store = self.clone(); + let rebalance_id = Arc::clone(&rebalance_id); async move { - store + let run_guard = store + .rebalance_run_guard(rebalance_id.as_ref(), "rebalance object migration") + .await?; + let result = store + .clone() .rebalance_object(src_pool_idx, bucket, rd, expected_bucket_incarnation_id) - .await + .await; + drop(run_guard); + result } }; // Route delete-marker migration through the store layer so it lands on the // cross-pool target (excluding the source pool), not back onto the source set. let mut delete_marker = |bucket: String, object: String, opts: ObjectOptions| { let store = self.clone(); - async move { store.delete_object(&bucket, &object, opts).await } + let rebalance_id = Arc::clone(&rebalance_id); + async move { + let run_guard = store + .rebalance_run_guard(rebalance_id.as_ref(), "rebalance delete-marker migration") + .await?; + let result = store.delete_object(&bucket, &object, opts).await; + drop(run_guard); + result + } }; let result = migrate_entry_version( &RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()), @@ -280,7 +304,10 @@ impl ECStore { error = %err, "Deferred rebalance entry after transient migration failure" ); - if let Err(stats_err) = self.update_rebalance_last_error(pool_index, deferred_error.clone()).await { + if let Err(stats_err) = self + .update_rebalance_last_error(pool_index, deferred_error.clone(), rebalance_id.as_ref()) + .await + { error!( "rebalance_entry {} failed to record deferred transient failure for {}: {}", &bucket, &entry.name, stats_err @@ -295,7 +322,12 @@ impl ECStore { if !stats_updates.is_empty() && let Err(stats_err) = self - .update_pool_stats_batch(pool_index, bucket.clone(), stats_updates.as_slice()) + .update_pool_stats_batch_for_rebalance( + pool_index, + bucket.clone(), + stats_updates.as_slice(), + rebalance_id.as_ref(), + ) .await { error!( @@ -323,6 +355,7 @@ impl ECStore { bucket.as_str(), entry.name.as_str(), stats_updates.as_slice(), + rebalance_id.as_ref(), data_movement::cleanup_source_entry_if_unchanged( set.clone(), bucket.as_str(), @@ -397,8 +430,13 @@ impl ECStore { ); resolve_rebalance_stats_update_result( - self.update_pool_stats_batch(pool_index, bucket.clone(), stats_updates.as_slice()) - .await, + self.update_pool_stats_batch_for_rebalance( + pool_index, + bucket.clone(), + stats_updates.as_slice(), + rebalance_id.as_ref(), + ) + .await, pool_index, bucket.as_str(), entry.name.as_str(), @@ -419,8 +457,9 @@ impl ECStore { data_movement::migrate_object(self, pool_idx, bucket, rd, expected_bucket_incarnation_id, "rebalance_object").await } - async fn update_rebalance_last_error(&self, pool_idx: usize, message: String) -> Result<()> { + async fn update_rebalance_last_error(&self, pool_idx: usize, message: String, expected_id: &str) -> Result<()> { let mut rebalance_meta = self.rebalance_meta.write().await; + super::control::ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "record rebalance last error")?; let Some(meta) = rebalance_meta.as_mut() else { return Err(rebalance_metadata_not_initialized_error("record rebalance last error")); }; @@ -441,6 +480,7 @@ impl ECStore { rx: CancellationToken, bucket: String, pool_index: usize, + rebalance_id: Arc, ) -> Result { ensure_valid_rebalance_pool_index(self.pools.len(), pool_index)?; @@ -473,6 +513,7 @@ impl ECStore { let bucket_configs = bucket_configs.clone(); let entry_tasks = entry_tasks.clone(); let entry_workers = entry_workers.clone(); + let rebalance_id = Arc::clone(&rebalance_id); move |entry: MetaCacheEntry| { let this = this.clone(); let bucket = bucket.clone(); @@ -482,6 +523,7 @@ impl ECStore { let bucket_configs = bucket_configs.clone(); let entry_tasks = entry_tasks.clone(); let entry_workers = entry_workers.clone(); + let rebalance_id = Arc::clone(&rebalance_id); Box::pin(async move { if callback_rx.is_cancelled() { return; @@ -534,7 +576,9 @@ impl ECStore { state = "task_started", "Started rebalance entry task" ); - let result = this.rebalance_entry(bucket, pool_index, entry, set, bucket_configs).await; + let result = this + .rebalance_entry(bucket, pool_index, entry, set, bucket_configs, rebalance_id) + .await; if let Err(err) = &result { error!("rebalance_entry: rebalance entry failed: {err}"); let mut first_err = entry_error.lock().await; @@ -679,7 +723,7 @@ mod tests { let finish_store = Arc::clone(&store); let finish = tokio::spawn(async move { finish_store - .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&version], async move { + .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&version], "", async move { cleanup_released.await.expect("cleanup release sender should remain alive"); Ok(ObjectInfo::default()) }) @@ -726,7 +770,7 @@ mod tests { meta.as_mut().expect("rebalance metadata should exist").pool_stats[0].bytes = 0; } let warning_result = store - .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], async { + .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], "", async { Err(Error::SlowDown.into()) }) .await @@ -743,7 +787,7 @@ mod tests { meta.as_mut().expect("rebalance metadata should exist").pool_stats[0].bytes = 0; } let deferred = store - .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], async { + .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], "", async { Err(data_movement::SourceCleanupError::SourceChanged) }) .await diff --git a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs index e52694b23..889d92795 100644 --- a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs @@ -66,10 +66,12 @@ use rustfs_filemeta::{FileInfo, MetaCacheEntry}; use rustfs_rio::Index; use s3s::dto::ReplicationConfiguration; use serde::Serialize; +use std::future::Future; use std::io::Cursor; use std::sync::Arc; use std::sync::Mutex; use std::sync::atomic::{AtomicUsize, Ordering}; +use std::task::{Context, Poll}; use time::OffsetDateTime; use tokio::sync::mpsc; use tokio::time::Duration; @@ -2711,9 +2713,83 @@ async fn test_start_rebalance_for_id_rejects_stopped_metadata() { assert!(err.to_string().contains("was stopped before start")); } +#[test] +fn test_stopped_activation_state_prevents_worker_token_commit() { + let mut meta = RebalanceMeta { + id: "rebalance-a".to_string(), + stopped_at: Some(OffsetDateTime::now_utc()), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + + 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"); +} + #[tokio::test] -async fn test_stop_at_activation_barrier_prevents_worker_token_commit() { - let meta = Arc::new(tokio::sync::RwLock::new(RebalanceMeta { +async fn test_old_worker_cannot_mutate_replacement_rebalance_state() { + let meta = RebalanceMeta { + id: "rebalance-b".to_string(), + pool_stats: vec![RebalanceStats { + participating: true, + buckets: vec!["bucket-a".to_string()], + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + let store = test_store_with_rebalance_meta(meta); + let mut fi = FileInfo::default(); + fi.size = 128; + + for err in [ + store + .next_rebal_bucket(0, "rebalance-a") + .await + .expect_err("old worker must not read replacement work"), + store + .bucket_rebalance_done(0, "bucket-a".to_string(), "rebalance-a") + .await + .expect_err("old worker must not complete replacement bucket"), + store + .update_pool_stats_batch_for_rebalance(0, "bucket-a".to_string(), &[&fi], "rebalance-a") + .await + .expect_err("old worker must not update replacement stats"), + store + .check_if_rebalance_done(0, "rebalance-a") + .await + .expect_err("old worker must not complete replacement pool"), + store + .save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a") + .await + .expect_err("old save task must not persist replacement metadata"), + ] { + assert!(err.to_string().contains("stale rebalance worker rejected")); + } + + let meta = store.rebalance_meta.read().await; + let meta = meta.as_ref().expect("replacement metadata should remain present"); + assert_eq!(meta.id, "rebalance-b"); + assert!(meta.pool_stats[0].rebalanced_buckets.is_empty()); + assert_eq!(meta.pool_stats[0].bytes, 0); + assert_eq!(meta.pool_stats[0].info.status, RebalStatus::Started); +} + +#[tokio::test] +async fn test_stop_waits_for_active_rebalance_side_effect_guard() { + let meta = RebalanceMeta { id: "rebalance-a".to_string(), pool_stats: vec![RebalanceStats { participating: true, @@ -2724,28 +2800,38 @@ async fn test_stop_at_activation_barrier_prevents_worker_token_commit() { ..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 store = test_store_with_rebalance_meta(meta); + let run_guard = store + .rebalance_run_guard("rebalance-a", "test side effect") + .await + .expect("active run should admit the side effect"); + let mut stop = Box::pin(store.stop_rebalance_for_id(Some("rebalance-a"))); + let mut context = Context::from_waker(futures::task::noop_waker_ref()); - 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; - }); + assert!(matches!(stop.as_mut().poll(&mut context), Poll::Pending)); + assert!( + store + .rebalance_meta + .read() + .await + .as_ref() + .is_some_and(|meta| meta.stopped_at.is_none()), + "stop must not change run state while a fenced side effect is active" + ); - 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"); + drop(run_guard); + stop.await + .expect_err("empty test store should fail only after committing the local stop state"); + assert!( + store + .rebalance_meta + .read() + .await + .as_ref() + .is_some_and(|meta| meta.stopped_at.is_some()), + "stop should commit after the production side-effect fence is released" + ); } fn test_store_with_rebalance_meta(meta: RebalanceMeta) -> Arc { diff --git a/crates/ecstore/src/services/rebalance/runtime.rs b/crates/ecstore/src/services/rebalance/runtime.rs index 72aafe071..5bf8692f8 100644 --- a/crates/ecstore/src/services/rebalance/runtime.rs +++ b/crates/ecstore/src/services/rebalance/runtime.rs @@ -79,12 +79,12 @@ impl ECStore { state = "starting", "Starting rebalance" ); - let expected_id = { + let expected_id: Arc = { let rebalance_meta = self.rebalance_meta.read().await; - rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.clone() + Arc::from(rebalance_meta.as_ref().ok_or(Error::ConfigNotFound)?.id.as_str()) }; let pool = clone_first_arc(self.pools.as_slice(), "start_rebalance: no pools available")?; - let activation_fence = match self.fence_rebalance_worker_activation(pool, &expected_id).await? { + let activation_fence = match self.fence_rebalance_worker_activation(pool, expected_id.as_ref()).await? { RebalanceWorkerActivationFence::Ready(fence) => fence, RebalanceWorkerActivationFence::NotStartedTerminal => return Ok(()), }; @@ -122,7 +122,7 @@ impl ECStore { meta_to_save = Some(meta.clone()); } activation_fence.ensure_held()?; - activation_outcome = commit_local_rebalance_worker_activation(meta, &expected_id, cancel_tx)?; + activation_outcome = commit_local_rebalance_worker_activation(meta, expected_id.as_ref(), cancel_tx)?; drop(rebalance_meta); } @@ -198,9 +198,10 @@ impl ECStore { let pool_idx = idx; let store = self.clone(); let rx_clone = rx.clone(); + let worker_id = Arc::clone(&expected_id); workers_started += 1; tokio::spawn(async move { - if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx).await { + if let Err(err) = store.rebalance_buckets(rx_clone, pool_idx, worker_id).await { error!( event = EVENT_REBALANCE_STATE, component = LOG_COMPONENT_ECSTORE, @@ -247,13 +248,14 @@ impl ECStore { } #[tracing::instrument(skip(self, rx))] - async fn rebalance_buckets(self: &Arc, rx: CancellationToken, pool_index: usize) -> Result<()> { + async fn rebalance_buckets(self: &Arc, rx: CancellationToken, pool_index: usize, rebalance_id: Arc) -> Result<()> { ensure_valid_rebalance_pool_index(self.pools.len(), pool_index)?; let (done_tx, mut done_rx) = tokio::sync::mpsc::channel::>(1); // Save rebalance metadata periodically let store = self.clone(); + let save_rebalance_id = Arc::clone(&rebalance_id); let save_task = tokio::spawn(async move { let mut timer = tokio::time::interval_at(Instant::now() + Duration::from_secs(30), Duration::from_secs(10)); let mut msg: String; @@ -267,6 +269,11 @@ impl ECStore { let terminal_event = classify_rebalance_terminal_event(result, now); msg = terminal_event.message().to_string(); let mut rebalance_meta = store.rebalance_meta.write().await; + super::control::ensure_rebalance_run_id( + rebalance_meta.as_ref(), + save_rebalance_id.as_ref(), + "apply rebalance terminal event", + )?; if let Some(meta) = rebalance_meta.as_mut() { let meta_stopped = meta.stopped_at.is_some(); if let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) { @@ -315,7 +322,10 @@ impl ECStore { } } - if let Err(err) = store.save_rebalance_stats(pool_index, RebalSaveOpt::Stats).await { + if let Err(err) = store + .save_rebalance_stats_for_id(pool_index, RebalSaveOpt::Stats, save_rebalance_id.as_ref()) + .await + { let wrapped = Error::other(format!("rebalance save_task stats save failed for pool {pool_index}: {err}")); error!("{} err: {:?}", msg, wrapped); if quit { @@ -381,7 +391,7 @@ impl ECStore { break; } - let next_bucket = match self.next_rebal_bucket(pool_index).await { + let next_bucket = match self.next_rebal_bucket(pool_index, rebalance_id.as_ref()).await { Ok(bucket) => bucket, Err(err) => { error!( @@ -413,7 +423,8 @@ impl ECStore { ); let outcome = match resolve_rebalance_bucket_result( - self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index).await, + self.rebalance_bucket(rx.clone(), bucket.clone(), pool_index, Arc::clone(&rebalance_id)) + .await, pool_index, &bucket, ) { @@ -476,7 +487,7 @@ impl ECStore { "Deferred rebalance bucket after transient object failures" ); if let Err(err) = self - .defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone()) + .defer_rebalance_bucket(pool_index, bucket.clone(), last_error.clone(), rebalance_id.as_ref()) .await { error!( @@ -540,7 +551,7 @@ impl ECStore { "Completed rebalance bucket" ); source_cleanup_deferred_attempts.remove(&bucket); - if let Err(err) = self.bucket_rebalance_done(pool_index, bucket).await { + if let Err(err) = self.bucket_rebalance_done(pool_index, bucket, rebalance_id.as_ref()).await { error!( event = EVENT_REBALANCE_BUCKET, component = LOG_COMPONENT_ECSTORE, @@ -601,8 +612,9 @@ impl ECStore { final_result } - pub(super) async fn check_if_rebalance_done(&self, pool_index: usize) -> bool { + pub(super) async fn check_if_rebalance_done(&self, pool_index: usize, expected_id: &str) -> Result { let mut rebalance_meta = self.rebalance_meta.write().await; + super::control::ensure_rebalance_worker_active(rebalance_meta.as_ref(), expected_id, "check rebalance completion")?; if let Some(meta) = rebalance_meta.as_mut() && let Some(pool_stat) = meta.pool_stats.get_mut(pool_index) @@ -617,7 +629,7 @@ impl ECStore { state = "already_completed", "Rebalance pool is already completed" ); - return true; + return Ok(true); } // Mark pool rebalance as done only after it reaches the PercentFreeGoal. @@ -647,19 +659,30 @@ impl ECStore { percent_free = pfi, "Marked rebalance pool completed" ); - return true; + return Ok(true); } } - false + Ok(false) } } impl ECStore { #[tracing::instrument(skip(self))] pub async fn save_rebalance_stats(&self, pool_idx: usize, opt: RebalSaveOpt) -> Result<()> { + self.save_rebalance_stats_inner(pool_idx, opt, None).await + } + + pub(crate) async fn save_rebalance_stats_for_id(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: &str) -> Result<()> { + self.save_rebalance_stats_inner(pool_idx, opt, Some(expected_id)).await + } + + async fn save_rebalance_stats_inner(&self, pool_idx: usize, opt: RebalSaveOpt, expected_id: Option<&str>) -> Result<()> { let meta_to_save = { let mut rebalance_meta = self.rebalance_meta.write().await; + if let Some(expected_id) = expected_id { + super::control::ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "save rebalance stats")?; + } let Some(meta) = rebalance_meta.as_mut() else { return Ok(()); }; @@ -681,10 +704,14 @@ impl ECStore { "Rebalance metadata save requested" ); let stage = format!("save_rebalance_stats for pool {pool_idx} opt {opt:?}"); - resolve_rebalance_meta_save_result( - self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await, - stage.as_str(), - )?; + let save_result = match expected_id { + Some(expected_id) => { + self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, stage.as_str(), expected_id) + .await + } + None => self.save_rebalance_meta_with_merge(pool, &meta_to_save, stage.as_str()).await, + }; + resolve_rebalance_meta_save_result(save_result, stage.as_str())?; Ok(()) } diff --git a/crates/ecstore/src/services/rebalance/types.rs b/crates/ecstore/src/services/rebalance/types.rs index 5f79275dc..a0b04bd4d 100644 --- a/crates/ecstore/src/services/rebalance/types.rs +++ b/crates/ecstore/src/services/rebalance/types.rs @@ -144,6 +144,8 @@ pub struct RebalanceMeta { #[serde(skip)] pub cancel: Option, // To be invoked on rebalance-stop #[serde(skip)] + pub activation_gate: std::sync::Arc>, + #[serde(skip)] pub last_refreshed_at: Option, #[serde(rename = "stopTs")] pub stopped_at: Option, // Time when rebalance-stop was issued