diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 4fc043286..c2b671e6b 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -27,7 +27,7 @@ use super::worker::{ }; use super::{ EVENT_REBALANCE_BUCKET, EVENT_REBALANCE_ENTRY, EVENT_REBALANCE_STATE, LOG_COMPONENT_ECSTORE, LOG_SUBSYSTEM_REBALANCE, - REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome, + ObjectInfo, REBALANCE_DEFERRED_ENTRY_ERROR_PREFIX, RebalanceBucketConfigs, RebalanceBucketOutcome, RebalanceEntryOutcome, }; use crate::core::pools::ListCallback; use crate::data_movement; @@ -37,13 +37,52 @@ use crate::object_api::{GetObjectReader, ObjectOptions}; use crate::set_disk::SetDisks; use crate::storage_api_contracts::object::ObjectOperations as _; use crate::store::ECStore; -use rustfs_filemeta::MetaCacheEntry; +use rustfs_filemeta::{FileInfo, MetaCacheEntry}; use std::sync::Arc; use time::OffsetDateTime; use tokio_util::sync::CancellationToken; use tracing::{debug, error, warn}; impl ECStore { + async fn finish_rebalance_entry_after_cleanup( + &self, + pool_index: usize, + bucket: &str, + object: &str, + stats_updates: &[&FileInfo], + cleanup: impl std::future::Future>, + ) -> Result> { + // Persisted stats can complete a pool on restart, so source cleanup must resolve first. + let cleanup_warning = resolve_rebalance_entry_cleanup_delete_result(cleanup.await, bucket, object)?; + if let Some(message) = cleanup_warning.as_ref() + && let Err(err) = self + .record_rebalance_cleanup_warning(pool_index, bucket, object, message.clone()) + .await + { + error!( + event = EVENT_REBALANCE_ENTRY, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REBALANCE, + pool_index, + bucket, + object, + stage = "cleanup_source", + error = ?err, + "Failed to record rebalance source cleanup warning" + ); + } + + resolve_rebalance_stats_update_result( + self.update_pool_stats_batch(pool_index, bucket.to_string(), stats_updates) + .await, + pool_index, + bucket, + object, + )?; + + Ok(cleanup_warning) + } + #[allow(unused_assignments)] #[tracing::instrument(skip(self, set))] async fn rebalance_entry( @@ -260,28 +299,23 @@ impl ECStore { } } - resolve_rebalance_stats_update_result( - self.update_pool_stats_batch(pool_index, bucket.clone(), stats_updates.as_slice()) - .await, - pool_index, - bucket.as_str(), - entry.name.as_str(), - )?; - if should_cleanup_rebalance_source_entry(rebalanced, fivs.versions.len(), expired) { - let cleanup_warning = resolve_rebalance_entry_cleanup_delete_result( - data_movement::cleanup_source_entry_if_unchanged( - set.clone(), + let cleanup_warning = self + .finish_rebalance_entry_after_cleanup( + pool_index, bucket.as_str(), entry.name.as_str(), - &fivs, - &cleanup_preflight_allowed_missing, - "rebalance", + stats_updates.as_slice(), + data_movement::cleanup_source_entry_if_unchanged( + set.clone(), + bucket.as_str(), + entry.name.as_str(), + &fivs, + &cleanup_preflight_allowed_missing, + "rebalance", + ), ) - .await, - bucket.as_str(), - entry.name.as_str(), - )?; + .await?; if let Some(message) = cleanup_warning { warn!( event = EVENT_REBALANCE_ENTRY, @@ -295,22 +329,6 @@ impl ECStore { error = %message, "Ignored rebalance source cleanup failure" ); - if let Err(err) = self - .record_rebalance_cleanup_warning(pool_index, bucket.as_str(), entry.name.as_str(), message) - .await - { - error!( - event = EVENT_REBALANCE_ENTRY, - component = LOG_COMPONENT_ECSTORE, - subsystem = LOG_SUBSYSTEM_REBALANCE, - pool_index, - bucket = %bucket, - object = %entry.name, - stage = "cleanup_source", - error = ?err, - "Failed to record rebalance source cleanup warning" - ); - } } else { debug!( event = EVENT_REBALANCE_ENTRY, @@ -337,6 +355,14 @@ impl ECStore { state = "source_retained", "Rebalance source object retained" ); + + resolve_rebalance_stats_update_result( + self.update_pool_stats_batch(pool_index, bucket.clone(), stats_updates.as_slice()) + .await, + pool_index, + bucket.as_str(), + entry.name.as_str(), + )?; } Ok(RebalanceEntryOutcome::Completed) @@ -548,3 +574,105 @@ impl ECStore { Ok(RebalanceBucketOutcome::Completed) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::services::rebalance::{RebalStatus, RebalanceInfo, RebalanceMeta, RebalanceStats}; + use rustfs_filemeta::FileInfo; + use time::OffsetDateTime; + + #[tokio::test] + async fn rebalance_stats_wait_for_source_cleanup_result() { + let endpoint_pools: crate::layout::endpoints::EndpointServerPools = Vec::new().into(); + let store = Arc::new(ECStore { + id: uuid::Uuid::new_v4(), + disk_map: std::collections::HashMap::new(), + pools: Vec::new(), + peer_sys: crate::cluster::rpc::S3PeerSys::new(&endpoint_pools), + pool_meta: tokio::sync::RwLock::new(crate::core::pools::PoolMeta::default()), + rebalance_meta: tokio::sync::RwLock::new(Some(RebalanceMeta { + pool_stats: vec![RebalanceStats { + participating: true, + info: RebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }], + ..Default::default() + })), + decommission_cancelers: tokio::sync::RwLock::new(Vec::new()), + start_gate: tokio::sync::Mutex::new(()), + pool_meta_save_gate: tokio::sync::Mutex::new(()), + ctx: crate::runtime::instance::bootstrap_ctx(), + }); + let mut version = FileInfo::new("object.bin", 4, 2); + version.name = "object.bin".to_string(); + version.size = 128; + version.is_latest = true; + let warning_version = version.clone(); + let (release_cleanup, cleanup_released) = tokio::sync::oneshot::channel(); + + 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 { + cleanup_released.await.expect("cleanup release sender should remain alive"); + Ok(ObjectInfo::default()) + }) + .await + }); + + tokio::task::yield_now().await; + assert_eq!( + store + .rebalance_meta + .read() + .await + .as_ref() + .expect("rebalance metadata should exist") + .pool_stats[0] + .bytes, + 0, + "stats must not become visible before source cleanup resolves" + ); + + release_cleanup.send(()).expect("cleanup waiter should remain alive"); + assert!( + finish + .await + .expect("finish task should not panic") + .expect("finish should succeed") + .is_none(), + "successful cleanup should not produce a warning" + ); + assert!( + store + .rebalance_meta + .read() + .await + .as_ref() + .expect("rebalance metadata should exist") + .pool_stats[0] + .bytes + > 0, + "stats should become visible after source cleanup resolves" + ); + + { + let mut meta = store.rebalance_meta.write().await; + meta.as_mut().expect("rebalance metadata should exist").pool_stats[0].bytes = 0; + } + let warning = store + .finish_rebalance_entry_after_cleanup(0, "bucket", "object.bin", &[&warning_version], async { Err(Error::SlowDown) }) + .await + .expect("cleanup warnings should not fail the completed migration"); + assert!(warning.is_some(), "cleanup failure should return a warning"); + let meta = store.rebalance_meta.read().await; + let pool_stats = &meta.as_ref().expect("rebalance metadata should exist").pool_stats[0]; + assert_eq!(pool_stats.cleanup_warnings.count, 1, "cleanup warning must block pool completion"); + assert!(pool_stats.bytes > 0, "completed migration bytes should still be recorded"); + } +}