From 06d00e250291521a6ad83795013b84ee0c9ad248 Mon Sep 17 00:00:00 2001 From: overtrue Date: Sun, 23 Aug 2026 11:03:38 +0800 Subject: [PATCH] test(ecstore): cover consumed free-version race --- crates/ecstore/src/store/init.rs | 78 +++++++++++++++++++++++++++++ crates/ecstore/src/store/object.rs | 80 ++++++++++++++++++++++++++++++ 2 files changed, 158 insertions(+) diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index b889001d6..b97f4339e 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -4050,6 +4050,84 @@ mod tests { shutdown.cancel(); } + #[cfg(feature = "test-util")] + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + #[serial_test::serial(storage_class_env)] + async fn decommission_entry_allows_free_version_consumed_before_source_lock() { + let temp_dir = tempfile::tempdir().expect("create consumed free-version store dir"); + let (ctx, store, shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "decommission-free-consumed", &[4, 4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let bucket = format!("decom-free-consumed-{}", uuid::Uuid::new_v4()); + let object = "free-consumed-object"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create consumed free-version bucket"); + let (_, free_version) = seed_transitioned_free_version(&ctx, &store, &bucket, object).await; + + mark_test_pool_decommissioning(&store, 0).await; + let source_set = store.pools[0].get_disks_by_key(object); + let barrier = crate::store::object::DecommissionFreeVersionSourceRaceBarrier::install(&bucket, object); + let decommission = tokio::spawn({ + let store = store.clone(); + let bucket = bucket.clone(); + let source_set = source_set.clone(); + async move { + store + .decommission_entry_for_test_with_bucket_incarnation( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + bucket, + source_set, + ) + .await + } + }); + + barrier.wait_until_paused().await; + store.pools[0] + .delete_object( + &bucket, + object, + ObjectOptions { + versioned: true, + version_id: Some(free_version.to_string()), + incl_free_versions: true, + ..Default::default() + }, + ) + .await + .expect("lifecycle should consume the source free version before decommission locks it"); + assert!( + source_set + .load_file_info_versions_exact(&bucket, object) + .await + .expect("consumed source metadata should remain readable") + .is_none(), + "the lifecycle delete should remove the source free version" + ); + + barrier.release(); + decommission + .await + .expect("decommission task should join") + .expect("a concurrently consumed free version should not fail source cleanup"); + assert!( + store.pools[1] + .get_disks_by_key(object) + .load_file_info_versions_exact(&bucket, object) + .await + .expect("target metadata should remain readable") + .is_none(), + "an already consumed free version should not be recreated on the target" + ); + shutdown.cancel(); + } + #[cfg(feature = "test-util")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial_test::serial(storage_class_env)] diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index dafcb56a0..9ee4afcc7 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -490,6 +490,82 @@ fn decommission_mutation_fence_for_test( .map(|hook| hook.fence.clone()) } +#[cfg(test)] +struct DecommissionFreeVersionSourceRaceState { + bucket: String, + object: String, + arrived: tokio::sync::Notify, + release: tokio::sync::Notify, +} + +#[cfg(test)] +pub(crate) struct DecommissionFreeVersionSourceRaceBarrier { + state: Arc, +} + +#[cfg(test)] +static DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER: std::sync::OnceLock< + std::sync::Mutex>>, +> = std::sync::OnceLock::new(); + +#[cfg(test)] +impl DecommissionFreeVersionSourceRaceBarrier { + pub(crate) fn install(bucket: &str, object: &str) -> Self { + let state = Arc::new(DecommissionFreeVersionSourceRaceState { + bucket: bucket.to_string(), + object: object.to_string(), + arrived: tokio::sync::Notify::new(), + release: tokio::sync::Notify::new(), + }); + let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission free-version source race barrier should not poison"); + assert!(slot.is_none(), "decommission free-version source race 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("decommission should pause before acquiring the free-version source lock"); + } + + pub(crate) fn release(&self) { + self.state.release.notify_one(); + } +} + +#[cfg(test)] +impl Drop for DecommissionFreeVersionSourceRaceBarrier { + fn drop(&mut self) { + self.state.release.notify_one(); + let mut slot = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission free-version source race barrier should not poison"); + if slot.as_ref().is_some_and(|state| Arc::ptr_eq(state, &self.state)) { + *slot = None; + } + } +} + +#[cfg(test)] +async fn pause_decommission_free_version_before_source_lock(bucket: &str, object: &str) { + let state = DECOMMISSION_FREE_VERSION_SOURCE_RACE_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission free-version source race barrier 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; + } +} + pub(crate) struct SourceCleanupMutationFence { guard: ObjectLockDiagGuard, source_lock_covered: bool, @@ -2309,6 +2385,10 @@ impl ECStore { &object, )? }; + #[cfg(test)] + if is_free_version { + pause_decommission_free_version_before_source_lock(bucket, logical_object).await; + } let _object_guards = self .acquire_data_movement_object_write_locks(bucket, &object, opts.src_pool_idx, idx, &mut opts) .await?;