From 8a77034463312e84426889a133a2d40388dd3a66 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 02:01:55 +0800 Subject: [PATCH] test(ecstore): exercise decommission delete fences --- crates/ecstore/src/core/pools.rs | 16 ++ crates/ecstore/src/set_disk/ops/object.rs | 23 ++ crates/ecstore/src/store/init.rs | 259 +++++++++++++++++----- crates/ecstore/src/store/object.rs | 121 ++++++++++ 4 files changed, 365 insertions(+), 54 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 1d5b1cd3e..08bd28a29 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -3383,6 +3383,22 @@ impl ECStore { Ok(()) } + #[cfg(test)] + pub(crate) async fn decommission_entry_for_test( + self: &Arc, + idx: usize, + entry: MetaCacheEntry, + bucket: String, + set: Arc, + ) -> Result<()> { + let worker_permit = Arc::new(Semaphore::new(1)) + .acquire_owned() + .await + .map_err(|err| Error::other(format!("decommission test worker permit acquire failed: {err}")))?; + self.decommission_entry(CancellationToken::new(), idx, entry, bucket, set, worker_permit, None, None, None, None) + .await + } + #[tracing::instrument(skip(self, rx))] async fn decommission_pool( self: &Arc, diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 73c7f2e16..6ef3c2ce4 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -2483,6 +2483,7 @@ impl SetDisks { }) .await?, ); + notify_put_object_commit_namespace_acquired(bucket, object); } #[cfg(not(any(test, feature = "test-util")))] { @@ -4630,6 +4631,7 @@ struct PutObjectCommitBarrierState { arrived: tokio::sync::Notify, release: tokio::sync::Notify, namespace_pending: tokio::sync::Notify, + namespace_acquired: std::sync::atomic::AtomicBool, } #[cfg(any(test, feature = "test-util"))] @@ -4651,6 +4653,7 @@ impl PutObjectCommitBarrier { arrived: tokio::sync::Notify::new(), release: tokio::sync::Notify::new(), namespace_pending: tokio::sync::Notify::new(), + namespace_acquired: std::sync::atomic::AtomicBool::new(false), }); let mut slot = PUT_OBJECT_COMMIT_BARRIER .get_or_init(|| std::sync::Mutex::new(Vec::new())) @@ -4685,6 +4688,10 @@ impl PutObjectCommitBarrier { .await .expect("put object should wait for the namespace lock after leaving the commit barrier"); } + + pub fn namespace_acquired(&self) -> bool { + self.state.namespace_acquired.load(std::sync::atomic::Ordering::Acquire) + } } #[cfg(any(test, feature = "test-util"))] @@ -4741,6 +4748,22 @@ fn notify_put_object_commit_namespace_pending(bucket: &str, object: &str) { } } +#[cfg(any(test, feature = "test-util"))] +fn notify_put_object_commit_namespace_acquired(bucket: &str, object: &str) { + let barrier = PUT_OBJECT_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("put object commit barrier mutex should not poison") + .iter() + .find(|barrier| { + barrier.bucket == bucket && barrier.object == object && barrier.pause == PutObjectCommitPause::BeforeNamespace + }) + .cloned(); + if let Some(barrier) = barrier { + barrier.namespace_acquired.store(true, std::sync::atomic::Ordering::Release); + } +} + #[cfg(test)] struct DeleteObjectCommitBarrierState { bucket: String, diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 6207bc72b..35d9c4cc2 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -612,7 +612,7 @@ mod tests { use rustfs_config::server_config::KVS; #[cfg(feature = "test-util")] use rustfs_filemeta::{FileInfo, FileMeta}; - use rustfs_filemeta::{FileInfoVersions, ObjectPartInfo}; + use rustfs_filemeta::{FileInfoVersions, MetaCacheEntry, ObjectPartInfo}; #[cfg(feature = "test-util")] use rustfs_protos::{TIER_MUTATION_RPC_PROTOCOL_VERSION, TierMutationRpcPhase}; use rustfs_rio::{Checksum, ChecksumType}; @@ -2827,21 +2827,30 @@ mod tests { #[tokio::test] #[serial_test::serial(storage_class_env)] - async fn delete_waits_for_decommission_commit_then_removes_every_copy() { + async fn decommission_entry_carries_migration_and_cleanup_mutation_fences() { let temp_dir = tempfile::tempdir().expect("create decommission delete-fence store dir"); - let (_ctx, store, shutdown) = - without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-delete-fence", &[4, 4])).await; + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store_with_layout( + temp_dir.path(), + "decommission-delete-fence", + &[(2, 4), (1, 4)], + CancellationToken::new(), + )) + .await; crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; let bucket = format!("decommission-delete-fence-{}", uuid::Uuid::new_v4()); - let object = "object.bin"; + let object = (0..128) + .map(|index| format!("object-{index}.bin")) + .find(|candidate| store.pools[0].get_disks_by_key(candidate).set_index == 1) + .expect("the deterministic object search should select source set 1"); + let source_body = b"source generation".to_vec(); store .make_bucket(&bucket, &MakeBucketOptions::default()) .await .expect("create decommission delete-fence bucket"); - let mut source = PutObjReader::from_vec(b"source generation".to_vec()); + let mut source = PutObjReader::from_vec(source_body.clone()); store.pools[0] - .put_object(&bucket, object, &mut source, &ObjectOptions::default()) + .put_object(&bucket, &object, &mut source, &ObjectOptions::default()) .await .expect("write source object to the pool being decommissioned"); { @@ -2855,77 +2864,219 @@ mod tests { let barrier = crate::set_disk::PutObjectCommitBarrier::install( &bucket, - object, + &object, crate::set_disk::PutObjectCommitPause::BeforeNamespace, ); - let migration_store = Arc::clone(&store); - let migration_bucket = bucket.clone(); - let migration = tokio::spawn(async move { - let source_reader = migration_store.pools[0] - .get_object_reader( - &migration_bucket, - object, - None, - HeaderMap::new(), - &ObjectOptions { - no_lock: true, - data_movement: true, - raw_data_movement_read: true, + let cleanup_barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(&bucket, &object); + let source_set = store.pools[0].get_disks_by_key(&object); + assert_eq!(source_set.set_index, 1, "the source entry must exercise the non-fixed set cleanup lock"); + let worker_store = Arc::clone(&store); + let worker_bucket = bucket.clone(); + let worker_object = object.clone(); + let worker = tokio::spawn(async move { + worker_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: worker_object, ..Default::default() }, + worker_bucket, + source_set, ) - .await?; - crate::data_movement::migrate_decommission_object( - migration_store, - 0, - migration_bucket, - source_reader, - None, - "test_decommission_delete_fence", - ) - .await + .await }); barrier.wait_until_paused().await; let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket); let delete_store = Arc::clone(&store); let delete_bucket = bucket.clone(); + let delete_object = object.clone(); let delete = tokio::spawn(async move { delete_store - .delete_object(&delete_bucket, object, ObjectOptions::default()) + .delete_object(&delete_bucket, &delete_object, ObjectOptions::default()) .await }); delete_barrier.wait_until_paused().await; delete_barrier.release_and_wait_until_namespace_pending().await; assert!( - !delete.is_finished(), - "DELETE must wait while the decommission source generation is being committed" + !delete_barrier.namespace_acquired() && !delete.is_finished(), + "DELETE must remain before namespace acquisition behind the decommission worker's target-commit mutation fence" + ); + delete.abort(); + assert!( + delete + .await + .expect_err("the blocked DELETE should be canceled") + .is_cancelled(), + "the competing DELETE must remain cancelable while blocked" ); barrier.release(); - migration - .await - .expect("decommission migration task should join") - .expect("decommission migration should commit before DELETE"); - delete - .await - .expect("DELETE task should join") - .expect("DELETE should remove the committed migration generation"); + cleanup_barrier.wait_until_paused().await; + drop(barrier); - for pool in &store.pools { - let err = pool - .get_object_info(&bucket, object, &ObjectOptions::default()) + let fixed_set = Arc::clone(&store.pools[0].disk_set[0]); + let fixed_mutation_barrier = crate::set_disk::PutObjectCommitBarrier::install( + &bucket, + &object, + crate::set_disk::PutObjectCommitPause::BeforeNamespace, + ); + let mutation_bucket = bucket.clone(); + let mutation_object = object.clone(); + let fixed_mutation = tokio::spawn(async move { + let mut reader = PutObjReader::from_vec(b"fixed-domain replacement".to_vec()); + fixed_set + .put_object(&mutation_bucket, &mutation_object, &mut reader, &ObjectOptions::default()) .await - .expect_err("DELETE must remove the source and migrated target copies"); - assert!( - matches!(err, StorageError::ObjectNotFound(_, _)), - "unexpected post-delete pool result: {err:?}" - ); - } - store - .get_object_info(&bucket, object, &ObjectOptions::default()) + }); + fixed_mutation_barrier.wait_until_paused().await; + fixed_mutation_barrier.release_and_wait_until_namespace_pending().await; + assert!( + !fixed_mutation_barrier.namespace_acquired() && !fixed_mutation.is_finished(), + "the source cleanup must retain the fixed mutation fence before the set-0 mutation acquires its namespace" + ); + fixed_mutation.abort(); + assert!( + fixed_mutation + .await + .expect_err("the fixed-domain mutation should be canceled") + .is_cancelled(), + "the competing fixed-domain mutation must remain cancelable while blocked" + ); + drop(fixed_mutation_barrier); + + cleanup_barrier.release(); + worker .await - .expect_err("the migrated source generation must not become visible again"); + .expect("decommission entry worker should join") + .expect("decommission entry should migrate and clean its source"); + + store.pools[0] + .get_object_info(&bucket, &object, &ObjectOptions::default()) + .await + .expect_err("the real decommission entry must clean the source generation"); + let mut target = store.pools[1] + .get_object_reader(&bucket, &object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("the real decommission entry must commit the target generation"); + let mut actual = Vec::new(); + target + .stream + .read_to_end(&mut actual) + .await + .expect("read the migrated target generation"); + assert_eq!(actual, source_body, "the decommission worker must preserve the migrated object body"); + + shutdown.cancel(); + } + + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn batch_delete_real_path_preserves_source_pool_errors_in_any_pool_order() { + let temp_dir = tempfile::tempdir().expect("create batch delete pool-error store dir"); + let (_ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "batch-delete-pool-errors", &[4, 4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + for source_pool_idx in [0, 1] { + { + let mut pool_meta = store.pool_meta.write().await; + for pool in &mut pool_meta.pools { + pool.decommission = None; + } + } + + let bucket = format!("batch-delete-pool-errors-{source_pool_idx}-{}", uuid::Uuid::new_v4()); + let object_names = vec![ + format!("third-{source_pool_idx}.bin"), + format!("first-{source_pool_idx}.bin"), + format!("second-{source_pool_idx}.bin"), + ]; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create batch delete pool-error bucket"); + for pool in &store.pools { + for object_name in &object_names { + let mut reader = PutObjReader::from_vec(format!("pool {} {object_name}", pool.pool_idx).into_bytes()); + pool.put_object(&bucket, object_name, &mut reader, &ObjectOptions::default()) + .await + .expect("seed each object in both the source and active pools"); + } + } + { + let mut pool_meta = store.pool_meta.write().await; + pool_meta.pools[source_pool_idx].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + } + assert!( + store.is_suspended(source_pool_idx).await, + "the injected error pool must be the decommission source" + ); + + let expected_errors = vec![ + StorageError::ErasureWriteQuorum, + StorageError::NamespaceLockQuorumUnavailable { + mode: "delete_objects_commit", + bucket: bucket.clone(), + object: object_names[1].clone(), + required: 3, + achieved: 2, + }, + StorageError::ErasureWriteQuorum, + ]; + let injection = crate::store::object::BatchDeletePoolErrorInjection::install( + &bucket, + source_pool_idx, + object_names.iter().cloned().zip(expected_errors.iter().cloned()).collect(), + ); + let requests = object_names + .iter() + .map(|object_name| ObjectToDelete { + object_name: object_name.clone(), + ..Default::default() + }) + .collect(); + + let (deleted, errors) = store.delete_objects(&bucket, requests, ObjectOptions::default()).await; + + assert_eq!( + injection.observed(), + object_names.len(), + "the source pool must first complete every real delete" + ); + assert_eq!( + errors, + expected_errors.iter().cloned().map(Some).collect::>(), + "a successful pool must not clear a source pool failure at any request index" + ); + assert_eq!( + deleted.iter().map(|object| object.object_name.as_str()).collect::>(), + object_names.iter().map(String::as_str).collect::>(), + "DeleteObjects must preserve request index mapping while aggregating pool failures" + ); + assert!( + deleted.iter().all(|object| object.found), + "the injected source results must retain real delete success data" + ); + + for pool in &store.pools { + for object_name in &object_names { + let error = pool + .get_object_info(&bucket, object_name, &ObjectOptions::default()) + .await + .expect_err("both the active and source pool delete calls must execute"); + assert!( + matches!(error, StorageError::ObjectNotFound(_, _)), + "unexpected residual object: {error:?}" + ); + } + } + drop(injection); + } shutdown.cancel(); } diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index ac644695f..384956bc5 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -838,6 +838,7 @@ struct DeleteAfterObjectLockSnapshotBarrierState { arrived: tokio::sync::Notify, release: tokio::sync::Notify, namespace_pending: tokio::sync::Notify, + namespace_acquired: AtomicBool, } #[cfg(test)] @@ -858,6 +859,7 @@ impl DeleteAfterObjectLockSnapshotBarrier { arrived: tokio::sync::Notify::new(), release: tokio::sync::Notify::new(), namespace_pending: tokio::sync::Notify::new(), + namespace_acquired: AtomicBool::new(false), }); let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) @@ -883,6 +885,10 @@ impl DeleteAfterObjectLockSnapshotBarrier { .await .expect("delete should proceed to its namespace lock after leaving the snapshot barrier"); } + + pub(crate) fn namespace_acquired(&self) -> bool { + self.state.namespace_acquired.load(Ordering::Acquire) + } } #[cfg(test)] @@ -914,6 +920,20 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) { } } +#[cfg(test)] +fn notify_delete_namespace_acquired(bucket: &str) { + let state = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("delete snapshot barrier mutex should not poison") + .as_ref() + .filter(|state| state.bucket == bucket) + .cloned(); + if let Some(state) = state { + state.namespace_acquired.store(true, Ordering::Release); + } +} + #[cfg(test)] struct VersionedDeleteMarkerCommitBarrierState { bucket: String, @@ -1048,6 +1068,88 @@ fn batch_delete_targets_pool(creates_latest_marker: bool, marker_target_pool_idx !creates_latest_marker || marker_target_pool_idx == Some(pool_idx) } +#[cfg(test)] +struct BatchDeletePoolErrorInjectionState { + bucket: String, + pool_idx: usize, + errors: std::collections::HashMap, + observed: std::sync::atomic::AtomicUsize, +} + +#[cfg(test)] +pub(crate) struct BatchDeletePoolErrorInjection { + state: Arc, +} + +#[cfg(test)] +static BATCH_DELETE_POOL_ERROR_INJECTION: std::sync::OnceLock>>> = + std::sync::OnceLock::new(); + +#[cfg(test)] +impl BatchDeletePoolErrorInjection { + pub(crate) fn install(bucket: &str, pool_idx: usize, errors: Vec<(String, Error)>) -> Self { + let state = Arc::new(BatchDeletePoolErrorInjectionState { + bucket: bucket.to_string(), + pool_idx, + errors: errors.into_iter().collect(), + observed: std::sync::atomic::AtomicUsize::new(0), + }); + let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("batch delete pool error injection mutex should not poison"); + assert!(slot.is_none(), "batch delete pool error injection must be unique"); + *slot = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) fn observed(&self) -> usize { + self.state.observed.load(Ordering::Acquire) + } +} + +#[cfg(test)] +impl Drop for BatchDeletePoolErrorInjection { + fn drop(&mut self) { + let mut slot = BATCH_DELETE_POOL_ERROR_INJECTION + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("batch delete pool error injection mutex should not poison"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +fn inject_batch_delete_pool_errors( + bucket: &str, + pool_idx: usize, + object_names: &[String], + result: &mut (Vec, Vec>), +) { + let state = BATCH_DELETE_POOL_ERROR_INJECTION + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("batch delete pool error injection mutex should not poison") + .as_ref() + .filter(|state| state.bucket == bucket && state.pool_idx == pool_idx) + .cloned(); + let Some(state) = state else { + return; + }; + + for (idx, object_name) in object_names.iter().enumerate() { + let Some(error) = state.errors.get(object_name) else { + continue; + }; + if result.1[idx].is_none() && result.0[idx].found { + result.1[idx] = Some(error.clone()); + state.observed.fetch_add(1, Ordering::AcqRel); + } + } +} + fn resolve_batch_delete_pool_results<'a>( initial_error: Option, pool_results: impl IntoIterator)>, @@ -2674,6 +2776,10 @@ impl ECStore { } else { None }; + #[cfg(test)] + if _object_lock_guard.is_some() { + notify_delete_namespace_acquired(bucket); + } if let Some(trigger) = opts.lifecycle_delete_all.as_ref() { let configs = delete_all_configs.as_ref().ok_or(StorageError::PreconditionFailed)?; let expected_bucket_incarnation_id = opts.expected_bucket_incarnation_id.ok_or(StorageError::PreconditionFailed)?; @@ -3006,6 +3112,10 @@ impl ECStore { Ok(guards) => guards, Err(err) => return return_batch_delete_lock_error(objects.as_slice(), err), }; + #[cfg(test)] + if !_object_lock_guards.is_empty() { + notify_delete_namespace_acquired(bucket); + } let delete_config_snapshot = opts .delete_replication_config_snapshot @@ -3057,7 +3167,18 @@ impl ECStore { let pool_opts = opts.clone(); futures.push(async move { + #[cfg(test)] + let pool_object_names = pool_objects + .iter() + .map(|object| object.object_name.clone()) + .collect::>(); let result = pool.delete_objects(bucket, pool_objects, pool_opts).await; + #[cfg(test)] + let result = { + let mut result = result; + inject_batch_delete_pool_errors(bucket, pool.pool_idx, &pool_object_names, &mut result); + result + }; (object_indices, result) }); }