fix(heal): recover replacement after transient disk errors (#7059)

* fix(heal): recover replacement after transient disk errors

* fix(heal): satisfy replacement status clippy lint
This commit is contained in:
Henry Guo
2026-09-03 07:03:01 +08:00
committed by GitHub
parent 07dce1cab0
commit 98f7e63396
5 changed files with 618 additions and 29 deletions
@@ -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<String>),
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ReplacementScenario {
Baseline,
MidRebuildIoFault,
}
struct MountNamespaceGuard {
mounts: Vec<PathBuf>,
}
@@ -102,6 +111,19 @@ mod tests {
impl FaultableBlockMount {
fn mount(target: &Path, image_root: &Path, label: &str) -> Result<Self, Box<dyn Error + Send + Sync>> {
Self::mount_with_live_recovery(target, image_root, label, false)
}
fn mount_live_recovery(target: &Path, image_root: &Path, label: &str) -> Result<Self, Box<dyn Error + Send + Sync>> {
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<Self, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
fn restore_linear_table(&self) -> Result<(), Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
self.restore_linear_table()
}
fn cleanup(&mut self) -> Result<(), Box<dyn Error + Send + Sync>> {
let mut first_error: Option<Box<dyn Error + Send + Sync>> = 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<Vec<BaselineVersion>, Box<dyn Error + Send + Sync>> {
async fn seed_baseline(
client: &Client,
target_disk: &Path,
extra_object_count: usize,
) -> Result<Vec<BaselineVersion>, Box<dyn Error + Send + Sync>> {
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::<Result<Vec<_>, Box<dyn Error + Send + Sync>>>()?;
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<Vec<(String, String)>, Box<dyn Error + Send + Sync>> {
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<Option<String>, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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::<Vec<_>>();
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<dyn Error + Send + Sync>> {
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::<rustfs_filemeta::Error>(),
@@ -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<dyn Error + Send + Sync>> {
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<Vec<u32>, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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<String, Box<dyn Error + Send + Sync>> {
let partial_timeout_secs = std::env::var("RUSTFS_HEAL_DISK_IO_PARTIAL_TIMEOUT_SECS")
.ok()
.and_then(|value| value.parse::<u64>().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<dyn Error + Send + Sync>>(())
}
.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<dyn Error + Send + Sync>> {
async fn run_replacement_e2e(
parity: usize,
test_name: &str,
scenario: ReplacementScenario,
) -> Result<(), Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
run_replacement_e2e(
4,
"replacement_privileged_e2e_test::tests::test_privileged_3x4_auto_replacement_recovers_from_mid_rebuild_eio",
ReplacementScenario::MidRebuildIoFault,
)
.await
}
+91 -3
View File
@@ -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() {
+15 -6
View File
@@ -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<ResumeManager> {
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<String>, 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?;
}
+44
View File
@@ -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");
+24 -7
View File
@@ -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");