From 01a0a8b66bccc17e3e47aa77512188af8a3905af Mon Sep 17 00:00:00 2001 From: overtrue Date: Fri, 21 Aug 2026 23:16:21 +0800 Subject: [PATCH] fix(ecstore): fence deletes against decommission commits --- crates/ecstore/src/core/pools.rs | 2 +- crates/ecstore/src/data_movement/mod.rs | 53 ++++++++++++ crates/ecstore/src/store/init.rs | 105 ++++++++++++++++++++++++ crates/ecstore/src/store/object.rs | 75 +++++++++++++++-- 4 files changed, 227 insertions(+), 8 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index fb983dc0c..27631ceab 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -4314,7 +4314,7 @@ impl ECStore { ) -> Result<()> { warn!("decommission_object: start {} {}", &bucket, &rd.object_info.name); let object_name = rd.object_info.name.clone(); - let result = data_movement::migrate_object( + let result = data_movement::migrate_decommission_object( self, pool_idx, bucket.clone(), diff --git a/crates/ecstore/src/data_movement/mod.rs b/crates/ecstore/src/data_movement/mod.rs index 1f4daa703..f76796bc6 100644 --- a/crates/ecstore/src/data_movement/mod.rs +++ b/crates/ecstore/src/data_movement/mod.rs @@ -27,6 +27,7 @@ use crate::storage_api_contracts::{ object::{HTTPPreconditions, ObjectOperations as _}, }; use crate::store::ECStore; +use crate::store::ObjectLockDiagGuard; use bytes::Bytes; use rustfs_filemeta::{FileInfo, FileInfoVersions, ObjectPartInfo}; use rustfs_rio::{EtagResolvable, HashReader, HashReaderDetector, Index, TryGetIndex}; @@ -1330,6 +1331,37 @@ fn data_movement_part_upload_failure_stage(err: &Error) -> &'static str { } } +pub(crate) async fn migrate_decommission_object( + store: Arc, + pool_idx: usize, + bucket: String, + rd: GetObjectReader, + source_bucket_incarnation_id: Option, + op_label: &str, +) -> Result<()> { + let source = rd.object_info.clone(); + let _mutation_fence = store + .acquire_decommission_object_mutation_fence(&bucket, &source.name) + .await?; + let current = find_data_movement_target_info(store.as_ref(), pool_idx, &bucket, &source) + .await? + .ok_or(Error::FileNotFound)?; + if !is_equivalent_data_movement_object_identity(&source, ¤t, true, false) { + return Err(Error::FileNotFound); + } + + migrate_object_inner( + store, + pool_idx, + bucket, + rd, + source_bucket_incarnation_id, + op_label, + Some(&_mutation_fence), + ) + .await +} + pub(crate) async fn migrate_object( store: Arc, pool_idx: usize, @@ -1337,6 +1369,18 @@ pub(crate) async fn migrate_object( rd: GetObjectReader, source_bucket_incarnation_id: Option, op_label: &str, +) -> Result<()> { + migrate_object_inner(store, pool_idx, bucket, rd, source_bucket_incarnation_id, op_label, None).await +} + +async fn migrate_object_inner( + store: Arc, + pool_idx: usize, + bucket: String, + rd: GetObjectReader, + source_bucket_incarnation_id: Option, + op_label: &str, + mutation_fence: Option<&ObjectLockDiagGuard>, ) -> Result<()> { let object_info = rd.object_info.clone(); let has_part_checksums = object_info @@ -1349,6 +1393,9 @@ pub(crate) async fn migrate_object( if should_use_multipart_data_movement(&object_info, has_part_checksums) { let mut new_multipart_opts = data_movement_new_multipart_opts(&object_info, pool_idx); new_multipart_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id; + if let Some(fence) = mutation_fence { + fence.add_namespace_lock_fence(&mut new_multipart_opts); + } let (res, target_pool_idx, expected_bucket_incarnation_id) = match store .handle_new_multipart_upload_with_pool_idx(&bucket, &object_info.name, &new_multipart_opts) .await @@ -1445,6 +1492,9 @@ pub(crate) async fn migrate_object( ) })?; complete_multipart_opts.expected_bucket_incarnation_id = expected_bucket_incarnation_id; + if let Some(fence) = mutation_fence { + fence.add_namespace_lock_fence(&mut complete_multipart_opts); + } if let Err(err) = store .clone() .complete_multipart_upload_for_data_movement( @@ -1608,6 +1658,9 @@ pub(crate) async fn migrate_object( let mut put_opts = data_movement_put_object_opts(&object_info, pool_idx); put_opts.expected_bucket_incarnation_id = source_bucket_incarnation_id; + if let Some(fence) = mutation_fence { + fence.add_namespace_lock_fence(&mut put_opts); + } let (target_pool_idx, put_result) = store .put_object_for_data_movement(&bucket, &object_info.name, &mut data, &put_opts) .await diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index e0452c5bd..2bac2d094 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -2752,6 +2752,111 @@ mod tests { .expect_err("suspended delete must remove the requested UUID version"); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn delete_waits_for_decommission_commit_then_removes_every_copy() { + 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; + 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"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create decommission delete-fence 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 source object to the pool being decommissioned"); + { + let mut pool_meta = store.pool_meta.write().await; + pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + } + assert!(store.is_suspended(0).await, "pool 0 must be a suspended decommission source"); + + let barrier = crate::set_disk::PutObjectCommitBarrier::install( + &bucket, + 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, + ..Default::default() + }, + ) + .await?; + crate::data_movement::migrate_decommission_object( + migration_store, + 0, + migration_bucket, + source_reader, + None, + "test_decommission_delete_fence", + ) + .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 mut delete = tokio::spawn(async move { + delete_store + .delete_object(&delete_bucket, object, ObjectOptions::default()) + .await + }); + delete_barrier.wait_until_paused().await; + delete_barrier.release(); + assert!( + tokio::time::timeout(Duration::from_millis(100), &mut delete).await.is_err(), + "DELETE must wait while the decommission source generation is being committed" + ); + + 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"); + + for pool in &store.pools { + let err = pool + .get_object_info(&bucket, object, &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()) + .await + .expect_err("the migrated source generation must not become visible again"); + + 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 312a25662..98b5ad29a 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -393,6 +393,14 @@ impl ObjectLockDiagGuard { pub(crate) fn is_lock_lost(&self) -> bool { self.guard.is_lock_lost() } + + pub(crate) fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) { + opts.no_lock = true; + opts.ensure_namespace_lock_fence(); + if let Some(signal) = self.lock_lost_signal() { + opts.add_namespace_lock_lost_signal(signal); + } + } } /// Opaque write-lock guard for the RestoreObject accept path; see @@ -410,10 +418,7 @@ impl RestoreAcceptGuard { } pub fn add_namespace_lock_fence(&self, opts: &mut ObjectOptions) { - opts.ensure_namespace_lock_fence(); - if let Some(signal) = self.0.lock_lost_signal() { - opts.add_namespace_lock_lost_signal(signal); - } + self.0.add_namespace_lock_fence(opts); } } @@ -913,6 +918,16 @@ fn writer_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions lookup_opts } +fn delete_pool_lookup_opts(opts: &ObjectOptions, no_lock: bool) -> ObjectOptions { + let mut lookup_opts = writer_pool_lookup_opts(opts, no_lock); + lookup_opts.skip_decommissioned = opts.data_movement; + lookup_opts +} + +fn should_delete_from_all_pools(opts: &ObjectOptions, pool_count: usize) -> bool { + pool_count > 0 && (!opts.versioned && !opts.version_suspended || opts.version_id.is_some()) +} + fn transition_restore_pool_opts(opts: &ObjectOptions) -> ObjectOptions { let mut lookup_opts = opts.clone(); lookup_opts.skip_decommissioned = true; @@ -1533,6 +1548,22 @@ impl ECStore { ))) } + pub(crate) async fn acquire_decommission_object_mutation_fence( + &self, + bucket: &str, + object: &str, + ) -> Result { + if self.ctx.lock_manager().is_disabled() { + return Err(Error::other("decommission object migration requires namespace locking")); + } + + let object = encode_dir_object(object); + let mut opts = ObjectOptions::default(); + 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")) + } + pub(crate) async fn acquire_all_object_read_locks( &self, op: &'static str, @@ -2479,7 +2510,7 @@ impl ECStore { return Ok(ObjectInfo::default()); } - let gopts = writer_pool_lookup_opts(&opts, true); + let gopts = delete_pool_lookup_opts(&opts, true); if opts.data_movement { let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await; @@ -2622,7 +2653,7 @@ impl ECStore { None }; - if !errs.is_empty() && !opts.versioned && !opts.version_suspended { + if should_delete_from_all_pools(&opts, errs.len()) { let mut obj = match self.delete_object_from_all_pools(bucket, object, &opts, errs).await { Ok(obj) => obj, Err(err) => { @@ -2640,7 +2671,7 @@ impl ECStore { } for pool in self.pools.iter() { - if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await { + if self.is_pool_rebalancing(pool.pool_idx).await { continue; } @@ -4433,6 +4464,36 @@ mod tests { assert_eq!(lookup_opts.version_id.as_deref(), Some("vid-1")); } + #[test] + fn ordinary_delete_lookup_includes_decommission_source_and_skips_rebalance_source() { + let lookup_opts = delete_pool_lookup_opts(&ObjectOptions::default(), true); + + assert!(lookup_opts.no_lock); + assert!(!lookup_opts.skip_decommissioned); + assert!(lookup_opts.skip_rebalancing); + } + + #[test] + fn delete_fans_out_for_unversioned_and_explicit_version_mutations() { + assert!(should_delete_from_all_pools(&ObjectOptions::default(), 1)); + assert!(should_delete_from_all_pools( + &ObjectOptions { + versioned: true, + version_id: Some(uuid::Uuid::new_v4().to_string()), + ..Default::default() + }, + 2, + )); + assert!(!should_delete_from_all_pools( + &ObjectOptions { + versioned: true, + ..Default::default() + }, + 1, + )); + assert!(!should_delete_from_all_pools(&ObjectOptions::default(), 0)); + } + #[test] fn data_movement_pool_lookup_opts_keeps_no_lock_for_tiered_moves() { let lookup_opts = data_movement_pool_lookup_opts(