diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index dd82512a5..572d2bcf3 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -16,13 +16,13 @@ #[cfg(test)] mod tests { - use crate::chaos::signed_admin_post; + use crate::chaos::{VersionShardCensus, census_object_version_on_disk, signed_admin_post}; use crate::common::{RustFSTestClusterEnvironment, RustFSTestEnvironment, init_logging}; use aws_sdk_s3::primitives::ByteStream; use std::collections::HashSet; use std::error::Error; use std::path::{Path, PathBuf}; - use tokio::time::{Duration, sleep, timeout}; + use tokio::time::{Duration, Instant, sleep, timeout}; use tracing::info; fn has_file_under(path: &Path) -> bool { @@ -48,6 +48,53 @@ mod tests { disk.join(bucket).join(key).join("xl.meta").is_file() } + fn matching_manifest_count( + disk: &Path, + bucket: &str, + expected_manifests: &[(String, VersionShardCensus)], + ) -> Result> { + let mut matching = 0; + for (key, expected) in expected_manifests { + let actual = census_object_version_on_disk(disk, bucket, key, None)?; + if actual.matches_manifest(expected) { + matching += 1; + } + } + Ok(matching) + } + + fn metadata_count(disk: &Path, bucket: &str, expected_manifests: &[(String, VersionShardCensus)]) -> usize { + expected_manifests + .iter() + .filter(|(key, _)| object_metadata_exists_on_disk(disk, bucket, key)) + .count() + } + + fn heal_task_status_diagnostic(body: &str) -> String { + let Ok(status) = serde_json::from_str::(body) else { + return body.to_string(); + }; + let items = status["items"].as_array(); + let mut unresolved_states = HashSet::new(); + for item in items.into_iter().flatten() { + for drive in item["after"]["drives"].as_array().into_iter().flatten() { + if let Some(state) = drive["state"].as_str() + && state != "ok" + { + unresolved_states.insert(state.to_string()); + } + } + } + let mut unresolved_states = unresolved_states.into_iter().collect::>(); + unresolved_states.sort(); + format!( + "summary={:?}, detail={:?}, item_count={}, unresolved_drive_states={unresolved_states:?}", + status["summary"].as_str(), + status["detail"].as_str(), + items.map_or(0, Vec::len) + ) + } + async fn assert_object_body(env: &RustFSTestEnvironment, bucket: &str, key: &str, expected: &[u8]) { let client = env.create_s3_client(); let response = client @@ -329,35 +376,60 @@ mod tests { } #[tokio::test(flavor = "multi_thread")] - async fn test_cluster_root_heal_rebuilds_replaced_remote_disk() -> Result<(), Box> { + async fn test_cluster_root_heal_resumes_replaced_remote_disk_after_node_restart() -> Result<(), Box> + { init_logging(); - info!("Root recursive heal should rebuild data on a remote node after its disk is replaced and the node rejoins"); + info!("Root recursive heal should resume after its replacement target restarts during a partial rebuild"); let mut cluster = RustFSTestClusterEnvironment::new(4).await?; cluster.set_env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true"); cluster.set_env("RUSTFS_HEAL_ENABLED", "true"); - cluster.set_env("RUSTFS_SCANNER_ENABLED", "true"); + cluster.set_env("RUSTFS_HEAL_AUTO_HEAL_ENABLE", "false"); + cluster.set_env("RUSTFS_SCANNER_ENABLED", "false"); + cluster.set_env("RUSTFS_HEAL_MAX_CONCURRENT_HEALS", "1"); + cluster.set_env("RUSTFS_HEAL_MAX_CONCURRENT_PER_SET", "1"); + cluster.set_env("RUSTFS_HEAL_PAGE_OBJECT_CONCURRENCY", "1"); + cluster.set_env("RUSTFS_HEAL_PAGE_PARALLEL_ENABLE", "false"); + cluster.set_env("RUST_LOG", "rustfs::heal::task=info,rustfs=error"); cluster.start().await?; let clients = cluster.create_all_clients()?; - let bucket = "heal-replaced-remote-disk"; + let bucket = "heal-restart-during-rebuild"; clients[0].create_bucket().bucket(bucket).send().await?; - let online_key = "cluster/online-before-replacement.bin"; - let online_body = b"object written while all cluster nodes are online".to_vec(); - clients[0] - .put_object() - .bucket(bucket) - .key(online_key) - .body(ByteStream::from(online_body.clone())) - .send() - .await?; - let replaced_disk = PathBuf::from(&cluster.nodes[1].data_dir); - assert!( - object_metadata_exists_on_disk(&replaced_disk, bucket, online_key), - "node 1 should contain metadata before disk replacement" - ); + let online_object_count = std::env::var("RUSTFS_HEAL_CHAOS_OBJECT_COUNT") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(24) + .clamp(8, 64); + let object_size_bytes = std::env::var("RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(4 * 1024 * 1024) + .clamp(1024 * 1024, 16 * 1024 * 1024); + let expected_body = vec![0x5a; object_size_bytes]; + let mut expected_manifests = Vec::with_capacity(online_object_count); + for index in 0..online_object_count { + let key = format!("cluster/online/object-{index:04}.bin"); + timeout( + Duration::from_secs(30), + clients[0] + .put_object() + .bucket(bucket) + .key(&key) + .body(ByteStream::from(expected_body.clone())) + .send(), + ) + .await??; + let census = census_object_version_on_disk(&replaced_disk, bucket, &key, None)?; + assert!(census.is_complete(), "node 1 should hold a complete baseline shard for {key}: {census:?}"); + assert!( + !census.expected_part_numbers.is_empty(), + "chaos objects must use physical part shards rather than inline data: {census:?}" + ); + expected_manifests.push((key, census)); + } cluster.stop_node(1)?; std::fs::remove_dir_all(&replaced_disk)?; @@ -365,18 +437,48 @@ mod tests { assert!(!has_file_under(&replaced_disk), "replacement disk must start empty"); let outage_key = "cluster/written-while-node-down.bin"; - let outage_body = b"object written while one remote node is offline".to_vec(); - timeout(Duration::from_secs(30), async { + timeout( + Duration::from_secs(30), clients[0] .put_object() .bucket(bucket) .key(outage_key) - .body(ByteStream::from(outage_body.clone())) - .send() - .await - }) + .body(ByteStream::from(expected_body.clone())) + .send(), + ) .await??; + let mut outage_peer_erasure_indices = HashSet::new(); + for (node_index, node) in cluster.nodes.iter().enumerate() { + if node_index == 1 { + continue; + } + let census = census_object_version_on_disk(Path::new(&node.data_dir), bucket, outage_key, None)?; + assert!( + census.is_complete(), + "online node {node_index} must hold a complete outage-object shard: {census:?}" + ); + let erasure_index = census + .erasure_index + .ok_or_else(|| format!("online node {node_index} outage-object shard has no erasure index: {census:?}"))?; + assert!( + (1..=cluster.nodes.len()).contains(&erasure_index), + "online node {node_index} outage-object erasure index is out of range: {census:?}" + ); + assert!( + outage_peer_erasure_indices.insert(erasure_index), + "outage-object erasure index {erasure_index} is duplicated across online nodes" + ); + } + assert_eq!( + outage_peer_erasure_indices.len(), + cluster.nodes.len().saturating_sub(1), + "every online node must contribute one unique outage-object erasure index" + ); + let expected_outage_target_erasure_index = (1..=cluster.nodes.len()) + .find(|index| !outage_peer_erasure_indices.contains(index)) + .ok_or("online outage-object shards leave no erasure index for the replacement target")?; + cluster.start_node(1).await?; let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url); @@ -399,47 +501,174 @@ mod tests { serde_json::Value::Bool(true), "cluster heal status should recover before root heal starts: {recovered}" ); + assert_eq!( + matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?, + 0, + "auto heal is disabled, so the replacement target must remain empty before the explicit root heal" + ); + assert!( + !census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?.has_xl_meta, + "the object written during the outage must be absent before the explicit root heal" + ); let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#; let heal_url = format!("{}/rustfs/admin/v3/heal/?forceStart=true", cluster.nodes[0].url); - signed_admin_post(&heal_url, Some(heal_body), &cluster.access_key, &cluster.secret_key).await?; + let heal_start_body = signed_admin_post(&heal_url, Some(heal_body), &cluster.access_key, &cluster.secret_key).await?; + let heal_start: serde_json::Value = serde_json::from_str(&heal_start_body) + .map_err(|err| format!("heal start response is not JSON ({err}): {heal_start_body}"))?; + let client_token = heal_start["clientToken"] + .as_str() + .filter(|token| !token.is_empty()) + .ok_or_else(|| format!("heal start response has no client token: {heal_start}"))?; + let task_status_url = format!("{}/rustfs/admin/v3/heal/?clientToken={client_token}", cluster.nodes[0].url); + + let partial_timeout_secs = std::env::var("RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(60); + let partial_deadline = Instant::now() + Duration::from_secs(partial_timeout_secs); + loop { + let status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?; + let active_status: serde_json::Value = serde_json::from_str(&status_body) + .map_err(|err| format!("background heal status is not JSON ({err}): {status_body}"))?; + let operations = &active_status["healOperations"]; + let admin_active = operations["activeBySource"]["admin"].as_u64().is_some_and(|count| count > 0) + || operations["retryingBySource"]["admin"] + .as_u64() + .is_some_and(|count| count > 0); + let active = active_status["state"].as_str() == Some("active") + && (operations["activeTasks"].as_u64().is_some_and(|count| count > 0) + || operations["retryingTasks"].as_u64().is_some_and(|count| count > 0)) + && admin_active; + if active { + break; + } + if Instant::now() >= partial_deadline { + return Err(format!("root heal never became active within {partial_timeout_secs}s: {active_status}").into()); + } + sleep(Duration::from_millis(50)).await; + } + let partial_count = loop { + let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?; + if matching > 0 && matching < expected_manifests.len() { + break matching; + } + if matching == expected_manifests.len() { + return Err(format!( + "root heal rebuilt all {} baseline objects before the target could be interrupted", + expected_manifests.len() + ) + .into()); + } + if Instant::now() >= partial_deadline { + return Err(format!( + "root heal made no observable partial progress on the replacement target within {partial_timeout_secs}s" + ) + .into()); + } + sleep(Duration::from_millis(10)).await; + }; + + cluster.stop_node(1)?; + let stopped_count = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?; + assert!( + stopped_count > 0 && stopped_count < expected_manifests.len(), + "the target must stop after a partial rebuild, observed before stop={partial_count}, after stop={stopped_count}, total={}", + expected_manifests.len() + ); + cluster.start_node(1).await?; - let expected_objects = [(online_key, online_body.as_slice()), (outage_key, outage_body.as_slice())]; - let mut remaining_rebuild_keys: HashSet<&str> = expected_objects.iter().map(|(key, _)| *key).collect(); let heal_timeout_secs = std::env::var("RUSTFS_HEAL_REPLACED_DISK_TIMEOUT_SECS") .ok() .and_then(|value| value.parse::().ok()) - .unwrap_or(90); - - for _ in 0..heal_timeout_secs { - for (key, body) in &expected_objects { - let response = clients[0].get_object().bucket(bucket).key(*key).send().await?; - let actual = response.body.collect().await?.into_bytes(); - assert_eq!(actual.as_ref(), *body, "object body changed for {key}"); - } - - if !remaining_rebuild_keys.is_empty() { - let rebuilt = remaining_rebuild_keys - .iter() - .copied() - .filter(|key| object_metadata_exists_on_disk(&replaced_disk, bucket, key)) - .collect::>(); - for key in rebuilt { - let _ = remaining_rebuild_keys.remove(key); + .unwrap_or(180); + let heal_deadline = Instant::now() + Duration::from_secs(heal_timeout_secs); + loop { + if metadata_count(&replaced_disk, bucket, &expected_manifests) == expected_manifests.len() + && object_metadata_exists_on_disk(&replaced_disk, bucket, outage_key) + { + let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?; + let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?; + if matching == expected_manifests.len() && outage_census.is_complete() { + break; } } - - if remaining_rebuild_keys.is_empty() { - return Ok(()); + if Instant::now() >= heal_deadline { + let matching = matching_manifest_count(&replaced_disk, bucket, &expected_manifests)?; + let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?; + let final_status = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key) + .await + .unwrap_or_else(|err| format!("status request failed: {err}")); + let task_status = match signed_admin_post(&task_status_url, None, &cluster.access_key, &cluster.secret_key).await + { + Ok(body) => heal_task_status_diagnostic(&body), + Err(err) => format!("task status request failed: {err}"), + }; + return Err(format!( + "root heal did not resume after target restart within {heal_timeout_secs}s: baseline={matching}/{}, outage={outage_census:?}, status={final_status}, task_status={task_status}", + expected_manifests.len() + ) + .into()); } - - sleep(Duration::from_secs(1)).await; + sleep(Duration::from_millis(250)).await; } - Err(format!( - "admin deep heal did not rebuild replaced remote disk metadata for {remaining_rebuild_keys:?} within timeout" - ) - .into()) + for (key, expected) in &expected_manifests { + let actual = census_object_version_on_disk(&replaced_disk, bucket, key, None)?; + assert!( + actual.matches_manifest(expected), + "rebuilt target shard differs from its baseline for {key}: {actual:?}" + ); + } + let outage_census = census_object_version_on_disk(&replaced_disk, bucket, outage_key, None)?; + assert!( + outage_census.is_complete(), + "outage object must have a complete target shard: {outage_census:?}" + ); + assert_eq!( + outage_census.erasure_index, + Some(expected_outage_target_erasure_index), + "the outage object must be rebuilt into its own missing erasure slot" + ); + + let target_client = cluster.create_s3_client(1)?; + for (key, _) in &expected_manifests { + let response = target_client.get_object().bucket(bucket).key(key).send().await?; + let actual = response.body.collect().await?.into_bytes(); + assert_eq!(actual.as_ref(), expected_body.as_slice(), "object body changed for {key}"); + } + let response = target_client.get_object().bucket(bucket).key(outage_key).send().await?; + let actual = response.body.collect().await?.into_bytes(); + assert_eq!(actual.as_ref(), expected_body.as_slice(), "object body changed for {outage_key}"); + + let terminal_deadline = Instant::now() + Duration::from_secs(30); + loop { + let status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?; + let status: serde_json::Value = serde_json::from_str(&status_body) + .map_err(|err| format!("background heal status is not JSON ({err}): {status_body}"))?; + let operations = &status["healOperations"]; + let terminal = status["clusterStatusComplete"] == serde_json::Value::Bool(true) + && status["state"].as_str() == Some("idle") + && operations["queueLength"].as_u64() == Some(0) + && operations["activeTasks"].as_u64() == Some(0) + && operations["retryingTasks"].as_u64() == Some(0); + if terminal { + break; + } + if Instant::now() >= terminal_deadline { + return Err(format!("heal data rebuilt but operations did not converge to terminal idle: {status}").into()); + } + sleep(Duration::from_millis(250)).await; + } + + let task_status_body = signed_admin_post(&task_status_url, None, &cluster.access_key, &cluster.secret_key).await?; + let task_status: serde_json::Value = serde_json::from_str(&task_status_body) + .map_err(|err| format!("heal task status is not JSON ({err}): {task_status_body}"))?; + if task_status["summary"].as_str() != Some("finished") { + return Err(format!("heal data rebuilt but task did not finish successfully: {task_status}").into()); + } + + Ok(()) } /// Issue #5850: `background-heal/status` must answer while a peer is down. diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index c59fc16e2..ab4f21f0b 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -36,6 +36,20 @@ const EVENT_HEAL_OBJECT_RENAME: &str = "heal_object_rename"; const HEAL_RENAME_INCOMPLETE: &str = "heal rename incomplete"; const READ_REPAIR_DATA_PHASE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60 * 60); +fn heal_drive_state_for_error(error: &DiskError) -> DriveState { + match error { + DiskError::DiskNotFound | DiskError::RemoteClientUnavailable(_) => DriveState::Offline, + DiskError::FaultyDisk | DiskError::FaultyRemoteDisk => DriveState::Faulty, + DiskError::FileNotFound + | DiskError::FileVersionNotFound + | DiskError::VolumeNotFound + | DiskError::PartMissingOrCorrupt + | DiskError::OutdatedXLMeta => DriveState::Missing, + DiskError::FileCorrupt => DriveState::Corrupt, + _ => DriveState::Unknown(error.to_string()), + } +} + #[cfg(test)] static HEAL_RENAME_FAILURES: std::sync::Mutex> = std::sync::Mutex::new(Vec::new()); @@ -892,16 +906,7 @@ impl SetDisks { } let drive_state = match reason { - Some(err) => match err { - DiskError::DiskNotFound => DriveState::Offline.to_string(), - DiskError::FileNotFound - | DiskError::FileVersionNotFound - | DiskError::VolumeNotFound - | DiskError::PartMissingOrCorrupt - | DiskError::OutdatedXLMeta => DriveState::Missing.to_string(), - DiskError::FileCorrupt => DriveState::Corrupt.to_string(), - _ => DriveState::Unknown(err.to_string()).to_string(), - }, + Some(err) => heal_drive_state_for_error(&err).to_string(), None => DriveState::Ok.to_string(), }; result.before.drives.push(HealDriveInfo { @@ -2673,6 +2678,17 @@ mod heal_result_report_tests { assert!(!super::metadata_less_part_file("xl.meta")); } + #[test] + fn unavailable_heal_errors_use_stable_drive_states() { + for error in [DiskError::FaultyDisk, DiskError::FaultyRemoteDisk] { + assert_eq!(super::heal_drive_state_for_error(&error).to_string(), DriveState::Faulty.to_string()); + } + assert_eq!( + super::heal_drive_state_for_error(&DiskError::RemoteClientUnavailable("peer restarting".to_string())).to_string(), + DriveState::Offline.to_string() + ); + } + #[test] fn read_repair_commit_fingerprint_tracks_commit_identity_only() { let data_dir = Uuid::parse_str("11111111-1111-1111-1111-111111111111").expect("data dir should parse"); diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index b39d92162..c95770834 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -24,7 +24,7 @@ use crate::heal::{ use crate::{Error, Result}; use metrics::{counter, histogram}; use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit}; -use rustfs_heal_contracts::heal_channel::{HealOpts, HealRequestSource, HealScanMode}; +use rustfs_heal_contracts::heal_channel::{DriveState, HealOpts, HealRequestSource, HealScanMode}; use rustfs_madmin::heal_commands::HealResultItem; use rustfs_utils::path::SLASH_SEPARATOR; use serde::{Deserialize, Serialize}; diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index 37459cac2..d44df549e 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -16,6 +16,22 @@ use super::*; use crate::heal::progress::{add_bytes, increment_counter, stable_generation}; use crate::heal::utils::format_set_disk_id; +fn unavailable_recreate_error(result: &HealResultItem, opts: &HealOpts) -> Option { + if opts.dry_run || !opts.recreate { + return None; + } + + let mut offline = false; + for drive in &result.after.drives { + if drive.state == DriveState::Faulty.to_str() { + return Some(Error::Disk(DiskError::FaultyDisk)); + } + offline |= drive.state == DriveState::Offline.to_str(); + } + + offline.then_some(Error::Disk(DiskError::DiskNotFound)) +} + impl HealTask { pub(super) async fn heal_bucket(&self, bucket: &str) -> Result<()> { debug!( @@ -335,13 +351,16 @@ impl HealTask { ) .await { - Ok((result, None)) => { - telemetry_unknown |= !increment_counter(&mut healed); - telemetry_unknown |= - !add_bytes(&mut bytes, u64::try_from(result.object_size).unwrap_or(u64::MAX)); - self.record_result_item(result).await; - None - } + Ok((result, None)) => match unavailable_recreate_error(&result, &heal_opts) { + Some(error) => Some(error), + None => { + telemetry_unknown |= !increment_counter(&mut healed); + telemetry_unknown |= + !add_bytes(&mut bytes, u64::try_from(result.object_size).unwrap_or(u64::MAX)); + self.record_result_item(result).await; + None + } + }, Ok((_, Some(err))) if is_missing_object_dir_heal_result(object, &err) => { telemetry_unknown |= !increment_counter(&mut healed); debug!( diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 430ad8fa9..186043906 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -706,11 +706,28 @@ enum MockHealObjectOutcome { OkWithOtherError(&'static str), ErrOther(&'static str), DanglingGraceDeferred, + UnavailableDrive(DriveState), RetryableReadQuorum, RetryableSlowDown, PermanentOther(&'static str), } +fn unavailable_drive_heal_result(state: DriveState) -> (HealResultItem, Option) { + ( + HealResultItem { + after: Infos { + drives: vec![HealDriveInfo { + endpoint: "remote-target".to_string(), + state: state.to_string(), + ..Default::default() + }], + }, + ..Default::default() + }, + None, + ) +} + #[derive(Clone, Copy)] enum MockObjectExists { Exists(bool), @@ -813,6 +830,7 @@ impl HealStorageAPI for MockStorage { "dangling object deletion deferred by heal grace window; retry_after_secs=3599; grace_secs=3600", ))), )), + MockHealObjectOutcome::UnavailableDrive(state) => Ok(unavailable_drive_heal_result(state)), MockHealObjectOutcome::RetryableReadQuorum => Err(Error::Storage(EcstoreError::InsufficientReadQuorum( bucket.to_string(), object.to_string(), @@ -833,6 +851,7 @@ impl HealStorageAPI for MockStorage { "dangling object deletion deferred by heal grace window; retry_after_secs=3599; grace_secs=3600", ))), )), + MockHealObjectOutcome::UnavailableDrive(state) => Ok(unavailable_drive_heal_result(state)), MockHealObjectOutcome::OkWithOtherError(message) => Ok((HealResultItem::default(), Some(Error::other(message)))), MockHealObjectOutcome::ErrOther(message) | MockHealObjectOutcome::PermanentOther(message) => { Err(Error::other(message)) @@ -1449,6 +1468,46 @@ async fn test_recursive_bucket_heal_retries_only_retryable_objects() { assert_eq!(progress.objects_failed, 0); } +#[tokio::test(start_paused = true)] +async fn recursive_bucket_heal_retries_when_recreate_target_is_unavailable() { + for state in [DriveState::Offline, DriveState::Faulty] { + let state_name = state.to_string(); + let storage = Arc::new(MockStorage::default()); + storage + .heal_object_outcomes + .lock() + .unwrap() + .insert("object-a".to_string(), VecDeque::from([MockHealObjectOutcome::UnavailableDrive(state)])); + let request = HealRequest::new( + HealType::Bucket { + bucket: "bucket-a".to_string(), + }, + HealOptions { + recursive: true, + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + let task = HealTask::from_request(request, storage.clone()); + + task.heal_bucket("bucket-a") + .await + .expect("an unavailable recreate target should be retried after it returns"); + + assert_eq!( + storage.heal_object_calls.lock().unwrap().as_slice(), + ["object-a".to_string(), "object-b".to_string(), "object-a".to_string()], + "unexpected calls for unavailable state {state_name}" + ); + let progress = task.get_progress().await; + assert_eq!(progress.objects_scanned, 2, "unexpected scanned count for state {state_name}"); + assert_eq!(progress.objects_healed, 2, "unexpected healed count for state {state_name}"); + assert_eq!(progress.objects_failed, 0, "unexpected failed count for state {state_name}"); + } +} + #[tokio::test(start_paused = true)] async fn recursive_bucket_heal_skips_dangling_delete_grace_without_batch_failure() { let storage = Arc::new(MockStorage::default());