diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 28fedd49a..54a795a48 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -2639,6 +2639,11 @@ impl ECStore { ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await) } + async fn ensure_decommission_rebalance_idle_after_refresh_under_start_gate(&self) -> Result<()> { + self.load_rebalance_meta_under_start_gate().await?; + ensure_decommission_not_rebalancing(self.is_rebalance_conflicting_with_decommission().await) + } + pub async fn status(&self, idx: usize) -> Result { let space_info = self.get_decommission_pool_space_info(idx).await?; @@ -4148,7 +4153,8 @@ impl ECStore { } let _start_guard = self.start_gate.lock().await; - self.ensure_decommission_rebalance_idle_after_refresh().await?; + self.ensure_decommission_rebalance_idle_after_refresh_under_start_gate() + .await?; let all_space_infos = self.get_decommission_all_pool_space_infos().await?; { diff --git a/crates/ecstore/src/services/notification_sys.rs b/crates/ecstore/src/services/notification_sys.rs index 72cc11d45..351352364 100644 --- a/crates/ecstore/src/services/notification_sys.rs +++ b/crates/ecstore/src/services/notification_sys.rs @@ -1121,15 +1121,19 @@ impl NotificationSys { } } - match store.stop_rebalance_for_id(expected_rebalance_id).await { + let local_rebalance_id = match expected_rebalance_id { + Some(expected_id) => Some(expected_id.to_owned()), + None => store.current_rebalance_id().await, + }; + match store.stop_rebalance_for_id(local_rebalance_id.as_deref()).await { Ok(_) => { - let save_result = match expected_rebalance_id { + let save_result = match local_rebalance_id.as_deref() { 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, + None => Ok(()), }; if let Err(err) = save_result { error!( diff --git a/crates/ecstore/src/services/rebalance/control.rs b/crates/ecstore/src/services/rebalance/control.rs index 7106188ac..8f5cfeb77 100644 --- a/crates/ecstore/src/services/rebalance/control.rs +++ b/crates/ecstore/src/services/rebalance/control.rs @@ -330,6 +330,10 @@ impl ECStore { #[tracing::instrument(skip_all)] pub async fn load_rebalance_meta(&self) -> Result<()> { let _start_guard = self.start_gate.lock().await; + self.load_rebalance_meta_under_start_gate().await + } + + pub(crate) async fn load_rebalance_meta_under_start_gate(&self) -> Result<()> { let mut meta = RebalanceMeta::new(); debug!( event = EVENT_REBALANCE_STATE, @@ -867,27 +871,31 @@ impl ECStore { Ok(()) } - pub async fn record_rebalance_stop_propagation(self: &Arc, record: RebalanceStopPropagationRecord) -> Result<()> { + pub async fn record_rebalance_stop_propagation( + self: &Arc, + expected_id: &str, + record: RebalanceStopPropagationRecord, + ) -> Result<()> { if !record.has_failures() { return Ok(()); } + let _start_guard = self.start_gate.lock().await; + let _activation_guard = self + .rebalance_activation_write_guard(Some(expected_id), "record rebalance stop propagation") + .await?; let encoded_error = encode_rebalance_stop_propagation_record(&record); let meta_to_save = { let mut rebalance_meta = self.rebalance_meta.write().await; + ensure_rebalance_run_id(rebalance_meta.as_ref(), expected_id, "record rebalance stop propagation")?; record_rebalance_stop_propagation_snapshot(rebalance_meta.as_mut(), encoded_error, OffsetDateTime::now_utc()) }; 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_for_id_with_merge( - pool, - &meta_to_save, - "record_rebalance_stop_propagation", - meta_to_save.id.as_str(), - ) - .await, + self.save_rebalance_meta_for_id_with_merge(pool, &meta_to_save, "record_rebalance_stop_propagation", expected_id) + .await, "record_rebalance_stop_propagation", )?; } diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 37a2ded82..84f0feb90 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -220,33 +220,22 @@ 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 { - let run_guard = store - .rebalance_run_guard(rebalance_id.as_ref(), "rebalance object migration") - .await?; - let result = store + store .clone() .rebalance_object(src_pool_idx, bucket, rd, expected_bucket_incarnation_id) - .await; - drop(run_guard); - result + .await } }; // 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(); - 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 - } + async move { store.delete_object(&bucket, &object, opts).await } }; + let run_guard = self + .rebalance_run_guard(rebalance_id.as_ref(), "rebalance version migration") + .await?; let result = migrate_entry_version( &RebalanceMigrationBackend::new(set.as_ref(), self.as_ref()), bucket.clone(), @@ -260,6 +249,7 @@ impl ECStore { &mut delete_marker, ) .await; + drop(run_guard); if result.ignored { if should_count_rebalance_version_complete(&result) { diff --git a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs index 889d92795..d9c16ee06 100644 --- a/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs +++ b/crates/ecstore/src/services/rebalance/rebalance_unit_tests.rs @@ -2775,6 +2775,10 @@ async fn test_old_worker_cannot_mutate_replacement_rebalance_state() { .save_rebalance_stats_for_id(0, RebalSaveOpt::Stats, "rebalance-a") .await .expect_err("old save task must not persist replacement metadata"), + store + .save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, "rebalance-a") + .await + .expect_err("old stop path must not persist replacement metadata"), ] { assert!(err.to_string().contains("stale rebalance worker rejected")); } @@ -2785,10 +2789,11 @@ async fn test_old_worker_cannot_mutate_replacement_rebalance_state() { 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); + assert!(meta.stopped_at.is_none()); } #[tokio::test] -async fn test_stop_waits_for_active_rebalance_side_effect_guard() { +async fn test_stop_waits_for_active_rebalance_migration_guard() { let meta = RebalanceMeta { id: "rebalance-a".to_string(), pool_stats: vec![RebalanceStats { @@ -2803,9 +2808,9 @@ async fn test_stop_waits_for_active_rebalance_side_effect_guard() { }; let store = test_store_with_rebalance_meta(meta); let run_guard = store - .rebalance_run_guard("rebalance-a", "test side effect") + .rebalance_run_guard("rebalance-a", "rebalance remote-tier migration") .await - .expect("active run should admit the side effect"); + .expect("active run should admit the remote-tier migration"); let mut stop = Box::pin(store.stop_rebalance_for_id(Some("rebalance-a"))); let mut context = Context::from_waker(futures::task::noop_waker_ref()); @@ -2834,6 +2839,52 @@ async fn test_stop_waits_for_active_rebalance_side_effect_guard() { ); } +#[tokio::test] +async fn test_rebalance_metadata_reload_under_start_gate_does_not_reacquire_gate() { + let store = test_store_with_rebalance_meta(RebalanceMeta::default()); + let _start_guard = store.start_gate.lock().await; + + let err = store + .load_rebalance_meta_under_start_gate() + .await + .expect_err("empty test store should reach the metadata load without waiting on start_gate again"); + + assert!(err.to_string().contains("no pools available")); +} + +#[tokio::test] +async fn test_stale_stop_propagation_cannot_mutate_replacement_rebalance() { + let meta = RebalanceMeta { + id: "rebalance-b".to_string(), + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + }; + let store = test_store_with_rebalance_meta(meta); + let record = RebalanceStopPropagationRecord { + stop_failures: vec!["old rebalance stop failed".to_string()], + ..Default::default() + }; + + let err = store + .record_rebalance_stop_propagation("rebalance-a", record) + .await + .expect_err("old propagation failure must not mutate 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.last_refreshed_at.is_none()); + assert!(meta.pool_stats[0].info.last_error.is_none()); +} + 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 5bf8692f8..b24b23812 100644 --- a/crates/ecstore/src/services/rebalance/runtime.rs +++ b/crates/ecstore/src/services/rebalance/runtime.rs @@ -673,7 +673,7 @@ impl ECStore { 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<()> { + pub 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 } diff --git a/rustfs/src/admin/handlers/rebalance.rs b/rustfs/src/admin/handlers/rebalance.rs index 81493dab5..775f099cd 100644 --- a/rustfs/src/admin/handlers/rebalance.rs +++ b/rustfs/src/admin/handlers/rebalance.rs @@ -177,12 +177,15 @@ async fn rollback_cluster_rebalance_start( terminal_reload_attempt_at: Some(terminal_reload_attempt_at), terminal_reload_failures: terminal_reload_failures.clone(), }; - store.record_rebalance_stop_propagation(record).await.map_err(|err| { - format!( - "cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}", - rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures) - ) - })?; + store + .record_rebalance_stop_propagation(rebalance_id, record) + .await + .map_err(|err| { + format!( + "cluster rebalance rollback for {rebalance_id} partial; failed to persist stop propagation: {err}; {}", + rebalance_rollback_failure_message(rebalance_id, &stop_failures, &terminal_reload_failures) + ) + })?; return Err(rebalance_rollback_failure_message( rebalance_id, &stop_failures, @@ -197,7 +200,7 @@ async fn rollback_cluster_rebalance_start( .await .map_err(|err| format!("local stop_rebalance rollback for {rebalance_id} failed: {err}"))?; store - .save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt) + .save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, rebalance_id) .await .map_err(|err| format!("local rollback stop metadata save for {rebalance_id} failed: {err}"))?; Ok(()) @@ -679,7 +682,7 @@ impl Operation for RebalanceStart { terminal_reload_attempt_at: Some(terminal_reload_attempt_at), terminal_reload_failures: terminal_reload_failures.clone(), }; - store.record_rebalance_stop_propagation(record).await.map_err(|err| { + store.record_rebalance_stop_propagation(&id, record).await.map_err(|err| { rebalance_internal_error(format!( "failed to persist rebalance local-start rollback propagation metadata: {err}" )) @@ -926,23 +929,26 @@ impl Operation for RebalanceStop { log_rebalance_request_rejected("stop", "rebalance_not_started", &request_id, &actor, &remote_addr); return Err(s3_error!(NoSuchResource, "pool rebalance is not started")); } + let Some(expected_rebalance_id) = expected_rebalance_id else { + return Err(s3_error!(InternalError, "active rebalance metadata has no activation id")); + }; let notification_sys = current_notification_system(); let stop_attempt_at = OffsetDateTime::now_utc(); let mut stop_failures = Vec::new(); if let Some(notification_sys) = notification_sys.as_ref() { stop_failures = notification_sys - .stop_rebalance_failures(expected_rebalance_id.as_deref()) + .stop_rebalance_failures(Some(expected_rebalance_id.as_str())) .await .map_err(|e| s3_error!(InternalError, "failed to stop rebalance via notification system: {}", e))?; } else { store - .stop_rebalance_for_id(expected_rebalance_id.as_deref()) + .stop_rebalance_for_id(Some(expected_rebalance_id.as_str())) .await .map_err(|e| s3_error!(InternalError, "failed to stop rebalance: {}", e))?; store - .save_rebalance_stats(usize::MAX, RebalSaveOpt::StoppedAt) + .save_rebalance_stats_for_id(usize::MAX, RebalSaveOpt::StoppedAt, expected_rebalance_id.as_str()) .await .map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop metadata: {}", e))?; } @@ -1007,7 +1013,7 @@ impl Operation for RebalanceStop { terminal_reload_failures: terminal_reload_failures.clone(), }; store - .record_rebalance_stop_propagation(record) + .record_rebalance_stop_propagation(expected_rebalance_id.as_str(), record) .await .map_err(|e| s3_error!(InternalError, "failed to persist rebalance stop propagation metadata: {}", e))?;