From 9fa1d3f58f4766cdbdba549d6fe9f879aa444721 Mon Sep 17 00:00:00 2001 From: overtrue Date: Tue, 8 Sep 2026 20:36:34 +0800 Subject: [PATCH] fix(heal): recover pool metadata during ordinary healing --- .../distributed_startup_regression_test.rs | 4 +- crates/ecstore/src/store/heal.rs | 229 +++++++++++++ crates/heal/src/heal/erasure_healer.rs | 307 ++++++++++++++---- crates/heal/src/heal/manager/tests.rs | 4 + crates/heal/src/heal/storage.rs | 13 + crates/heal/src/heal/task/heal_bucket.rs | 14 + crates/heal/src/heal/task/tests.rs | 205 ++++++++++++ 7 files changed, 703 insertions(+), 73 deletions(-) diff --git a/crates/e2e_test/src/distributed_startup_regression_test.rs b/crates/e2e_test/src/distributed_startup_regression_test.rs index 57674fbdf..25a331896 100644 --- a/crates/e2e_test/src/distributed_startup_regression_test.rs +++ b/crates/e2e_test/src/distributed_startup_regression_test.rs @@ -498,7 +498,9 @@ mod tests { .stderr(log) .spawn()?, ); - let status = tokio::time::timeout(Duration::from_secs(10), async { + // macOS evaluates each fresh binary copy before its capability hook can run. + let probe_timeout = if cfg!(target_os = "macos") { 60 } else { 10 }; + let status = tokio::time::timeout(Duration::from_secs(probe_timeout), async { loop { if let Some(status) = child.0.try_wait()? { return Ok::<_, std::io::Error>(status); diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index edd5bb673..90cf33ed5 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -18,6 +18,7 @@ use crate::services::rebalance::{REBAL_META_NAME, RebalStatus}; use crate::set_disk::get_lock_acquire_timeout; use crate::storage_api_contracts::heal::HealOperations as _; use crate::storage_api_contracts::namespace::NamespaceLocking as _; +use rustfs_heal_contracts::heal_channel::DriveState; use rustfs_lock::NamespaceLockGuard; use std::collections::BTreeSet; use tracing::trace; @@ -378,6 +379,58 @@ impl ECStore { Ok(result) } + /// Heal every pool metadata owner in the selected scope without allowing + /// one healthy pool to hide another pool's failed repair. + pub async fn heal_pool_metadata(&self, opts: &HealOpts) -> Result> { + let scopes = self.heal_erasure_set_scopes(opts).await?; + let mut results = Vec::new(); + for (pool_index, set_index) in scopes { + if !self.replacement_pool_metadata_applies(pool_index, set_index)? { + continue; + } + let set = &self.pools[pool_index].disk_set[set_index]; + let targets = set.set_endpoints.iter().map(ToString::to_string).collect::>(); + if targets.is_empty() + || targets.len() != set.set_drive_count + || targets.iter().collect::>().len() != targets.len() + { + return Err(Error::SlowDown); + } + // Administrative remove/no-lock options apply to user objects, + // never to the cluster's authoritative metadata transaction. + let metadata_opts = HealOpts { + dry_run: opts.dry_run, + scan_mode: opts.scan_mode, + pool: Some(pool_index), + set: Some(set_index), + ..Default::default() + }; + let (result, error) = self + .handle_heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, "", &metadata_opts) + .await?; + if let Some(error) = error { + return Err(error); + } + if !opts.dry_run { + let ok_state = DriveState::Ok.to_string(); + let complete = result.after.drives.len() == targets.len() + && targets.iter().all(|target| { + let mut outcomes = result.after.drives.iter().filter(|drive| drive.endpoint == *target); + outcomes.next().is_some_and(|drive| drive.state == ok_state) && outcomes.next().is_none() + }); + if !complete + || !set + .replacement_targets_have_version(RUSTFS_META_BUCKET, POOL_META_NAME, "", &targets) + .await? + { + return Err(Error::SlowDown); + } + } + results.push(result); + } + Ok(results) + } + /// Whether this replacement set owns the pool's metadata replica. /// /// Pool metadata follows normal object placement within each pool. A valid @@ -916,6 +969,182 @@ mod tests { assert!(store.replacement_pool_metadata_applies(store.pools.len(), 0).is_err()); } + #[tokio::test] + #[serial_test::serial] + async fn ordinary_pool_metadata_heal_repairs_each_owner_and_preserves_dry_run() { + let (_temp_dirs, store, _other_store) = test_two_pool_stores(None).await; + let first_missing = remove_pool_meta_shard(&store, 0).await; + let second_missing = remove_pool_meta_shard(&store, 1).await; + let destructive_options = HealOpts { + remove: true, + no_lock: true, + ..Default::default() + }; + let results = store + .heal_pool_metadata(&HealOpts { + dry_run: true, + ..destructive_options + }) + .await + .expect("dry-run should inspect both metadata owners without requiring a commit"); + assert_eq!(results.len(), 2); + assert!( + first_missing + .read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false) + .await + .is_err() + ); + assert!( + second_missing + .read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false) + .await + .is_err() + ); + + let lock = store.pools[0] + .new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME) + .await + .expect("metadata namespace lock should be available"); + let guard = lock + .get_read_lock(get_lock_acquire_timeout()) + .await + .expect("a metadata reader should hold the shared fence"); + let error = temp_env::async_with_vars( + [(rustfs_config::ENV_OBJECT_LOCK_ACQUIRE_TIMEOUT, Some("1"))], + store.heal_pool_metadata(&destructive_options), + ) + .await + .expect_err("administrative no-lock cannot bypass the metadata write fence"); + assert!(matches!(error, Error::Lock(rustfs_lock::LockError::Timeout { .. }))); + drop(guard); + let results = store + .heal_pool_metadata(&destructive_options) + .await + .expect("every metadata owner should be repaired"); + assert_eq!(results.len(), 2); + assert!(first_missing.read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false).await.is_ok()); + assert!( + second_missing + .read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false) + .await + .is_ok() + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_pool_metadata_heal_does_not_hide_missing_later_pool() { + let (_temp_dirs, store, _other_store) = test_two_pool_stores(None).await; + delete_config(store.pools[1].clone(), POOL_META_NAME) + .await + .expect("the second pool metadata replica should be removed"); + + let error = store + .heal_pool_metadata(&HealOpts::default()) + .await + .expect_err("the healthy first pool must not hide the second owner's missing replica"); + + assert!(!matches!(error, Error::NoHealRequired)); + let second_set = store.pools[1].get_disks_by_key(POOL_META_NAME); + for disk in second_set.disks.read().await.iter().flatten() { + assert!(disk.read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false).await.is_err()); + } + } + + #[tokio::test] + async fn ordinary_pool_metadata_heal_skips_only_valid_non_owner_sets() { + let mut store = minimal_heal_store().await; + store.ctx = Arc::new(InstanceContext::new()); + for algorithm in [ + crate::disk::format::DistributionAlgoVersion::V1, + crate::disk::format::DistributionAlgoVersion::V2, + crate::disk::format::DistributionAlgoVersion::V3, + ] { + let mut temp_dirs = Vec::new(); + for pool_index in 0..store.pools.len() { + let (dirs, mut pool) = + crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(Arc::clone(&store.ctx), pool_index).await; + temp_dirs.extend(dirs); + Arc::get_mut(&mut pool) + .expect("fixture pool should have one owner") + .distribution_algo = algorithm.clone(); + store.pools[pool_index] = pool; + } + for pool_index in 0..store.pools.len() { + let owner = (0..store.pools[pool_index].disk_set.len()) + .find(|set_index| { + store + .replacement_pool_metadata_applies(pool_index, *set_index) + .expect("valid metadata placement") + }) + .expect("every pool must have one metadata owner"); + let non_owner = 1 - owner; + assert!( + store + .heal_pool_metadata(&HealOpts { + pool: Some(pool_index), + set: Some(non_owner), + ..Default::default() + }) + .await + .expect("valid non-owner should need no metadata write") + .is_empty() + ); + assert!( + store + .heal_pool_metadata(&HealOpts { + pool: Some(pool_index), + set: Some(owner), + ..Default::default() + }) + .await + .is_err(), + "an owner with no authoritative metadata must fail" + ); + assert!( + store + .heal_pool_metadata(&HealOpts { + pool: Some(pool_index), + set: Some(2), + ..Default::default() + }) + .await + .is_err(), + "invalid sets cannot claim the non-owner exemption" + ); + } + } + assert!( + store + .heal_pool_metadata(&HealOpts { + pool: Some(2), + ..Default::default() + }) + .await + .is_err() + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn ordinary_pool_metadata_heal_requires_every_owner_endpoint() { + let (_temp_dirs, store, _other_store) = test_two_pool_stores(None).await; + let owner = store.pools[0].get_disks_by_key(POOL_META_NAME); + let offline_disk = owner.disks.write().await[0] + .take() + .expect("fixture owner disk should start online"); + + let result = store + .heal_pool_metadata(&HealOpts { + pool: Some(0), + ..Default::default() + }) + .await; + + assert!(result.is_err(), "a surviving metadata shard must not hide an offline owner endpoint"); + owner.disks.write().await[0] = Some(offline_disk); + } + 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] diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index df72f281c..63a0c868b 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -842,7 +842,7 @@ impl ErasureSetHealer { } if failed_objects == 0 && skipped_objects == 0 && failed_buckets == 0 { - self.heal_replacement_pool_metadata( + self.heal_pool_metadata( set_disk_id, &mut ErasureSetPassCounters { processed_objects: &mut processed_objects, @@ -941,25 +941,40 @@ impl ErasureSetHealer { Ok(()) } - async fn heal_replacement_pool_metadata( + async fn heal_pool_metadata( &self, set_disk_id: &str, counters: &mut ErasureSetPassCounters<'_>, resume_manager: &ResumeManager, checkpoint_manager: &CheckpointManager, ) -> Result<()> { - if self.replacement_task_id.is_none() { - return Ok(()); - } - if self.target_endpoints.is_empty() { - return Err(Error::TaskExecutionFailed { - message: "Replacement pool metadata heal requires target endpoints".to_string(), - }); - } - - if !self.storage.replacement_pool_metadata_applies(&self.heal_opts).await? { - return Ok(()); - } + let ordinary_opts = if self.replacement_task_id.is_none() { + let (pool_index, set_index) = crate::heal::utils::parse_set_disk_id(set_disk_id)?; + if self.heal_opts.pool.is_some_and(|pool| pool != pool_index) + || self.heal_opts.set.is_some_and(|set| set != set_index) + { + return Err(Error::TaskExecutionFailed { + message: format!("Pool metadata scope does not match resumed set {set_disk_id}"), + }); + } + Some(HealOpts { + dry_run: self.heal_opts.dry_run, + scan_mode: self.heal_opts.scan_mode, + pool: Some(pool_index), + set: Some(set_index), + ..Default::default() + }) + } else { + if self.target_endpoints.is_empty() { + return Err(Error::TaskExecutionFailed { + message: "Replacement pool metadata heal requires target endpoints".to_string(), + }); + } + if !self.storage.replacement_pool_metadata_applies(&self.heal_opts).await? { + return Ok(()); + } + None + }; let object_key = format!("{RUSTFS_META_BUCKET}/{POOL_META_NAME}"); let checkpoint_key = compose_key(&object_key, None); @@ -977,67 +992,88 @@ impl ErasureSetHealer { .set_current_item(Some(RUSTFS_META_BUCKET.to_string()), Some(POOL_META_NAME.to_string())) .await?; - let result = match self - .storage - .heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, &self.heal_opts) - .await - { - Ok((result, None)) if target_outcomes_complete(&result, &self.target_endpoints) => { - let object_size = result_object_size_u64(&result); - match self - .storage - .replacement_targets_have_version( - RUSTFS_META_BUCKET, - POOL_META_NAME, - None, - &self.heal_opts, - &self.target_endpoints, - ) - .await - { - Ok(true) => (object_size, Ok(())), - Ok(false) => ( - object_size, - Err(Error::transient_skip( - "Skipped replacement pool metadata heal because target readback did not confirm the committed version", - )), - ), - Err(err) => ( - object_size, - Err(Error::transient_skip(format!( - "Skipped replacement pool metadata heal because target readback failed: {err}" - ))), - ), + let result = if let Some(opts) = ordinary_opts { + match self.storage.heal_pool_metadata(&opts).await { + Ok(results) if results.is_empty() => return Ok(()), + Ok(results) => { + let [result] = results.as_slice() else { + return Err(Error::TaskExecutionFailed { + message: format!("Pool metadata returned multiple replicas for set {set_disk_id}"), + }); + }; + (result_object_size_u64(result), Ok(())) } + Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err), + Err(err) => match Self::classify_heal_object_error(&err) { + HealObjectOutcome::Absent | HealObjectOutcome::Transient => { + (0, Err(Error::transient_skip(format!("Pool metadata heal must be retried: {err}")))) + } + HealObjectOutcome::Failed => (0, Err(err)), + }, } - Ok((result, None)) => ( - result_object_size_u64(&result), - Err(Error::transient_skip( - "Skipped replacement pool metadata heal because a replacement target was not committed", - )), - ), - Ok((result, Some(err))) => { - let object_size = result_object_size_u64(&result); - match Self::classify_heal_object_error(&err) { + } else { + match self + .storage + .heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, &self.heal_opts) + .await + { + Ok((result, None)) if target_outcomes_complete(&result, &self.target_endpoints) => { + let object_size = result_object_size_u64(&result); + match self + .storage + .replacement_targets_have_version( + RUSTFS_META_BUCKET, + POOL_META_NAME, + None, + &self.heal_opts, + &self.target_endpoints, + ) + .await + { + Ok(true) => (object_size, Ok(())), + Ok(false) => ( + object_size, + Err(Error::transient_skip( + "Skipped replacement pool metadata heal because target readback did not confirm the committed version", + )), + ), + Err(err) => ( + object_size, + Err(Error::transient_skip(format!( + "Skipped replacement pool metadata heal because target readback failed: {err}" + ))), + ), + } + } + Ok((result, None)) => ( + result_object_size_u64(&result), + Err(Error::transient_skip( + "Skipped replacement pool metadata heal because a replacement target was not committed", + )), + ), + Ok((result, Some(err))) => { + let object_size = result_object_size_u64(&result); + match Self::classify_heal_object_error(&err) { + HealObjectOutcome::Absent | HealObjectOutcome::Transient => ( + object_size, + Err(Error::transient_skip(format!( + "Skipped replacement pool metadata heal due to transient error: {err}" + ))), + ), + HealObjectOutcome::Failed => (object_size, Err(err)), + } + } + Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err), + Err(err) => match Self::classify_heal_object_error(&err) { HealObjectOutcome::Absent | HealObjectOutcome::Transient => ( - object_size, + 0, Err(Error::transient_skip(format!( "Skipped replacement pool metadata heal due to transient error: {err}" ))), ), - HealObjectOutcome::Failed => (object_size, Err(err)), - } + HealObjectOutcome::Failed => (0, Err(err)), + }, } - Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err), - Err(err) => match Self::classify_heal_object_error(&err) { - HealObjectOutcome::Absent | HealObjectOutcome::Transient => ( - 0, - Err(Error::transient_skip(format!( - "Skipped replacement pool metadata heal due to transient error: {err}" - ))), - ), - HealObjectOutcome::Failed => (0, Err(err)), - }, }; let (object_size, result) = result; @@ -1056,7 +1092,7 @@ impl ErasureSetHealer { bucket = RUSTFS_META_BUCKET, object = POOL_META_NAME, state = "healed", - "Replacement pool metadata healed" + "Pool metadata healed" ); CheckpointObjectOutcome::Processed } @@ -1073,7 +1109,7 @@ impl ErasureSetHealer { object = POOL_META_NAME, state = "transient_skip", error = %message, - "Replacement pool metadata heal skipped due to transient error" + "Pool metadata heal skipped due to transient error" ); CheckpointObjectOutcome::Skipped } @@ -1090,7 +1126,7 @@ impl ErasureSetHealer { object = POOL_META_NAME, state = "failed", error = %err, - "Replacement pool metadata heal failed" + "Pool metadata heal failed" ); CheckpointObjectOutcome::Failed } @@ -1531,7 +1567,9 @@ impl ErasureSetHealer { ); CheckpointObjectOutcome::Processed } - Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err), + Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => { + return Err(err); + } Err(Error::TransientSkip { message }) => { telemetry_unknown |= !increment_counter(skipped_objects); telemetry_unknown |= !add_bytes(&mut bytes_processed, object_size); @@ -2060,6 +2098,8 @@ mod resume_loop_tests { /// Target-specific physical readback evidence per `compose_key`; the /// fake models a healthy backend unless a test explicitly revokes it. replacement_commit_evidence: Mutex>, + ordinary_pool_metadata_required: AtomicBool, + ordinary_pool_metadata_opts: Mutex>, pool_metadata_not_applicable: AtomicBool, fail_pool_metadata_scope: AtomicBool, lifecycle_expired: Mutex>, @@ -2177,6 +2217,20 @@ mod resume_loop_tests { async fn heal_format(&self, _dry: bool) -> Result<(HealResultItem, Option)> { Ok((HealResultItem::default(), None)) } + async fn heal_pool_metadata(&self, opts: &HealOpts) -> Result> { + if !self.ordinary_pool_metadata_required.load(Ordering::SeqCst) { + return Ok(Vec::new()); + } + self.ordinary_pool_metadata_opts.lock().expect("metadata options").push(*opts); + if !self.replacement_pool_metadata_applies(opts).await? { + return Ok(Vec::new()); + } + let (result, error) = self.heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, opts).await?; + if let Some(error) = error { + return Err(error); + } + Ok(vec![result]) + } async fn replacement_pool_metadata_applies(&self, opts: &HealOpts) -> Result { if self.fail_pool_metadata_scope.load(Ordering::SeqCst) { return Err(Error::other("injected pool metadata scope failure")); @@ -2679,6 +2733,115 @@ mod resume_loop_tests { assert!(state.completed, "successful data heal must be persisted before cleanup is attempted"); } + #[tokio::test] + async fn ordinary_set_heals_pool_metadata_without_replacement_generation_or_targets() { + let env = make_env().await; + env.storage.ordinary_pool_metadata_required.store(true, Ordering::SeqCst); + env.storage + .set_result(POOL_META_NAME, None, replacement_target_ok_result("metadata-disk", POOL_META_NAME)); + assert!(env.healer.replacement_task_id.is_none()); + assert!(env.healer.target_endpoints.is_empty()); + + env.healer + .execute_heal_with_resume(&[], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect("ordinary set recovery should repair metadata even without user buckets"); + + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + { + let opts = env.storage.ordinary_pool_metadata_opts.lock().expect("metadata options"); + assert_eq!(opts.len(), 1); + assert_eq!((opts[0].pool, opts[0].set), (Some(0), Some(0))); + } + let state = env.resume.get_state().await; + assert!(state.completed); + assert_eq!(state.successful_objects, 1, "metadata must enter durable completion counters"); + } + + #[tokio::test] + async fn ordinary_set_pool_metadata_respects_non_owner_and_dry_run() { + let mut env = make_env().await; + env.storage.ordinary_pool_metadata_required.store(true, Ordering::SeqCst); + env.storage.pool_metadata_not_applicable.store(true, Ordering::SeqCst); + env.healer.heal_opts.pool = Some(0); + env.healer.heal_opts.set = Some(1); + env.healer + .execute_heal_with_resume(&[], "pool_0_set_1", &env.resume, &env.checkpoint) + .await + .expect("a valid non-owner set must not invent a metadata replica"); + assert!(env.storage.calls().is_empty()); + assert_eq!(env.resume.get_state().await.successful_objects, 0); + + let mut env = make_env().await; + env.storage.ordinary_pool_metadata_required.store(true, Ordering::SeqCst); + env.healer.heal_opts.dry_run = true; + env.healer.heal_opts.remove = true; + env.healer.heal_opts.no_lock = true; + env.storage.set_replacement_commit_evidence(POOL_META_NAME, None, false); + env.healer + .execute_heal_with_resume(&[], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect("ordinary dry-run metadata work must not require a replacement commit"); + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + let opts = env.storage.ordinary_pool_metadata_opts.lock().expect("metadata options"); + assert_eq!(opts.len(), 1); + assert!(opts[0].dry_run); + assert!(!opts[0].remove); + assert!(!opts[0].no_lock); + } + + #[tokio::test] + async fn ordinary_set_missing_pool_metadata_preserves_retry_state() { + let env = make_env().await; + env.storage.ordinary_pool_metadata_required.store(true, Ordering::SeqCst); + env.storage.set_outcome(POOL_META_NAME, None, HealOutcome::FileNotFound); + + let error = env + .healer + .execute_heal_with_resume(&[], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect_err("missing required pool metadata must prevent ordinary set completion"); + + assert!(matches!(error, Error::TransientSkip { .. })); + let state = env.resume.get_state().await; + assert!(!state.completed); + assert_eq!(state.retry_count, 1); + assert!(CheckpointManager::has_checkpoint(&env.healer.disk, &env.task_id).await); + } + + #[tokio::test] + async fn ordinary_set_pool_metadata_timeout_keeps_control_error() { + let env = make_env().await; + env.storage.ordinary_pool_metadata_required.store(true, Ordering::SeqCst); + env.storage.set_outcome(POOL_META_NAME, None, HealOutcome::Timeout); + + let error = env + .healer + .execute_heal_with_resume(&[], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect_err("metadata timeout must abort the ordinary set pass"); + + assert!(matches!(error, Error::TaskTimeout)); + assert!(!env.resume.get_state().await.completed); + } + + #[tokio::test] + async fn ordinary_set_pool_metadata_rejects_mismatched_explicit_scope() { + let mut env = make_env().await; + env.storage.ordinary_pool_metadata_required.store(true, Ordering::SeqCst); + env.healer.heal_opts.pool = Some(1); + + let error = env + .healer + .execute_heal_with_resume(&[], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect_err("explicit metadata scope must agree with the resumed set"); + + assert!(matches!(error, Error::TaskExecutionFailed { .. })); + assert!(env.storage.calls().is_empty(), "scope mismatch must fail before metadata mutation"); + assert!(!env.resume.get_state().await.completed); + } + #[tokio::test] async fn replacement_completion_keeps_resume_artifacts_until_marker_cleanup() { let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; @@ -2820,7 +2983,7 @@ mod resume_loop_tests { let mut failed_objects = 0; let mut skipped_objects = 0; let error = healer - .heal_replacement_pool_metadata( + .heal_pool_metadata( "pool_0_set_0", &mut super::ErasureSetPassCounters { processed_objects: &mut processed_objects, diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index c065b33d3..3d20754fd 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -555,6 +555,10 @@ async fn completed_retention_scheduler_preserves_progress_aliases_and_atomic_han #[async_trait::async_trait] impl HealStorageAPI for MockStorage { + async fn heal_pool_metadata(&self, _opts: &HealOpts) -> Result> { + Ok(Vec::new()) + } + async fn get_object_meta(&self, _bucket: &str, _object: &str) -> Result> { Ok(None) } diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index 63844c247..600719dc8 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -436,6 +436,15 @@ pub trait HealStorageAPI: Send + Sync { Err(Error::other("target-scoped replacement format is unsupported")) } + /// Heal each pool metadata replica owned by the selected live scope. + /// + /// A successful result requires every applicable owner to finish; an empty + /// result is valid only for a known scope with no metadata replica. Backends + /// without pool metadata must explicitly implement that empty result. + async fn heal_pool_metadata(&self, _opts: &HealOpts) -> Result> { + Err(Error::other("pool metadata healing is unsupported")) + } + /// Whether the selected replacement set owns the pool metadata replica. /// /// Only a topology-aware backend may exempt a valid non-owner set. The @@ -1276,6 +1285,10 @@ impl HealStorageAPI for ECStoreHealStorage { .map_err(Error::Storage) } + async fn heal_pool_metadata(&self, opts: &HealOpts) -> Result> { + self.ecstore.heal_pool_metadata(opts).await.map_err(Error::Storage) + } + async fn replacement_pool_metadata_applies(&self, opts: &HealOpts) -> Result { let pool_index = opts .pool diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index d7465fdb5..6c711b97b 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -328,6 +328,20 @@ impl HealTask { } } + let metadata_opts = HealOpts { + dry_run: self.options.dry_run, + scan_mode: self.options.scan_mode, + pool: self.options.pool_index, + set: self.options.set_index, + ..Default::default() + }; + for result in self + .await_with_control(self.storage.heal_pool_metadata(&metadata_opts)) + .await? + { + self.record_result_item(result).await; + } + if failed > 0 { let failure = BatchHealFailure { scope: "cluster".to_string(), diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 048946f2a..55ffe38fa 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -1325,6 +1325,7 @@ struct MockStorage { retry_test_events: Mutex>, listed: Mutex, list_each_bucket: bool, + pool_metadata_required: bool, fail_second_listing_page: bool, recoverable_second_page_failures: Mutex>, listing_tokens: Mutex>>, @@ -1944,6 +1945,34 @@ impl HealStorageAPI for MockStorage { .collect()) } + async fn heal_pool_metadata(&self, opts: &HealOpts) -> Result> { + if !self.pool_metadata_required { + return Ok(Vec::new()); + } + let scopes = self.erasure_set_scopes.lock().expect("metadata scopes").clone(); + let scopes = if scopes.is_empty() { + vec![(opts.pool.unwrap_or(0), opts.set.unwrap_or(0))] + } else { + scopes + }; + let mut results = Vec::new(); + for (pool, set) in scopes { + let scoped_opts = HealOpts { + pool: Some(pool), + set: Some(set), + ..*opts + }; + let (result, error) = self + .heal_object(RUSTFS_META_BUCKET, crate::heal::POOL_META_NAME, None, &scoped_opts) + .await?; + if let Some(error) = error { + return Err(error); + } + results.push(result); + } + Ok(results) + } + async fn object_exists(&self, _bucket: &str, object: &str) -> Result { if let Some(result) = self.object_exists_by_name.lock().unwrap().get(object).copied() { return match result { @@ -2684,6 +2713,182 @@ async fn test_recursive_bucket_heal_treats_missing_continuation_token_as_end() { ); } +#[tokio::test] +async fn root_heal_restores_pool_metadata_without_user_buckets() { + let storage = Arc::new(MockStorage { + pool_metadata_required: true, + listed_buckets: Mutex::new(Some(Vec::new())), + ..Default::default() + }); + assert!(storage.pool_metadata_required); + let task = HealTask::from_request( + HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal), + storage.clone(), + ); + + task.execute().await.expect("root heal should restore required pool metadata"); + + assert_eq!( + storage.heal_object_calls.lock().expect("heal calls").as_slice(), + [crate::heal::POOL_META_NAME] + ); + assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); +} + +#[tokio::test] +async fn root_heal_pool_metadata_cannot_hide_a_later_owner_failure() { + let storage = Arc::new(MockStorage { + pool_metadata_required: true, + listed_buckets: Mutex::new(Some(Vec::new())), + erasure_set_scopes: Mutex::new(vec![(0, 0), (1, 1)]), + ..Default::default() + }); + storage.heal_object_outcomes.lock().expect("metadata outcomes").insert( + crate::heal::POOL_META_NAME.to_string(), + VecDeque::from([ + MockHealObjectOutcome::UnavailableDrive(DriveState::Ok), + MockHealObjectOutcome::OkWithReadQuorum, + ]), + ); + let task = HealTask::from_request( + HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal), + storage.clone(), + ); + + let error = task + .execute() + .await + .expect_err("one healthy owner cannot satisfy another owner's recovery"); + + assert!(matches!(error, Error::Storage(EcstoreError::InsufficientReadQuorum(_, _)))); + { + let opts = storage.object_heal_opts.lock().expect("owner options"); + assert_eq!( + opts.iter().map(|opts| (opts.pool, opts.set)).collect::>(), + vec![(Some(0), Some(0)), (Some(1), Some(1))] + ); + } + assert!(!matches!(task.get_status().await, HealTaskStatus::Completed)); +} + +#[tokio::test] +async fn root_heal_pool_metadata_does_not_inherit_remove_or_no_lock() { + let storage = Arc::new(MockStorage { + pool_metadata_required: true, + listed_buckets: Mutex::new(Some(Vec::new())), + ..Default::default() + }); + let task = HealTask::from_request( + HealRequest::new( + HealType::Cluster, + HealOptions { + remove_corrupted: true, + no_lock: true, + dry_run: true, + ..Default::default() + }, + HealPriority::Normal, + ), + storage.clone(), + ); + + task.execute() + .await + .expect("dry-run metadata inspection should be fenced and non-destructive"); + + let opts = storage.object_heal_opts.lock().expect("metadata options"); + assert_eq!(opts.len(), 1, "an empty user namespace must still inspect metadata"); + assert!(opts[0].dry_run); + assert!(!opts[0].remove); + assert!(!opts[0].no_lock); +} + +#[tokio::test] +async fn root_heal_pool_metadata_failure_does_not_prevent_user_repairs() { + let storage = Arc::new(MockStorage { + pool_metadata_required: true, + ..Default::default() + }); + storage.heal_object_outcomes.lock().expect("metadata outcome").insert( + crate::heal::POOL_META_NAME.to_string(), + VecDeque::from([MockHealObjectOutcome::OkWithReadQuorum]), + ); + let task = HealTask::from_request( + HealRequest::new( + HealType::Cluster, + HealOptions { + recursive: true, + ..Default::default() + }, + HealPriority::Normal, + ), + storage.clone(), + ); + + let error = task + .execute() + .await + .expect_err("unrecovered metadata must still fail root completion"); + + assert!(matches!(error, Error::Storage(EcstoreError::InsufficientReadQuorum(_, _)))); + assert_eq!( + storage.heal_object_calls.lock().expect("heal calls").as_slice(), + ["object-a", "object-b", crate::heal::POOL_META_NAME] + ); + assert!(!matches!(task.get_status().await, HealTaskStatus::Completed)); +} + +#[tokio::test] +async fn root_heal_pool_metadata_preserves_typed_quorum_failure() { + let storage = Arc::new(MockStorage { + pool_metadata_required: true, + listed_buckets: Mutex::new(Some(Vec::new())), + ..Default::default() + }); + storage.heal_object_outcomes.lock().expect("metadata outcome").insert( + crate::heal::POOL_META_NAME.to_string(), + VecDeque::from([MockHealObjectOutcome::OkWithReadQuorum]), + ); + let task = HealTask::from_request(HealRequest::new(HealType::Cluster, HealOptions::default(), HealPriority::Normal), storage); + + let error = task + .execute() + .await + .expect_err("metadata quorum failure must prevent root completion"); + + assert!(matches!(error, Error::Storage(EcstoreError::InsufficientReadQuorum(_, _)))); + assert!(!matches!(task.get_status().await, HealTaskStatus::Completed)); +} + +#[tokio::test(start_paused = true)] +async fn root_heal_pool_metadata_obeys_task_timeout() { + let storage = Arc::new(MockStorage { + pool_metadata_required: true, + listed_buckets: Mutex::new(Some(Vec::new())), + retry_test_delays: HashMap::from([(crate::heal::POOL_META_NAME.to_string(), Duration::from_secs(10))]), + ..Default::default() + }); + let task = HealTask::from_request( + HealRequest::new( + HealType::Cluster, + HealOptions { + timeout: Some(Duration::from_millis(10)), + ..Default::default() + }, + HealPriority::Normal, + ), + storage, + ); + + let error = task + .execute() + .await + .expect_err("metadata work must stay inside the root task budget"); + + assert!(matches!(error, Error::TaskTimeout)); + assert!(!matches!(task.get_status().await, HealTaskStatus::Completed)); +} + #[tokio::test] async fn test_cluster_heal_visits_bucket_objects() { let storage = Arc::new(MockStorage::default());