From ffe4085d9189a51a18e2c666b5e8351419330ab6 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sat, 22 Aug 2026 03:11:16 +0800 Subject: [PATCH] fix(ecstore): fence decommission commit loss --- crates/ecstore/src/data_movement/mod.rs | 57 ++- crates/ecstore/src/object_api/types.rs | 30 +- crates/ecstore/src/set_disk/mod.rs | 4 + crates/ecstore/src/set_disk/ops/multipart.rs | 95 ++++ crates/ecstore/src/set_disk/ops/object.rs | 8 +- crates/ecstore/src/store/init.rs | 493 +++++++++++++++++++ crates/ecstore/src/store/object.rs | 111 ++++- 7 files changed, 764 insertions(+), 34 deletions(-) diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 147e0846d..039b44d72 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -856,7 +856,6 @@ fn is_equivalent_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) fn is_superseding_unversioned_data_movement_object(source: &ObjectInfo, target: &ObjectInfo) -> bool { is_unversioned_data_movement_object(source) && is_unversioned_data_movement_object(target) - && !target.delete_marker && source .mod_time .zip(target.mod_time) @@ -3640,25 +3639,47 @@ mod tests { } #[test] - fn test_precondition_conflict_rejects_newer_delete_marker() { - let source = ObjectInfo { - size: 128, - etag: Some("etag-source".to_string()), - mod_time: Some(OffsetDateTime::UNIX_EPOCH), - ..Default::default() - }; - let target = ObjectInfo { - delete_marker: true, - etag: None, - mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND), - ..source.clone() - }; + fn test_precondition_conflict_accepts_only_newer_null_delete_marker() { + for version_id in [None, Some(Uuid::nil())] { + let source = ObjectInfo { + version_id, + size: 128, + etag: Some("etag-source".to_string()), + mod_time: Some(OffsetDateTime::UNIX_EPOCH), + ..Default::default() + }; + let target = ObjectInfo { + delete_marker: true, + etag: None, + mod_time: OffsetDateTime::UNIX_EPOCH.checked_add(time::Duration::SECOND), + ..source.clone() + }; - let should_resume = - resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(target)), &source, 0, 1) - .expect("delete marker conflict should be evaluated"); + assert!( + resolve_data_movement_overwrite_resume_result( + &Error::PreconditionFailed, + Ok(Some(target.clone())), + &source, + 0, + 1, + ) + .expect("newer null delete marker should be evaluated") + ); - assert!(!should_resume); + let mut same_time = target.clone(); + same_time.mod_time = source.mod_time; + assert!( + !resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(same_time)), &source, 0, 1,) + .expect("same-generation null delete marker should be rejected") + ); + + let mut versioned = target; + versioned.version_id = Some(Uuid::new_v4()); + assert!( + !resolve_data_movement_overwrite_resume_result(&Error::PreconditionFailed, Ok(Some(versioned)), &source, 0, 1,) + .expect("a UUID delete marker must not erase a null source version") + ); + } } #[test] diff --git a/crates/ecstore/src/object_api/types.rs b/crates/ecstore/src/object_api/types.rs index 1bbff7a7f..dd4443db4 100644 --- a/crates/ecstore/src/object_api/types.rs +++ b/crates/ecstore/src/object_api/types.rs @@ -24,7 +24,7 @@ use crate::storage_api_contracts::{ pub struct NamespaceLockFence { signals: Arc>>, #[cfg(test)] - forced_lost: Arc, + forced_lost: Arc>>, } impl Debug for NamespaceLockFence { @@ -40,13 +40,17 @@ impl NamespaceLockFence { Self { signals: Arc::default(), #[cfg(test)] - forced_lost: Arc::new(std::sync::atomic::AtomicBool::new(false)), + forced_lost: Arc::new(vec![Arc::new(std::sync::atomic::AtomicBool::new(false))]), } } pub(crate) fn is_lock_lost(&self) -> bool { #[cfg(test)] - if self.forced_lost.load(std::sync::atomic::Ordering::Acquire) { + if self + .forced_lost + .iter() + .any(|lost| lost.load(std::sync::atomic::Ordering::Acquire)) + { return true; } self.signals.iter().any(|signal| signal.is_lost()) @@ -57,27 +61,26 @@ impl NamespaceLockFence { } fn extend(&mut self, other: &Self) { - if Arc::ptr_eq(&self.signals, &other.signals) { - return; + if !Arc::ptr_eq(&self.signals, &other.signals) { + Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned()); } - Arc::make_mut(&mut self.signals).extend(other.signals.iter().cloned()); #[cfg(test)] - if other.forced_lost.load(std::sync::atomic::Ordering::Acquire) { - self.forced_lost.store(true, std::sync::atomic::Ordering::Release); + if !Arc::ptr_eq(&self.forced_lost, &other.forced_lost) { + Arc::make_mut(&mut self.forced_lost).extend(other.forced_lost.iter().cloned()); } } #[cfg(test)] pub(crate) fn lost_for_test() -> Self { let fence = Self::new(); - fence.forced_lost.store(true, std::sync::atomic::Ordering::Release); + fence.forced_lost[0].store(true, std::sync::atomic::Ordering::Release); fence } #[cfg(test)] pub(crate) fn loss_handle_for_test() -> (Self, Arc) { let fence = Self::new(); - (fence.clone(), Arc::clone(&fence.forced_lost)) + (fence.clone(), Arc::clone(&fence.forced_lost[0])) } } @@ -411,6 +414,13 @@ impl ObjectOptions { self.namespace_lock_fence.get_or_insert_with(NamespaceLockFence::new); } + #[cfg(test)] + pub(crate) fn add_namespace_lock_fence_for_test(&mut self, fence: &NamespaceLockFence) { + self.namespace_lock_fence + .get_or_insert_with(NamespaceLockFence::new) + .extend(fence); + } + pub(crate) fn ensure_lifecycle_delete_all_journal(&mut self) { self.lifecycle_delete_all_journal .get_or_insert_with(|| Arc::new(parking_lot::Mutex::new(LifecycleDeleteAllJournalState::default()))); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index cbf050a4b..b39240996 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -735,8 +735,12 @@ pub(crate) use core::io_primitives::disk_call_counters; mod ctx; mod metadata; mod ops; +#[cfg(test)] +pub(crate) use ops::multipart::NewMultipartUploadCommitObservation; #[cfg(any(test, feature = "test-util"))] pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause}; +#[cfg(test)] +pub(crate) use ops::object::DeleteObjectCommitBarrier; #[cfg(feature = "test-util")] pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier; pub(crate) use ops::object::body_cache_plaintext_len; diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 6af71db17..5aa59b158 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -32,6 +32,8 @@ use crate::crash_inject::{self, CrashPoint}; use crate::multipart_listing::paginate_multipart_listing; use futures::{StreamExt, stream}; use std::future::Future; +#[cfg(test)] +use std::sync::atomic::AtomicBool; #[cfg(any(test, feature = "test-util"))] use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; @@ -65,6 +67,7 @@ impl StaleMultipartCleanupGuard { #[cfg(any(test, feature = "test-util"))] #[derive(Clone, Copy, PartialEq, Eq)] pub enum MultipartCommitPause { + NewUploadBeforeLockLost, PutPartBeforeLockAcquire, PutPartBeforeLockLost, PutPartAfterRename, @@ -156,6 +159,72 @@ impl Drop for MultipartCommitBarrier { } } +#[cfg(test)] +struct NewMultipartUploadCommitObservationState { + bucket: String, + object: String, + committed: AtomicBool, +} + +#[cfg(test)] +pub(crate) struct NewMultipartUploadCommitObservation { + state: Arc, +} + +#[cfg(test)] +static NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +impl NewMultipartUploadCommitObservation { + pub(crate) fn install(bucket: &str, object: &str) -> Self { + let state = Arc::new(NewMultipartUploadCommitObservationState { + bucket: bucket.to_string(), + object: object.to_string(), + committed: AtomicBool::new(false), + }); + let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("new multipart upload commit observation mutex should not poison"); + assert!(slot.is_none(), "new multipart upload commit observation must be unique"); + *slot = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) fn committed(&self) -> bool { + self.state.committed.load(Ordering::Acquire) + } +} + +#[cfg(test)] +impl Drop for NewMultipartUploadCommitObservation { + fn drop(&mut self) { + let mut slot = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("new multipart upload commit observation mutex should not poison"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +fn observe_new_multipart_upload_commit(bucket: &str, object: &str) { + let state = NEW_MULTIPART_UPLOAD_COMMIT_OBSERVATION + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("new multipart upload commit observation mutex should not poison") + .as_ref() + .filter(|state| state.bucket == bucket && state.object == object) + .cloned(); + if let Some(state) = state { + state.committed.store(true, Ordering::Release); + } +} + #[cfg(any(test, feature = "test-util"))] async fn pause_multipart_commit(bucket: &str, object: &str, pause: MultipartCommitPause) { let barrier = { @@ -1615,6 +1684,30 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { let upload_path = Self::get_multipart_upload_dir(bucket, object, upload_uuid.as_str(), opts.data_movement); + #[cfg(any(test, feature = "test-util"))] + pause_multipart_commit(bucket, object, MultipartCommitPause::NewUploadBeforeLockLost).await; + if _object_lock_guard.as_ref().is_some_and(|guard| guard.is_lock_lost()) { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "new_multipart_upload_commit", + bucket: bucket.to_string(), + object: object.to_string(), + required: 1, + achieved: 0, + }); + } + if opts + .namespace_lock_fence + .as_ref() + .is_some_and(NamespaceLockFence::is_lock_lost) + { + return Err(StorageError::NamespaceLockQuorumUnavailable { + mode: "new_multipart_upload_outer_lock", + bucket: bucket.to_string(), + object: object.to_string(), + required: 1, + achieved: 0, + }); + } ensure_multipart_bucket_lifecycle_lock_held(bucket, object, opts)?; Self::write_unique_file_info( &shuffle_disks, @@ -1626,6 +1719,8 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { ) .await .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; + #[cfg(test)] + observe_new_multipart_upload_commit(bucket, object); // evalDisks diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 6ef3c2ce4..b57dbb7a8 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -4773,7 +4773,7 @@ struct DeleteObjectCommitBarrierState { } #[cfg(test)] -struct DeleteObjectCommitBarrier { +pub(crate) struct DeleteObjectCommitBarrier { state: Arc, } @@ -4783,7 +4783,7 @@ static DELETE_OBJECT_COMMIT_BARRIER: std::sync::OnceLock Self { + pub(crate) fn install(bucket: &str, object: &str) -> Self { let state = Arc::new(DeleteObjectCommitBarrierState { bucket: bucket.to_string(), object: object.to_string(), @@ -4799,13 +4799,13 @@ impl DeleteObjectCommitBarrier { Self { state } } - async fn wait_until_paused(&self) { + pub(crate) async fn wait_until_paused(&self) { tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified()) .await .expect("delete object should reach the deterministic commit barrier"); } - fn release(&self) { + pub(crate) fn release(&self) { self.state.release.notify_one(); } } diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 7fc9958fa..f9f822cab 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -1299,6 +1299,139 @@ mod tests { (source_version, expected_source_versions) } + async fn mark_test_pool_decommissioning(store: &Arc, pool_idx: usize) { + let mut pool_meta = store.pool_meta.write().await; + pool_meta.pools[pool_idx].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + } + + async fn write_decommission_test_multipart_source( + store: &Arc, + pool_idx: usize, + bucket: &str, + object: &str, + ) { + let pool = &store.pools[pool_idx]; + let upload = pool + .new_multipart_upload(bucket, object, &ObjectOptions::default()) + .await + .expect("create decommission multipart source upload"); + let first_part = vec![b'm'; 5 * 1024 * 1024]; + let second_part = b"decommission multipart tail".to_vec(); + let mut completed_parts = Vec::with_capacity(2); + for (part_number, body) in [(1, first_part), (2, second_part)] { + let mut reader = PutObjReader::from_vec(body); + let part = pool + .put_object_part(bucket, object, &upload.upload_id, part_number, &mut reader, &ObjectOptions::default()) + .await + .expect("write decommission multipart source part"); + completed_parts.push(crate::storage_api_contracts::multipart::CompletePart { + part_num: part.part_num, + etag: part.etag, + ..Default::default() + }); + } + pool.clone() + .complete_multipart_upload(bucket, object, &upload.upload_id, completed_parts, &ObjectOptions::default()) + .await + .expect("complete decommission multipart source object"); + } + + async fn assert_pool_object_present(pool: &Arc, bucket: &str, object: &str) { + pool.get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect("expected object generation must remain present"); + } + + async fn assert_pool_object_absent(pool: &Arc, bucket: &str, object: &str) { + let err = pool + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect_err("fenced decommission target must remain absent"); + assert!( + matches!(err, StorageError::ObjectNotFound(_, _) | StorageError::VersionNotFound(_, _, _)), + "unexpected fenced target result: {err:?}" + ); + } + + async fn write_suspended_decommission_source(store: &Arc, bucket: &str, object: &str) { + let mut reader = PutObjReader::from_vec(b"suspended source generation".to_vec()); + let source = store.pools[0] + .put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + version_suspended: true, + mod_time: Some(OffsetDateTime::UNIX_EPOCH), + ..Default::default() + }, + ) + .await + .expect("write suspended null source version"); + assert!( + source.version_id.is_none_or(|version_id| version_id.is_nil()), + "suspended source must use the null version identity" + ); + } + + async fn assert_suspended_null_source_present(store: &Arc, bucket: &str, object: &str) { + let versions = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("suspended source versions should be readable") + .expect("suspended source must exist before worker convergence"); + assert!( + versions + .versions + .iter() + .any(|version| !version.deleted && version.version_id.is_none_or(|version_id| version_id.is_nil())), + "the source pool must retain its null data version while DELETE owns the fixed fence" + ); + } + + async fn assert_suspended_decommission_converged(store: &Arc, bucket: &str, object: &str) { + let source_versions = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("source versions should remain readable after suspended convergence"); + assert!( + source_versions.is_none_or(|versions| versions.versions.is_empty()), + "worker convergence must remove only the decommissioned source null version" + ); + + let target_versions = store.pools[1] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("active target versions should be readable") + .expect("active target must retain the suspended DELETE marker"); + assert!( + matches!(target_versions.versions.as_slice(), [marker] if marker.deleted && marker.version_id.is_none_or(|version_id| version_id.is_nil())), + "active target must contain only its null delete marker: {target_versions:?}" + ); + + let err = store + .get_object_info( + bucket, + object, + &ObjectOptions { + version_suspended: true, + ..Default::default() + }, + ) + .await + .expect_err("the active null delete marker must hide the migrated source generation"); + assert!( + matches!(err, StorageError::ObjectNotFound(_, _)), + "unexpected suspended latest-object result: {err:?}" + ); + } + #[tokio::test] #[serial_test::serial(storage_class_env)] async fn tag_updates_skip_active_rebalance_source_pool() { @@ -2969,6 +3102,206 @@ mod tests { shutdown.cancel(); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn decommission_outer_fence_loss_blocks_target_put_commit() { + let temp_dir = tempfile::tempdir().expect("create decommission PUT fence-loss store dir"); + let (_ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-put-fence-loss", &[4, 4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("decommission-put-fence-loss-{}", uuid::Uuid::new_v4()); + let object = "ordinary.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create decommission PUT fence-loss bucket"); + let mut source = PutObjReader::from_vec(b"source generation".to_vec()); + store.pools[0] + .put_object(&bucket, object, &mut source, &ObjectOptions::default()) + .await + .expect("write decommission PUT source"); + mark_test_pool_decommissioning(&store, 0).await; + + let loss_hook = crate::store::object::DecommissionMutationFenceLossHook::install( + &bucket, + object, + crate::store::object::DecommissionMutationFenceTestPhase::Migration, + ); + let barrier = crate::set_disk::PutObjectCommitBarrier::install( + &bucket, + object, + crate::set_disk::PutObjectCommitPause::BeforeQuotaRename, + ); + let source_set = store.pools[0].get_disks_by_key(object); + let worker_store = Arc::clone(&store); + let worker_bucket = bucket.clone(); + let worker = tokio::spawn(async move { + worker_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + worker_bucket, + source_set, + ) + .await + }); + + barrier.wait_until_paused().await; + loss_hook.mark_lost(); + barrier.release(); + drop(barrier); + worker + .await + .expect("decommission PUT fence-loss worker should join") + .expect("a fenced migration failure should remain retryable at entry scope"); + + assert_pool_object_absent(&store.pools[1], &bucket, object).await; + assert_pool_object_present(&store.pools[0], &bucket, object).await; + shutdown.cancel(); + } + + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn decommission_outer_fence_loss_blocks_multipart_commits() { + let temp_dir = tempfile::tempdir().expect("create decommission multipart fence-loss store dir"); + let (_ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-multipart-fence-loss", &[4, 4])) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("decommission-multipart-fence-loss-{}", uuid::Uuid::new_v4()); + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create decommission multipart fence-loss bucket"); + for object in ["new-upload.bin", "complete.bin"] { + write_decommission_test_multipart_source(&store, 0, &bucket, object).await; + } + mark_test_pool_decommissioning(&store, 0).await; + + for (object, pause) in [ + ("new-upload.bin", crate::set_disk::MultipartCommitPause::NewUploadBeforeLockLost), + ("complete.bin", crate::set_disk::MultipartCommitPause::BeforeLockLost), + ] { + let loss_hook = crate::store::object::DecommissionMutationFenceLossHook::install( + &bucket, + object, + crate::store::object::DecommissionMutationFenceTestPhase::Migration, + ); + let commit_observation = (pause == crate::set_disk::MultipartCommitPause::NewUploadBeforeLockLost) + .then(|| crate::set_disk::NewMultipartUploadCommitObservation::install(&bucket, object)); + let barrier = crate::set_disk::MultipartCommitBarrier::install(&bucket, object, pause); + let source_set = store.pools[0].get_disks_by_key(object); + let worker_store = Arc::clone(&store); + let worker_bucket = bucket.clone(); + let worker = tokio::spawn(async move { + worker_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + worker_bucket, + source_set, + ) + .await + }); + + barrier.wait_until_paused().await; + loss_hook.mark_lost(); + barrier.release(); + drop(barrier); + worker + .await + .expect("decommission multipart fence-loss worker should join") + .expect("a fenced multipart migration failure should remain retryable at entry scope"); + + if let Some(commit_observation) = commit_observation { + assert!( + !commit_observation.committed(), + "new multipart upload metadata must not commit after the outer fence is lost" + ); + } + assert_pool_object_absent(&store.pools[1], &bucket, object).await; + assert_pool_object_present(&store.pools[0], &bucket, object).await; + let uploads = store.pools[1] + .list_multipart_uploads(&bucket, object, None, None, None, 100) + .await + .expect("list target multipart uploads after fenced migration"); + assert!(uploads.uploads.is_empty(), "fenced multipart migration must not retain target staging"); + } + + shutdown.cancel(); + } + + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn decommission_outer_fence_loss_blocks_source_cleanup_delete_commit() { + let temp_dir = tempfile::tempdir().expect("create decommission cleanup fence-loss store dir"); + let (_ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-cleanup-fence-loss", &[4, 4])) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("decommission-cleanup-fence-loss-{}", uuid::Uuid::new_v4()); + let object = "cleanup.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create decommission cleanup fence-loss bucket"); + let mut source = PutObjReader::from_vec(b"source generation".to_vec()); + store.pools[0] + .put_object(&bucket, object, &mut source, &ObjectOptions::default()) + .await + .expect("write decommission cleanup source"); + mark_test_pool_decommissioning(&store, 0).await; + + let loss_hook = crate::store::object::DecommissionMutationFenceLossHook::install( + &bucket, + object, + crate::store::object::DecommissionMutationFenceTestPhase::SourceCleanup, + ); + let barrier = crate::set_disk::DeleteObjectCommitBarrier::install(&bucket, object); + let source_set = store.pools[0].get_disks_by_key(object); + let worker_store = Arc::clone(&store); + let worker_bucket = bucket.clone(); + let worker = tokio::spawn(async move { + worker_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + worker_bucket, + source_set, + ) + .await + }); + + barrier.wait_until_paused().await; + loss_hook.mark_lost(); + barrier.release(); + drop(barrier); + let err = worker + .await + .expect("decommission cleanup fence-loss worker should join") + .expect_err("source cleanup must fail after its outer fence is lost"); + assert!( + err.to_string().contains("delete_object_commit"), + "cleanup failure must come from the delete commit fence: {err:?}" + ); + + assert_pool_object_present(&store.pools[0], &bucket, object).await; + assert_pool_object_present(&store.pools[1], &bucket, object).await; + shutdown.cancel(); + } + #[tokio::test] #[serial_test::serial(storage_class_env)] async fn reverse_decommission_reuses_fixed_target_fence_for_put_and_multipart() { @@ -3609,6 +3942,166 @@ mod tests { shutdown.cancel(); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn suspended_delete_marker_then_decommission_worker_converges_null_source() { + let temp_dir = tempfile::tempdir().expect("create suspended decommission DELETE store dir"); + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "suspended-decommission-delete-convergence", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("suspended-decommission-delete-convergence-{}", uuid::Uuid::new_v4()); + let object = "single.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create suspended decommission DELETE bucket"); + write_suspended_decommission_source(&store, &bucket, object).await; + mark_test_pool_decommissioning(&store, 0).await; + + let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object); + let delete_store = Arc::clone(&store); + let delete_bucket = bucket.clone(); + let delete = tokio::spawn(async move { + delete_store + .delete_object( + &delete_bucket, + object, + ObjectOptions { + version_suspended: true, + ..Default::default() + }, + ) + .await + }); + delete_barrier.wait_until_paused().await; + assert_suspended_null_source_present(&store, &bucket, object).await; + + let source_set = store.pools[0].get_disks_by_key(object); + let worker_store = Arc::clone(&store); + let worker_bucket = bucket.clone(); + let worker = tokio::spawn(async move { + worker_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + worker_bucket, + source_set, + ) + .await + }); + + delete_barrier.release(); + let marker = delete + .await + .expect("suspended DELETE task should join") + .expect("suspended DELETE should commit its active-pool marker"); + drop(delete_barrier); + assert!(marker.delete_marker, "suspended DELETE must create a marker"); + assert!( + marker.version_id.is_none_or(|version_id| version_id.is_nil()), + "suspended DELETE marker must keep the null version identity" + ); + worker + .await + .expect("suspended decommission worker should join") + .expect("worker must treat the newer active null marker as a completed migration"); + + assert_suspended_decommission_converged(&store, &bucket, object).await; + shutdown.cancel(); + } + + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn suspended_batch_delete_marker_then_decommission_worker_converges_null_source() { + let temp_dir = tempfile::tempdir().expect("create suspended batch decommission DELETE store dir"); + let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store( + temp_dir.path(), + "suspended-batch-decommission-delete-convergence", + &[4, 4], + )) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("suspended-batch-decommission-delete-convergence-{}", uuid::Uuid::new_v4()); + let object = "batch.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create suspended batch decommission DELETE bucket"); + write_suspended_decommission_source(&store, &bucket, object).await; + mark_test_pool_decommissioning(&store, 0).await; + + let delete_config_snapshot = + Arc::new(crate::bucket::replication::DeleteReplicationConfigSnapshot::from_configs_for_test( + s3s::dto::VersioningConfiguration { + status: Some(s3s::dto::BucketVersioningStatus::from_static(s3s::dto::BucketVersioningStatus::SUSPENDED)), + ..Default::default() + }, + None, + )); + let delete_barrier = crate::store::object::VersionedDeleteMarkerCommitBarrier::install(&bucket, object); + let delete_store = Arc::clone(&store); + let delete_bucket = bucket.clone(); + let delete = tokio::spawn(async move { + delete_store + .delete_objects( + &delete_bucket, + vec![ObjectToDelete { + object_name: object.to_string(), + ..Default::default() + }], + ObjectOptions { + delete_replication_config_snapshot: Some(delete_config_snapshot), + ..Default::default() + }, + ) + .await + }); + delete_barrier.wait_until_paused().await; + assert_suspended_null_source_present(&store, &bucket, object).await; + + let source_set = store.pools[0].get_disks_by_key(object); + let worker_store = Arc::clone(&store); + let worker_bucket = bucket.clone(); + let worker = tokio::spawn(async move { + worker_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + worker_bucket, + source_set, + ) + .await + }); + + delete_barrier.release(); + let (deleted, errors) = delete.await.expect("suspended batch DELETE task should join"); + drop(delete_barrier); + assert!(errors.iter().all(Option::is_none), "suspended batch DELETE should succeed: {errors:?}"); + assert!( + matches!(deleted.as_slice(), [marker] if marker.delete_marker && marker.delete_marker_version_id.is_none_or(|version_id| version_id.is_nil())), + "suspended batch DELETE must create one null marker: {deleted:?}" + ); + worker + .await + .expect("suspended batch decommission worker should join") + .expect("worker must treat the newer batch null marker as a completed migration"); + + assert_suspended_decommission_converged(&store, &bucket, object).await; + shutdown.cancel(); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index ccde2194b..c77f3dade 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -353,6 +353,8 @@ impl fmt::Display for ObjectLockDiagMode { pub(crate) struct ObjectLockDiagGuard { guard: rustfs_lock::NamespaceLockGuard, + #[cfg(test)] + test_namespace_lock_fence: Option, enabled: bool, op: &'static str, bucket: Option, @@ -374,6 +376,8 @@ impl ObjectLockDiagGuard { ) -> Self { Self { guard, + #[cfg(test)] + test_namespace_lock_fence: None, enabled, op, bucket, @@ -400,9 +404,92 @@ impl ObjectLockDiagGuard { if let Some(signal) = self.lock_lost_signal() { opts.add_namespace_lock_lost_signal(signal); } + #[cfg(test)] + if let Some(fence) = self.test_namespace_lock_fence.as_ref() { + opts.add_namespace_lock_fence_for_test(fence); + } } } +#[cfg(test)] +#[derive(Clone, Copy, PartialEq, Eq)] +pub(crate) enum DecommissionMutationFenceTestPhase { + Migration, + SourceCleanup, +} + +#[cfg(test)] +struct DecommissionMutationFenceLossState { + bucket: String, + object: String, + phase: DecommissionMutationFenceTestPhase, + fence: NamespaceLockFence, + loss_handle: Arc, +} + +#[cfg(test)] +pub(crate) struct DecommissionMutationFenceLossHook { + state: Arc, +} + +#[cfg(test)] +static DECOMMISSION_MUTATION_FENCE_LOSS_HOOK: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +impl DecommissionMutationFenceLossHook { + pub(crate) fn install(bucket: &str, object: &str, phase: DecommissionMutationFenceTestPhase) -> Self { + let (fence, loss_handle) = NamespaceLockFence::loss_handle_for_test(); + let state = Arc::new(DecommissionMutationFenceLossState { + bucket: bucket.to_string(), + object: object.to_string(), + phase, + fence, + loss_handle, + }); + let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission mutation fence loss hooks should not poison"); + assert!(slot.is_none(), "decommission mutation fence loss hook must be unique"); + *slot = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) fn mark_lost(&self) { + self.state.loss_handle.store(true, Ordering::Release); + } +} + +#[cfg(test)] +impl Drop for DecommissionMutationFenceLossHook { + fn drop(&mut self) { + let mut slot = DECOMMISSION_MUTATION_FENCE_LOSS_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission mutation fence loss hooks should not poison"); + if slot.as_ref().is_some_and(|hook| Arc::ptr_eq(hook, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +fn decommission_mutation_fence_for_test( + bucket: &str, + object: &str, + phase: DecommissionMutationFenceTestPhase, +) -> Option { + DECOMMISSION_MUTATION_FENCE_LOSS_HOOK + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission mutation fence loss hooks should not poison") + .as_ref() + .filter(|hook| hook.bucket == bucket && hook.object == object && hook.phase == phase) + .map(|hook| hook.fence.clone()) +} + pub(crate) struct SourceCleanupMutationFence { guard: ObjectLockDiagGuard, source_lock_covered: bool, @@ -1827,11 +1914,22 @@ impl ECStore { return Err(Error::other("decommission object migration requires namespace locking")); } + #[cfg(test)] + let test_namespace_lock_fence = + decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::Migration); let object = encode_dir_object(object); let mut opts = ObjectOptions::default(); - self.acquire_object_read_lock_if_needed("decommission_object", bucket, &object, &mut opts) + let guard = self + .acquire_object_read_lock_if_needed("decommission_object", bucket, &object, &mut opts) .await? - .ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence")) + .ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence"))?; + #[cfg(test)] + let guard = { + let mut guard = guard; + guard.test_namespace_lock_fence = test_namespace_lock_fence; + guard + }; + Ok(guard) } pub(super) fn apply_decommission_target_mutation_fence( @@ -1865,6 +1963,9 @@ impl ECStore { #[cfg(test)] crate::data_movement::notify_source_cleanup_mutation_fence_pending(bucket, object); + #[cfg(test)] + let test_namespace_lock_fence = + decommission_mutation_fence_for_test(bucket, object, DecommissionMutationFenceTestPhase::SourceCleanup); let object = encode_dir_object(object); let fixed_set = Arc::clone(&self.pools[0].disk_set[0]); // SetDisks namespaces include pool/set identity, so only the canonical @@ -1875,6 +1976,12 @@ impl ECStore { let guard = self .acquire_object_write_lock("decommission_source_cleanup", bucket, &object) .await?; + #[cfg(test)] + let guard = { + let mut guard = guard; + guard.test_namespace_lock_fence = test_namespace_lock_fence; + guard + }; Ok(SourceCleanupMutationFence { guard,