From 65e55c0f8ea2fcd57da9162e8ac22e5fcd5f8ef4 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 8 Aug 2026 21:14:26 +0800 Subject: [PATCH] fix(rebalance): fence batch delete source pools (#5846) --- crates/ecstore/src/store/init.rs | 75 ++++++++++++++++++++++++++++++ crates/ecstore/src/store/object.rs | 3 ++ 2 files changed, 78 insertions(+) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 9e94e6d59..a979328e7 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -1332,6 +1332,81 @@ mod tests { shutdown.cancel(); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn delete_objects_skips_active_rebalance_source_pool() { + let temp_dir = tempfile::tempdir().expect("create batch-delete writer-fencing store dir"); + let (_ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-delete-rebalance", &[4, 4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("batch-delete-rebalance-{}", uuid::Uuid::new_v4()); + let object = "delete-me.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create batch delete rebalance bucket"); + + let mut source_reader = PutObjReader::from_vec(b"source-body".to_vec()); + store.pools[0] + .put_object(&bucket, object, &mut source_reader, &ObjectOptions::default()) + .await + .expect("write object on active source pool"); + let mut target_reader = PutObjReader::from_vec(b"target-body".to_vec()); + store.pools[1] + .put_object(&bucket, object, &mut target_reader, &ObjectOptions::default()) + .await + .expect("write object on non-rebalancing target pool"); + + let mut pool_stats = vec![RebalanceStats::default(); store.pools.len()]; + pool_stats[0] = RebalanceStats { + participating: true, + info: RebalanceInfo { + start_time: Some(OffsetDateTime::now_utc()), + status: RebalStatus::Started, + ..Default::default() + }, + ..Default::default() + }; + *store.rebalance_meta.write().await = Some(RebalanceMeta { + id: uuid::Uuid::new_v4().to_string(), + pool_stats, + ..Default::default() + }); + assert!(store.is_pool_rebalancing(0).await, "pool 0 must be marked as an active rebalance source"); + + let (deleted, errs) = store + .delete_objects( + &bucket, + vec![crate::storage_api_contracts::object::ObjectToDelete { + object_name: object.to_string(), + ..Default::default() + }], + ObjectOptions::default(), + ) + .await; + assert!(matches!(errs.as_slice(), [None]), "batch delete must not fail: {errs:?}"); + assert!( + matches!(deleted.as_slice(), [deleted] if deleted.found && deleted.object_name == object), + "batch delete must report the non-rebalancing pool deletion" + ); + + store.pools[0] + .get_object_info(&bucket, object, &ObjectOptions::default()) + .await + .expect("active source pool object must not be deleted by DeleteObjects"); + let target_err = store.pools[1] + .get_object_info(&bucket, object, &ObjectOptions::default()) + .await + .expect_err("non-rebalancing target pool object must be deleted"); + assert!( + matches!(target_err, StorageError::ObjectNotFound(_, _)), + "target pool should report object not found after DeleteObjects, got {target_err:?}" + ); + + shutdown.cancel(); + } + #[tokio::test] #[serial_test::serial(storage_class_env)] async fn data_movement_conflicts_preserve_newer_target_and_abort_staging() { diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index fde1bd075..bf8c091b3 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -1932,6 +1932,9 @@ impl ECStore { let mut futures = Vec::with_capacity(self.pools.len()); for pool in self.pools.iter() { + if self.is_pool_rebalancing(pool.pool_idx).await { + continue; + } futures.push(pool.delete_objects(bucket, objects.clone(), opts.clone())); }