fix(heal): recover pool metadata during ordinary healing

This commit is contained in:
overtrue
2026-09-08 20:36:34 +08:00
parent c03d3cdd59
commit 9fa1d3f58f
7 changed files with 703 additions and 73 deletions
@@ -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);
+229
View File
@@ -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<Vec<HealResultItem>> {
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::<Vec<_>>();
if targets.is_empty()
|| targets.len() != set.set_drive_count
|| targets.iter().collect::<BTreeSet<_>>().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]
+235 -72
View File
@@ -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<HashMap<String, ReplacementCommitEvidence>>,
ordinary_pool_metadata_required: AtomicBool,
ordinary_pool_metadata_opts: Mutex<Vec<HealOpts>>,
pool_metadata_not_applicable: AtomicBool,
fail_pool_metadata_scope: AtomicBool,
lifecycle_expired: Mutex<HashSet<String>>,
@@ -2177,6 +2217,20 @@ mod resume_loop_tests {
async fn heal_format(&self, _dry: bool) -> Result<(HealResultItem, Option<Error>)> {
Ok((HealResultItem::default(), None))
}
async fn heal_pool_metadata(&self, opts: &HealOpts) -> Result<Vec<HealResultItem>> {
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<bool> {
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,
+4
View File
@@ -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<Vec<HealResultItem>> {
Ok(Vec::new())
}
async fn get_object_meta(&self, _bucket: &str, _object: &str) -> Result<Option<HealObjectInfo>> {
Ok(None)
}
+13
View File
@@ -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<Vec<HealResultItem>> {
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<Vec<HealResultItem>> {
self.ecstore.heal_pool_metadata(opts).await.map_err(Error::Storage)
}
async fn replacement_pool_metadata_applies(&self, opts: &HealOpts) -> Result<bool> {
let pool_index = opts
.pool
+14
View File
@@ -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(),
+205
View File
@@ -1325,6 +1325,7 @@ struct MockStorage {
retry_test_events: Mutex<Vec<String>>,
listed: Mutex<bool>,
list_each_bucket: bool,
pool_metadata_required: bool,
fail_second_listing_page: bool,
recoverable_second_page_failures: Mutex<Option<usize>>,
listing_tokens: Mutex<Vec<Option<String>>>,
@@ -1944,6 +1945,34 @@ impl HealStorageAPI for MockStorage {
.collect())
}
async fn heal_pool_metadata(&self, opts: &HealOpts) -> Result<Vec<HealResultItem>> {
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<bool> {
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<_>>(),
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());