diff --git a/crates/ecstore/src/disk/error.rs b/crates/ecstore/src/disk/error.rs index 881f93778..b987056a9 100644 --- a/crates/ecstore/src/disk/error.rs +++ b/crates/ecstore/src/disk/error.rs @@ -244,13 +244,16 @@ impl StdError for ConditionalFileNotCommittedError { } } -fn classify_internode_missing_error(error: &InternodeHttpError) -> Option { +fn classify_internode_disk_error(error: &InternodeHttpError) -> Option { 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 for DiskError { fn from(e: std::io::Error) -> Self { - if let Some(error) = e.get_ref().and_then(|source| source.downcast_ref::()) { - 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::()) + && let Some(classified) = classify_internode_disk_error(error) + { + return classified; } let e = match e.downcast::() { Ok(terminal_error) => { @@ -541,7 +541,7 @@ impl From for DiskError { && let Some(internode_error) = io_error .get_ref() .and_then(|source| source.downcast_ref::()) - && 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) ); } } diff --git a/crates/ecstore/src/io_support/shard_integrity.rs b/crates/ecstore/src/io_support/shard_integrity.rs index 9c8f9c417..deb41f382 100644 --- a/crates/ecstore/src/io_support/shard_integrity.rs +++ b/crates/ecstore/src/io_support/shard_integrity.rs @@ -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)); } diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index b9e766c2f..68cf4d582 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -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 diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 1f57c1142..906ff1eb3 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -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 { + 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 { + 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 { + 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] diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index 17c5fa2a9..c814aa114 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -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]) -> 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, @@ -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::()) + .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::()) + .expect("clone must retain typed internode source") + .is_remote_file_corrupt() + ); + } + server.abort(); + } + struct BlockedH2Server { url: String, shutdown: tokio::sync::oneshot::Sender<()>, diff --git a/rustfs/src/admin/handlers/audit_runtime_config.rs b/rustfs/src/admin/handlers/audit_runtime_config.rs index 1055d3bf3..ad58b726a 100644 --- a/rustfs/src/admin/handlers/audit_runtime_config.rs +++ b/rustfs/src/admin/handlers/audit_runtime_config.rs @@ -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}"), } } }) diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index 2b0bd9a68..a96f992db 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -1712,16 +1712,17 @@ fn response_with_status(status: StatusCode, message: impl Into) -> Respo } fn response_with_disk_error(error: &DiskError, message: impl Into) -> Response { - 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))