fix(rebalance): commit stats after source cleanup (#5795)

This commit is contained in:
cxymds
2026-08-07 17:18:15 +08:00
committed by GitHub
parent 8d582a096c
commit 58d4bdc79f
+164 -36
View File
@@ -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<Output = Result<ObjectInfo>>,
) -> Result<Option<String>> {
// 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");
}
}