diff --git a/crates/e2e_test/src/chaos.rs b/crates/e2e_test/src/chaos.rs index 0b7a8b8f4..920a87f0c 100644 --- a/crates/e2e_test/src/chaos.rs +++ b/crates/e2e_test/src/chaos.rs @@ -62,8 +62,8 @@ pub(crate) struct VersionShardCensus { pub data_dir: Option, pub erasure_index: Option, pub expected_part_numbers: BTreeSet, - pub present_part_numbers: BTreeSet, pub present_part_fingerprints: BTreeMap, + pub inline_data_fingerprint: Option, } #[derive(Clone, Debug, Eq, PartialEq)] @@ -74,7 +74,12 @@ pub(crate) struct PartShardFingerprint { impl VersionShardCensus { pub(crate) fn is_complete(&self) -> bool { - self.has_xl_meta && self.expected_part_numbers == self.present_part_numbers + self.has_xl_meta + && self.expected_part_numbers.len() == self.present_part_fingerprints.len() + && self + .expected_part_numbers + .iter() + .all(|part_number| self.present_part_fingerprints.contains_key(part_number)) } pub(crate) fn matches_manifest(&self, manifest: &Self) -> bool { @@ -85,6 +90,7 @@ impl VersionShardCensus { && self.erasure_index == manifest.erasure_index && self.expected_part_numbers == manifest.expected_part_numbers && self.present_part_fingerprints == manifest.present_part_fingerprints + && self.inline_data_fingerprint == manifest.inline_data_fingerprint } } @@ -93,6 +99,13 @@ fn sha256_hex(data: &[u8]) -> String { digest.iter().map(|byte| format!("{byte:02x}")).collect() } +fn shard_fingerprint(data: &[u8]) -> ChaosResult { + Ok(PartShardFingerprint { + size: u64::try_from(data.len())?, + sha256: sha256_hex(data), + }) +} + /// Single-node RustFS server with `disk_count` local volume directories that /// can be faulted individually while the server is running. pub struct DiskFaultHarness { @@ -301,8 +314,8 @@ pub(crate) fn census_object_version_on_disk( data_dir: None, erasure_index: None, expected_part_numbers: BTreeSet::new(), - present_part_numbers: BTreeSet::new(), present_part_fingerprints: BTreeMap::new(), + inline_data_fingerprint: None, }); } @@ -315,10 +328,10 @@ pub(crate) fn census_object_version_on_disk( }; let data_dir = file_info.data_dir.map(|id| id.to_string()); let erasure_index = Some(file_info.erasure.index); + let inline_data_fingerprint = file_info.data.as_deref().map(shard_fingerprint).transpose()?; let part_dir = data_dir.as_ref().map_or_else(|| object_dir.clone(), |id| object_dir.join(id)); - let (present_part_numbers, present_part_fingerprints) = match std::fs::read_dir(&part_dir) { + let present_part_fingerprints = match std::fs::read_dir(&part_dir) { Ok(entries) => { - let mut numbers = BTreeSet::new(); let mut fingerprints = BTreeMap::new(); for entry in entries { let entry = entry?; @@ -333,19 +346,12 @@ pub(crate) fn census_object_version_on_disk( else { continue; }; - numbers.insert(part_number); let data = std::fs::read(entry.path())?; - fingerprints.insert( - part_number, - PartShardFingerprint { - size: u64::try_from(data.len())?, - sha256: sha256_hex(&data), - }, - ); + fingerprints.insert(part_number, shard_fingerprint(&data)?); } - (numbers, fingerprints) + fingerprints } - Err(error) if error.kind() == std::io::ErrorKind::NotFound => (BTreeSet::new(), BTreeMap::new()), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => BTreeMap::new(), Err(error) => return Err(error.into()), }; @@ -355,8 +361,8 @@ pub(crate) fn census_object_version_on_disk( data_dir, erasure_index, expected_part_numbers, - present_part_numbers, present_part_fingerprints, + inline_data_fingerprint, }) } @@ -396,3 +402,43 @@ pub async fn signed_admin_post(url: &str, body: Option<&str>, access_key: &str, Ok(body) } + +#[cfg(test)] +mod tests { + use super::*; + + fn complete_census() -> VersionShardCensus { + VersionShardCensus { + version_id: Some("version".to_string()), + has_xl_meta: true, + data_dir: Some("data-dir".to_string()), + erasure_index: Some(3), + expected_part_numbers: BTreeSet::from([1]), + present_part_fingerprints: BTreeMap::from([(1, shard_fingerprint(b"part").unwrap())]), + inline_data_fingerprint: None, + } + } + + #[test] + fn shard_fingerprint_uses_physical_length_and_sha256() { + assert_eq!( + shard_fingerprint(b"abc").unwrap(), + PartShardFingerprint { + size: 3, + sha256: "ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad".to_string(), + } + ); + } + + #[test] + fn manifest_requires_matching_inline_payload() { + let mut expected = complete_census(); + expected.expected_part_numbers.clear(); + expected.present_part_fingerprints.clear(); + expected.inline_data_fingerprint = Some(shard_fingerprint(b"expected").unwrap()); + let mut changed = expected.clone(); + changed.inline_data_fingerprint = Some(shard_fingerprint(b"changed").unwrap()); + assert!(expected.matches_manifest(&expected)); + assert!(!changed.matches_manifest(&expected)); + } +} diff --git a/crates/e2e_test/src/reliability_disk_fault_test.rs b/crates/e2e_test/src/reliability_disk_fault_test.rs index 1aa1a0d70..3b2ae3e29 100644 --- a/crates/e2e_test/src/reliability_disk_fault_test.rs +++ b/crates/e2e_test/src/reliability_disk_fault_test.rs @@ -349,11 +349,32 @@ mod tests { .send() .await?; + let first_inline = client + .put_object() + .bucket(bucket) + .key("versions/inline.bin") + .body(ByteStream::from(payload(8 * 1024, 40))) + .send() + .await?; + let first_inline_version = first_inline + .version_id() + .ok_or("first inline PUT did not return a version ID")?; + let second_inline = client + .put_object() + .bucket(bucket) + .key("versions/inline.bin") + .body(ByteStream::from(payload(8 * 1024, 41))) + .send() + .await?; + let second_inline_version = second_inline + .version_id() + .ok_or("second inline PUT did not return a version ID")?; + let first = client .put_object() .bucket(bucket) .key(key) - .body(ByteStream::from(payload(256 * 1024, 41))) + .body(ByteStream::from(payload(128 * 1024, 41))) .send() .await?; let first_version = first.version_id().ok_or("first PUT did not return a version ID")?; @@ -361,16 +382,36 @@ mod tests { .put_object() .bucket(bucket) .key(key) - .body(ByteStream::from(payload(256 * 1024, 42))) + .body(ByteStream::from(payload(3 * 1024 * 1024, 42))) .send() .await?; let second_version = second.version_id().ok_or("second PUT did not return a version ID")?; let delete = client.delete_object().bucket(bucket).key(key).send().await?; let delete_version = delete.version_id().ok_or("delete marker did not return a version ID")?; + let first_inline_census = harness.census_object_version(0, bucket, "versions/inline.bin", Some(first_inline_version))?; + let second_inline_census = + harness.census_object_version(0, bucket, "versions/inline.bin", Some(second_inline_version))?; let first_census = harness.census_object_version(0, bucket, key, Some(first_version))?; + let first_other_disk_census = harness.census_object_version(1, bucket, key, Some(first_version))?; let second_census = harness.census_object_version(0, bucket, key, Some(second_version))?; let delete_census = harness.census_object_version(0, bucket, key, Some(delete_version))?; + assert!( + first_inline_census.is_complete() && second_inline_census.is_complete(), + "inline version physical census is incomplete: first={first_inline_census:?} second={second_inline_census:?}" + ); + assert!( + first_inline_census.present_part_fingerprints.is_empty() && second_inline_census.present_part_fingerprints.is_empty(), + "inline versions must not select external shard files: first={first_inline_census:?} second={second_inline_census:?}" + ); + assert!( + first_inline_census.inline_data_fingerprint.is_some() && second_inline_census.inline_data_fingerprint.is_some(), + "inline versions must fingerprint payload bytes stored in xl.meta" + ); + assert_ne!( + first_inline_census.inline_data_fingerprint, second_inline_census.inline_data_fingerprint, + "same-size inline versions with different payloads must retain distinct xl.meta fingerprints" + ); assert!( first_census.is_complete(), "first version physical census is incomplete: {first_census:?}" @@ -379,6 +420,14 @@ mod tests { second_census.is_complete(), "second version physical census is incomplete: {second_census:?}" ); + assert!( + first_other_disk_census.is_complete(), + "first version physical census on the second disk is incomplete: {first_other_disk_census:?}" + ); + assert_ne!( + first_census.erasure_index, first_other_disk_census.erasure_index, + "physical census must preserve each disk's erasure index" + ); assert_ne!( first_census.data_dir, second_census.data_dir, "distinct object versions must select distinct physical data directories" @@ -387,6 +436,24 @@ mod tests { first_census.expected_part_numbers, second_census.expected_part_numbers, "same single-part shape should expose the same part numbers" ); + let first_part = first_census + .present_part_fingerprints + .values() + .next() + .ok_or("first version did not expose a physical part fingerprint")?; + let second_part = second_census + .present_part_fingerprints + .values() + .next() + .ok_or("second version did not expose a physical part fingerprint")?; + assert_ne!( + first_part.size, second_part.size, + "different shard lengths must retain their physical sizes" + ); + assert_ne!( + first_part.sha256, second_part.sha256, + "different shard contents must retain their physical hashes" + ); assert!( delete_census.is_complete(), "delete marker physical census is incomplete: {delete_census:?}" @@ -396,7 +463,7 @@ mod tests { "delete marker must not declare object shards: {delete_census:?}" ); assert!( - delete_census.present_part_numbers.is_empty(), + delete_census.present_part_fingerprints.is_empty(), "delete marker must not select stale object shards: {delete_census:?}" ); Ok(()) diff --git a/crates/e2e_test/src/replacement_privileged_e2e_test.rs b/crates/e2e_test/src/replacement_privileged_e2e_test.rs index ab12c2eef..67c1b0097 100644 --- a/crates/e2e_test/src/replacement_privileged_e2e_test.rs +++ b/crates/e2e_test/src/replacement_privileged_e2e_test.rs @@ -59,6 +59,13 @@ mod tests { expected: VersionShardCensus, } + #[derive(Debug, Eq, PartialEq)] + enum CompletionSample { + Pending, + Ready, + CompletedWithIncomplete(BTreeSet), + } + struct MountNamespaceGuard { mounts: Vec, } @@ -75,7 +82,7 @@ mod tests { impl MountNamespaceGuard { fn new() -> Result> { verify_isolated_mount_namespace()?; - run_command("mount", ["--make-rprivate", "/"])?; + run_command("mount", &["--make-rprivate", "/"])?; Ok(Self { mounts: Vec::new() }) } @@ -108,15 +115,15 @@ mod tests { return Err("losetup --find --show returned an empty loop device".into()); } - run_command_dynamic("mkfs.ext4", &["-F", &loop_device])?; + 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"); let mapper = format!("/dev/mapper/{dm_name}"); - run_command_dynamic("dmsetup", &["create", &dm_name, "--table", &table])?; + run_command("dmsetup", &["create", &dm_name, "--table", &table])?; let target_arg = path_to_string(target, "faultable mount target")?; - run_command_dynamic("mount", &[&mapper, &target_arg])?; + run_command("mount", &[&mapper, &target_arg])?; Ok(Self { target: target.to_path_buf(), @@ -131,17 +138,17 @@ mod tests { fn make_unavailable(&self) -> Result<(), Box> { let sectors = run_command_stdout("blockdev", &["--getsz", &self.loop_device])?; let error_table = format!("0 {sectors} error"); - run_command_dynamic("dmsetup", &["suspend", &self.dm_name])?; - run_command_dynamic("dmsetup", &["load", &self.dm_name, "--table", &error_table])?; - run_command_dynamic("dmsetup", &["resume", &self.dm_name]) + run_command("dmsetup", &["suspend", &self.dm_name])?; + run_command("dmsetup", &["load", &self.dm_name, "--table", &error_table])?; + run_command("dmsetup", &["resume", &self.dm_name]) } fn restore_available(&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_dynamic("dmsetup", &["suspend", &self.dm_name])?; - run_command_dynamic("dmsetup", &["load", &self.dm_name, "--table", &linear_table])?; - run_command_dynamic("dmsetup", &["resume", &self.dm_name]) + run_command("dmsetup", &["suspend", &self.dm_name])?; + run_command("dmsetup", &["load", &self.dm_name, "--table", &linear_table])?; + run_command("dmsetup", &["resume", &self.dm_name]) } fn cleanup(&mut self) -> Result<(), Box> { @@ -157,14 +164,14 @@ mod tests { } } if self.dm_created { - if let Err(error) = run_command_dynamic("dmsetup", &["remove", "-f", &self.dm_name]) { + if let Err(error) = run_command("dmsetup", &["remove", "-f", &self.dm_name]) { first_error.get_or_insert(error); } else { self.dm_created = false; } } if !self.loop_device.is_empty() { - if let Err(error) = run_command_dynamic("losetup", &["-d", &self.loop_device]) { + if let Err(error) = run_command("losetup", &["-d", &self.loop_device]) { first_error.get_or_insert(error); } else { self.loop_device.clear(); @@ -188,14 +195,10 @@ mod tests { } } - fn run_command(program: &str, args: [&str; N]) -> Result<(), Box> { - run_command_dynamic(program, &args) - } - - fn run_command_dynamic(program: &str, args: &[&str]) -> Result<(), Box> { + fn checked_command_output(program: &str, args: &[&str]) -> Result> { let output = Command::new(program).args(args).output()?; if output.status.success() { - return Ok(()); + return Ok(output); } Err(format!( "{program} {} failed with status {}: stdout={} stderr={}", @@ -207,19 +210,14 @@ mod tests { .into()) } + fn run_command(program: &str, args: &[&str]) -> Result<(), Box> { + checked_command_output(program, args).map(drop) + } + fn run_command_stdout(program: &str, args: &[&str]) -> Result> { - let output = Command::new(program).args(args).output()?; - if output.status.success() { - return Ok(String::from_utf8_lossy(&output.stdout).trim().to_string()); - } - Err(format!( - "{program} {} failed with status {}: stdout={} stderr={}", - args.join(" "), - output.status, - String::from_utf8_lossy(&output.stdout), - String::from_utf8_lossy(&output.stderr) - ) - .into()) + Ok(String::from_utf8(checked_command_output(program, args)?.stdout)? + .trim() + .to_string()) } fn path_to_string(path: &Path, label: &str) -> Result> { @@ -260,14 +258,14 @@ mod tests { let target = target .to_str() .ok_or_else(|| format!("tmpfs target path is not UTF-8: {target:?}"))?; - run_command("mount", ["-t", "tmpfs", "-o", MOUNT_SIZE, label, target]) + run_command("mount", &["-t", "tmpfs", "-o", MOUNT_SIZE, label, target]) } fn detach_mount(target: &Path) -> Result<(), Box> { let target = target .to_str() .ok_or_else(|| format!("umount target path is not UTF-8: {target:?}"))?; - run_command("umount", [target]) + run_command("umount", &[target]) } fn privileged_run_enabled() -> Result> { @@ -423,6 +421,10 @@ mod tests { .await?; versions.push((versioned_bucket, "history/object.bin", deleted.version_id().map(str::to_owned), None)); + let (version_id, body_sha256) = + put_object_version(client, versioned_bucket, "history/inline.bin", payload(8 * 1024, 9)).await?; + versions.push((versioned_bucket, "history/inline.bin", version_id, Some(body_sha256))); + let (version_id, body_sha256) = put_multipart_version( client, versioned_bucket, @@ -436,7 +438,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))); - versions + let 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())?; @@ -451,7 +453,15 @@ mod tests { expected, }) }) - .collect() + .collect::, Box>>()?; + let inline = versions + .iter() + .find(|version| version.key == "history/inline.bin") + .ok_or("inline replacement baseline was not recorded")?; + if inline.expected.inline_data_fingerprint.is_none() || !inline.expected.present_part_fingerprints.is_empty() { + return Err(format!("inline replacement baseline lacks xl.meta payload evidence: {:?}", inline.expected).into()); + } + Ok(versions) } async fn verify_bodies(client: &Client, versions: &[BaselineVersion]) -> Result<(), Box> { @@ -597,6 +607,33 @@ mod tests { Ok(String::from_utf8_lossy(&log[start..]).into_owned()) } + fn live_disk_loss_scan_completed(log: &str, target_disk: &Path) -> bool { + let target = target_disk.to_string_lossy(); + let mut saw_live_loss = false; + for line in log.lines() { + if line.contains("Heal auto-scan disk inspection failed") + && line.contains("check_failed") + && line.contains(target.as_ref()) + { + saw_live_loss = true; + continue; + } + if saw_live_loss && (line.contains("Heal auto disk scanner idle") || line.contains("Heal auto-scan cycle completed")) + { + return true; + } + } + false + } + + fn live_disk_loss_scan_completed_from_path( + log_path: &Path, + start_offset: u64, + target_disk: &Path, + ) -> Result> { + Ok(live_disk_loss_scan_completed(&log_from_offset(log_path, start_offset)?, target_disk)) + } + async fn wait_for_live_disk_loss_observation( log_path: &Path, target_disk: &Path, @@ -605,20 +642,12 @@ mod tests { ) -> Result<(), Box> { let deadline = Instant::now() + Duration::from_secs(timeout_secs); let mut tick = interval(Duration::from_secs(1)); - let target = target_disk.to_string_lossy(); loop { - let log = log_from_offset(log_path, start_offset)?; - let mut saw_live_loss = false; - for line in log.lines() { - if line.contains("check_failed") && line.contains(target.as_ref()) { - saw_live_loss = true; - continue; - } - if saw_live_loss && line.contains("Heal auto disk scanner idle") { - return Ok(()); - } + if live_disk_loss_scan_completed_from_path(log_path, start_offset, target_disk)? { + return Ok(()); } if Instant::now() >= deadline { + let log = log_from_offset(log_path, start_offset)?; return Err(format!( "scanner did not finish a live target-loss scan for {target_disk:?} within {timeout_secs}s; log tail:\n{}", log_tail(&log) @@ -629,14 +658,10 @@ mod tests { } } - fn require_definitive_replacement_status( - status: &serde_json::Value, - context: &str, - ) -> Result<(), Box> { - if status["cluster"]["definitive"].as_bool().unwrap_or(false) { - return Ok(()); - } - Err(format!("{context} requires a definitive cluster replacement status: {status}").into()) + fn cluster_status_is_definitive(status: &serde_json::Value) -> Result> { + status["cluster"]["definitive"] + .as_bool() + .ok_or_else(|| format!("replacement recovery status omitted cluster.definitive: {status}").into()) } async fn assert_no_replacement_status_records( @@ -644,8 +669,17 @@ mod tests { target_disk: &Path, ) -> Result<(), Box> { let status = replacement_status(cluster).await?; - require_definitive_replacement_status(&status, "live missing replacement status check")?; - let states = target_record_states(&status, target_disk); + assert_no_replacement_status_records_in_status(&status, target_disk) + } + + fn assert_no_replacement_status_records_in_status( + status: &serde_json::Value, + target_disk: &Path, + ) -> Result<(), Box> { + if !cluster_status_is_definitive(status)? { + return Err(format!("live missing replacement status check requires a definitive cluster status: {status}").into()); + } + let states = target_record_states(status, target_disk); if states.is_empty() { return Ok(()); } @@ -655,11 +689,6 @@ mod tests { .into()) } - fn target_record_has_state(status: &serde_json::Value, target_disk: &Path, states: &[&str]) -> bool { - let present_states = target_record_states(status, target_disk); - states.iter().any(|state| present_states.contains(*state)) - } - fn target_record_states(status: &serde_json::Value, target_disk: &Path) -> BTreeSet { let target = target_disk.to_string_lossy(); status["cluster"]["records"] @@ -694,6 +723,57 @@ mod tests { Ok(missing) } + fn replacement_completion_state( + status: &serde_json::Value, + target_disk: &Path, + missing: BTreeSet, + ) -> Result> { + if !cluster_status_is_definitive(status)? { + return Ok(CompletionSample::Pending); + } + if !target_record_states(status, target_disk).contains("completed") { + return Ok(CompletionSample::Pending); + } + if missing.is_empty() { + return Ok(CompletionSample::Ready); + } + Ok(CompletionSample::CompletedWithIncomplete(missing)) + } + + async fn sample_replacement_completion( + target_disk: &Path, + census: C, + status: S, + ) -> Result> + where + C: FnOnce() -> Result, Box>, + S: FnOnce() -> F, + F: std::future::Future>>, + { + let missing = census()?; + let status = status().await?; + replacement_completion_state(&status, target_disk, missing) + } + + async fn confirm_replacement_completion( + target_disk: &Path, + mut census: C, + mut status: S, + ) -> Result> + where + C: FnMut() -> Result, Box>, + S: FnMut() -> F, + F: std::future::Future>>, + { + let missing = census()?; + let status = status().await?; + let first = replacement_completion_state(&status, target_disk, missing)?; + if matches!(first, CompletionSample::CompletedWithIncomplete(_)) { + return replacement_completion_state(&status, target_disk, census()?); + } + Ok(first) + } + async fn wait_for_completed_replacement_with_census( cluster: &RustFSTestClusterEnvironment, target_disk: &Path, @@ -703,28 +783,25 @@ mod tests { let deadline = Instant::now() + Duration::from_secs(timeout_secs); let mut tick = interval(Duration::from_secs(1)); loop { - let status = replacement_status(cluster).await?; - let missing = incomplete_versions(target_disk, versions)?; - if require_definitive_replacement_status(&status, "replacement completion poll").is_err() { - if Instant::now() >= deadline { + match confirm_replacement_completion( + target_disk, + || incomplete_versions(target_disk, versions), + || replacement_status(cluster), + ) + .await? + { + CompletionSample::Ready => return Ok(()), + CompletionSample::CompletedWithIncomplete(confirmed_missing) => { return Err(format!( - "replacement recovery status never became definitive within {timeout_secs}s while waiting for physical census; latest status: {status}; missing: {missing:?}" - ) - .into()); + "replacement status remained completed across two incomplete physical censuses: {confirmed_missing:?}" + ) + .into()); } - tick.tick().await; - continue; - } - if target_record_has_state(&status, target_disk, &["completed"]) { - if !missing.is_empty() { - return Err(format!( - "replacement status reached completed before target physical census matched baseline: {missing:?}; status: {status}" - ) - .into()); - } - return Ok(()); + CompletionSample::Pending => {} } if Instant::now() >= deadline { + let missing = incomplete_versions(target_disk, versions)?; + let status = replacement_status(cluster).await?; return Err(format!( "replacement target did not reach completed with matching physical census within {timeout_secs}s: missing={missing:?}; status={status}" ) @@ -812,6 +889,175 @@ mod tests { Ok(()) } + #[test] + fn live_loss_barrier_requires_scanner_failure_after_log_offset() -> Result<(), Box> { + let target = Path::new("/mnt/target"); + assert!(live_disk_loss_scan_completed( + "Heal auto-scan disk inspection failed endpoint=/mnt/target disk_state=check_failed\nHeal auto-scan cycle completed", + target + )); + assert!(live_disk_loss_scan_completed( + "Heal auto-scan disk inspection failed endpoint=/mnt/target disk_state=check_failed\nHeal auto disk scanner idle", + target + )); + assert!(!live_disk_loss_scan_completed( + "Heal auto disk scanner idle\nHeal auto-scan disk inspection failed endpoint=/mnt/target disk_state=check_failed", + target + )); + assert!(!live_disk_loss_scan_completed( + "event=disk_health_check_failed endpoint=/mnt/target disk_state=check_failed\nHeal auto disk scanner idle", + target + )); + assert!(!live_disk_loss_scan_completed( + "Heal auto-scan disk inspection failed endpoint=/mnt/other disk_state=check_failed\nHeal auto disk scanner idle", + target + )); + let path = std::env::temp_dir().join(format!("rustfs-replacement-scan-{}.log", std::process::id())); + let stale = + "Heal auto-scan disk inspection failed endpoint=/mnt/target disk_state=check_failed\nHeal auto disk scanner idle\n"; + fs::write(&path, stale)?; + let offset = log_len(&path)?; + assert!(!live_disk_loss_scan_completed_from_path(&path, offset, target)?); + let fresh = + "Heal auto-scan disk inspection failed endpoint=/mnt/target disk_state=check_failed\nHeal auto disk scanner idle\n"; + fs::write(&path, format!("{stale}{fresh}"))?; + assert!(live_disk_loss_scan_completed_from_path(&path, offset, target)?); + fs::remove_file(path)?; + Ok(()) + } + + #[test] + fn completion_requires_definitive_status_and_prior_census_match() { + let target = Path::new("/mnt/target"); + let non_definitive = serde_json::json!({ + "cluster": {"definitive": false, "records": [{"state": "completed", "targetSlots": ["/mnt/target"]}]} + }); + assert_eq!( + replacement_completion_state(&non_definitive, target, BTreeSet::new()).unwrap(), + CompletionSample::Pending + ); + + let omitted = serde_json::json!({ + "cluster": {"records": [{"state": "completed", "targetSlots": ["/mnt/target"]}]} + }); + assert!(replacement_completion_state(&omitted, target, BTreeSet::new()).is_err()); + + let definitive = serde_json::json!({ + "cluster": {"definitive": true, "records": [{"state": "completed", "targetSlots": ["/mnt/target"]}]} + }); + assert_eq!( + replacement_completion_state(&definitive, target, BTreeSet::from(["missing".to_string()])).unwrap(), + CompletionSample::CompletedWithIncomplete(BTreeSet::from(["missing".to_string()])) + ); + assert_eq!( + replacement_completion_state(&definitive, target, BTreeSet::new()).unwrap(), + CompletionSample::Ready + ); + } + + #[tokio::test] + async fn completion_poll_samples_census_before_status() { + let order = std::rc::Rc::new(std::cell::RefCell::new(Vec::new())); + let census_order = order.clone(); + let status_order = order.clone(); + let sample = sample_replacement_completion( + Path::new("/mnt/target"), + move || { + census_order.borrow_mut().push("census"); + Ok::<_, Box>(BTreeSet::new()) + }, + move || async move { + status_order.borrow_mut().push("status"); + Ok::<_, Box>(serde_json::json!({ + "cluster": {"definitive": true, "records": [{"state": "completed", "targetSlots": ["/mnt/target"]}]} + })) + }, + ) + .await + .unwrap(); + assert_eq!(sample, CompletionSample::Ready); + assert_eq!(*order.borrow(), ["census", "status"]); + } + + #[tokio::test] + async fn completed_status_confirms_a_stale_incomplete_census() { + let samples = std::rc::Rc::new(std::cell::RefCell::new(std::collections::VecDeque::from([ + BTreeSet::from(["missing".to_string()]), + BTreeSet::new(), + ]))); + let census_samples = samples.clone(); + let result = confirm_replacement_completion( + Path::new("/mnt/target"), + move || { + census_samples + .borrow_mut() + .pop_front() + .ok_or_else(|| "missing census sample".into()) + }, + || async { + Ok(serde_json::json!({ + "cluster": {"definitive": true, "records": [{"state": "completed", "targetSlots": ["/mnt/target"]}]} + })) + }, + ) + .await + .unwrap(); + assert_eq!(result, CompletionSample::Ready); + assert!(samples.borrow().is_empty()); + + let status_samples = std::rc::Rc::new(std::cell::RefCell::new(std::collections::VecDeque::from([ + serde_json::json!({ + "cluster": {"definitive": true, "records": [{"state": "completed", "targetSlots": ["/mnt/target"]}]} + }), + serde_json::json!({ + "cluster": {"definitive": false, "records": []} + }), + ]))); + let persistent = std::rc::Rc::new(std::cell::RefCell::new(std::collections::VecDeque::from([ + BTreeSet::from(["missing".to_string()]), + BTreeSet::from(["still-missing".to_string()]), + ]))); + let census_samples = persistent.clone(); + let statuses = status_samples.clone(); + let result = confirm_replacement_completion( + Path::new("/mnt/target"), + move || { + census_samples + .borrow_mut() + .pop_front() + .ok_or_else(|| "missing census sample".into()) + }, + move || { + let statuses = statuses.clone(); + async move { + statuses + .borrow_mut() + .pop_front() + .ok_or_else(|| "missing status sample".into()) + } + }, + ) + .await + .unwrap(); + assert_eq!( + result, + CompletionSample::CompletedWithIncomplete(BTreeSet::from(["still-missing".to_string()])) + ); + assert!(persistent.borrow().is_empty()); + assert_eq!(status_samples.borrow().len(), 1); + } + + #[test] + fn absent_status_requires_definitive_empty_records() { + let target = Path::new("/mnt/target"); + let non_definitive = serde_json::json!({"cluster": {"definitive": false, "records": []}}); + assert!(assert_no_replacement_status_records_in_status(&non_definitive, target).is_err()); + let omitted = serde_json::json!({"cluster": {"records": []}}); + assert!(assert_no_replacement_status_records_in_status(&omitted, target).is_err()); + let definitive = serde_json::json!({"cluster": {"definitive": true, "records": []}}); + assert!(assert_no_replacement_status_records_in_status(&definitive, target).is_ok()); + } + /// 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")]