diff --git a/crates/e2e_test/src/replacement_privileged_e2e_test.rs b/crates/e2e_test/src/replacement_privileged_e2e_test.rs index d533e20d0..60cfccf68 100644 --- a/crates/e2e_test/src/replacement_privileged_e2e_test.rs +++ b/crates/e2e_test/src/replacement_privileged_e2e_test.rs @@ -35,7 +35,8 @@ mod tests { use std::fs; use std::path::{Path, PathBuf}; use std::process::Command; - use tokio::time::{Duration, Instant, interval}; + use tokio::net::TcpStream; + use tokio::time::{Duration, Instant, interval, sleep, timeout}; use tracing::info; const ENABLE_ENV: &str = "RUSTFS_PRIVILEGED_REPLACEMENT_E2E"; @@ -48,6 +49,8 @@ mod tests { const REPLACEMENT_INTENT_SUFFIX: &str = "_ahm_replacement_intent.json"; const REPLACEMENT_COMPLETION_PROOF_SUFFIX: &str = "_ahm_replacement_completion_proof.json"; const RESUME_CHECKPOINT_SUFFIX: &str = "_ahm_checkpoint.json"; + const FAULT_WINDOW_OBJECT_COUNT: usize = 24; + const FAULT_WINDOW_OBJECT_BYTES: usize = 32 * 1024 * 1024; #[derive(Debug)] struct BaselineVersion { @@ -65,6 +68,12 @@ mod tests { CompletedWithIncomplete(BTreeSet), } + #[derive(Clone, Copy, Debug, Eq, PartialEq)] + enum ReplacementScenario { + Baseline, + MidRebuildIoFault, + } + struct MountNamespaceGuard { mounts: Vec, } @@ -102,6 +111,19 @@ mod tests { impl FaultableBlockMount { fn mount(target: &Path, image_root: &Path, label: &str) -> Result> { + Self::mount_with_live_recovery(target, image_root, label, false) + } + + fn mount_live_recovery(target: &Path, image_root: &Path, label: &str) -> Result> { + Self::mount_with_live_recovery(target, image_root, label, true) + } + + fn mount_with_live_recovery( + target: &Path, + image_root: &Path, + label: &str, + live_recovery: bool, + ) -> Result> { fs::create_dir_all(image_root)?; let image = image_root.join(format!("{label}.img")); let file = fs::File::create(&image)?; @@ -114,7 +136,15 @@ mod tests { return Err("losetup --find --show returned an empty loop device".into()); } - run_command("mkfs.ext4", &["-F", &loop_device])?; + if live_recovery { + // Keep the filesystem and RustFS' persistent root descriptor attached + // across the transient all-block EIO. A journaling ext4 abort requires + // an unmount to recover, which would test process/disk reattachment + // instead of live I/O recovery. + run_command("mkfs.ext4", &["-F", "-O", "^has_journal", &loop_device])?; + } else { + run_command("mkfs.ext4", &["-F", &loop_device])?; + } let sectors = run_command_stdout("blockdev", &["--getsz", &loop_device])?; let dm_name = format!("rustfs_e2e_{label}_{}", std::process::id()); let table = format!("0 {sectors} linear {loop_device} 0"); @@ -122,7 +152,11 @@ mod tests { run_command("dmsetup", &["create", &dm_name, "--table", &table])?; let target_arg = path_to_string(target, "faultable mount target")?; - run_command("mount", &[&mapper, &target_arg])?; + if live_recovery { + run_command("mount", &["-o", "errors=continue", &mapper, &target_arg])?; + } else { + run_command("mount", &[&mapper, &target_arg])?; + } Ok(Self { target: target.to_path_buf(), @@ -159,18 +193,25 @@ mod tests { Ok(()) } - fn restore_available(&self) -> Result<(), Box> { + fn restore_linear_table(&self) -> Result<(), Box> { let sectors = run_command_stdout("blockdev", &["--getsz", &self.loop_device])?; let linear_table = format!("0 {sectors} linear {} 0", self.loop_device); - run_command("dmsetup", &["suspend", &self.dm_name])?; + // An ext4 journal abort can leave the mounted filesystem internally + // read-only. Avoid dmsetup's filesystem freeze/flush in that state; + // all I/O sent to the error target has already completed with EIO. + run_command("dmsetup", &["suspend", "--noflush", &self.dm_name])?; run_command("dmsetup", &["load", &self.dm_name, "--table", &linear_table])?; run_command("dmsetup", &["resume", &self.dm_name]) } + fn restore_available(&self) -> Result<(), Box> { + self.restore_linear_table() + } + fn cleanup(&mut self) -> Result<(), Box> { let mut first_error: Option> = None; if self.dm_created { - let _ = self.restore_available(); + let _ = self.restore_linear_table(); } if self.mounted { if let Err(error) = detach_mount(&self.target) { @@ -469,7 +510,11 @@ mod tests { Ok((completed.version_id().map(str::to_owned), digest)) } - async fn seed_baseline(client: &Client, target_disk: &Path) -> Result, Box> { + async fn seed_baseline( + client: &Client, + target_disk: &Path, + extra_object_count: usize, + ) -> Result, Box> { let plain_bucket = "priv-replacement-plain"; let versioned_bucket = "priv-replacement-versions"; let null_bucket = "priv-replacement-null"; @@ -532,7 +577,7 @@ mod tests { put_object_version(client, null_bucket, "null/current.bin", payload(512 * 1024, 8)).await?; versions.push((null_bucket, "null/current.bin", version_id, Some(body_sha256))); - let versions = versions + let mut versions = versions .into_iter() .map(|(bucket, key, version_id, body_sha256)| { let expected = census_object_version_on_disk(target_disk, bucket, key, version_id.as_deref())?; @@ -548,6 +593,23 @@ mod tests { }) }) .collect::, Box>>()?; + for index in 0..extra_object_count { + let key = format!("fault-window/object-{index:04}.bin"); + let seed = u8::try_from(index + 32)?; + let (version_id, body_sha256) = + put_object_version(client, plain_bucket, &key, payload(FAULT_WINDOW_OBJECT_BYTES, seed)).await?; + let expected = census_object_version_on_disk(target_disk, plain_bucket, &key, version_id.as_deref())?; + if !expected.is_complete() { + return Err(format!("fault-window baseline census is incomplete for {plain_bucket}/{key}: {expected:?}").into()); + } + versions.push(BaselineVersion { + bucket: plain_bucket.to_string(), + key, + version_id, + body_sha256: Some(body_sha256), + expected, + }); + } let inline = versions .iter() .find(|version| version.key == "history/inline.bin") @@ -739,6 +801,98 @@ mod tests { .collect() } + fn target_record_details( + status: &serde_json::Value, + target_disk: &Path, + ) -> Result, Box> { + let target = target_disk.to_string_lossy(); + let records = status["cluster"]["records"] + .as_array() + .ok_or_else(|| format!("replacement recovery status omitted cluster.records: {status}"))?; + records + .iter() + .filter(|record| { + record["targetSlots"] + .as_array() + .into_iter() + .flatten() + .filter_map(serde_json::Value::as_str) + .any(|slot| slot.contains(target.as_ref())) + }) + .map(|record| { + let task_id = record["taskId"] + .as_str() + .filter(|task_id| !task_id.is_empty()) + .ok_or_else(|| format!("replacement recovery record omitted taskId: {record}"))?; + let state = record["state"] + .as_str() + .filter(|state| !state.is_empty()) + .ok_or_else(|| format!("replacement recovery record omitted state: {record}"))?; + Ok((task_id.to_string(), state.to_string())) + }) + .collect() + } + + fn running_target_generation( + status: &serde_json::Value, + target_disk: &Path, + ) -> Result, Box> { + if !cluster_status_is_definitive(status)? { + return Ok(None); + } + let records = target_record_details(status, target_disk)?; + if records.len() == 1 && records[0].1 == "running" { + return Ok(Some(records[0].0.clone())); + } + Ok(None) + } + + fn assert_target_generation_nonterminal( + status: &serde_json::Value, + target_disk: &Path, + expected_task_id: &str, + ) -> Result<(), Box> { + if !cluster_status_is_definitive(status)? { + return Err(format!("replacement recovery became non-definitive during target EIO: {status}").into()); + } + let records = target_record_details(status, target_disk)?; + let matching = records + .iter() + .filter(|(task_id, _)| task_id == expected_task_id) + .collect::>(); + if matching.len() != 1 { + return Err(format!( + "replacement generation {expected_task_id} must remain uniquely observable during target EIO: {records:?}" + ) + .into()); + } + match matching[0].1.as_str() { + "waiting_for_replacement" | "running" | "incomplete" => Ok(()), + state => Err(format!( + "replacement generation {expected_task_id} reached invalid state {state:?} during target EIO: {status}" + ) + .into()), + } + } + + fn assert_target_generation_completed( + status: &serde_json::Value, + target_disk: &Path, + expected_task_id: &str, + ) -> Result<(), Box> { + if !cluster_status_is_definitive(status)? { + return Err(format!("completed replacement recovery status is non-definitive: {status}").into()); + } + let records = target_record_details(status, target_disk)?; + if records == [(expected_task_id.to_string(), "completed".to_string())] { + return Ok(()); + } + Err( + format!("replacement generation {expected_task_id} did not retain its identity through EIO recovery: {records:?}") + .into(), + ) + } + fn is_transient_recovery_version_absence(error: &(dyn Error + 'static)) -> bool { matches!( error.downcast_ref::(), @@ -775,6 +929,165 @@ mod tests { Ok(missing) } + async fn wait_for_partial_replacement<'a>( + cluster: &RustFSTestClusterEnvironment, + target_disk: &Path, + versions: &'a [BaselineVersion], + timeout_secs: u64, + ) -> Result<(usize, &'a BaselineVersion, String), Box> { + let deadline = Instant::now() + Duration::from_secs(timeout_secs); + loop { + let missing = incomplete_versions(target_disk, versions)?; + let completed = versions.len().saturating_sub(missing.len()); + if completed > 0 && completed < versions.len() { + let witness = versions.iter().find(|version| { + census_object_version_on_disk(target_disk, &version.bucket, &version.key, version.version_id.as_deref()) + .is_ok_and(|actual| actual.matches_manifest(&version.expected)) + }); + if let Some(witness) = witness { + let status = replacement_status(cluster).await?; + if let Some(task_id) = running_target_generation(&status, target_disk)? { + return Ok((completed, witness, task_id)); + } + } + } + if completed == versions.len() { + return Err(format!( + "replacement rebuilt all {} baseline versions before target EIO could be injected", + versions.len() + ) + .into()); + } + if Instant::now() >= deadline { + let status = replacement_status(cluster).await?; + return Err(format!( + "replacement made no observable running partial progress within {timeout_secs}s: completed={completed}/{} status={status}", + versions.len() + ) + .into()); + } + sleep(Duration::from_millis(10)).await; + } + } + + fn cluster_process_ids(cluster: &RustFSTestClusterEnvironment) -> Result, Box> { + cluster + .nodes + .iter() + .enumerate() + .map(|(index, node)| { + node.process + .as_ref() + .map(std::process::Child::id) + .ok_or_else(|| format!("cluster node {index} process is not running").into()) + }) + .collect() + } + + async fn assert_cluster_processes_and_listeners_unchanged( + cluster: &mut RustFSTestClusterEnvironment, + expected_pids: &[u32], + ) -> Result<(), Box> { + if cluster.nodes.len() != expected_pids.len() { + return Err("cluster node count changed during target EIO".into()); + } + for (index, (node, expected_pid)) in cluster.nodes.iter_mut().zip(expected_pids).enumerate() { + let process = node + .process + .as_mut() + .ok_or_else(|| format!("cluster node {index} process disappeared during target EIO"))?; + if process.id() != *expected_pid { + return Err(format!( + "cluster node {index} PID changed during target EIO: expected {expected_pid}, got {}", + process.id() + ) + .into()); + } + if let Some(status) = process.try_wait()? { + return Err(format!("cluster node {index} exited during target EIO with {status}").into()); + } + match timeout(Duration::from_secs(2), TcpStream::connect(&node.address)).await { + Ok(Ok(stream)) => drop(stream), + Ok(Err(error)) => { + return Err(format!("cluster node {index} TCP listener failed during target EIO: {error}").into()); + } + Err(_) => return Err(format!("cluster node {index} TCP listener timed out during target EIO").into()), + } + } + Ok(()) + } + + async fn exercise_mid_rebuild_io_fault( + cluster: &mut RustFSTestClusterEnvironment, + replacement_mount: &FaultableBlockMount, + target_disk: &Path, + versions: &[BaselineVersion], + ) -> Result> { + let partial_timeout_secs = std::env::var("RUSTFS_HEAL_DISK_IO_PARTIAL_TIMEOUT_SECS") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(120); + let (partial_count, witness, task_id) = + wait_for_partial_replacement(cluster, target_disk, versions, partial_timeout_secs).await?; + let expected_pids = cluster_process_ids(cluster)?; + + replacement_mount + .make_unavailable() + .map_err(|error| format!("failed to install dm-error on the active replacement: {error}"))?; + let fault_result = async { + replacement_mount + .verify_raw_io_is_unavailable() + .map_err(|error| format!("active replacement dm-error was not proven by direct I/O: {error}"))?; + assert_cluster_processes_and_listeners_unchanged(cluster, &expected_pids).await?; + + let observation_deadline = Instant::now() + Duration::from_secs(2); + loop { + let status = timeout(Duration::from_secs(5), replacement_status(cluster)) + .await + .map_err(|_| "replacement recovery status timed out during target EIO")??; + assert_target_generation_nonterminal(&status, target_disk, &task_id)?; + assert_cluster_processes_and_listeners_unchanged(cluster, &expected_pids).await?; + if Instant::now() >= observation_deadline { + break; + } + sleep(Duration::from_millis(100)).await; + } + Ok::<(), Box>(()) + } + .await; + let restore_result = replacement_mount + .restore_available() + .map_err(|error| format!("failed to restore the active replacement after dm-error: {error}")); + if let Err(error) = fault_result { + if let Err(restore_error) = restore_result { + info!(%restore_error, "replacement restore also failed while preserving target EIO failure"); + } + return Err(error); + } + restore_result?; + + assert_cluster_processes_and_listeners_unchanged(cluster, &expected_pids).await?; + let actual = census_object_version_on_disk(target_disk, &witness.bucket, &witness.key, witness.version_id.as_deref())?; + if !actual.matches_manifest(&witness.expected) { + return Err(format!( + "witnessed replacement shard did not survive target EIO for {}/{}@{:?}: {actual:?}", + witness.bucket, witness.key, witness.version_id + ) + .into()); + } + let completed_after_restore = versions + .len() + .saturating_sub(incomplete_versions(target_disk, versions)?.len()); + if completed_after_restore < partial_count { + return Err(format!( + "replacement progress regressed across target EIO: before={partial_count}, after={completed_after_restore}" + ) + .into()); + } + + Ok(task_id) + } + fn replacement_completion_state( status: &serde_json::Value, target_disk: &Path, @@ -863,7 +1176,11 @@ mod tests { } } - async fn run_replacement_e2e(parity: usize, test_name: &str) -> Result<(), Box> { + async fn run_replacement_e2e( + parity: usize, + test_name: &str, + scenario: ReplacementScenario, + ) -> Result<(), Box> { init_logging(); if !privileged_run_enabled()? { return Ok(()); @@ -900,7 +1217,10 @@ mod tests { } } let mut target_mount = target_mount.ok_or("target drive was not mounted with the faultable block fixture")?; - let mut replacement_mount = ZramBlockMount::reserve(&target_disk)?; + let mut zram_replacement = match scenario { + ReplacementScenario::Baseline => Some(ZramBlockMount::reserve(&target_disk)?), + ReplacementScenario::MidRebuildIoFault => None, + }; cluster.set_env("RUSTFS_HEAL_ENABLED", "true"); cluster.set_env("RUSTFS_SCANNER_ENABLED", "true"); @@ -908,13 +1228,21 @@ mod tests { cluster.set_env("RUSTFS_SCANNER_CYCLE", "1"); cluster.set_env("RUSTFS_SCANNER_START_DELAY_SECS", "0"); cluster.set_env("RUSTFS_STORAGE_CLASS_STANDARD", format!("EC:{parity}")); + if scenario == ReplacementScenario::MidRebuildIoFault { + cluster.set_env("RUSTFS_HEAL_PAGE_OBJECT_CONCURRENCY", "1"); + cluster.set_env("RUSTFS_HEAL_PAGE_PARALLEL_ENABLE", "false"); + } for node_index in 0..cluster.nodes.len() { cluster.set_node_env(node_index, "RUST_LOG", "rustfs=info,rustfs::heal::manager=debug,rustfs_notify=debug")?; } cluster.start().await?; let clients = cluster.create_all_clients()?; - let versions = seed_baseline(&clients[0], &target_disk) + let extra_object_count = match scenario { + ReplacementScenario::Baseline => 0, + ReplacementScenario::MidRebuildIoFault => FAULT_WINDOW_OBJECT_COUNT, + }; + let versions = seed_baseline(&clients[0], &target_disk, extra_object_count) .await .map_err(|error| format!("pre-fault baseline seeding failed: {error}"))?; verify_bodies(&clients[0], &versions) @@ -935,7 +1263,20 @@ mod tests { cluster.stop_node_gracefully(TARGET_NODE).await?; target_mount.cleanup()?; - replacement_mount.mount_target()?; + let mut faultable_replacement = match scenario { + ReplacementScenario::Baseline => { + zram_replacement + .as_mut() + .ok_or("baseline replacement zram was not reserved")? + .mount_target()?; + None + } + ReplacementScenario::MidRebuildIoFault => Some(FaultableBlockMount::mount_live_recovery( + &target_disk, + &image_root, + &format!("p{parity}_replacement_node{TARGET_NODE}_drive{TARGET_DRIVE}"), + )?), + }; let missing_before_restart = incomplete_versions(&target_disk, &versions)?; assert_eq!( missing_before_restart.len(), @@ -945,12 +1286,28 @@ mod tests { cluster.start_node(TARGET_NODE).await?; let recovery_result = async { + let faulted_task_id = match faultable_replacement.as_ref() { + Some(replacement) => { + Some(exercise_mid_rebuild_io_fault(&mut cluster, replacement, &target_disk, &versions).await?) + } + None => None, + }; wait_for_completed_replacement_with_census(&cluster, &target_disk, &versions, 420).await?; + if let Some(task_id) = faulted_task_id { + let status = replacement_status(&cluster).await?; + assert_target_generation_completed(&status, &target_disk, &task_id)?; + } verify_bodies(&clients[0], &versions).await } .await; let stop_result = cluster.stop_node_gracefully(TARGET_NODE).await; - let replacement_cleanup_result = replacement_mount.cleanup(); + let replacement_cleanup_result = match faultable_replacement.as_mut() { + Some(replacement) => replacement.cleanup(), + None => zram_replacement + .as_mut() + .ok_or("baseline replacement zram disappeared before cleanup")? + .cleanup(), + }; if let Err(error) = recovery_result { if let Err(stop_error) = stop_result { @@ -1097,6 +1454,65 @@ mod tests { assert_eq!(status_samples.borrow().len(), 1); } + #[test] + fn target_eio_status_preserves_one_nonterminal_generation() { + let target = Path::new("/mnt/target"); + for state in ["waiting_for_replacement", "running", "incomplete"] { + let status = serde_json::json!({ + "cluster": { + "definitive": true, + "records": [{ + "taskId": "generation-a", + "state": state, + "targetSlots": ["http://127.0.0.1:9000/mnt/target"] + }] + } + }); + assert!(assert_target_generation_nonterminal(&status, target, "generation-a").is_ok()); + } + + let running = serde_json::json!({ + "cluster": { + "definitive": true, + "records": [{ + "taskId": "generation-a", + "state": "running", + "targetSlots": ["/mnt/target"] + }] + } + }); + assert_eq!(running_target_generation(&running, target).unwrap().as_deref(), Some("generation-a")); + } + + #[test] + fn target_eio_status_rejects_false_or_replaced_completion() { + let target = Path::new("/mnt/target"); + let completed = serde_json::json!({ + "cluster": { + "definitive": true, + "records": [{ + "taskId": "generation-a", + "state": "completed", + "targetSlots": ["/mnt/target"] + }] + } + }); + assert!(assert_target_generation_nonterminal(&completed, target, "generation-a").is_err()); + assert!(assert_target_generation_completed(&completed, target, "generation-a").is_ok()); + assert!(assert_target_generation_completed(&completed, target, "generation-b").is_err()); + + let duplicate = serde_json::json!({ + "cluster": { + "definitive": true, + "records": [ + {"taskId": "generation-a", "state": "running", "targetSlots": ["/mnt/target"]}, + {"taskId": "generation-a", "state": "incomplete", "targetSlots": ["/mnt/target"]} + ] + } + }); + assert!(assert_target_generation_nonterminal(&duplicate, target, "generation-a").is_err()); + } + #[test] fn absent_status_requires_definitive_empty_records() { let target = Path::new("/mnt/target"); @@ -1117,6 +1533,7 @@ mod tests { run_replacement_e2e( 4, "replacement_privileged_e2e_test::tests::test_privileged_3x4_auto_replacement_rebuilds_ec8_plus_4_without_admin_heal", + ReplacementScenario::Baseline, ) .await } @@ -1130,6 +1547,20 @@ mod tests { run_replacement_e2e( 6, "replacement_privileged_e2e_test::tests::test_privileged_3x4_auto_replacement_rebuilds_ec6_plus_6_without_admin_heal", + ReplacementScenario::Baseline, + ) + .await + } + + /// Linux mount namespaces are per-thread; keep mount setup and process + /// spawning on one OS thread so child RustFS nodes inherit the test mounts. + #[tokio::test(flavor = "current_thread")] + #[ignore = "requires Linux root/CAP_SYS_ADMIN and RUSTFS_PRIVILEGED_REPLACEMENT_E2E=1"] + async fn test_privileged_3x4_auto_replacement_recovers_from_mid_rebuild_eio() -> Result<(), Box> { + run_replacement_e2e( + 4, + "replacement_privileged_e2e_test::tests::test_privileged_3x4_auto_replacement_recovers_from_mid_rebuild_eio", + ReplacementScenario::MidRebuildIoFault, ) .await } diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index ed8de1da8..b56dca504 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -111,9 +111,12 @@ impl ECStore { write_state.ensure_write_safe("heal format fence failed")?; // Metadata fence order is part of the decommission/rebalance protocol: - // pool.bin must always be acquired before rebalance.bin. + // pool.bin must always be acquired before rebalance.bin. A read guard is + // sufficient to freeze the pool snapshot and lets replacement format + // recovery coexist with ordinary Heal's capacity fence. The rebalance + // write guard still serializes format repair and excludes transitions. let pool_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; - let pool_guard = pool_lock.get_write_lock(get_lock_acquire_timeout()).await?; + let pool_guard = pool_lock.get_read_lock(get_lock_acquire_timeout()).await?; let rebalance_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?; let rebalance_guard = rebalance_lock.get_write_lock(get_lock_acquire_timeout()).await?; @@ -1102,7 +1105,8 @@ mod tests { pool_index: usize, disk_index: usize, ) -> String { - let target = store.pools[pool_index].endpoints.endpoints.as_ref()[disk_index].to_string(); + let endpoint = store.pools[pool_index].endpoints.endpoints.as_ref()[disk_index].clone(); + let target = endpoint.to_string(); let format_path = heal_test_format_path(temp_dir, pool_index, disk_index); tokio::fs::remove_file(&format_path) .await @@ -1112,6 +1116,22 @@ mod tests { .await .expect("replacement target format path should be inspectable") ); + let replacement = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("unformatted replacement disk should open"); + let set = &store.pools[pool_index].disk_set[0]; + let target_slot = set + .set_endpoints + .iter() + .position(|candidate| candidate == &endpoint) + .expect("replacement endpoint should belong to its set"); + set.disks.write().await[target_slot] = Some(replacement); target } @@ -1124,6 +1144,74 @@ mod tests { ); } + #[tokio::test] + #[serial_test::serial] + async fn replacement_format_heal_coexists_with_capacity_read_fence() { + let (temp_dir, store, shutdown) = multi_pool_heal_store().await; + let target = remove_heal_test_format(&temp_dir, &store, 0, 3).await; + let capacity_guard = store + .acquire_external_decommission_capacity_fence(&[0], "heal") + .await + .expect("ordinary heal capacity fence should be acquired"); + + let (_, err) = + tokio::time::timeout(std::time::Duration::from_secs(2), store.heal_replacement_format(false, 0, 0, &[target])) + .await + .expect("replacement format heal must not wait on an existing capacity read fence") + .expect("replacement format heal should complete"); + + assert!(err.is_none(), "replacement format heal should succeed: {err:?}"); + assert!( + tokio::fs::try_exists(heal_test_format_path(&temp_dir, 0, 3)) + .await + .expect("replacement target format path should be inspectable"), + "replacement format heal should restore the target while the capacity read fence is held" + ); + drop(capacity_guard); + shutdown.cancel(); + } + + #[tokio::test] + #[serial_test::serial] + async fn replacement_format_heal_remains_blocked_by_pool_metadata_writer() { + let (temp_dir, store, shutdown) = multi_pool_heal_store().await; + let target = remove_heal_test_format(&temp_dir, &store, 0, 3).await; + let pool_lock = store.pools[0] + .new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME) + .await + .expect("pool metadata lock should be created"); + let pool_writer = pool_lock + .get_write_lock(get_lock_acquire_timeout()) + .await + .expect("pool metadata writer should be acquired"); + + assert!( + tokio::time::timeout( + std::time::Duration::from_secs(1), + store.heal_replacement_format(false, 0, 0, std::slice::from_ref(&target)), + ) + .await + .is_err(), + "replacement format heal must wait for a pool metadata writer" + ); + assert_heal_test_format_missing(&temp_dir, 0, 3).await; + + drop(pool_writer); + let (_, err) = + tokio::time::timeout(std::time::Duration::from_secs(30), store.heal_replacement_format(false, 0, 0, &[target])) + .await + .expect("replacement format heal should resume after the writer releases") + .expect("replacement format heal should complete after the writer releases"); + assert!(err.is_none(), "replacement format heal should succeed: {err:?}"); + assert!( + tokio::fs::try_exists(heal_test_format_path(&temp_dir, 0, 3)) + .await + .expect("replacement target format path should be inspectable"), + "replacement format heal should restore the target after the writer releases" + ); + shutdown.cancel(); + } + #[tokio::test] #[serial_test::serial] async fn full_format_heal_preblocked_pool_metadata_never_writes_format() { diff --git a/crates/heal/src/heal/task/heal_erasure_set.rs b/crates/heal/src/heal/task/heal_erasure_set.rs index e9478c721..cc1d638ba 100644 --- a/crates/heal/src/heal/task/heal_erasure_set.rs +++ b/crates/heal/src/heal/task/heal_erasure_set.rs @@ -13,6 +13,13 @@ // limitations under the License. /// erasure-set heal: drives the ErasureSetHealer across the set's buckets use super::*; +use crate::heal::DiskStore; + +pub(super) async fn load_verified_replacement_resume(disk: &DiskStore, task_id: &str) -> Result { + let resume_manager = ResumeManager::load_replacement_intent(disk.clone(), task_id).await?; + resume_manager.ensure_replacement_completion_proof().await?; + Ok(resume_manager) +} impl HealTask { pub(super) async fn heal_erasure_set(&self, buckets: Vec, set_disk_id: String) -> Result<()> { @@ -445,14 +452,16 @@ impl HealTask { self.verify_replacement_identity_fence(expected_identities, &set_disk_id, "marker completion") .await?; } + let verified_replacement_resume = if let Some((disk, _, _)) = replacement_resume.as_ref() { + Some((disk.clone(), load_verified_replacement_resume(disk, &self.id).await?)) + } else { + None + }; super::super::clear_healing_markers_after_verified(&self.heal_endpoints, &healing_marker).await?; - if let Some((disk, resume_manager, _)) = replacement_resume.as_ref() { + if let Some((disk, resume_manager)) = verified_replacement_resume { resume_manager.mark_replacement_cleanup_pending().await?; - if CheckpointManager::has_checkpoint(disk, &self.id).await { - CheckpointManager::load_from_disk(disk.clone(), &self.id) - .await? - .cleanup() - .await?; + if CheckpointManager::has_checkpoint(&disk, &self.id).await { + CheckpointManager::load_from_disk(disk, &self.id).await?.cleanup().await?; } resume_manager.cleanup().await?; } diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index d2846dc5c..1734d0c7a 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -469,6 +469,50 @@ async fn cleanup_pending_recovery_removes_checkpoint_without_rebuild_work() { assert!(!*storage.listed.lock().unwrap()); } +#[tokio::test] +async fn replacement_cleanup_reloads_verified_state_from_survivor_anchor() { + let temp = TempDir::new().expect("temporary resume disk directory should be created"); + let anchor = make_resume_disk(&temp).await; + let task_id = crate::heal::resume::ResumeUtils::generate_task_id(); + let stale_resume = ResumeManager::new_replacement_intent( + anchor.clone(), + task_id.clone(), + "pool_0_set_0".to_string(), + vec!["bucket-a".to_string()], + vec!["replacement-a".to_string()], + vec![replacement_identity("replacement-a", "device-a", "filesystem-a")], + ) + .await + .expect("replacement intent should persist on the survivor anchor"); + let rebuilding_resume = ResumeManager::load_replacement_intent(anchor.clone(), &task_id) + .await + .expect("the inner healer should load the persisted replacement intent"); + rebuilding_resume + .mark_replacement_completed_and_verified() + .await + .expect("the inner healer should persist verified completion"); + + assert!( + stale_resume.mark_replacement_cleanup_pending().await.is_err(), + "the original in-memory manager must remain stale after another manager persists verification" + ); + let verified_resume = super::heal_erasure_set::load_verified_replacement_resume(&anchor, &task_id) + .await + .expect("cleanup should reload and verify durable replacement completion"); + verified_resume + .mark_replacement_cleanup_pending() + .await + .expect("the freshly loaded verified state should enter cleanup"); + + let state = ResumeManager::load_replacement_intent(anchor, &task_id) + .await + .expect("cleanup-pending state should remain durable") + .get_state() + .await; + assert!(state.completed); + assert_eq!(state.replacement_phase, ReplacementPhase::CleanupPending); +} + #[tokio::test] async fn verified_recovery_keeps_state_when_marker_clear_fails() { let temp = TempDir::new().expect("temporary resume disk directory should be created"); diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index cbfdaf2fd..f10ad9123 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -650,14 +650,15 @@ fn aggregate_replacement_recovery_cluster_status( reason.get_or_insert("peer_replacement_recovery_status_not_definitive"); } - let first_records = snapshots - .first() + // Replacement recovery state is anchored on one surviving disk. Nodes + // without that disk legitimately report a definitive empty snapshot, so + // only non-empty snapshots must agree with each other. + let mut non_empty_records = snapshots + .iter() + .filter(|snapshot| !snapshot.records.is_empty()) .map(|snapshot| canonical_replacement_records(&snapshot.records)); - if let Some(first_records) = &first_records - && snapshots - .iter() - .skip(1) - .any(|snapshot| canonical_replacement_records(&snapshot.records) != *first_records) + if let Some(first_records) = non_empty_records.next() + && non_empty_records.any(|records| records != first_records) { reason.get_or_insert("peer_replacement_recovery_status_conflict"); } @@ -1627,6 +1628,22 @@ mod tests { assert_eq!(cluster.records.len(), 2); } + #[test] + fn replacement_recovery_cluster_accepts_definitive_empty_peers() { + let local = rustfs_heal::ReplacementRecoverySnapshot { + records: Vec::new(), + definitive: true, + reason: None, + }; + let peer = replacement_snapshot("11111111-1111-4111-8111-111111111111"); + let cluster = aggregate_replacement_recovery_cluster_status(vec![local], vec![Ok(Some(peer))], 2, true); + + assert!(cluster.definitive); + assert!(cluster.reason.is_none()); + assert_eq!(cluster.records.len(), 1); + assert_eq!(cluster.records[0].state, rustfs_heal::ReplacementRecoveryState::Completed); + } + #[test] fn replacement_recovery_cluster_rejects_conflicting_records_inside_a_snapshot() { let mut local = replacement_snapshot("11111111-1111-4111-8111-111111111111");