From 1bf27f9a2d1e717e9ab2d371894961ccdf591e0a Mon Sep 17 00:00:00 2001 From: overtrue Date: Fri, 21 Aug 2026 23:53:57 +0800 Subject: [PATCH] fix(ecstore): preserve decommission target write locks --- crates/ecstore/src/store/init.rs | 147 ++++++++++++++++++++++++++++- crates/ecstore/src/store/object.rs | 12 ++- 2 files changed, 155 insertions(+), 4 deletions(-) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 2bac2d094..475852576 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -2817,15 +2817,15 @@ mod tests { 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 { + let delete = tokio::spawn(async move { delete_store .delete_object(&delete_bucket, object, ObjectOptions::default()) .await }); delete_barrier.wait_until_paused().await; - delete_barrier.release(); + delete_barrier.release_and_wait_until_namespace_pending().await; assert!( - tokio::time::timeout(Duration::from_millis(100), &mut delete).await.is_err(), + !delete.is_finished(), "DELETE must wait while the decommission source generation is being committed" ); @@ -2857,6 +2857,147 @@ mod tests { shutdown.cancel(); } + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn versioned_delete_waits_for_decommission_commit_then_publishes_marker() { + 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])) + .await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + + let bucket = format!("versioned-decommission-delete-fence-{}", uuid::Uuid::new_v4()); + let object = "object.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create versioned decommission delete-fence bucket"); + let mut source = PutObjReader::from_vec(b"source generation".to_vec()); + let source_info = store.pools[0] + .put_object( + &bucket, + object, + &mut source, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .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 mut pool_meta = store.pool_meta.write().await; + pool_meta.pools[0].decommission = Some(PoolDecommissionInfo { + start_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }); + } + + 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 { + versioned: true, + version_id: Some(source_version.to_string()), + 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_versioned_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 delete = tokio::spawn(async move { + delete_store + .delete_object( + &delete_bucket, + object, + ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + }); + delete_barrier.wait_until_paused().await; + delete_barrier.release_and_wait_until_namespace_pending().await; + assert!( + !delete.is_finished(), + "versioned DELETE must wait while the source version is being committed" + ); + + barrier.release(); + migration + .await + .expect("versioned decommission migration task should join") + .expect("versioned decommission migration should commit before DELETE"); + let marker = delete + .await + .expect("versioned DELETE task should join") + .expect("versioned DELETE should publish a delete marker after migration"); + assert!(marker.delete_marker, "versioned DELETE must publish a delete marker"); + assert!( + marker.version_id.is_some_and(|version_id| !version_id.is_nil()), + "the delete marker must have a non-nil version ID" + ); + + let err = store + .get_object_info( + &bucket, + object, + &ObjectOptions { + versioned: true, + ..Default::default() + }, + ) + .await + .expect_err("the post-migration delete marker must hide the migrated version"); + assert!( + matches!(err, StorageError::ObjectNotFound(_, _)), + "unexpected latest-version result: {err:?}" + ); + store + .get_object_info( + &bucket, + object, + &ObjectOptions { + versioned: true, + version_id: Some(source_version.to_string()), + ..Default::default() + }, + ) + .await + .expect("the migrated source version must remain addressable below the delete marker"); + + 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 98b5ad29a..e687b0b0e 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -395,7 +395,6 @@ impl ObjectLockDiagGuard { } 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); @@ -818,6 +817,7 @@ struct DeleteAfterObjectLockSnapshotBarrierState { bucket: String, arrived: tokio::sync::Notify, release: tokio::sync::Notify, + namespace_pending: tokio::sync::Notify, } #[cfg(test)] @@ -837,6 +837,7 @@ impl DeleteAfterObjectLockSnapshotBarrier { bucket: bucket.to_string(), arrived: tokio::sync::Notify::new(), release: tokio::sync::Notify::new(), + namespace_pending: tokio::sync::Notify::new(), }); let mut slot = DELETE_AFTER_OBJECT_LOCK_SNAPSHOT_BARRIER .get_or_init(|| std::sync::Mutex::new(None)) @@ -854,6 +855,14 @@ impl DeleteAfterObjectLockSnapshotBarrier { pub(crate) fn release(&self) { self.state.release.notify_one(); } + + pub(crate) async fn release_and_wait_until_namespace_pending(&self) { + let namespace_pending = self.state.namespace_pending.notified(); + self.release(); + tokio::time::timeout(Duration::from_secs(5), namespace_pending) + .await + .expect("delete should proceed to its namespace lock after leaving the snapshot barrier"); + } } #[cfg(test)] @@ -881,6 +890,7 @@ async fn pause_delete_after_object_lock_snapshot(bucket: &str) { if let Some(state) = state { state.arrived.notify_one(); state.release.notified().await; + state.namespace_pending.notify_one(); } }