From 03aecc5c3eb1ccdde62e661aa179f390d0ac074e Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 1 Sep 2026 18:31:21 +0800 Subject: [PATCH] fix(heal): avoid pool metadata lock recursion (#6991) --- crates/ecstore/src/core/pools.rs | 36 +++- crates/ecstore/src/store/heal.rs | 276 +++++++++++++++++++++++++++++-- 2 files changed, 297 insertions(+), 15 deletions(-) diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index d69648af7..0bc96dabb 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -8018,6 +8018,19 @@ impl ECStore { write_state: &mut PoolMetaWriteState, operation: &str, ) -> Result<(rustfs_lock::NamespaceLockGuard, PoolMeta)> { + self.acquire_pool_meta_write_guard_with_lock_error(write_state, operation, activation_pool_meta_lock_error) + .await + } + + async fn acquire_pool_meta_write_guard_with_lock_error( + &self, + write_state: &mut PoolMetaWriteState, + operation: &str, + map_lock_error: F, + ) -> Result<(rustfs_lock::NamespaceLockGuard, PoolMeta)> + where + F: FnOnce(rustfs_lock::LockError) -> Error, + { write_state.ensure_write_safe(operation)?; load_pool_meta_identity_observing(self.pools.clone(), write_state).await?; let pool = self.pools.first().cloned().ok_or_else(|| { @@ -8031,7 +8044,7 @@ impl ECStore { let pool_meta_guard = pool_meta_lock .get_write_lock(get_lock_acquire_timeout()) .await - .map_err(activation_pool_meta_lock_error)?; + .map_err(map_lock_error)?; let selection = load_pool_meta_replicas_observing(self.pools.clone(), true, write_state).await?; write_state.observe_replicas(selection.replica_state); write_state.ensure_write_safe(operation)?; @@ -8087,6 +8100,27 @@ impl ECStore { Ok((pool_meta_guard, has_active_source)) } + /// Fence healing of the pool metadata object itself without recursively + /// acquiring its namespace lock through the ordinary capacity probe. The + /// caller must retain the returned write guard through every admitted + /// repair; admission results preserve the input target order. + pub(crate) async fn acquire_pool_meta_object_heal_fence( + &self, + target_pool_indices: &[usize], + ) -> Result<(rustfs_lock::NamespaceLockGuard, Vec>)> { + let mut save_guard = self.pool_meta_save_gate.lock().await; + let (pool_meta_guard, snapshot) = self + .acquire_pool_meta_write_guard_with_lock_error(&mut save_guard, "pool metadata heal admission failed", Error::from) + .await?; + let admissions = target_pool_indices + .iter() + .copied() + .map(|target_pool_index| ensure_external_decommission_target_admission(&snapshot, target_pool_index, "heal")) + .collect(); + drop(save_guard); + Ok((pool_meta_guard, admissions)) + } + pub(crate) async fn acquire_decommission_capacity_release_fence_with_active_source( &self, ) -> Result<(rustfs_lock::NamespaceLockGuard, bool)> { diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index d96765f7a..ed8de1da8 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -36,6 +36,10 @@ fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error { ) } +fn is_pool_meta_object(bucket: &str, object: &str) -> bool { + bucket == RUSTFS_META_BUCKET && object == POOL_META_NAME +} + #[derive(Debug, Clone, Copy)] enum HealFormatPoolSkip { Completed, @@ -492,7 +496,7 @@ impl ECStore { #[cfg(test)] let store_id = self.id; - let mut futures = Vec::with_capacity(pools.len()); + let mut heal_pools = Vec::with_capacity(pools.len()); for pool in pools.iter() { let suspended_complete = { let pool_meta = self.pool_meta.read().await; @@ -520,19 +524,58 @@ impl ECStore { } continue; } - let pool_idx = pool.pool_idx; - let pool = Arc::clone(pool); - let pool_object = object.clone(); - let opts = *opts; - futures.push( - self.run_external_decommission_capacity_heal(pool_idx, bucket, &object, opts, move |opts| async move { - #[cfg(test)] - crate::core::pools::notify_decommission_external_heal_operation_started(store_id); - pool.heal_object(bucket, &pool_object, version_id, &opts).await - }), - ); + heal_pools.push(Arc::clone(pool)); } - let results = join_all(futures).await; + let results = if is_pool_meta_object(bucket, &object) && !opts.no_lock && !heal_pools.is_empty() { + let target_pool_indices = heal_pools.iter().map(|pool| pool.pool_idx).collect::>(); + match self.acquire_pool_meta_object_heal_fence(&target_pool_indices).await { + Ok((pool_meta_guard, admissions)) => { + let fixed_set = self.pools.first().and_then(|pool| pool.disk_set.first()).cloned(); + let futures = heal_pools.iter().zip(admissions).map(|(pool, admission)| { + let pool = Arc::clone(pool); + let pool_object = object.clone(); + let fixed_set = fixed_set.clone(); + let mut opts = *opts; + async move { + admission?; + let fixed_set = fixed_set.ok_or_else(|| Error::other("pool metadata heal requires a fixed set"))?; + let target_set = pool.get_disks_for_heal_object(&pool_object, &opts)?; + opts.no_lock = fixed_set.shares_namespace_lock_domain(&target_set).await; + #[cfg(test)] + if !opts.no_lock { + crate::core::pools::notify_decommission_external_heal_target_lock_attempted(); + } + #[cfg(test)] + crate::core::pools::notify_decommission_external_heal_operation_started(store_id); + pool.heal_object(bucket, &pool_object, version_id, &opts).await + } + }); + let results = join_all(futures).await; + drop(pool_meta_guard); + results + } + Err(err) => (0..heal_pools.len()).map(|_| Err(err.clone())).collect(), + } + } else { + let mut futures = Vec::with_capacity(heal_pools.len()); + for pool in heal_pools { + let pool_idx = pool.pool_idx; + let pool_object = object.clone(); + let opts = *opts; + futures.push(self.run_external_decommission_capacity_heal( + pool_idx, + bucket, + &object, + opts, + move |opts| async move { + #[cfg(test)] + crate::core::pools::notify_decommission_external_heal_operation_started(store_id); + pool.heal_object(bucket, &pool_object, version_id, &opts).await + }, + )); + } + join_all(futures).await + }; let mut errs = Vec::with_capacity(self.pools.len()); let mut ress = Vec::with_capacity(self.pools.len()); @@ -631,7 +674,7 @@ mod tests { }; use crate::core::sets::HealFormatAfterSaveBarrier; use crate::disk::error::Result as DiskResult; - use crate::disk::{DeleteOptions, DiskOption, FORMAT_CONFIG_FILE, format::FormatV3, new_disk}; + use crate::disk::{DeleteOptions, DiskOption, DiskStore, FORMAT_CONFIG_FILE, format::FormatV3, new_disk}; use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::runtime::instance::InstanceContext; use crate::services::rebalance::{ @@ -784,6 +827,30 @@ mod tests { } } + async fn remove_pool_meta_shard(store: &ECStore, pool_idx: usize) -> DiskStore { + let target_set = store.pools[pool_idx].get_disks_by_key(POOL_META_NAME); + let missing_disk = target_set.disks.read().await[0] + .clone() + .expect("pool metadata fixture disk should be online"); + missing_disk + .delete( + RUSTFS_META_BUCKET, + POOL_META_NAME, + DeleteOptions { + recursive: true, + immediate: true, + ..Default::default() + }, + ) + .await + .expect("one pool metadata shard should be removable"); + assert!( + missing_disk.read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false).await.is_err(), + "pool metadata fixture must start with one missing shard" + ); + missing_disk + } + #[tokio::test] async fn heal_erasure_set_scopes_follow_requested_pool_and_set() { let store = minimal_heal_store().await; @@ -1243,6 +1310,121 @@ mod tests { ); } + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_read_repair_reuses_its_write_fence() { + let (_temp_dirs, store, _other_store) = test_two_pool_stores(None).await; + let missing_disk = remove_pool_meta_shard(&store, 0).await; + + let (result, err) = tokio::time::timeout( + std::time::Duration::from_secs(30), + store.handle_heal_object( + RUSTFS_META_BUCKET, + POOL_META_NAME, + "", + &HealOpts { + read_repair: true, + pool: Some(0), + ..Default::default() + }, + ), + ) + .await + .expect("pool metadata read repair must not wait on its own capacity fence") + .expect("pool metadata read repair should complete"); + + assert!(err.is_none(), "pool metadata read repair should succeed: {err:?}"); + assert_eq!(result.object, POOL_META_NAME); + assert!( + missing_disk.read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false).await.is_ok(), + "pool metadata read repair should restore the missing shard" + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_heal_preserves_typed_lock_timeout() { + let (_temp_dirs, store, _other_store) = test_two_pool_stores(None).await; + let lock = store.pools[0] + .new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME) + .await + .expect("pool metadata lock should be created"); + let guard = lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .expect("pool metadata lock should be acquired"); + + let (_, err) = temp_env::async_with_vars( + [(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], + store.handle_heal_object( + RUSTFS_META_BUCKET, + POOL_META_NAME, + "", + &HealOpts { + pool: Some(0), + ..Default::default() + }, + ), + ) + .await + .expect("pool metadata lock timeout should be mapped into the heal result"); + + assert!( + matches!(err, Some(Error::Lock(rustfs_lock::LockError::Timeout { .. }))), + "pool metadata heal must preserve the recoverable lock timeout: {err:?}" + ); + drop(guard); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_neighbor_keeps_ordinary_object_locking() { + let (_temp_dirs, store, _other_store) = test_two_pool_stores(None).await; + let object = "pool.bin.backup"; + save_config(store.pools[0].clone(), object, b"neighbor metadata".to_vec()) + .await + .expect("neighbor metadata fixture should be written"); + + let target_set = store.pools[0].get_disks_by_key(object); + let lock = target_set + .new_ns_lock(RUSTFS_META_BUCKET, object) + .await + .expect("neighbor metadata lock should be created"); + let guard = lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .expect("neighbor metadata lock should be acquired"); + let barrier = DecommissionCapacityLockOrderBarrier::install(store.id, store.id); + let heal_store = Arc::clone(&store); + let mut heal = tokio::spawn(async move { + heal_store + .handle_heal_object( + RUSTFS_META_BUCKET, + object, + "", + &HealOpts { + pool: Some(0), + ..Default::default() + }, + ) + .await + }); + + barrier.wait_until_external_heal_target_lock_attempted().await; + tokio::task::yield_now().await; + assert!(!heal.is_finished(), "neighbor metadata heal must retain ordinary object locking"); + drop(guard); + + let (result, err) = tokio::time::timeout(std::time::Duration::from_secs(30), &mut heal) + .await + .expect("neighbor metadata heal should finish after the object lock is released") + .expect("neighbor metadata heal task should not panic") + .expect("neighbor metadata heal should complete"); + assert!(err.is_none(), "neighbor metadata heal should succeed: {err:?}"); + assert_eq!(result.object, object); + drop(barrier); + } + #[tokio::test] #[serial_test::serial] async fn targeted_heal_is_blocked_by_exact_fit_decommission_reservation() { @@ -1322,6 +1504,72 @@ mod tests { missing_disk.read_xl(&bucket, object, false).await.is_err(), "capacity-blocked targeted heal must not rewrite the missing shard" ); + + let (_, err) = tokio::time::timeout( + std::time::Duration::from_secs(30), + store.handle_heal_object( + RUSTFS_META_BUCKET, + POOL_META_NAME, + "", + &HealOpts { + pool: Some(1), + ..Default::default() + }, + ), + ) + .await + .expect("pool metadata admission should not recurse on its namespace lock") + .expect("pool metadata capacity rejection should be mapped"); + assert!( + matches!(err, Some(Error::SlowDown)), + "pool metadata heal must preserve target reservation admission: {err:?}" + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn pool_meta_heal_keeps_per_target_capacity_admission() { + let (_temp_dirs, store, _other_store) = test_three_pool_stores_with_isolated_node_contexts(None).await; + let layout = DecommissionErasureLayout { data: 1, parity: 0 }; + set_decommission_capacity_info_overrides_for_test( + store.id, + vec![vec![ + DecommissionPoolCapacityInfo::for_test(0, layout, 0, 30, 30), + DecommissionPoolCapacityInfo::for_test(1, layout, 60, 60, 0), + DecommissionPoolCapacityInfo::for_test(2, layout, 60, 60, 0), + ]], + ); + store + .save_current_pool_meta_for_decommission_start(&[0], Vec::new()) + .await + .expect("pool metadata reservation should activate"); + { + 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("pool metadata reservation should be durable"); + assert_eq!(reservation.targets.len(), 1); + assert_eq!(reservation.targets[0].pool_index, 1); + } + + let missing_disk = remove_pool_meta_shard(&store, 2).await; + + let (result, err) = tokio::time::timeout( + std::time::Duration::from_secs(30), + store.handle_heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, "", &HealOpts::default()), + ) + .await + .expect("unscoped pool metadata heal should complete") + .expect("unscoped pool metadata heal should return a mapped result"); + + assert!(err.is_none(), "an admitted target should let unscoped metadata heal succeed: {err:?}"); + assert_eq!(result.object, POOL_META_NAME); + assert!( + missing_disk.read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false).await.is_ok(), + "unscoped metadata heal should repair the admitted target while another target is reserved" + ); } #[tokio::test]