fix(heal): preserve remote corruption and reject incomplete deep scans (#8032)

* fix(heal): reject truncated encoded shards

* test(heal): cover truncated shards across integrity modes (#8035)

Co-authored-by: zhi22915 <qiuzgang@gmail.com>

* fix(heal): preserve remote corruption and reject incomplete deep scans

* test(admin): tolerate config lock timeout while polling

* fix(test): use admin storage error facade

---------

Co-authored-by: Hauser <housemecn@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
Co-authored-by: overtrue <anzhengchao@gmail.com>
This commit is contained in:
cxymds
2026-09-21 10:33:03 +08:00
committed by GitHub
parent 42dd76d3b2
commit a7c875b48f
7 changed files with 667 additions and 85 deletions
+17 -18
View File
@@ -244,13 +244,16 @@ impl StdError for ConditionalFileNotCommittedError {
}
}
fn classify_internode_missing_error(error: &InternodeHttpError) -> Option<DiskError> {
fn classify_internode_disk_error(error: &InternodeHttpError) -> Option<DiskError> {
if error.is_remote_file_not_found() {
return Some(DiskError::FileNotFound);
}
if error.is_remote_volume_not_found() {
return Some(DiskError::VolumeNotFound);
}
if error.is_remote_file_corrupt() {
return Some(DiskError::FileCorrupt);
}
None
}
@@ -526,13 +529,10 @@ fn io_error_chain_contains_kind(io_error: &std::io::Error, kind: std::io::ErrorK
impl From<std::io::Error> for DiskError {
fn from(e: std::io::Error) -> Self {
if let Some(error) = e.get_ref().and_then(|source| source.downcast_ref::<InternodeHttpError>()) {
if error.is_remote_file_not_found() {
return DiskError::FileNotFound;
}
if error.is_remote_volume_not_found() {
return DiskError::VolumeNotFound;
}
if let Some(error) = e.get_ref().and_then(|source| source.downcast_ref::<InternodeHttpError>())
&& let Some(classified) = classify_internode_disk_error(error)
{
return classified;
}
let e = match e.downcast::<TerminalReadError>() {
Ok(terminal_error) => {
@@ -541,7 +541,7 @@ impl From<std::io::Error> for DiskError {
&& let Some(internode_error) = io_error
.get_ref()
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
&& let Some(classified) = classify_internode_missing_error(internode_error)
&& let Some(classified) = classify_internode_disk_error(internode_error)
{
return classified;
}
@@ -963,6 +963,7 @@ mod tests {
for (remote_error, expected) in [
(rustfs_rio::new_test_remote_file_not_found_http_io_error(), DiskError::FileNotFound),
(rustfs_rio::new_test_remote_volume_not_found_http_io_error(), DiskError::VolumeNotFound),
(rustfs_rio::new_test_remote_file_corrupt_http_io_error(), DiskError::FileCorrupt),
] {
let wrapped = terminal_read_error_to_io(DiskError::Io(remote_error));
assert_eq!(DiskError::from(wrapped), expected);
@@ -1533,25 +1534,23 @@ mod tests {
}
#[test]
fn test_internode_missing_errors_preserve_disk_error_types() {
fn test_internode_disk_errors_preserve_disk_error_types() {
let file_missing = DiskError::from(rustfs_rio::new_test_remote_file_not_found_http_io_error());
let volume_missing = DiskError::from(rustfs_rio::new_test_remote_volume_not_found_http_io_error());
let file_corrupt = DiskError::from(rustfs_rio::new_test_remote_file_corrupt_http_io_error());
let unmarked_server_error = DiskError::from(rustfs_rio::new_test_internode_http_io_error(
rustfs_rio::InternodeHttpErrorKind::HttpStatus(http::StatusCode::INTERNAL_SERVER_ERROR),
));
assert_eq!(file_missing, DiskError::FileNotFound);
assert_eq!(volume_missing, DiskError::VolumeNotFound);
assert_eq!(file_corrupt, DiskError::FileCorrupt);
assert!(matches!(unmarked_server_error, DiskError::Io(_)));
for missing in [file_missing, volume_missing] {
assert_eq!(missing.clone(), missing);
for error in [file_missing, volume_missing, file_corrupt] {
assert_eq!(error.clone(), error);
assert_eq!(
crate::disk::error_reduce::reduce_write_quorum_errs(
&[Some(missing.clone()), Some(missing.clone()), None],
&[],
2
),
Some(missing)
crate::disk::error_reduce::reduce_write_quorum_errs(&[Some(error.clone()), Some(error.clone()), None], &[], 2),
Some(error)
);
}
}
@@ -587,10 +587,8 @@ pub(crate) async fn verify_deep_parts(
return Err(DiskError::FileCorrupt);
}
let part_status = statuses.get_mut(&part_index).ok_or(DiskError::FileCorrupt)?;
// Deep verification reads the complete encoded shard. `shard_file_offset`
// is a range-read helper and can stop before the final encoded byte for
// a partial data block, which would let a one-byte tail truncation pass
// as healthy. Use the physical shard length for the integrity proof.
// Deep verification must consume every encoded stripe, including a
// partial final stripe, before certifying the shard's integrity.
let length = usize::try_from(erasure.shard_file_size(part.size as i64)).map_err(|_| DiskError::FileCorrupt)?;
let mut readers = Vec::with_capacity(disks.len());
for (index, disk) in disks.iter().enumerate() {
@@ -635,6 +633,9 @@ pub(crate) async fn verify_deep_parts(
part_status[index] = CHECK_PART_FILE_CORRUPT;
}
Err(error) => {
if !matches!(error, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::FileCorrupt) {
return Err(error);
}
readers.push(None);
part_status[index] = conv_part_err_to_int(&Some(error));
}
+70
View File
@@ -6473,6 +6473,28 @@ async fn verify_inline_part_bitrot(meta: &FileInfo) -> disk::error::Result<()> {
.map_err(|_| DiskError::FileCorrupt)
}
fn validate_deep_scan_results(results: &[usize], expected_parts: usize) -> disk::error::Result<()> {
if results.len() != expected_parts {
return Err(DiskError::other(format!(
"incomplete deep scan result: expected {expected_parts} parts, received {}",
results.len()
)));
}
for (part, status) in results.iter().enumerate() {
match *status {
CHECK_PART_SUCCESS | CHECK_PART_FILE_NOT_FOUND | CHECK_PART_FILE_CORRUPT => {}
CHECK_PART_DISK_NOT_FOUND => return Err(DiskError::DiskNotFound),
crate::disk::CHECK_PART_VOLUME_NOT_FOUND => return Err(DiskError::VolumeNotFound),
_ => {
return Err(DiskError::other(format!(
"incomplete deep scan result: part {part} has unverified status {status}"
)));
}
}
}
Ok(())
}
/// disks_with_all_partsv2 is a corrected version based on Go implementation.
/// It sets partsMetadata and onlineDisks when xl.meta is inexistant/corrupted or outdated.
/// It also checks if the status of each part (corrupted, missing, ok) in each drive.
@@ -6693,9 +6715,13 @@ async fn disks_with_all_parts(
// it needs healing too.
match disk.verify_file(bucket, object, meta).await {
Ok(v) => {
validate_deep_scan_results(&v.results, latest_meta.parts.len())?;
verify_resp = v;
}
Err(err) => {
if !matches!(err, DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::FileCorrupt) {
return Err(err);
}
debug!(
event = EVENT_SET_DISK_HEAL,
component = LOG_COMPONENT_ECSTORE,
@@ -11109,6 +11135,50 @@ mod tests {
assert_eq!(conv_part_err_to_int(&Some(other_err)), CHECK_PART_UNKNOWN); // Other errors should return UNKNOWN, not SUCCESS
}
#[test]
fn deep_scan_results_accept_only_complete_verified_or_repairable_parts() {
for (results, expected_parts) in [
(vec![], 0),
(vec![CHECK_PART_SUCCESS], 1),
(vec![CHECK_PART_FILE_NOT_FOUND], 1),
(vec![CHECK_PART_FILE_CORRUPT], 1),
(vec![CHECK_PART_SUCCESS, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_FILE_CORRUPT], 3),
] {
validate_deep_scan_results(&results, expected_parts)
.expect("complete results must allow healthy or repairable parts");
}
}
#[test]
fn deep_scan_results_reject_unknown_and_incomplete_observations() {
for (results, expected_parts) in [
(vec![CHECK_PART_UNKNOWN], 1),
(vec![usize::MAX], 1),
(vec![CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN], 2),
(vec![], 1),
(vec![CHECK_PART_SUCCESS], 2),
(vec![CHECK_PART_SUCCESS, CHECK_PART_SUCCESS], 1),
(vec![CHECK_PART_SUCCESS], 0),
] {
let error =
validate_deep_scan_results(&results, expected_parts).expect_err("unverified parts must not become healthy");
assert!(matches!(error, DiskError::Io(_)), "an incomplete observation is not proven corruption");
assert!(error.to_string().contains("incomplete deep scan result"));
}
}
#[test]
fn deep_scan_results_preserve_disk_and_volume_failures() {
for (status, expected_error) in [
(CHECK_PART_DISK_NOT_FOUND, DiskError::DiskNotFound),
(CHECK_PART_VOLUME_NOT_FOUND, DiskError::VolumeNotFound),
] {
let error = validate_deep_scan_results(&[CHECK_PART_SUCCESS, status], 2)
.expect_err("an unavailable part cannot certify a healthy disk");
assert_eq!(error, expected_error);
}
}
#[test]
fn test_has_part_err() {
// Test checking for part errors
+444 -54
View File
@@ -4378,63 +4378,243 @@ mod heal_result_report_tests {
);
}
#[tokio::test]
async fn deep_heal_rebuilds_h2_near_tail_truncated_part() {
let (temp_dirs, disks, set) = hermetic_set_disks_for_pool_with_default_parity_isolated(16, 0, 4).await;
let bucket = "deep-heal-h2-near-tail-truncation";
let object = "object.bin";
for disk in &disks {
disk.make_volume(bucket).await.expect("bucket volume should be created");
#[derive(Debug)]
struct FailedIntegrityReadTransport {
error: DiskError,
reads: std::sync::atomic::AtomicUsize,
}
#[async_trait::async_trait]
impl crate::cluster::rpc::internode_data_transport::InternodeDataTransport for FailedIntegrityReadTransport {
async fn open_read(
&self,
_request: crate::cluster::rpc::internode_data_transport::ReadStreamRequest,
) -> crate::disk::error::Result<crate::disk::FileReader> {
self.reads.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let error = match &self.error {
DiskError::FileCorrupt => DiskError::from(rustfs_rio::new_test_remote_file_corrupt_http_io_error()),
DiskError::FileNotFound => DiskError::from(rustfs_rio::new_test_remote_file_not_found_http_io_error()),
DiskError::VolumeNotFound => DiskError::from(rustfs_rio::new_test_remote_volume_not_found_http_io_error()),
error => error.clone(),
};
Err(error)
}
let expected_payload = vec![0x6d; 5 * 1024 * 1024 + 123];
async fn open_write(
&self,
_request: crate::cluster::rpc::internode_data_transport::WriteStreamRequest,
) -> crate::disk::error::Result<crate::disk::FileWriter> {
panic!("an integrity scan must not write remote data");
}
async fn open_walk_dir(
&self,
_request: crate::cluster::rpc::internode_data_transport::WalkDirStreamRequest,
) -> crate::disk::error::Result<crate::disk::FileReader> {
panic!("an object integrity scan must not walk remote directories");
}
fn name(&self) -> &'static str {
"failed-integrity-read"
}
fn capabilities(&self) -> crate::cluster::rpc::internode_data_transport::InternodeDataTransportCapabilities {
crate::cluster::rpc::internode_data_transport::InternodeDataTransportCapabilities::tcp_http()
}
}
async fn assert_remote_integrity_read_failure(expected_error: DiskError) {
use crate::disk::{CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN};
use crate::object_api::ShardIntegrityWriteMode;
let (dirs, disks, set) = hermetic_set_disks_for_pool_with_default_parity_isolated(4, 0, 2).await;
let bucket = "remote-integrity-read-failure";
let object = "object.bin";
for disk in &disks {
disk.make_volume(bucket).await.expect("create fixture bucket");
}
set.put_object(
bucket,
object,
&mut PutObjReader::from_vec(expected_payload.clone()),
&mut PutObjReader::from_vec(vec![0x6d; 1024 * 1024 + 123]),
&ObjectOptions {
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Protected),
no_lock: true,
..Default::default()
},
)
.await
.expect("source object should be written before shard truncation");
let source = disks[2]
.read_version("", bucket, object, "", &ReadOptions::default())
.await
.expect("source metadata should be readable");
let data_dir = source.data_dir.expect("non-inline source should have a data directory");
let truncated_part = temp_dirs[1]
.path()
.join(bucket)
.join(object)
.join(data_dir.to_string())
.join("part.1");
let original_len = tokio::fs::metadata(&truncated_part)
.await
.expect("target shard should exist")
.len();
assert!(original_len > 1, "test shard must be large enough to truncate");
let file = tokio::fs::OpenOptions::new()
.write(true)
.open(&truncated_part)
.await
.expect("target shard should be writable");
file.set_len(original_len - 1)
.await
.expect("target shard should be truncated");
.expect("write protected external fixture");
let mut reader = set
.get_object_reader(bucket, object, None, Default::default(), &ObjectOptions::default())
.await
.expect("GET should remain readable after a one-byte shard truncation");
let mut read_back = Vec::new();
tokio::io::copy(&mut reader, &mut read_back)
.await
.expect("GET should reconstruct the truncated shard through EC");
assert_eq!(read_back, expected_payload, "EC GET must preserve the object bytes");
let mut files = Vec::new();
let mut snapshots = Vec::new();
for (dir, disk) in dirs.iter().zip(&disks) {
let file = disk
.read_version("", bucket, object, "", &ReadOptions::default())
.await
.expect("read fixture metadata");
assert!(file.data.is_none(), "the remote test must open an external shard");
let object_dir = dir.path().join(bucket).join(object);
let part_path = object_dir
.join(file.data_dir.expect("external shard directory").to_string())
.join("part.1");
for path in [part_path, object_dir.join("xl.meta")] {
let bytes = tokio::fs::read(&path).await.expect("snapshot fixture before scan");
snapshots.push((path, bytes));
}
files.push(file);
}
let transport = Arc::new(FailedIntegrityReadTransport {
error: expected_error.clone(),
reads: std::sync::atomic::AtomicUsize::new(0),
});
let endpoint = Endpoint {
url: url::Url::parse("http://remote-integrity.invalid:9000/data/disk0").expect("fixture endpoint"),
is_local: false,
pool_idx: 0,
set_idx: 0,
disk_idx: 0,
};
let remote = crate::cluster::rpc::RemoteDisk::new(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
transport.clone(),
)
.await
.expect("construct remote disk with failing data transport");
let mut verification_disks: Vec<_> = disks.iter().cloned().map(Some).collect();
verification_disks[0] = Some(Arc::new(crate::disk::Disk::Remote(Box::new(remote))));
let mut statuses = [(0, vec![CHECK_PART_UNKNOWN; disks.len()])].into_iter().collect();
let result = crate::io_support::shard_integrity::verify_deep_parts(
&files,
&verification_disks,
&files[0],
bucket,
object,
&mut statuses,
)
.await;
assert!(transport.reads.load(std::sync::atomic::Ordering::Relaxed) > 0);
if matches!(expected_error, DiskError::FileCorrupt | DiskError::FileNotFound) {
result.expect("typed remote payload damage must identify a repairable shard");
let expected_status = if expected_error == DiskError::FileCorrupt {
CHECK_PART_FILE_CORRUPT
} else {
CHECK_PART_FILE_NOT_FOUND
};
assert_eq!(statuses[&0][0], expected_status);
assert!(statuses[&0][1..].iter().all(|status| *status == CHECK_PART_SUCCESS));
let (heal, metadata, reason) = super::should_heal_object_on_disk(&None, &[statuses[&0][0]], &files[0], &files[0]);
assert!(heal, "remote payload damage must schedule repair: {expected_error:?}");
assert!(!metadata);
let expected_reason = if expected_error == DiskError::FileCorrupt {
DiskError::FileCorrupt
} else {
DiskError::PartMissingOrCorrupt
};
assert_eq!(reason, Some(expected_reason));
} else {
let error = result.expect_err("unavailable payload storage must abort health verification");
if expected_error.is_internode_http_status(500) {
assert!(error.is_internode_http_status(500), "preserve the original transport error: {error:?}");
} else {
assert_eq!(error, expected_error, "preserve the original storage failure");
}
assert_eq!(statuses[&0][0], CHECK_PART_UNKNOWN, "failed verification cannot certify the shard");
}
for (path, before) in snapshots {
assert_eq!(tokio::fs::read(path).await.expect("read fixture after scan"), before);
}
}
let (result, error) = set
#[tokio::test]
async fn deep_scan_remote_payload_damage_schedules_repair() {
for error in [DiskError::FileCorrupt, DiskError::FileNotFound] {
assert_remote_integrity_read_failure(error).await;
}
}
#[tokio::test]
async fn deep_scan_remote_storage_failure_cannot_certify_healthy() {
for error in [
DiskError::from(rustfs_rio::new_test_internode_http_io_error(
rustfs_rio::InternodeHttpErrorKind::HttpStatus(reqwest::StatusCode::INTERNAL_SERVER_ERROR),
)),
DiskError::DiskNotFound,
DiskError::VolumeNotFound,
DiskError::FileAccessDenied,
] {
assert_remote_integrity_read_failure(error).await;
}
}
#[cfg(unix)]
#[tokio::test]
async fn deep_heal_legacy_verification_error_cannot_certify_healthy() {
use crate::object_api::ShardIntegrityWriteMode;
let (dirs, disks, set) = hermetic_set_disks_isolated(4).await;
let bucket = "legacy-deep-scan-verification-error";
let object = "object.bin";
for disk in &disks {
disk.make_volume(bucket).await.expect("create legacy fixture bucket");
}
set.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![0x6d; 1024 * 1024 + 123]),
&ObjectOptions {
shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Legacy),
no_lock: true,
..Default::default()
},
)
.await
.expect("write external legacy fixture");
let mut files = Vec::new();
let mut part_paths = Vec::new();
let mut preserved = Vec::new();
for (index, (dir, disk)) in dirs.iter().zip(&disks).enumerate() {
let file = disk
.read_version("", bucket, object, "", &ReadOptions::default())
.await
.expect("read legacy fixture metadata");
assert!(file.data.is_none(), "legacy verification must open the physical shard");
assert!(file.parts.iter().all(|part| part.integrity.is_none()));
let object_dir = dir.path().join(bucket).join(object);
let part = object_dir
.join(file.data_dir.expect("external legacy shard directory").to_string())
.join("part.1");
let metadata = object_dir.join("xl.meta");
let bytes = tokio::fs::read(&metadata).await.expect("snapshot legacy metadata");
preserved.push((metadata, bytes));
if index != 0 {
let bytes = tokio::fs::read(&part).await.expect("snapshot healthy legacy shard");
preserved.push((part.clone(), bytes));
}
part_paths.push(part);
files.push(file);
}
let broken_part = &part_paths[0];
tokio::fs::remove_file(broken_part)
.await
.expect("remove fixture shard before symlink fault");
std::os::unix::fs::symlink("part.1", broken_part).expect("create self-referential shard symlink");
let open_error = tokio::fs::File::open(broken_part)
.await
.expect_err("self-referential symlink must fail to open");
assert_eq!(open_error.raw_os_error(), Some(libc::ELOOP));
let verification_error = disks[0]
.verify_file(bucket, object, &files[0])
.await
.expect_err("legacy disk verifier must reject the symlink path");
assert_eq!(verification_error, DiskError::InvalidPath);
let error = set
.heal_object(
bucket,
object,
@@ -4446,19 +4626,229 @@ mod heal_result_report_tests {
},
)
.await
.expect("deep heal should finish after a one-byte H2 tail truncation");
assert!(error.is_none(), "deep heal should recover the H2 truncated shard: {error:?}");
assert_eq!(result.drives_healed(), Some(1));
assert_eq!(result.before.drives[1].state, DriveState::Corrupt.to_string());
.expect_err("an unverified legacy part must fail the complete heal operation");
assert_eq!(error, DiskError::InvalidPath, "complete heal must preserve the verifier's path error");
assert_eq!(
tokio::fs::metadata(&truncated_part)
tokio::fs::read_link(broken_part)
.await
.expect("repaired H2 shard should exist")
.len(),
original_len,
"deep heal must restore the complete H2 shard length"
.expect("failed heal must preserve the symlink fault"),
std::path::Path::new("part.1")
);
for (path, before) in preserved {
assert_eq!(tokio::fs::read(path).await.expect("read unchanged legacy fixture"), before);
}
}
async fn assert_h2_truncated_shards(
mode: crate::object_api::ShardIntegrityWriteMode,
coding_indexes: &[usize],
retained_len: fn(usize, usize) -> usize,
) {
use crate::object_api::ShardIntegrityWriteMode;
let (temp_dirs, disks, set) = hermetic_set_disks_for_pool_with_default_parity_isolated(16, 0, 4).await;
let bucket = "deep-heal-h2-truncation";
let object = "object.bin";
for disk in &disks {
disk.make_volume(bucket).await.expect("create bucket volume");
}
let expected_payload: Vec<_> = (0..5 * 1024 * 1024 + 123)
.map(|index| u8::try_from((index / 4096 + index) % 251).expect("payload byte fits"))
.collect();
set.put_object(
bucket,
object,
&mut PutObjReader::from_vec(expected_payload.clone()),
&ObjectOptions {
shard_integrity_write_mode: Some(mode),
no_lock: true,
..Default::default()
},
)
.await
.expect("commit all source shards");
let mut files = Vec::new();
let mut paths = Vec::new();
let mut original_shards = Vec::new();
let mut original_metadata = Vec::new();
let mut damaged = Vec::new();
for (slot, disk) in disks.iter().enumerate() {
let source = disk
.read_version("", bucket, object, "", &ReadOptions::default())
.await
.expect("read source metadata");
assert_eq!((source.erasure.data_blocks, source.erasure.parity_blocks), (12, 4));
assert_eq!(source.parts[0].integrity.is_some(), mode == ShardIntegrityWriteMode::Protected);
let directory = temp_dirs[slot].path().join(bucket).join(object);
let path = directory
.join(source.data_dir.expect("external shard data directory").to_string())
.join("part.1");
let original = tokio::fs::read(&path).await.expect("read original shard");
original_metadata.push(tokio::fs::read(directory.join("xl.meta")).await.expect("read xl.meta"));
if coding_indexes.contains(&source.erasure.index) {
let erasure = crate::erasure::coding::Erasure::try_new_with_options(
12,
4,
source.erasure.block_size,
source.uses_legacy_checksum,
)
.expect("valid erasure layout");
let stripe_size = erasure.shard_size() + source.erasure.get_checksum_info(1).algorithm.size();
let keep = retained_len(original.len(), stripe_size);
assert!(keep < original.len(), "fixture must truncate the shard");
tokio::fs::OpenOptions::new()
.write(true)
.open(&path)
.await
.expect("open target shard")
.set_len(u64::try_from(keep).expect("shard length fits"))
.await
.expect("truncate target shard");
damaged.push((slot, keep));
}
paths.push(path);
original_shards.push(original);
files.push(source);
}
assert_eq!(damaged.len(), coding_indexes.len());
if mode == ShardIntegrityWriteMode::Protected {
let mut statuses = [(0, vec![crate::disk::CHECK_PART_UNKNOWN; disks.len()])]
.into_iter()
.collect();
let online_disks: Vec<_> = disks.iter().cloned().map(Some).collect();
crate::io_support::shard_integrity::verify_deep_parts(
&files,
&online_disks,
&files[0],
bucket,
object,
&mut statuses,
)
.await
.expect("verify protected shards directly");
for (slot, status) in statuses[&0].iter().enumerate() {
let expected = if damaged.iter().any(|(index, _)| *index == slot) {
crate::disk::CHECK_PART_FILE_CORRUPT
} else {
crate::disk::CHECK_PART_SUCCESS
};
assert_eq!(*status, expected, "protected shard {slot}");
}
}
let opts = HealOpts {
no_lock: true,
scan_mode: HealScanMode::Deep,
..Default::default()
};
let (dry_run, error) = set
.heal_object(bucket, object, "", &HealOpts { dry_run: true, ..opts })
.await
.expect("dry-run should report corruption");
assert!(error.is_none(), "dry-run error: {error:?}");
assert_eq!(dry_run.drives_healed(), Some(0));
for (slot, keep) in &damaged {
assert_eq!(dry_run.after.drives[*slot].state, DriveState::Corrupt.to_string());
assert_eq!(
tokio::fs::read(&paths[*slot]).await.expect("dry-run preserved shard"),
original_shards[*slot][..*keep]
);
}
let recoverable = damaged.len() <= 4;
if recoverable {
let mut reader = set
.get_object_reader(bucket, object, None, Default::default(), &ObjectOptions::default())
.await
.expect("degraded GET should remain readable");
let mut read_back = Vec::new();
tokio::io::copy(&mut reader, &mut read_back)
.await
.expect("read complete degraded GET");
assert_eq!(read_back, expected_payload);
}
let (result, error) = set
.heal_object(bucket, object, "", &opts)
.await
.expect("heal should report outcome");
if recoverable {
assert!(error.is_none(), "heal error: {error:?}");
assert_eq!(result.drives_healed(), Some(damaged.len()));
for (slot, _) in &damaged {
assert_eq!(result.before.drives[*slot].state, DriveState::Corrupt.to_string());
}
} else {
assert_eq!(error, Some(DiskError::ErasureReadQuorum));
assert_eq!(result.drives_healed(), Some(0));
}
for (slot, path) in paths.iter().enumerate() {
let expected = if !recoverable && let Some((_, keep)) = damaged.iter().find(|(index, _)| *index == slot) {
&original_shards[slot][..*keep]
} else {
original_shards[slot].as_slice()
};
assert_eq!(
tokio::fs::read(path).await.expect("read shard after heal"),
expected,
"physical shard {slot}"
);
if !recoverable {
let metadata_path = temp_dirs[slot].path().join(bucket).join(object).join("xl.meta");
assert_eq!(tokio::fs::read(metadata_path).await.expect("preserved metadata"), original_metadata[slot]);
}
}
if recoverable {
let mut reader = set
.get_object_reader(bucket, object, None, Default::default(), &ObjectOptions::default())
.await
.expect("GET after repair");
let mut read_back = Vec::new();
tokio::io::copy(&mut reader, &mut read_back)
.await
.expect("read repaired object");
assert_eq!(read_back, expected_payload);
let (healthy, error) = set.heal_object(bucket, object, "", &opts).await.expect("second deep scan");
assert!(error.is_none());
assert_eq!(healthy.drives_healed(), Some(0), "repair must converge");
assert_eq!(healthy.integrity_verified, mode == ShardIntegrityWriteMode::Protected);
}
}
#[tokio::test]
async fn deep_heal_rebuilds_h2_near_tail_truncated_part() {
use crate::object_api::ShardIntegrityWriteMode::{Legacy, Protected};
for mode in [Legacy, Protected] {
for coding_index in [1, 13] {
assert_h2_truncated_shards(mode, &[coding_index], |length, _| length - 1).await;
}
}
}
#[tokio::test]
async fn deep_heal_rebuilds_h2_half_truncated_part() {
use crate::object_api::ShardIntegrityWriteMode::{Legacy, Protected};
for mode in [Legacy, Protected] {
for coding_index in [1, 13] {
assert_h2_truncated_shards(mode, &[coding_index], |length, _| length / 2).await;
}
}
}
#[tokio::test]
async fn deep_heal_rebuilds_h2_stripe_boundary_truncated_parts_at_read_quorum() {
use crate::object_api::ShardIntegrityWriteMode::{Legacy, Protected};
for mode in [Legacy, Protected] {
assert_h2_truncated_shards(mode, &[1, 2, 13, 16], |_, stripe| stripe).await;
}
}
#[tokio::test]
async fn deep_heal_rejects_h2_truncated_parts_below_read_quorum() {
use crate::object_api::ShardIntegrityWriteMode::{Legacy, Protected};
for mode in [Legacy, Protected] {
assert_h2_truncated_shards(mode, &[1, 2, 3, 13, 16], |length, _| length / 2).await;
}
}
#[tokio::test]
+112
View File
@@ -57,6 +57,7 @@ const EXCESSIVE_EMPTY_CHUNKS_ERROR: &str = "HTTP body returned too many empty ch
pub const INTERNODE_DISK_ERROR_HEADER: &str = "x-rustfs-disk-error";
pub const INTERNODE_FILE_NOT_FOUND: &str = "file-not-found";
pub const INTERNODE_VOLUME_NOT_FOUND: &str = "volume-not-found";
pub const INTERNODE_FILE_CORRUPT: &str = "file-corrupt";
#[derive(Debug, Clone, Copy, Eq, PartialEq)]
pub enum InternodeHttpErrorKind {
@@ -185,6 +186,7 @@ pub struct InternodeHttpError {
enum RemoteDiskErrorKind {
FileNotFound,
VolumeNotFound,
FileCorrupt,
}
impl std::fmt::Debug for InternodeHttpError {
@@ -215,6 +217,10 @@ impl InternodeHttpError {
self.remote_disk_error == Some(RemoteDiskErrorKind::VolumeNotFound)
}
pub fn is_remote_file_corrupt(&self) -> bool {
self.remote_disk_error == Some(RemoteDiskErrorKind::FileCorrupt)
}
fn new(kind: InternodeHttpErrorKind, context: InternodeHttpRequestContext) -> Self {
Self {
kind,
@@ -328,6 +334,20 @@ pub fn new_test_remote_volume_not_found_http_io_error() -> io::Error {
.into_io_error()
}
#[doc(hidden)]
pub fn new_test_remote_file_corrupt_http_io_error() -> io::Error {
InternodeHttpError::with_remote_disk_error(
InternodeHttpErrorKind::HttpStatus(reqwest::StatusCode::INTERNAL_SERVER_ERROR),
InternodeHttpRequestContext {
method: "GET".to_string(),
target: READ_FILE_STREAM_PATH.to_string(),
operation: Some(INTERNODE_OPERATION_READ_FILE_STREAM),
},
RemoteDiskErrorKind::FileCorrupt,
)
.into_io_error()
}
fn add_root_certificates_from_der(builder: reqwest::ClientBuilder, certs_der: &[Vec<u8>]) -> reqwest::ClientBuilder {
let mut b = builder;
for der in certs_der {
@@ -850,6 +870,7 @@ fn classify_http_response(
let remote_disk_error = match headers.get(INTERNODE_DISK_ERROR_HEADER).and_then(|value| value.to_str().ok()) {
Some(INTERNODE_FILE_NOT_FOUND) => Some(RemoteDiskErrorKind::FileNotFound),
Some(INTERNODE_VOLUME_NOT_FOUND) => Some(RemoteDiskErrorKind::VolumeNotFound),
Some(INTERNODE_FILE_CORRUPT) => Some(RemoteDiskErrorKind::FileCorrupt),
_ => None,
};
ClassifiedHttpResponse { kind, remote_disk_error }
@@ -2047,6 +2068,7 @@ mod tests {
assert_eq!(INTERNODE_DISK_ERROR_HEADER, "x-rustfs-disk-error");
assert_eq!(INTERNODE_FILE_NOT_FOUND, "file-not-found");
assert_eq!(INTERNODE_VOLUME_NOT_FOUND, "volume-not-found");
assert_eq!(INTERNODE_FILE_CORRUPT, "file-corrupt");
}
#[test]
@@ -2090,6 +2112,48 @@ mod tests {
assert!(wrong_operation.remote_disk_error.is_none());
}
#[test]
fn classify_http_response_scopes_corrupt_disk_error_to_read_failures() {
let mut headers = HeaderMap::new();
headers.insert(INTERNODE_DISK_ERROR_HEADER, INTERNODE_FILE_CORRUPT.parse().expect("valid error token"));
let classified = classify_http_response(
reqwest::StatusCode::INTERNAL_SERVER_ERROR,
&headers,
Some(INTERNODE_OPERATION_READ_FILE_STREAM),
);
assert_eq!(classified.remote_disk_error, Some(RemoteDiskErrorKind::FileCorrupt));
for status in [reqwest::StatusCode::NOT_FOUND, reqwest::StatusCode::SERVICE_UNAVAILABLE] {
assert!(
classify_http_response(status, &headers, Some(INTERNODE_OPERATION_READ_FILE_STREAM))
.remote_disk_error
.is_none()
);
}
for operation in [
None,
Some(INTERNODE_OPERATION_WALK_DIR),
Some(INTERNODE_OPERATION_PUT_FILE_STREAM),
] {
assert!(
classify_http_response(reqwest::StatusCode::INTERNAL_SERVER_ERROR, &headers, operation)
.remote_disk_error
.is_none()
);
}
headers.insert(INTERNODE_DISK_ERROR_HEADER, "unknown-disk-error".parse().expect("valid unknown token"));
for headers in [&headers, &HeaderMap::new()] {
assert!(
classify_http_response(
reqwest::StatusCode::INTERNAL_SERVER_ERROR,
headers,
Some(INTERNODE_OPERATION_READ_FILE_STREAM),
)
.remote_disk_error
.is_none()
);
}
}
#[derive(Clone, Default)]
struct TestState {
head_count: Arc<AtomicUsize>,
@@ -2158,6 +2222,16 @@ mod tests {
let addr = listener.local_addr().expect("listener local address should be available");
let app = Router::new()
.route("/stream", get(get_stream).head(reject_head).put(accept_put))
.route(
READ_FILE_STREAM_PATH,
get(|| async {
(
StatusCode::INTERNAL_SERVER_ERROR,
[(INTERNODE_DISK_ERROR_HEADER, INTERNODE_FILE_CORRUPT)],
"read file err file corrupt",
)
}),
)
.route(WALK_DIR_PATH, get(get_stream))
.route("/reject-put", get(get_stream).put(reject_put))
.route("/stall", get(get_stalling_stream))
@@ -2172,6 +2246,44 @@ mod tests {
Some((format!("http://{addr}/stream"), handle))
}
#[tokio::test]
async fn http_readers_preserve_remote_file_corruption() {
let (url, server) = start_test_server(TestState::default())
.await
.expect("corruption regression server must bind");
let url = format!("{}{READ_FILE_STREAM_PATH}", url.strip_suffix("/stream").expect("test server URL suffix"));
let byte_error = HttpReader::new(url.clone(), Method::GET, HeaderMap::new(), None)
.await
.err()
.expect("byte reader must reject corrupt shard response");
let chunk_error = HttpChunkReader::new_with_stall_timeout(url, Method::GET, HeaderMap::new(), None, None)
.await
.err()
.expect("chunk reader must reject corrupt shard response");
for error in [byte_error, chunk_error] {
let source = error
.get_ref()
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
.expect("HTTP error must retain typed internode source");
assert!(source.is_remote_file_corrupt());
assert!(!source.is_remote_file_not_found());
assert!(!source.is_remote_volume_not_found());
assert_eq!(
source.kind(),
InternodeHttpErrorKind::HttpStatus(reqwest::StatusCode::INTERNAL_SERVER_ERROR)
);
let cloned = clone_internode_http_io_error(&error).expect("typed transport error must be cloneable");
assert!(
cloned
.get_ref()
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
.expect("clone must retain typed internode source")
.is_remote_file_corrupt()
);
}
server.abort();
}
struct BlockedH2Server {
url: String,
shutdown: tokio::sync::oneshot::Sender<()>,
@@ -189,6 +189,7 @@ mod tests {
use crate::admin::handlers::target_descriptor::admin_target_spec_from_builtin;
use crate::admin::runtime_sources::{IamInterface, KmsInterface};
use crate::admin::storage_api::config::save_admin_server_config;
use crate::admin::storage_api::error::StorageError;
use rustfs_config::audit::AUDIT_WEBHOOK_SUB_SYS;
use rustfs_config::server_config::KVS;
use rustfs_config::{ENABLE_KEY, EnableState, SCANNER_CYCLE, SCANNER_SUB_SYS, WEBHOOK_ENDPOINT, WEBHOOK_QUEUE_DIR};
@@ -231,11 +232,17 @@ mod tests {
let mut poll = tokio::time::interval(Duration::from_millis(10));
loop {
poll.tick().await;
let config = read_admin_config_without_migrate(store.clone())
.await
.expect("read persisted server config");
if config.0.get(subsystem).is_some_and(|targets| targets.contains_key(target)) {
return;
match read_admin_config_without_migrate(store.clone()).await {
Ok(config) => {
if config.0.get(subsystem).is_some_and(|targets| targets.contains_key(target)) {
return;
}
}
Err(StorageError::Lock(rustfs_lock::LockError::Timeout { .. })) => {
// The writer may still be committing the config snapshot. Retry after
// the object lock is released instead of failing the polling helper.
}
Err(error) => panic!("read persisted server config: {error}"),
}
}
})
+7 -4
View File
@@ -1712,16 +1712,17 @@ fn response_with_status(status: StatusCode, message: impl Into<String>) -> Respo
}
fn response_with_disk_error(error: &DiskError, message: impl Into<String>) -> Response<Body> {
let missing = match error {
let disk_error = match error {
DiskError::FileNotFound => Some(rustfs_rio::INTERNODE_FILE_NOT_FOUND),
DiskError::VolumeNotFound => Some(rustfs_rio::INTERNODE_VOLUME_NOT_FOUND),
DiskError::FileCorrupt => Some(rustfs_rio::INTERNODE_FILE_CORRUPT),
_ => None,
};
let mut response = response_with_status(StatusCode::INTERNAL_SERVER_ERROR, message);
if let Some(missing) = missing {
if let Some(disk_error) = disk_error {
response
.headers_mut()
.insert(rustfs_rio::INTERNODE_DISK_ERROR_HEADER, HeaderValue::from_static(missing));
.insert(rustfs_rio::INTERNODE_DISK_ERROR_HEADER, HeaderValue::from_static(disk_error));
}
response
}
@@ -3087,12 +3088,14 @@ mod tests {
}
#[test]
fn read_file_error_response_marks_only_missing_disk_errors() {
fn read_file_error_response_preserves_typed_disk_errors() {
for (error, expected) in [
(DiskError::FileNotFound, rustfs_rio::INTERNODE_FILE_NOT_FOUND),
(DiskError::VolumeNotFound, rustfs_rio::INTERNODE_VOLUME_NOT_FOUND),
(DiskError::FileCorrupt, rustfs_rio::INTERNODE_FILE_CORRUPT),
] {
let response = response_with_disk_error(&error, error.to_string());
assert_eq!(response.status(), StatusCode::INTERNAL_SERVER_ERROR);
assert_eq!(
response.headers().get(rustfs_rio::INTERNODE_DISK_ERROR_HEADER),
Some(&HeaderValue::from_static(expected))