Files
rustfs/crates/e2e_test/src/replacement_privileged_e2e_test.rs
T
Henry Guo 98f7e63396 fix(heal): recover replacement after transient disk errors (#7059)
* fix(heal): recover replacement after transient disk errors

* fix(heal): satisfy replacement status clippy lint
2026-09-03 07:03:01 +08:00

1568 lines
65 KiB
Rust

// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Privileged Linux replacement-heal E2E.
//!
//! These ignored tests exercise the documented drive-swap path from rustfs#5869:
//! a 3-node x 4-drive cluster, a restarted target node whose absent endpoint is
//! observed by the scanner, a blank replacement mounted back at the same
//! endpoint, and automatic recovery without an Admin deep-heal request. The final
//! assertion is a per-version physical census of the replacement target's
//! `xl.meta` and `part.N` files.
#[cfg(all(test, target_os = "linux"))]
mod tests {
use crate::chaos::{VersionShardCensus, census_object_version_on_disk};
use crate::common::{ClusterTopology, RustFSTestClusterEnvironment, admin_request, init_logging};
use aws_sdk_s3::Client;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, VersioningConfiguration};
use http::Method;
use sha2::{Digest, Sha256};
use std::collections::BTreeSet;
use std::error::Error;
use std::fs;
use std::path::{Path, PathBuf};
use std::process::Command;
use tokio::net::TcpStream;
use tokio::time::{Duration, Instant, interval, sleep, timeout};
use tracing::info;
const ENABLE_ENV: &str = "RUSTFS_PRIVILEGED_REPLACEMENT_E2E";
const NAMESPACE_ENV: &str = "RUSTFS_PRIVILEGED_REPLACEMENT_E2E_IN_NAMESPACE";
const LOG_DIR_ENV: &str = "RUSTFS_PRIVILEGED_REPLACEMENT_LOG_DIR";
const TARGET_NODE: usize = 1;
const TARGET_DRIVE: usize = 0;
const MOUNT_SIZE: &str = "size=128m,mode=0700";
const REPLACEMENT_RECOVERY_DIR: &str = ".rustfs.sys/buckets/ahm-replacement";
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 {
bucket: String,
key: String,
version_id: Option<String>,
body_sha256: Option<String>,
expected: VersionShardCensus,
}
#[derive(Debug, Eq, PartialEq)]
enum CompletionSample {
Pending,
Ready,
CompletedWithIncomplete(BTreeSet<String>),
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ReplacementScenario {
Baseline,
MidRebuildIoFault,
}
struct MountNamespaceGuard {
mounts: Vec<PathBuf>,
}
struct FaultableBlockMount {
target: PathBuf,
image: PathBuf,
loop_device: String,
dm_name: String,
mounted: bool,
dm_created: bool,
}
impl MountNamespaceGuard {
fn new() -> Result<Self, Box<dyn Error + Send + Sync>> {
verify_isolated_mount_namespace()?;
run_command("mount", &["--make-rprivate", "/"])?;
Ok(Self { mounts: Vec::new() })
}
fn mount_tmpfs(&mut self, target: &Path, label: &str) -> Result<(), Box<dyn Error + Send + Sync>> {
mount_tmpfs(target, label)?;
self.mounts.push(target.to_path_buf());
Ok(())
}
}
impl Drop for MountNamespaceGuard {
fn drop(&mut self) {
for mount in self.mounts.iter().rev() {
let _ = detach_mount(mount);
}
}
}
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)?;
file.set_len(256 * 1024 * 1024)?;
drop(file);
let image_arg = path_to_string(&image, "loop image")?;
let loop_device = run_command_stdout("losetup", &["--find", "--show", &image_arg])?;
if loop_device.is_empty() {
return Err("losetup --find --show returned an empty loop device".into());
}
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");
let mapper = format!("/dev/mapper/{dm_name}");
run_command("dmsetup", &["create", &dm_name, "--table", &table])?;
let target_arg = path_to_string(target, "faultable mount target")?;
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(),
image,
loop_device,
dm_name,
mounted: true,
dm_created: true,
})
}
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("dmsetup", &["suspend", &self.dm_name])?;
run_command("dmsetup", &["load", &self.dm_name, "--table", &error_table])?;
run_command("dmsetup", &["resume", &self.dm_name])
}
fn verify_raw_io_is_unavailable(&self) -> Result<(), Box<dyn Error + Send + Sync>> {
let mapper = format!("/dev/mapper/{}", self.dm_name);
let output = Command::new("dd")
.env("LC_ALL", "C")
.arg(format!("if={mapper}"))
.args(["of=/dev/null", "bs=4096", "count=1", "iflag=direct", "status=none"])
.output()?;
if output.status.success() {
return Err(format!("dm-error target unexpectedly allowed a raw read from {mapper}").into());
}
let stderr = String::from_utf8_lossy(&output.stderr);
if !stderr.contains("Input/output error") {
return Err(format!("raw read from dm-error target failed unexpectedly: {stderr}").into());
}
Ok(())
}
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);
// 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_linear_table();
}
if self.mounted {
if let Err(error) = detach_mount(&self.target) {
first_error.get_or_insert(error);
} else {
self.mounted = false;
}
}
if self.dm_created {
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("losetup", &["-d", &self.loop_device]) {
first_error.get_or_insert(error);
} else {
self.loop_device.clear();
}
}
if self.image.exists()
&& let Err(error) = fs::remove_file(&self.image)
{
first_error.get_or_insert(error.into());
}
if let Some(error) = first_error {
return Err(error);
}
Ok(())
}
}
impl Drop for FaultableBlockMount {
fn drop(&mut self) {
let _ = self.cleanup();
}
}
struct ZramBlockMount {
target: PathBuf,
device: String,
mounted: bool,
}
impl ZramBlockMount {
fn reserve(target: &Path) -> Result<Self, Box<dyn Error + Send + Sync>> {
if !Path::new("/dev/zram-control").exists() {
run_command("modprobe", &["zram"])?;
}
let device = run_command_stdout("zramctl", &["--find", "--size", "256M"])?;
if device.is_empty() {
return Err("zramctl --find --size returned an empty device".into());
}
Ok(Self {
target: target.to_path_buf(),
device,
mounted: false,
})
}
fn mount_target(&mut self) -> Result<(), Box<dyn Error + Send + Sync>> {
let result = (|| {
run_command("mkfs.ext4", &["-F", &self.device])?;
let target_arg = path_to_string(&self.target, "zram replacement mount target")?;
run_command("mount", &[&self.device, &target_arg])
})();
if let Err(error) = result {
let _ = self.cleanup();
return Err(error);
}
self.mounted = true;
Ok(())
}
fn cleanup(&mut self) -> Result<(), Box<dyn Error + Send + Sync>> {
let mut first_error: Option<Box<dyn Error + Send + Sync>> = None;
if self.mounted {
if let Err(error) = detach_mount(&self.target) {
first_error.get_or_insert(error);
} else {
self.mounted = false;
}
}
if !self.device.is_empty() {
if let Err(error) = run_command("zramctl", &["--reset", &self.device]) {
first_error.get_or_insert(error);
} else {
self.device.clear();
}
}
if let Some(error) = first_error {
return Err(error);
}
Ok(())
}
}
impl Drop for ZramBlockMount {
fn drop(&mut self) {
let _ = self.cleanup();
}
}
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(output);
}
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())
}
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>> {
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>> {
path.to_str()
.map(str::to_owned)
.ok_or_else(|| format!("{label} path is not UTF-8: {path:?}").into())
}
fn mount_namespace_link(proc_entry: &str) -> Result<String, Box<dyn Error + Send + Sync>> {
Ok(fs::read_link(format!("/proc/{proc_entry}/ns/mnt"))?
.to_string_lossy()
.into_owned())
}
fn parent_pid() -> Result<String, Box<dyn Error + Send + Sync>> {
let status = fs::read_to_string("/proc/self/status")?;
for line in status.lines() {
if let Some(ppid) = line.strip_prefix("PPid:") {
return Ok(ppid.trim().to_string());
}
}
Err("/proc/self/status does not contain PPid".into())
}
fn verify_isolated_mount_namespace() -> Result<(), Box<dyn Error + Send + Sync>> {
let current = mount_namespace_link("self")?;
let parent = mount_namespace_link(&parent_pid()?)?;
if current == parent {
return Err(format!(
"privileged replacement E2E must run in a private mount namespace before mounting test drives; current namespace {current} still matches parent"
)
.into());
}
Ok(())
}
fn mount_tmpfs(target: &Path, label: &str) -> Result<(), Box<dyn Error + Send + Sync>> {
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])
}
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])
}
fn privileged_run_enabled() -> Result<bool, Box<dyn Error + Send + Sync>> {
let enabled = std::env::var(ENABLE_ENV)
.ok()
.is_some_and(|value| matches!(value.as_str(), "1" | "true" | "TRUE" | "yes" | "YES"));
if !enabled {
info!("{ENABLE_ENV}=1 is not set; privileged replacement E2E is skipped");
return Ok(false);
}
Ok(true)
}
fn run_current_test_in_mount_namespace(test_name: &str) -> Result<(), Box<dyn Error + Send + Sync>> {
let test_binary = std::env::current_exe()?;
let status = Command::new("unshare")
.arg("--mount")
.arg("--propagation")
.arg("private")
.arg("--")
.arg(test_binary)
.arg("--exact")
.arg(test_name)
.arg("--ignored")
.arg("--nocapture")
.env(NAMESPACE_ENV, "1")
.status()?;
if status.success() {
return Ok(());
}
Err(format!("{ENABLE_ENV}=1 requires root or CAP_SYS_ADMIN; unshare exited with status {status}").into())
}
fn replacement_node_log_path(
cluster_temp_dir: &str,
parity: usize,
node_index: usize,
) -> Result<PathBuf, Box<dyn Error + Send + Sync>> {
let log_dir = std::env::var_os(LOG_DIR_ENV)
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from(cluster_temp_dir));
fs::create_dir_all(&log_dir)?;
Ok(log_dir.join(format!("replacement-ec{parity}-node{node_index}-{}.log", std::process::id())))
}
fn payload(len: usize, seed: u8) -> Vec<u8> {
let mut next = seed;
(0..len)
.map(|_| {
let byte = next;
next = next.wrapping_add(31);
byte
})
.collect()
}
fn sha256_hex(data: &[u8]) -> String {
let digest = Sha256::digest(data);
digest.iter().map(|byte| format!("{byte:02x}")).collect()
}
async fn put_object_version(
client: &Client,
bucket: &str,
key: &str,
body: Vec<u8>,
) -> Result<(Option<String>, String), Box<dyn Error + Send + Sync>> {
let digest = sha256_hex(&body);
let output = client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from(body))
.send()
.await?;
Ok((output.version_id().map(str::to_owned), digest))
}
async fn put_multipart_version(
client: &Client,
bucket: &str,
key: &str,
parts: Vec<Vec<u8>>,
) -> Result<(Option<String>, String), Box<dyn Error + Send + Sync>> {
let body = parts.iter().flatten().copied().collect::<Vec<_>>();
let digest = sha256_hex(&body);
let create = client.create_multipart_upload().bucket(bucket).key(key).send().await?;
let upload_id = create
.upload_id()
.ok_or("create_multipart_upload returned no upload id")?
.to_string();
let mut completed_parts = Vec::with_capacity(parts.len());
for (index, part) in parts.into_iter().enumerate() {
let part_number = i32::try_from(index + 1)?;
let uploaded = client
.upload_part()
.bucket(bucket)
.key(key)
.upload_id(&upload_id)
.part_number(part_number)
.body(ByteStream::from(part))
.send()
.await?;
completed_parts.push(
CompletedPart::builder()
.part_number(part_number)
.e_tag(uploaded.e_tag().ok_or("upload_part returned no etag")?)
.build(),
);
}
let completed = client
.complete_multipart_upload()
.bucket(bucket)
.key(key)
.upload_id(upload_id)
.multipart_upload(CompletedMultipartUpload::builder().set_parts(Some(completed_parts)).build())
.send()
.await?;
Ok((completed.version_id().map(str::to_owned), digest))
}
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";
client.create_bucket().bucket(plain_bucket).send().await?;
client.create_bucket().bucket(versioned_bucket).send().await?;
client.create_bucket().bucket(null_bucket).send().await?;
client
.put_bucket_versioning()
.bucket(versioned_bucket)
.versioning_configuration(
VersioningConfiguration::builder()
.status(BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await?;
let mut versions = Vec::new();
let (version_id, body_sha256) =
put_object_version(client, plain_bucket, "objects/single-part.bin", payload(512 * 1024, 1)).await?;
versions.push((plain_bucket, "objects/single-part.bin", version_id, Some(body_sha256)));
let (version_id, body_sha256) = put_multipart_version(
client,
plain_bucket,
"objects/multipart.bin",
vec![payload(5 * 1024 * 1024, 2), payload(1024 * 1024, 3)],
)
.await?;
versions.push((plain_bucket, "objects/multipart.bin", version_id, Some(body_sha256)));
let (version_id, body_sha256) =
put_object_version(client, versioned_bucket, "history/object.bin", payload(512 * 1024, 4)).await?;
versions.push((versioned_bucket, "history/object.bin", version_id, Some(body_sha256)));
let (version_id, body_sha256) =
put_object_version(client, versioned_bucket, "history/object.bin", payload(768 * 1024, 5)).await?;
versions.push((versioned_bucket, "history/object.bin", version_id, Some(body_sha256)));
let deleted = client
.delete_object()
.bucket(versioned_bucket)
.key("history/object.bin")
.send()
.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,
"history/multipart.bin",
vec![payload(5 * 1024 * 1024, 6), payload(2 * 1024 * 1024, 7)],
)
.await?;
versions.push((versioned_bucket, "history/multipart.bin", version_id, Some(body_sha256)));
let (version_id, body_sha256) =
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 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())?;
if !expected.is_complete() {
return Err(format!("baseline census is incomplete for {bucket}/{key}@{version_id:?}: {expected:?}").into());
}
Ok(BaselineVersion {
bucket: bucket.to_string(),
key: key.to_string(),
version_id,
body_sha256,
expected,
})
})
.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")
.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>> {
for version in versions {
let Some(expected_sha256) = &version.body_sha256 else {
continue;
};
let mut request = client.get_object().bucket(&version.bucket).key(&version.key);
if let Some(version_id) = &version.version_id {
request = request.version_id(version_id);
}
let response = request.send().await.map_err(|error| {
format!("body GET failed for {}/{}@{:?}: {error}", version.bucket, version.key, version.version_id)
})?;
let body = response
.body
.collect()
.await
.map_err(|error| {
format!(
"body stream failed for {}/{}@{:?}: {error}",
version.bucket, version.key, version.version_id
)
})?
.into_bytes();
assert_eq!(
sha256_hex(&body),
*expected_sha256,
"body hash changed for {}/{}@{:?}",
version.bucket,
version.key,
version.version_id
);
}
Ok(())
}
async fn replacement_status(
cluster: &RustFSTestClusterEnvironment,
) -> Result<serde_json::Value, Box<dyn Error + Send + Sync>> {
let (status, body) = admin_request(
&cluster.nodes[0].url,
Method::GET,
"/rustfs/admin/v4/heal/replacement-recovery",
None,
&cluster.access_key,
&cluster.secret_key,
)
.await?;
if !status.is_success() {
return Err(format!("replacement recovery status failed: {status} {body}").into());
}
Ok(serde_json::from_str(&body)?)
}
fn replacement_artifact_targets_disk(
artifact: &serde_json::Value,
target_disk: &Path,
) -> Result<bool, Box<dyn Error + Send + Sync>> {
let target = target_disk.to_string_lossy();
Ok(artifact["replacement_targets"]
.as_array()
.ok_or("replacement artifact has no replacement_targets")?
.iter()
.any(|slot| slot.as_str().is_some_and(|slot| slot.contains(target.as_ref()))))
}
fn assert_no_replacement_admission_artifacts(
cluster: &RustFSTestClusterEnvironment,
target_disk: &Path,
) -> Result<(), Box<dyn Error + Send + Sync>> {
for node in &cluster.nodes {
for drive in &node.data_dirs {
let drive = Path::new(drive);
let recovery_dir = drive.join(REPLACEMENT_RECOVERY_DIR);
match fs::read_dir(&recovery_dir) {
Ok(entries) => {
for entry in entries {
let entry = entry?;
let file_name = entry.file_name().to_string_lossy().into_owned();
if file_name.ends_with(REPLACEMENT_INTENT_SUFFIX)
|| file_name.ends_with(REPLACEMENT_COMPLETION_PROOF_SUFFIX)
{
let artifact: serde_json::Value = serde_json::from_slice(&fs::read(entry.path())?)?;
if replacement_artifact_targets_disk(&artifact, target_disk)? {
return Err(format!(
"absent replacement target was admitted before it became ready: {:?}",
entry.path()
)
.into());
}
continue;
}
if file_name.ends_with(RESUME_CHECKPOINT_SUFFIX) {
return Err(format!(
"replacement created checkpoint while target was absent: {:?}",
entry.path()
)
.into());
}
}
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(_) if drive == target_disk => {}
Err(error) => {
return Err(format!("failed to read replacement recovery dir {recovery_dir:?}: {error}").into());
}
}
let bucket_meta_dir = drive.join(".rustfs.sys").join("buckets");
let entries = match fs::read_dir(&bucket_meta_dir) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue,
Err(_) if drive == target_disk => continue,
Err(error) => return Err(format!("failed to read bucket metadata dir {bucket_meta_dir:?}: {error}").into()),
};
for entry in entries {
let entry = entry?;
let file_name = entry.file_name().to_string_lossy().into_owned();
if file_name.ends_with(RESUME_CHECKPOINT_SUFFIX) {
return Err(format!("replacement created checkpoint while target was absent: {:?}", entry.path()).into());
}
}
}
}
let marker = target_disk.join(".rustfs.sys").join("healing.bin");
if marker.exists() {
return Err(format!("replacement created healing marker while target was absent: {marker:?}").into());
}
Ok(())
}
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(
cluster: &RustFSTestClusterEnvironment,
target_disk: &Path,
) -> Result<(), Box<dyn Error + Send + Sync>> {
let status = replacement_status(cluster).await?;
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(());
}
Err(format!(
"live missing replacement target must not have durable recovery records; observed states {states:?} in status {status}"
)
.into())
}
fn target_record_states(status: &serde_json::Value, target_disk: &Path) -> BTreeSet<String> {
let target = target_disk.to_string_lossy();
status["cluster"]["records"]
.as_array()
.into_iter()
.flatten()
.filter_map(|record| {
let state = record["state"].as_str()?;
let target_matches = record["targetSlots"]
.as_array()
.into_iter()
.flatten()
.filter_map(serde_json::Value::as_str)
.any(|slot| slot.contains(target.as_ref()));
target_matches.then(|| state.to_string())
})
.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>(),
Some(rustfs_filemeta::Error::FileVersionNotFound)
)
}
fn incomplete_versions(
target_disk: &Path,
versions: &[BaselineVersion],
) -> Result<BTreeSet<String>, Box<dyn Error + Send + Sync>> {
let mut missing = BTreeSet::new();
for version in versions {
let actual =
match census_object_version_on_disk(target_disk, &version.bucket, &version.key, version.version_id.as_deref()) {
Ok(actual) => actual,
// During replacement recovery, xl.meta may arrive before this
// particular historical version. The generic census helper
// correctly reports that as an error; this progress poll must
// instead wait for the version to be restored.
Err(error) if is_transient_recovery_version_absence(error.as_ref()) => {
missing.insert(format!(
"{}/{}@{:?}: version metadata not yet present on replacement",
version.bucket, version.key, version.version_id
));
continue;
}
Err(error) => return Err(error),
};
if !actual.matches_manifest(&version.expected) {
missing.insert(format!("{}/{}@{:?}: {actual:?}", version.bucket, version.key, version.version_id));
}
}
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,
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,
versions: &[BaselineVersion],
timeout_secs: u64,
) -> Result<(), Box<dyn Error + Send + Sync>> {
let deadline = Instant::now() + Duration::from_secs(timeout_secs);
let mut tick = interval(Duration::from_secs(1));
loop {
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 status remained completed across two incomplete physical censuses: {confirmed_missing:?}"
)
.into());
}
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}"
)
.into());
}
tick.tick().await;
}
}
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(());
}
if std::env::var_os(NAMESPACE_ENV).is_none() {
return run_current_test_in_mount_namespace(test_name);
}
verify_isolated_mount_namespace()?;
let mut mount_ns = MountNamespaceGuard::new()?;
let mut cluster = RustFSTestClusterEnvironment::with_topology(ClusterTopology::single_pool_multidrive(3, 4)).await?;
for node_index in 0..cluster.nodes.len() {
let node_log_path = replacement_node_log_path(&cluster.temp_dir, parity, node_index)?;
cluster.set_node_capture_log_path(node_index, node_log_path.to_string_lossy())?;
}
let target_disk = PathBuf::from(&cluster.nodes[TARGET_NODE].data_dirs[TARGET_DRIVE]);
// The blank target uses a temporary zram block device, so the
// replacement readiness fence sees no root or sibling alias.
cluster.extra_env.retain(|(key, _)| key != "RUSTFS_UNSAFE_BYPASS_DISK_CHECK");
let image_root = PathBuf::from(&cluster.temp_dir).join("replacement-block-images");
let mut target_mount = None;
for (node_index, node) in cluster.nodes.iter().enumerate() {
for (drive_index, drive) in node.data_dirs.iter().enumerate() {
let drive = Path::new(drive);
if node_index == TARGET_NODE && drive_index == TARGET_DRIVE {
target_mount = Some(FaultableBlockMount::mount(
drive,
&image_root,
&format!("p{parity}_node{node_index}_drive{drive_index}"),
)?);
} else {
mount_ns.mount_tmpfs(drive, &format!("rustfs-e2e-p{parity}-node{node_index}-drive{drive_index}"))?;
}
}
}
let mut target_mount = target_mount.ok_or("target drive was not mounted with the faultable block fixture")?;
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");
cluster.set_env("RUSTFS_HEAL_INTERVAL_SECS", "1");
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 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)
.await
.map_err(|error| format!("pre-fault body verification failed: {error}"))?;
target_mount
.make_unavailable()
.map_err(|error| format!("failed to install the dm-error target: {error}"))?;
target_mount
.verify_raw_io_is_unavailable()
.map_err(|error| format!("dm-error target was not proven by a direct raw read: {error}"))?;
assert_no_replacement_status_records(&cluster, &target_disk)
.await
.map_err(|error| format!("live-fault replacement status check failed: {error}"))?;
assert_no_replacement_admission_artifacts(&cluster, &target_disk)
.map_err(|error| format!("live-fault replacement artifact check failed: {error}"))?;
cluster.stop_node_gracefully(TARGET_NODE).await?;
target_mount.cleanup()?;
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(),
versions.len(),
"blank replacement must start without any baseline version shards"
);
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 = 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 {
info!(%stop_error, "replacement target stop failed while preserving recovery failure");
}
if let Err(cleanup_error) = replacement_cleanup_result {
info!(%cleanup_error, "replacement zram cleanup failed while preserving recovery failure");
}
return Err(error);
}
stop_result?;
replacement_cleanup_result?;
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
);
}
#[test]
fn recovery_census_only_treats_missing_version_as_transient() {
let missing_version: Box<dyn Error + Send + Sync> = Box::new(rustfs_filemeta::Error::FileVersionNotFound);
let missing_file: Box<dyn Error + Send + Sync> = Box::new(rustfs_filemeta::Error::FileNotFound);
assert!(is_transient_recovery_version_absence(missing_version.as_ref()));
assert!(!is_transient_recovery_version_absence(missing_file.as_ref()));
}
#[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 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");
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")]
#[ignore = "requires Linux root/CAP_SYS_ADMIN and RUSTFS_PRIVILEGED_REPLACEMENT_E2E=1"]
async fn test_privileged_3x4_auto_replacement_rebuilds_ec8_plus_4_without_admin_heal()
-> Result<(), Box<dyn Error + Send + Sync>> {
run_replacement_e2e(
4,
"replacement_privileged_e2e_test::tests::test_privileged_3x4_auto_replacement_rebuilds_ec8_plus_4_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_rebuilds_ec6_plus_6_without_admin_heal()
-> Result<(), Box<dyn Error + Send + Sync>> {
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
}
}