diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 9c166e251..aacb97047 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -3416,6 +3416,9 @@ impl ECStore { ) .await?; + let source_cleanup_mutation_fence = self + .acquire_decommission_source_cleanup_fence(bucket.as_str(), entry.name.as_str(), set.as_ref()) + .await?; let cleanup_result = data_movement::cleanup_source_entry_if_unchanged( set.clone(), bucket.as_str(), @@ -3427,6 +3430,7 @@ impl ECStore { lifecycle_guard: bucket_incarnation_fence .as_ref() .and_then(|guard| guard.namespace_lock_guard()), + object_mutation_fence: Some(&source_cleanup_mutation_fence), }, "decommission", ) diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index f76796bc6..3d7b34fab 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -26,8 +26,7 @@ use crate::storage_api_contracts::{ namespace::NamespaceLocking as _, object::{HTTPPreconditions, ObjectOperations as _}, }; -use crate::store::ECStore; -use crate::store::ObjectLockDiagGuard; +use crate::store::{ECStore, ObjectLockDiagGuard, SourceCleanupMutationFence}; use bytes::Bytes; use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo}; use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex}; @@ -1029,6 +1028,7 @@ pub(crate) enum SourceCleanupError { pub(crate) struct SourceCleanupBucketFence<'a> { pub(crate) expected_incarnation_id: Option, pub(crate) lifecycle_guard: Option<&'a rustfs_lock::NamespaceLockGuard>, + pub(crate) object_mutation_fence: Option<&'a SourceCleanupMutationFence>, } fn ensure_source_cleanup_versions_match( @@ -1066,7 +1066,9 @@ pub(crate) async fn ensure_source_cleanup_versions_unchanged( struct SourceCleanupDeleteBarrierState { bucket: String, object: String, + fence_pending: tokio::sync::Notify, arrived: tokio::sync::Notify, + is_paused: AtomicBool, release: tokio::sync::Notify, } @@ -1080,7 +1082,7 @@ pub(crate) struct SourceCleanupDeleteBarrier { } #[cfg(test)] -static SOURCE_CLEANUP_DELETE_BARRIER: std::sync::OnceLock>>> = +static SOURCE_CLEANUP_DELETE_BARRIERS: std::sync::OnceLock>>> = std::sync::OnceLock::new(); #[cfg(test)] @@ -1093,15 +1095,22 @@ impl SourceCleanupDeleteBarrier { let state = Arc::new(SourceCleanupDeleteBarrierState { bucket: bucket.to_string(), object: object.to_string(), + fence_pending: tokio::sync::Notify::new(), arrived: tokio::sync::Notify::new(), + is_paused: AtomicBool::new(false), release: tokio::sync::Notify::new(), }); - let mut slot = SOURCE_CLEANUP_DELETE_BARRIER - .get_or_init(|| std::sync::Mutex::new(None)) + let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS + .get_or_init(|| std::sync::Mutex::new(Vec::new())) .lock() .expect("source cleanup delete barrier mutex should not poison"); - assert!(slot.is_none(), "source cleanup delete barrier must be unique"); - *slot = Some(Arc::clone(&state)); + assert!( + !barriers + .iter() + .any(|barrier| barrier.bucket == bucket && barrier.object == object), + "source cleanup delete barrier must be unique per object" + ); + barriers.push(Arc::clone(&state)); Self { state } } @@ -1111,35 +1120,58 @@ impl SourceCleanupDeleteBarrier { .expect("source cleanup should reach the pre-delete barrier"); } + pub(crate) async fn wait_until_fence_pending(&self) { + tokio::time::timeout(StdDuration::from_secs(30), self.state.fence_pending.notified()) + .await + .expect("source cleanup should attempt the fixed mutation fence"); + } + + pub(crate) fn is_paused(&self) -> bool { + self.state.is_paused.load(Ordering::Acquire) + } + pub(crate) fn release(&self) { self.state.release.notify_one(); } } +#[cfg(test)] +pub(crate) fn notify_source_cleanup_mutation_fence_pending(bucket: &str, object: &str) { + let barrier = SOURCE_CLEANUP_DELETE_BARRIERS + .get_or_init(|| std::sync::Mutex::new(Vec::new())) + .lock() + .expect("source cleanup delete barrier mutex should not poison") + .iter() + .find(|barrier| barrier.bucket == bucket && barrier.object == object) + .cloned(); + if let Some(barrier) = barrier { + barrier.fence_pending.notify_one(); + } +} + #[cfg(test)] impl Drop for SourceCleanupDeleteBarrier { fn drop(&mut self) { self.state.release.notify_one(); - let mut slot = SOURCE_CLEANUP_DELETE_BARRIER - .get_or_init(|| std::sync::Mutex::new(None)) + let mut barriers = SOURCE_CLEANUP_DELETE_BARRIERS + .get_or_init(|| std::sync::Mutex::new(Vec::new())) .lock() .expect("source cleanup delete barrier mutex should not poison"); - if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { - *slot = None; - } + barriers.retain(|state| !Arc::ptr_eq(state, &self.state)); } } #[cfg(test)] async fn pause_source_cleanup_before_delete(bucket: &str, object: &str) { - let barrier = SOURCE_CLEANUP_DELETE_BARRIER - .get_or_init(|| std::sync::Mutex::new(None)) + let barrier = SOURCE_CLEANUP_DELETE_BARRIERS + .get_or_init(|| std::sync::Mutex::new(Vec::new())) .lock() .expect("source cleanup delete barrier mutex should not poison") - .as_ref() - .filter(|barrier| barrier.bucket == bucket && barrier.object == object) + .iter() + .find(|barrier| barrier.bucket == bucket && barrier.object == object) .cloned(); if let Some(barrier) = barrier { + barrier.is_paused.store(true, Ordering::Release); barrier.arrived.notify_one(); barrier.release.notified().await; } @@ -1155,11 +1187,20 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( op_label: &str, ) -> std::result::Result { let cleanup_key = encode_dir_object(object); - let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?; - let _guard = ns_lock - .get_write_lock(get_lock_acquire_timeout()) - .await - .map_err(Error::from)?; + let source_guard = if bucket_fence + .object_mutation_fence + .is_some_and(SourceCleanupMutationFence::source_lock_covered) + { + None + } else { + let ns_lock = set.new_ns_lock(bucket, cleanup_key.as_str()).await?; + Some( + ns_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .map_err(Error::from)?, + ) + }; if bucket_fence .lifecycle_guard @@ -1169,6 +1210,14 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( "{op_label}: bucket incarnation fence was lost before source cleanup" )))); } + if bucket_fence + .object_mutation_fence + .is_some_and(SourceCleanupMutationFence::is_lock_lost) + { + return Err(SourceCleanupError::Storage(Error::other(format!( + "{op_label}: object mutation fence was lost before source cleanup" + )))); + } ensure_source_cleanup_versions_unchanged(set.clone(), bucket, object, expected, allowed_missing, op_label).await?; @@ -1183,7 +1232,12 @@ pub(crate) async fn cleanup_source_entry_if_unchanged( expected_bucket_incarnation_id: bucket_fence.expected_incarnation_id, ..Default::default() }; - opts.add_namespace_lock_guard(&_guard); + if let Some(source_guard) = source_guard.as_ref() { + opts.add_namespace_lock_guard(source_guard); + } + if let Some(object_mutation_fence) = bucket_fence.object_mutation_fence { + object_mutation_fence.add_namespace_lock_fence(&mut opts); + } if let Some(bucket_lifecycle_guard) = bucket_fence.lifecycle_guard { opts.add_bucket_lifecycle_lock_guard(bucket_lifecycle_guard); } diff --git a/crates/ecstore/src/services/rebalance/entry.rs b/crates/ecstore/src/services/rebalance/entry.rs index 764a68500..e9e1e2343 100644 --- a/crates/ecstore/src/services/rebalance/entry.rs +++ b/crates/ecstore/src/services/rebalance/entry.rs @@ -334,6 +334,7 @@ impl ECStore { lifecycle_guard: bucket_incarnation_fence .as_ref() .and_then(|guard| guard.namespace_lock_guard()), + ..Default::default() }, "rebalance", ), diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index cb8f70766..73c7f2e16 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -11636,6 +11636,7 @@ mod transition_upload_integrity_tests { crate::data_movement::SourceCleanupBucketFence { expected_incarnation_id: None, lifecycle_guard: Some(&bucket_guard), + ..Default::default() }, "test_data_movement", ) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 475852576..5c581a275 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -2859,7 +2859,7 @@ mod tests { #[tokio::test] #[serial_test::serial(storage_class_env)] - async fn versioned_delete_waits_for_decommission_commit_then_publishes_marker() { + async fn versioned_delete_marker_survives_decommission_source_cleanup() { let temp_dir = tempfile::tempdir().expect("create versioned decommission delete-fence store dir"); let (_ctx, store, shutdown) = without_storage_class_env(build_isolated_test_store(temp_dir.path(), "versioned-decommission-delete-fence", &[4, 4])) @@ -2886,6 +2886,12 @@ mod tests { .await .expect("write source version to the pool being decommissioned"); let source_version = source_info.version_id.expect("versioned source must have a version ID"); + let expected_source_versions = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(&bucket, object) + .await + .expect("source versions should be readable before migration") + .expect("source versions should exist before migration"); { let mut pool_meta = store.pool_meta.write().await; pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { @@ -2929,8 +2935,13 @@ mod tests { .await }); barrier.wait_until_paused().await; + barrier.release(); + migration + .await + .expect("versioned decommission migration task should join") + .expect("versioned decommission migration should commit before DELETE"); - let delete_barrier = crate::store::object::DeleteAfterObjectLockSnapshotBarrier::install(&bucket); + 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 { @@ -2946,17 +2957,46 @@ mod tests { .await }); delete_barrier.wait_until_paused().await; - delete_barrier.release_and_wait_until_namespace_pending().await; + let cleanup_set = store.pools[0].get_disks_by_key(object); + crate::data_movement::ensure_source_cleanup_versions_unchanged( + Arc::clone(&cleanup_set), + &bucket, + object, + &expected_source_versions, + &[], + "test_versioned_decommission_delete_fence", + ) + .await + .expect("the committed delete marker must not be published to the suspended source pool"); + + let cleanup_delete_barrier = crate::data_movement::SourceCleanupDeleteBarrier::install(&bucket, object); + let cleanup_store = Arc::clone(&store); + let cleanup_bucket = bucket.clone(); + let cleanup = tokio::spawn(async move { + let mutation_fence = cleanup_store + .acquire_decommission_source_cleanup_fence(&cleanup_bucket, object, cleanup_set.as_ref()) + .await?; + crate::data_movement::cleanup_source_entry_if_unchanged( + cleanup_set, + &cleanup_bucket, + object, + &expected_source_versions, + &[], + crate::data_movement::SourceCleanupBucketFence { + object_mutation_fence: Some(&mutation_fence), + ..Default::default() + }, + "test_versioned_decommission_delete_fence", + ) + .await + }); + cleanup_delete_barrier.wait_until_fence_pending().await; assert!( - !delete.is_finished(), - "versioned DELETE must wait while the source version is being committed" + !cleanup_delete_barrier.is_paused(), + "source cleanup must wait for the versioned DELETE mutation fence" ); - barrier.release(); - migration - .await - .expect("versioned decommission migration task should join") - .expect("versioned decommission migration should commit before DELETE"); + delete_barrier.release(); let marker = delete .await .expect("versioned DELETE task should join") @@ -2967,6 +3007,13 @@ mod tests { "the delete marker must have a non-nil version ID" ); + cleanup_delete_barrier.wait_until_paused().await; + cleanup_delete_barrier.release(); + cleanup + .await + .expect("source cleanup task should join") + .expect("source cleanup should preserve the active-pool delete marker"); + let err = store .get_object_info( &bucket, @@ -2994,6 +3041,18 @@ mod tests { ) .await .expect("the migrated source version must remain addressable below the delete marker"); + store.pools[0] + .get_object_info( + &bucket, + object, + &ObjectOptions { + versioned: true, + version_id: Some(source_version.to_string()), + ..Default::default() + }, + ) + .await + .expect_err("source cleanup must remove the decommissioned source versions"); shutdown.cancel(); } diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 4b8a9cb71..3a298267a 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -151,7 +151,7 @@ pub(crate) mod init_format; pub(crate) mod list_objects; mod multipart; mod object; -pub(crate) use object::ObjectLockDiagGuard; +pub(crate) use object::{ObjectLockDiagGuard, SourceCleanupMutationFence}; pub use object::{ PrepareSelectObjectSnapshotError, PreparedGetObjectReader, SelectObjectSnapshot, SelectObjectSnapshotReadError, SnapshotConsistencyError, diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index e687b0b0e..e913cfe2d 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -36,7 +36,7 @@ use crate::bucket::replication::ReplicationObjectBridge; use crate::disk::OldCurrentSize; use crate::object_api::{NamespaceLockFence, ObjectLockConfigSnapshot}; use crate::set_disk::{ - get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold, + SetDisks, get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold, is_lock_optimization_enabled, is_object_lock_diag_enabled, }; use crate::storage_api_contracts::{ @@ -402,6 +402,25 @@ impl ObjectLockDiagGuard { } } +pub(crate) struct SourceCleanupMutationFence { + guard: ObjectLockDiagGuard, + source_lock_covered: bool, +} + +impl SourceCleanupMutationFence { + pub(crate) fn source_lock_covered(&self) -> bool { + self.source_lock_covered + } + + pub(crate) fn is_lock_lost(&self) -> bool { + self.guard.is_lock_lost() + } + + pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) { + self.guard.add_namespace_lock_fence(opts); + } +} + /// Opaque write-lock guard for the RestoreObject accept path; see /// [`ECStore::acquire_restore_accept_guard`]. Deliberately not a general lock /// API — it only exists so the accept path's restore-status compare-and-set @@ -894,6 +913,82 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) { } } +#[cfg(test)] +struct VersionedDeleteMarkerCommitBarrierState { + bucket: String, + object: String, + arrived: tokio::sync::Notify, + release: tokio::sync::Notify, +} + +#[cfg(test)] +pub(crate) struct VersionedDeleteMarkerCommitBarrier { + state: Arc, +} + +#[cfg(test)] +static VERSIONED_DELETE_MARKER_COMMIT_BARRIER: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +impl VersionedDeleteMarkerCommitBarrier { + pub(crate) fn install(bucket: &str, object: &str) -> Self { + let state = Arc::new(VersionedDeleteMarkerCommitBarrierState { + bucket: bucket.to_string(), + object: object.to_string(), + arrived: tokio::sync::Notify::new(), + release: tokio::sync::Notify::new(), + }); + let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("versioned delete-marker commit barrier mutex should not poison"); + assert!(slot.is_none(), "versioned delete-marker commit barrier must be unique"); + *slot = Some(Arc::clone(&state)); + Self { state } + } + + pub(crate) async fn wait_until_paused(&self) { + tokio::time::timeout(Duration::from_secs(30), self.state.arrived.notified()) + .await + .expect("versioned DELETE should reach the post-marker-commit barrier"); + } + + pub(crate) fn release(&self) { + self.state.release.notify_one(); + } +} + +#[cfg(test)] +impl Drop for VersionedDeleteMarkerCommitBarrier { + fn drop(&mut self) { + self.state.release.notify_one(); + let mut slot = VERSIONED_DELETE_MARKER_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("versioned delete-marker commit barrier mutex should not poison"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +async fn pause_versioned_delete_marker_after_commit(bucket: &str, object: &str) { + let state = VERSIONED_DELETE_MARKER_COMMIT_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("versioned delete-marker commit barrier mutex should not poison") + .as_ref() + .filter(|state| state.bucket == bucket && state.object == object) + .cloned(); + if let Some(state) = state { + state.arrived.notify_one(); + state.release.notified().await; + } +} + /// Whether a delete-time lookup miss on a directory key should trigger an orphan /// empty-directory tree purge (issue #4189). /// @@ -1574,6 +1669,34 @@ impl ECStore { .ok_or_else(|| Error::other("decommission object migration failed to acquire its namespace fence")) } + pub(crate) async fn acquire_decommission_source_cleanup_fence( + &self, + bucket: &str, + object: &str, + source_set: &SetDisks, + ) -> Result { + if self.ctx.lock_manager().is_disabled() { + return Err(Error::other("decommission source cleanup requires namespace locking")); + } + + #[cfg(test)] + crate::data_movement::notify_source_cleanup_mutation_fence_pending(bucket, object); + let object = encode_dir_object(object); + let distributed = self.ctx.is_dist_erasure().await; + let fixed_set = Arc::clone(&self.pools[0].disk_set[0]); + let source_lock_covered = !distributed || same_distributed_lock_domain(&fixed_set.lockers, &source_set.lockers); + // Lock order: fixed store mutation domain first; source cleanup takes its + // hashed source-domain lock second only when this guard does not cover it. + let guard = self + .acquire_object_write_lock("decommission_source_cleanup", bucket, &object) + .await?; + + Ok(SourceCleanupMutationFence { + guard, + source_lock_covered, + }) + } + pub(crate) async fn acquire_all_object_read_locks( &self, op: &'static str, @@ -2681,12 +2804,14 @@ impl ECStore { } for pool in self.pools.iter() { - if self.is_pool_rebalancing(pool.pool_idx).await { + if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await { continue; } match pool.delete_object(bucket, object, opts.clone()).await { Ok(res) => { + #[cfg(test)] + pause_versioned_delete_marker_after_commit(bucket, object).await; if let (Some(api), Some(je)) = (tier_journal_api.as_ref(), journal_entry.as_ref()) { commit_prepared_tier_delete_journal_entry(api, je).await; } @@ -4481,6 +4606,16 @@ mod tests { assert!(lookup_opts.no_lock); assert!(!lookup_opts.skip_decommissioned); assert!(lookup_opts.skip_rebalancing); + + let explicit_version = delete_pool_lookup_opts( + &ObjectOptions { + versioned: true, + version_id: Some(uuid::Uuid::new_v4().to_string()), + ..Default::default() + }, + true, + ); + assert!(!explicit_version.skip_decommissioned); } #[test]