test(heal): strengthen replacement e2e evidence (#5956)

* test(heal): strengthen replacement e2e evidence

* test(heal): fix replacement e2e barriers

Accept the real post-fault scanner failure-to-idle sequence as the live disk loss barrier, and preserve the first definitive completed status while only resampling the physical census for premature-completion confirmation.

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
Zhengchao An
2026-08-12 16:14:52 +08:00
committed by GitHub
parent 3ebb426abe
commit 493a2cc1ba
3 changed files with 456 additions and 97 deletions
+62 -16
View File
@@ -62,8 +62,8 @@ pub(crate) struct VersionShardCensus {
pub data_dir: Option<String>,
pub erasure_index: Option<usize>,
pub expected_part_numbers: BTreeSet<usize>,
pub present_part_numbers: BTreeSet<usize>,
pub present_part_fingerprints: BTreeMap<usize, PartShardFingerprint>,
pub inline_data_fingerprint: Option<PartShardFingerprint>,
}
#[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<PartShardFingerprint> {
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));
}
}
@@ -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(())
@@ -59,6 +59,13 @@ mod tests {
expected: VersionShardCensus,
}
#[derive(Debug, Eq, PartialEq)]
enum CompletionSample {
Pending,
Ready,
CompletedWithIncomplete(BTreeSet<String>),
}
struct MountNamespaceGuard {
mounts: Vec<PathBuf>,
}
@@ -75,7 +82,7 @@ mod tests {
impl MountNamespaceGuard {
fn new() -> Result<Self, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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<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_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<dyn Error + Send + Sync>> {
@@ -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<const N: usize>(program: &str, args: [&str; N]) -> Result<(), Box<dyn Error + Send + Sync>> {
run_command_dynamic(program, &args)
}
fn run_command_dynamic(program: &str, args: &[&str]) -> Result<(), Box<dyn Error + Send + Sync>> {
fn checked_command_output(program: &str, args: &[&str]) -> Result<std::process::Output, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
checked_command_output(program, args).map(drop)
}
fn run_command_stdout(program: &str, args: &[&str]) -> Result<String, Box<dyn Error + Send + Sync>> {
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<String, Box<dyn Error + Send + Sync>> {
@@ -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<dyn Error + Send + Sync>> {
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<bool, Box<dyn Error + Send + Sync>> {
@@ -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::<Result<Vec<_>, Box<dyn Error + Send + Sync>>>()?;
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<dyn Error + Send + Sync>> {
@@ -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<bool, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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<bool, Box<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>> {
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<String> {
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<String>,
) -> Result<CompletionSample, Box<dyn Error + Send + Sync>> {
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<C, S, F>(
target_disk: &Path,
census: C,
status: S,
) -> Result<CompletionSample, Box<dyn Error + Send + Sync>>
where
C: FnOnce() -> Result<BTreeSet<String>, Box<dyn Error + Send + Sync>>,
S: FnOnce() -> F,
F: std::future::Future<Output = Result<serde_json::Value, Box<dyn Error + Send + Sync>>>,
{
let missing = census()?;
let status = status().await?;
replacement_completion_state(&status, target_disk, missing)
}
async fn confirm_replacement_completion<C, S, F>(
target_disk: &Path,
mut census: C,
mut status: S,
) -> Result<CompletionSample, Box<dyn Error + Send + Sync>>
where
C: FnMut() -> Result<BTreeSet<String>, Box<dyn Error + Send + Sync>>,
S: FnMut() -> F,
F: std::future::Future<Output = Result<serde_json::Value, Box<dyn Error + Send + Sync>>>,
{
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<dyn Error + Send + Sync>> {
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<dyn Error + Send + Sync>>(BTreeSet::new())
},
move || async move {
status_order.borrow_mut().push("status");
Ok::<_, Box<dyn Error + Send + Sync>>(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")]