From 2861e15d04de94f4c0925b8dd87b86b0863b601e Mon Sep 17 00:00:00 2001 From: cxymds Date: Wed, 2 Sep 2026 16:21:41 +0800 Subject: [PATCH] fix(ecstore): close decommission recovery races --- crates/ecstore/src/core/pools.rs | 58 ++++- crates/ecstore/src/core/pools_test.rs | 295 +++++++++++++++++++++--- crates/ecstore/src/store/multipart.rs | 8 +- scripts/error-other-format-baseline.txt | 4 +- 4 files changed, 322 insertions(+), 43 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index d5e9b3b85..31905d0ed 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -2190,7 +2190,7 @@ fn release_decommission_target_inflight( } else { 0 }; - let changed = temporary_mutation_changed || cleared_pending > 0; + let changed = released > 0 || temporary_mutation_changed || cleared_pending > 0; if changed { renew_decommission_capacity_reservation(reservation, now, true); pool.last_update = now; @@ -8114,6 +8114,9 @@ struct DecommissionCapacityLockOrderBarrierState { target_gate_retry_entered: tokio::sync::Notify, target_gate_retry_entries: AtomicUsize, target_gate_exact_reloads: AtomicUsize, + target_gate_acquire_pause_target: AtomicUsize, + target_gate_acquire_entered: tokio::sync::Notify, + target_gate_acquire_release: tokio::sync::Notify, cancel_before_start_entered: tokio::sync::Notify, cancel_before_start_release: tokio::sync::Notify, cancel_before_start_paused: AtomicBool, @@ -8150,6 +8153,9 @@ impl DecommissionCapacityLockOrderBarrier { target_gate_retry_entered: tokio::sync::Notify::new(), target_gate_retry_entries: AtomicUsize::new(0), target_gate_exact_reloads: AtomicUsize::new(0), + target_gate_acquire_pause_target: AtomicUsize::new(usize::MAX), + target_gate_acquire_entered: tokio::sync::Notify::new(), + target_gate_acquire_release: tokio::sync::Notify::new(), cancel_before_start_entered: tokio::sync::Notify::new(), cancel_before_start_release: tokio::sync::Notify::new(), cancel_before_start_paused: AtomicBool::new(false), @@ -8238,6 +8244,28 @@ impl DecommissionCapacityLockOrderBarrier { self.state.target_gate_exact_reloads.load(Ordering::Acquire) } + #[cfg(feature = "test-util")] + pub(crate) fn pause_target_gate_acquire(&self, target_pool_index: usize) { + self.state + .target_gate_acquire_pause_target + .store(target_pool_index, Ordering::Release); + } + + #[cfg(feature = "test-util")] + pub(crate) async fn wait_until_target_gate_acquire_paused(&self) { + tokio::time::timeout(std::time::Duration::from_secs(30), self.state.target_gate_acquire_entered.notified()) + .await + .expect("decommission capacity mutation should pause before acquiring its target gate"); + } + + #[cfg(feature = "test-util")] + pub(crate) fn release_target_gate_acquire(&self) { + self.state + .target_gate_acquire_pause_target + .store(usize::MAX, Ordering::Release); + self.state.target_gate_acquire_release.notify_one(); + } + pub(crate) fn pause_cancel_before_start(&self) { self.state.cancel_before_start_paused.store(true, Ordering::Release); } @@ -8255,6 +8283,7 @@ impl DecommissionCapacityLockOrderBarrier { #[cfg(feature = "test-util")] pub(crate) fn release_owner(&self) { self.state.owner_release.notify_one(); + self.state.target_gate_acquire_release.notify_one(); } #[cfg(feature = "test-util")] @@ -8318,6 +8347,24 @@ async fn pause_decommission_capacity_before_owner_write(store_id: uuid::Uuid) { } } +#[cfg(test)] +async fn pause_decommission_capacity_before_target_gate_acquire(store_id: uuid::Uuid, target_pool_index: usize) { + let barrier = DECOMMISSION_CAPACITY_LOCK_ORDER_BARRIER + .get_or_init(|| std::sync::Mutex::new(None)) + .lock() + .expect("decommission capacity lock-order barrier should not be poisoned") + .as_ref() + .filter(|state| { + state.owner_store_id == store_id + && state.target_gate_acquire_pause_target.load(Ordering::Acquire) == target_pool_index + }) + .cloned(); + if let Some(barrier) = barrier { + barrier.target_gate_acquire_entered.notify_one(); + barrier.target_gate_acquire_release.notified().await; + } +} + #[cfg(test)] async fn pause_decommission_cancel_before_start_gate(store_id: uuid::Uuid) { let barrier = DECOMMISSION_CAPACITY_LOCK_ORDER_BARRIER @@ -8894,6 +8941,8 @@ impl ECStore { })?; let object = format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/{target_pool_index}"); let target_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, &object).await?; + #[cfg(test)] + pause_decommission_capacity_before_target_gate_acquire(self.id, target_pool_index).await; match target_lock .get_write_lock_quiet(DECOMMISSION_CAPACITY_TARGET_LOCK_TIMEOUT) .await @@ -10010,7 +10059,7 @@ impl ECStore { )?; } true - } else if expected_target_physical_bytes == 0 { + } else if temporary { record_decommission_target_inflight( &mut snapshot, source_pool_index, @@ -17940,6 +17989,11 @@ mod tests { assert!(is_decommission_target_capacity_error(&wrap(Error::StorageFull))); assert!(!is_decommission_target_capacity_error(&wrap(Error::SlowDown))); + let gate_busy = decommission_capacity_blocked_error(format!( + "{DECOMMISSION_CAPACITY_TARGET_GATE_BUSY_PREFIX}7{DECOMMISSION_CAPACITY_TARGET_GATE_BUSY_SUFFIX}" + )); + assert_eq!(decommission_capacity_target_gate_busy_index(&wrap(gate_busy)), Some(7)); + // Cleanup safety: a not-found surfacing from inside a stage is the same // condition as one surfacing directly, so the source entry stays // eligible for cleanup. diff --git a/crates/ecstore/src/core/pools_test.rs b/crates/ecstore/src/core/pools_test.rs index 0201feaca..59427eb49 100644 --- a/crates/ecstore/src/core/pools_test.rs +++ b/crates/ecstore/src/core/pools_test.rs @@ -250,6 +250,27 @@ mod decommission_lock_order_tests { .expect("large-stack ecstore test thread should complete"); } + fn run_large_stack_current_thread_async_test(name: &str, case: C) + where + C: FnOnce() -> F + Send + 'static, + F: Future + 'static, + { + const STACK_SIZE: usize = 32 * 1024 * 1024; + std::thread::Builder::new() + .name(name.to_string()) + .stack_size(STACK_SIZE) + .spawn(move || { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("large-stack current-thread ecstore test runtime should build"); + runtime.block_on(case()); + }) + .expect("large-stack current-thread ecstore test should spawn") + .join() + .expect("large-stack current-thread ecstore test should complete"); + } + #[derive(Clone, Copy)] enum ExternalObjectMutation { Put, @@ -1311,9 +1332,16 @@ mod decommission_lock_order_tests { assert_same_object_copy_fences_capacity_lease_loss(true).await; } - #[tokio::test] + #[test] #[serial_test::serial] - async fn data_movement_multipart_restart_cleans_exact_zero_delta_upload_without_capacity_ledger_state() { + fn data_movement_multipart_restart_cleans_exact_zero_delta_upload_with_recovery_marker() { + run_large_stack_async_test( + "multipart-zero-delta-restart", + data_movement_multipart_restart_cleans_exact_zero_delta_upload_with_recovery_marker_case, + ); + } + + async fn data_movement_multipart_restart_cleans_exact_zero_delta_upload_with_recovery_marker_case() { let (_temp_dirs, store, _other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await; let bucket = test_bucket("multipart-zero-delta-restart"); let object = "staged-without-statfs-delta.bin"; @@ -1418,8 +1446,9 @@ mod decommission_lock_order_tests { .as_ref() .and_then(|info| info.capacity_reservation.as_ref()) .expect("zero-delta reservation should remain active"); - assert_eq!(reservation.pending_target_physical_bytes, 0); + assert_eq!(reservation.pending_target_physical_bytes, 1); assert_eq!(reservation.inflight_target_physical_bytes, 0); + assert_eq!(reservation.targets[0].pending_mutation_id, owner.mutation_id); assert_eq!(reservation.targets[0].temporary_mutations.len(), 1); assert_eq!(reservation.targets[0].temporary_mutations[0].mutation_id, owner.mutation_id.unwrap()); assert_eq!(reservation.targets[0].temporary_mutations[0].physical_bytes, 0); @@ -1465,6 +1494,181 @@ mod decommission_lock_order_tests { assert_eq!(store.data_movement_multipart_discovery_count_for_test(), 1); } + #[tokio::test] + #[serial_test::serial] + async fn multipart_cleanup_proves_target_absence_after_acquiring_capacity_gate() { + let (_temp_dirs, store, _other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await; + let bucket = test_bucket("multipart-cleanup-target-race"); + let object = "published-between-cleanup-proof-and-finalize.bin"; + let body = vec![0x5a; 64 * 1024]; + let version_id = uuid::Uuid::new_v4(); + let mod_time = time::OffsetDateTime::UNIX_EPOCH + time::Duration::seconds(19); + let upload_identity = format!("v1:{version_id}:{}", mod_time.unix_timestamp_nanos()); + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("create the cleanup target-race bucket"); + let incarnation = store + .bucket_incarnation_id(&bucket) + .await + .expect("load the cleanup target-race bucket incarnation"); + let layout = DecommissionErasureLayout { data: 1, parity: 0 }; + set_decommission_capacity_info_overrides_for_test( + store.id, + vec![vec![ + DecommissionPoolCapacityInfo::for_test(0, layout, 0, body.len(), body.len()), + DecommissionPoolCapacityInfo::for_test(1, layout, 0, body.len().saturating_mul(2), body.len().saturating_mul(2)), + DecommissionPoolCapacityInfo::for_test(2, layout, body.len().saturating_mul(2), body.len().saturating_mul(2), 0), + ]], + ); + store + .save_current_pool_meta_for_decommission_start(&[0], Vec::new()) + .await + .expect("activate the cleanup target-race reservation"); + let target_pool_index = store.pool_meta.read().await.pools[0] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .and_then(|reservation| reservation.targets.first()) + .expect("cleanup target-race fixture should reserve one target") + .pool_index; + assert_eq!(target_pool_index, 2); + let owner = decommission_capacity_owner(&*store.pool_meta.read().await).with_mutation_id(uuid::Uuid::new_v4()); + let mut upload_opts = ObjectOptions { + data_movement: true, + src_pool_idx: 0, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(mod_time), + expected_bucket_incarnation_id: Some(incarnation), + ..Default::default() + }; + rustfs_utils::http::insert_str( + &mut upload_opts.user_defined, + rustfs_utils::http::SUFFIX_DATA_MOVEMENT_UPLOAD, + upload_identity.clone(), + ); + owner.apply_to(&mut upload_opts); + let mut cleanup_opts = ObjectOptions { + data_movement: true, + src_pool_idx: 0, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(mod_time), + expected_bucket_incarnation_id: Some(incarnation), + ..Default::default() + }; + owner.apply_to(&mut cleanup_opts); + + let staging_store = Arc::clone(&store); + let staging_bucket = bucket.clone(); + store + .run_decommission_capacity_temporary_mutation(target_pool_index, Some(owner), None, || async move { + new_multipart_upload(&staging_store, target_pool_index, &staging_bucket, object, upload_opts).await?; + Err::<(), crate::error::Error>(crate::error::Error::OperationCanceled) + }) + .await + .expect_err("the fixture should fail after staging its recoverable upload"); + store + .run_decommission_capacity_temporary_mutation(target_pool_index, Some(owner), Some(body.len()), || async { + Err::<(), crate::error::Error>(crate::error::Error::OperationCanceled) + }) + .await + .expect_err("the fixture should retain an exact pending capacity intent"); + { + let pool_meta = store.pool_meta.read().await; + let reservation = pool_meta.pools[0] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .expect("cleanup target-race reservation should remain active"); + assert!(reservation.pending_target_physical_bytes > 0); + assert_eq!(reservation.targets[0].pending_mutation_id, owner.mutation_id); + assert!( + reservation.targets[0] + .temporary_mutations + .iter() + .any(|mutation| mutation.mutation_id == owner.mutation_id.unwrap()) + ); + } + + let target_lock = store + .new_ns_lock( + RUSTFS_META_BUCKET, + &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/{target_pool_index}"), + ) + .await + .expect("create the cleanup target-race gate"); + let target_guard = target_lock + .get_write_lock(Duration::from_secs(30)) + .await + .expect("hold the cleanup target-race gate"); + let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id); + barrier.disable_owner_pause(); + barrier.pause_target_gate_acquire(target_pool_index); + let mut cleanup = tokio::spawn({ + let cleanup_store = Arc::clone(&store); + let cleanup_bucket = bucket.clone(); + let cleanup_identity = upload_identity.clone(); + async move { + cleanup_store + .reconcile_multipart_uploads_for_data_movement( + target_pool_index, + &cleanup_bucket, + object, + &cleanup_identity, + &cleanup_opts, + ) + .await + } + }); + tokio::select! { + _ = barrier.wait_until_target_gate_acquire_paused() => {} + result = &mut cleanup => panic!("cleanup finished before its target-gate proof point: {result:?}"), + } + + let mut target_data = PutObjReader::from_vec(body); + let published = store.pools[target_pool_index] + .put_object( + &bucket, + object, + &mut target_data, + &ObjectOptions { + data_movement: true, + src_pool_idx: 0, + versioned: true, + version_id: Some(version_id.to_string()), + mod_time: Some(mod_time), + expected_bucket_incarnation_id: Some(incarnation), + ..Default::default() + }, + ) + .await + .expect("publish the exact owned target while cleanup waits for its capacity gate"); + assert!(crate::data_movement::is_owned_data_movement_target(&published)); + drop(target_guard); + barrier.release_target_gate_acquire(); + tokio::time::timeout(Duration::from_secs(30), cleanup) + .await + .expect("cleanup should finish after acquiring its target gate") + .expect("cleanup task should not panic") + .expect("cleanup should preserve the now-published target intent"); + drop(barrier); + + let pool_meta = store.pool_meta.read().await; + let reservation = pool_meta.pools[0] + .decommission + .as_ref() + .and_then(|info| info.capacity_reservation.as_ref()) + .expect("cleanup target-race reservation should remain active"); + assert!( + reservation.pending_target_physical_bytes > 0, + "cleanup must not clear capacity for a target published before its gated proof" + ); + assert_eq!(reservation.targets[0].pending_mutation_id, owner.mutation_id); + assert!(reservation.targets[0].temporary_mutations.is_empty()); + } + #[test] #[serial_test::serial] fn data_movement_multipart_restart_reconciles_published_capacity_before_new_upload() { @@ -4245,9 +4449,16 @@ mod decommission_lock_order_tests { } } - #[tokio::test] + #[test] #[serial_test::serial] - async fn data_movement_equivalent_target_reconciles_published_capacity_after_restart() { + fn data_movement_equivalent_target_reconciles_published_capacity_after_restart() { + run_large_stack_current_thread_async_test( + "equivalent-target-capacity-restart", + data_movement_equivalent_target_reconciles_published_capacity_after_restart_case, + ); + } + + async fn data_movement_equivalent_target_reconciles_published_capacity_after_restart_case() { let (_temp_dirs, store, other_store) = test_three_pool_stores_with_three_disk_sets_with_isolated_node_contexts(None).await; let bucket = test_bucket("equivalent-target"); @@ -4478,38 +4689,50 @@ mod decommission_lock_order_tests { crate::error::is_err_object_not_found(&interleaved_target_err) || crate::error::is_err_version_not_found(&interleaved_target_err) ); - let retry_reader = other_store.pools[0] - .get_object_reader( - &bucket, - object, - None, - HeaderMap::new(), - &ObjectOptions { - versioned: true, - version_id: Some(source_version), - no_lock: true, - data_movement: true, - raw_data_movement_read: true, - ..Default::default() - }, - ) + let target_lock = other_store + .new_ns_lock(RUSTFS_META_BUCKET, &format!("{DECOMMISSION_CAPACITY_TARGET_LOCK_PREFIX}/2")) .await - .expect("restart retry should reread the exact source version"); - tokio::time::timeout( - Duration::from_secs(30), - data_movement::migrate_decommission_object( - Arc::clone(&other_store), - 0, - bucket.clone(), - retry_reader, - Some(incarnation), - "equivalent_target_reconcile_retry", - Some(retry_owner), - ), - ) - .await - .expect("equivalent-target retry should not deadlock") - .expect("equivalent target retry should reconcile the pending capacity intent"); + .expect("create the equivalent-target retry gate"); + let target_guard = target_lock + .get_write_lock(Duration::from_secs(30)) + .await + .expect("hold the equivalent-target gate before restart retry"); + let retry_observer = DecommissionCapacityLockOrderBarrier::install(other_store.id, other_store.id); + retry_observer.disable_owner_pause(); + let retry_set = other_store.pools[0].get_disks_by_key(object); + let mut retry = tokio::spawn({ + let retry_store = Arc::clone(&other_store); + let retry_bucket = bucket.clone(); + async move { + retry_store + .decommission_entry_for_test( + 0, + MetaCacheEntry { + name: object.to_string(), + ..Default::default() + }, + retry_bucket, + retry_set, + ) + .await + } + }); + tokio::select! { + _ = retry_observer.wait_until_target_gate_retry() => {} + result = &mut retry => panic!("equivalent-target retry finished before target contention: {result:?}"), + } + drop(target_guard); + tokio::time::timeout(Duration::from_secs(30), retry) + .await + .expect("equivalent-target retry should not deadlock after permit transfer") + .expect("equivalent-target retry task should not panic") + .expect("equivalent target retry should reconcile the pending capacity intent"); + assert_eq!( + retry_observer.target_gate_exact_reloads(), + 1, + "equivalent reconciliation must consume its exact target permit without re-entering busy recovery" + ); + drop(retry_observer); let mut reconciled = crate::core::pools::PoolMeta::default(); reconciled diff --git a/crates/ecstore/src/store/multipart.rs b/crates/ecstore/src/store/multipart.rs index ea2413b66..afe68f1c1 100644 --- a/crates/ecstore/src/store/multipart.rs +++ b/crates/ecstore/src/store/multipart.rs @@ -810,9 +810,11 @@ impl ECStore { .await? .contains(&target_pool_idx) { - return Err(Error::other(format!( - "data movement multipart cleanup target pool {target_pool_idx} is outside its capacity reservation" - ))); + return Err(Error::DecommissionCapacityBlocked { + message: format!( + "data movement multipart cleanup target pool {target_pool_idx} is outside its capacity reservation" + ), + }); } if !self .has_decommission_capacity_temporary_mutation_state(target_pool_idx, owner) diff --git a/scripts/error-other-format-baseline.txt b/scripts/error-other-format-baseline.txt index e41af7bfd..b412c0fc8 100644 --- a/scripts/error-other-format-baseline.txt +++ b/scripts/error-other-format-baseline.txt @@ -27,7 +27,7 @@ 6|crates/ecstore/src/config/com.rs 14|crates/ecstore/src/config/storageclass.rs 182|crates/ecstore/src/core/pools.rs -8|crates/ecstore/src/data_movement/mod.rs +7|crates/ecstore/src/data_movement/mod.rs 2|crates/ecstore/src/data_usage/local_snapshot.rs 12|crates/ecstore/src/data_usage/mod.rs 5|crates/ecstore/src/disk/local.rs @@ -70,5 +70,5 @@ 12|crates/ecstore/src/store/init.rs 2|crates/ecstore/src/store/init_format.rs 3|crates/ecstore/src/store/multipart.rs -7|crates/ecstore/src/store/object.rs +6|crates/ecstore/src/store/object.rs 5|crates/ecstore/src/store/rebalance/support.rs