From f21b06dfd2194d3a82eee002dbe29ef0ee4bf702 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sun, 13 Sep 2026 19:39:02 +0800 Subject: [PATCH] fix(heal): rebuild truncated xl.meta from healthy quorum (#7730) * fix(heal): rebuild truncated xl.meta from healthy quorum * fix(test): pass topology to heal overlap RPC regression * fix(test): drive heal admission alongside partial PUT Poll the partial PUT and its mock heal receiver together, bound their handshake, and retain the existing repair-scope assertions. * fix(test): prepare durable MRF fixtures and Linux heal stack * fix(test): drive tier cleanup recovery after deferred attempts --------- Co-authored-by: Hauser --- Cargo.lock | 1 + .../src/ecstore_validation_blackbox.rs | 55 ++- crates/ecstore/src/set_disk/mod.rs | 24 ++ crates/ecstore/src/store/init.rs | 45 +-- crates/filemeta/src/filemeta.rs | 88 +++++ crates/filemeta/src/filemeta/codec.rs | 32 +- .../heal/tests/heal_truncated_xlmeta_test.rs | 322 ++++++++++++++++++ crates/heal/tests/mrf_partial_write_test.rs | 11 +- crates/heal/tests/storage_api.rs | 4 +- crates/scanner/Cargo.toml | 1 + .../src/scanner_folder/tests/mrf_ownership.rs | 16 +- 11 files changed, 526 insertions(+), 73 deletions(-) create mode 100644 crates/heal/tests/heal_truncated_xlmeta_test.rs diff --git a/Cargo.lock b/Cargo.lock index 499e6731c..699320263 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -10752,6 +10752,7 @@ dependencies = [ "rustfs-s3-types", "rustfs-scanner-metrics", "rustfs-storage-api", + "rustfs-test-utils", "rustfs-utils", "s3s", "serde", diff --git a/crates/ecstore/src/ecstore_validation_blackbox.rs b/crates/ecstore/src/ecstore_validation_blackbox.rs index b89caafe3..fbac4cd7f 100644 --- a/crates/ecstore/src/ecstore_validation_blackbox.rs +++ b/crates/ecstore/src/ecstore_validation_blackbox.rs @@ -374,9 +374,9 @@ async fn blackbox_heal_requests_preserve_repair_scope() { let mut heal_rx = rustfs_heal_contracts::heal_channel::init_heal_channel() .expect("this must be the only ecstore test that owns the heal channel receiver"); - // Ordinary PUTs use the same admission channel as read repair. A single - // rename target failure still satisfies write quorum, so the committed - // version must be queued for convergence without delaying the PUT ACK. + // Without a durable MRF consumer, partial PUTs fall back to the heal + // admission channel and wait for its receipt before acknowledging the write. + // Drive the test receiver alongside the PUT so neither waits on the other. let (_put_dirs, put_set) = make_local_set_disks(4, 2).await; let put_bucket = "bb-put-partial-convergence"; let put_object = "object.bin"; @@ -389,42 +389,33 @@ async fn blackbox_heal_requests_preserve_repair_scope() { disks[0].take() }; let mut put_reader = PutObjReader::from_vec(vec![0x42; BLOCK_SIZE_V2 + 1024]); - let committed = put_set - .put_object( - put_bucket, - put_object, - &mut put_reader, - &ObjectOptions { - no_lock: true, - versioned: true, - ..Default::default() + let (committed, request) = tokio::time::timeout(std::time::Duration::from_secs(30), async { + tokio::join!( + async { + put_set + .put_object( + put_bucket, + put_object, + &mut put_reader, + &ObjectOptions { + no_lock: true, + versioned: true, + ..Default::default() + }, + ) + .await + .expect("partial ordinary PUT should succeed at write quorum") }, + receive_matching_heal(&mut heal_rx, put_bucket, put_object), ) - .await - .expect("partial ordinary PUT should succeed at write quorum"); + }) + .await + .expect("partial ordinary PUT and repair admission should complete together"); let committed_version = committed .version_id .expect("versioned PUT should return a version id") .to_string(); - let request = tokio::time::timeout(std::time::Duration::from_secs(30), async { - loop { - match heal_rx.recv().await.expect("heal channel should stay open") { - HealChannelCommand::Start { request, response_tx } - if request.bucket == put_bucket && request.object_prefix.as_deref() == Some(put_object) => - { - let _ = response_tx.send(Ok(HealAdmissionResult::Accepted)); - break request; - } - HealChannelCommand::Start { response_tx, .. } => { - let _ = response_tx.send(Ok(HealAdmissionResult::Accepted)); - } - _ => {} - } - } - }) - .await - .expect("partial ordinary PUT should enqueue convergence heal"); assert_eq!(request.object_version_id.as_deref(), Some(committed_version.as_str())); assert_eq!(request.pool_index, Some(0)); assert_eq!(request.set_index, Some(0)); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 7b8821efd..363f261d5 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -10818,6 +10818,30 @@ mod tests { assert_eq!(reason, Some(DiskError::FileCorrupt)); } + #[test] + fn metadata_io_failures_never_authorize_heal_overwrite() { + let meta = FileInfo::default(); + let io_errors = [ + std::io::Error::new(std::io::ErrorKind::PermissionDenied, "metadata access denied"), + std::io::Error::other("transient metadata read failure"), + ]; + let mut errors: Vec<_> = io_errors + .into_iter() + .map(|error| DiskError::from(rustfs_filemeta::Error::Io(error))) + .collect(); + #[cfg(unix)] + errors.push(DiskError::from(rustfs_filemeta::Error::Io(std::io::Error::from_raw_os_error(libc::EIO)))); + errors.push(DiskError::Timeout); + for error in errors { + assert_ne!(error, DiskError::FileCorrupt); + let (heal, metadata, reason) = + should_heal_object_on_disk(&Some(error.clone()), &[CHECK_PART_FILE_CORRUPT], &meta, &meta); + assert!(!heal, "an I/O failure must not authorize overwriting metadata: {error}"); + assert!(!metadata); + assert_eq!(reason, Some(error)); + } + } + #[tokio::test] async fn test_get_disks_info_preserves_runtime_state_for_suspect_and_offline_disks() { let format = FormatV3::new(1, 3); diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index b7df96fa5..a95871805 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -3002,33 +3002,36 @@ mod tests { .expect("rebind the same remote destination after restart"); } let set = store.pools[0].get_disks_by_key(object); + // A deferred first cleanup must retain its durable owner + // until a later recovery scan can retry the operation. + backend.set_remove_failure(true); ExpiryState::resize_workers(1, Arc::clone(&store)).await; let recovered = recover_tier_free_versions(Arc::clone(&store), 100, None, None) .await .expect("recover persisted cleanup owner"); assert!(recovered.enqueued >= 1); - tokio::time::timeout(Duration::from_secs(30), async { - loop { - let versions = set - .load_file_info_versions_exact(&bucket, object) - .await - .expect("read cleanup progress") - .expect("new object must survive cleanup"); - if versions - .versions - .iter() - .chain(versions.free_versions.iter()) - .all(|fi| !fi.tier_free_version()) - { - break; - } - tokio::task::yield_now().await; - } - }) - .await - .expect("cleanup must converge"); + wait_for_expiry_workers_idle(&store).await; + assert!(backend.contains(&remote).await, "failed cleanup must retain remote bytes"); + assert_eq!(backend.remove_count().await, removed_before); + backend.set_remove_failure(false); + // This fixture starts expiry workers without the runtime's + // recovery loop, so drive its durable rescan explicitly. + wait_for_tier_free_version_recovery(Arc::clone(&store), &backend, removed_before + 1).await; + let versions = set + .load_file_info_versions_exact(&bucket, object) + .await + .expect("read cleanup progress") + .expect("new object must survive cleanup"); + assert!( + versions + .versions + .iter() + .chain(versions.free_versions.iter()) + .all(|fi| !fi.tier_free_version()), + "cleanup must remove its owner: {state:?}, suspended={suspended}, copy={self_copy}" + ); assert!(!backend.contains(&remote).await); - assert_eq!(backend.remove_count().await, removed_before + 1, "one remote DELETE per owner"); + assert_eq!(backend.remove_count().await, removed_before + 1, "one successful remote DELETE per owner"); assert_eq!(backend.remove_versions().await.last(), Some(&(remote.clone(), version.to_string()))); let mut reader = store .get_object_reader(&bucket, object, None, HeaderMap::new(), &options) diff --git a/crates/filemeta/src/filemeta.rs b/crates/filemeta/src/filemeta.rs index 114c19634..14a19bfaa 100644 --- a/crates/filemeta/src/filemeta.rs +++ b/crates/filemeta/src/filemeta.rs @@ -2139,6 +2139,94 @@ mod test { } } + #[test] + fn truncated_xlmeta_framing_is_file_corrupt() { + let buf = FileMeta::default() + .marshal_msg() + .expect("serialize metadata without inline data"); + FileMeta::load(&buf).expect("complete metadata must decode"); + for cut in 0..buf.len() { + assert_eq!( + FileMeta::load(&buf[..cut]).expect_err("every incomplete metadata frame must fail"), + Error::FileCorrupt, + "truncation at byte {cut} must remain repairable" + ); + } + } + + #[test] + fn truncated_xlmeta_index_and_format_reads_are_file_corrupt() { + let buf = FileMeta::default() + .marshal_msg() + .expect("serialize metadata without inline data"); + FileMeta::is_indexed_meta(&buf).expect("complete indexed metadata must decode"); + FileMeta::read_format_versions(&buf).expect("complete format header must decode"); + for cut in 0..buf.len() { + assert_eq!( + FileMeta::is_indexed_meta(&buf[..cut]).expect_err("incomplete index must fail"), + Error::FileCorrupt, + "index truncation at byte {cut}" + ); + if cut < buf.len() - 5 { + assert_eq!( + FileMeta::read_format_versions(&buf[..cut]).expect_err("incomplete metadata block must fail"), + Error::FileCorrupt, + "format truncation at byte {cut}" + ); + } + } + for cut in 0..5 { + assert_eq!( + FileMeta::read_bytes_header(&buf[8..8 + cut]).expect_err("incomplete bin32 prefix must fail"), + Error::FileCorrupt + ); + } + } + + #[test] + fn malformed_xlmeta_framing_is_file_corrupt() { + let original = FileMeta::default() + .marshal_msg() + .expect("serialize metadata without inline data"); + for offset in [8, original.len() - 5] { + let mut buf = original.clone(); + buf[offset] = 0xc0; // nil cannot encode a bin length or a CRC integer. + assert_eq!(FileMeta::load(&buf).expect_err("invalid framing marker"), Error::FileCorrupt); + assert_eq!( + FileMeta::is_indexed_meta(&buf).expect_err("invalid index framing marker"), + Error::FileCorrupt + ); + if offset == 8 { + assert_eq!( + FileMeta::read_format_versions(&buf).expect_err("invalid format framing marker"), + Error::FileCorrupt + ); + assert_eq!( + FileMeta::read_bytes_header(&buf[8..]).expect_err("invalid bin framing marker"), + Error::FileCorrupt + ); + } + } + } + + #[test] + fn unsupported_xlmeta_versions_are_not_classified_as_corruption() { + let mut buf = FileMeta::default().marshal_msg().expect("serialize metadata"); + buf[4..6].copy_from_slice(&(XL_FILE_VERSION_MAJOR + 1).to_le_bytes()); + assert_ne!(FileMeta::load(&buf).expect_err("unsupported file version"), Error::FileCorrupt); + for (header_ver, meta_ver) in [ + (XL_HEADER_VERSION + 1, XL_META_VERSION), + (XL_HEADER_VERSION, XL_META_VERSION + 1), + ] { + let mut meta = Vec::new(); + rmp::encode::write_uint(&mut meta, u64::from(header_ver)).expect("write header version"); + rmp::encode::write_uint(&mut meta, u64::from(meta_ver)).expect("write metadata version"); + rmp::encode::write_uint(&mut meta, 0).expect("write empty version count"); + let buf = build_xl_buffer(&meta); + assert_ne!(FileMeta::load(&buf).expect_err("unsupported schema version"), Error::FileCorrupt); + } + } + /// Regression test for rustfs/rustfs#2715: a corrupted version count in /// xl.meta must yield a decode error instead of sizing a huge allocation /// from the bogus count (which aborts the whole process). diff --git a/crates/filemeta/src/filemeta/codec.rs b/crates/filemeta/src/filemeta/codec.rs index 86af49b97..851785ddb 100644 --- a/crates/filemeta/src/filemeta/codec.rs +++ b/crates/filemeta/src/filemeta/codec.rs @@ -31,12 +31,12 @@ impl FileMeta { pub fn read_format_versions(buf: &[u8]) -> Result<(u16, u16, u8, u8)> { let (buf, major, minor) = Self::check_xl2_v1(buf)?; if buf.len() < 5 { - return Err(Error::other("insufficient data for metadata length prefix")); + return Err(Error::FileCorrupt); } let (mut size_buf, buf) = buf.split_at(5); - let bin_len = rmp::decode::read_bin_len(&mut size_buf)?; + let bin_len = rmp::decode::read_bin_len(&mut size_buf).map_err(|_| Error::FileCorrupt)?; if buf.len() < bin_len as usize { - return Err(Error::other("insufficient data for metadata")); + return Err(Error::FileCorrupt); } let (meta, _) = buf.split_at(bin_len as usize); let (_, header_ver, meta_ver, _) = Self::decode_xl_headers(meta)?; @@ -74,26 +74,26 @@ impl FileMeta { } if buf.len() < 5 { - return Err(Error::other("insufficient data for meta length")); + return Err(Error::FileCorrupt); } let (mut size_buf, buf) = buf.split_at(5); // Get meta data, buf = crc + data - let bin_len = rmp::decode::read_bin_len(&mut size_buf)?; + let bin_len = rmp::decode::read_bin_len(&mut size_buf).map_err(|_| Error::FileCorrupt)?; if buf.len() < bin_len as usize { - return Ok((&[], &[])); + return Err(Error::FileCorrupt); } let (meta, buf) = buf.split_at(bin_len as usize); if buf.len() < 5 { - return Err(Error::other("insufficient data for CRC")); + return Err(Error::FileCorrupt); } let (mut crc_buf, inline_data) = buf.split_at(5); // crc check - let crc = rmp::decode::read_u32(&mut crc_buf)?; + let crc = rmp::decode::read_u32(&mut crc_buf).map_err(|_| Error::FileCorrupt)?; let meta_crc = xxh64::xxh64(meta, XXHASH_SEED) as u32; if crc != meta_crc { @@ -119,13 +119,13 @@ impl FileMeta { // Fixed u32 pub fn read_bytes_header(buf: &[u8]) -> Result<(u32, &[u8])> { if buf.len() < 5 { - return Err(Error::other("insufficient data for bytes header")); + return Err(Error::FileCorrupt); } let (mut size_buf, remaining) = buf.split_at(5); // Get meta data, buf = crc + data - let bin_len = rmp::decode::read_bin_len(&mut size_buf)?; + let bin_len = rmp::decode::read_bin_len(&mut size_buf).map_err(|_| Error::FileCorrupt)?; Ok((bin_len, remaining)) } @@ -137,12 +137,14 @@ impl FileMeta { // check version, buf = buf[8..] let (buf, _, _) = Self::check_xl2_v1(buf)?; + // These bytes have already been read. Invalid framing is deterministic + // metadata damage; preserve FileCorrupt so quorum-backed heal can repair it. if buf.len() < 5 { error!( "insufficient data for metadata length prefix: expected at least 5 bytes, got {}", buf.len() ); - return Err(Error::other("insufficient data for metadata length prefix")); + return Err(Error::FileCorrupt); } let (mut size_buf, buf) = buf.split_at(5); @@ -150,25 +152,25 @@ impl FileMeta { // Get meta data, buf = crc + data let bin_len = rmp::decode::read_bin_len(&mut size_buf).map_err(|e| { error!("failed to read binary length for metadata: {}", e); - Error::other(format!("failed to read binary length for metadata: {e}")) + Error::FileCorrupt })?; if buf.len() < bin_len as usize { error!("insufficient data for metadata: expected {} bytes, got {} bytes", bin_len, buf.len()); - return Err(Error::other("insufficient data for metadata")); + return Err(Error::FileCorrupt); } let (meta, buf) = buf.split_at(bin_len as usize); if buf.len() < 5 { error!("insufficient data for CRC: expected 5 bytes, got {} bytes", buf.len()); - return Err(Error::other("insufficient data for CRC")); + return Err(Error::FileCorrupt); } let (mut crc_buf, buf) = buf.split_at(5); // crc check let crc = rmp::decode::read_u32(&mut crc_buf).map_err(|e| { error!("failed to read CRC value: {}", e); - Error::other(format!("failed to read CRC value: {e}")) + Error::FileCorrupt })?; let meta_crc = xxh64::xxh64(meta, XXHASH_SEED) as u32; diff --git a/crates/heal/tests/heal_truncated_xlmeta_test.rs b/crates/heal/tests/heal_truncated_xlmeta_test.rs new file mode 100644 index 000000000..e3de63233 --- /dev/null +++ b/crates/heal/tests/heal_truncated_xlmeta_test.rs @@ -0,0 +1,322 @@ +// 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. + +use http::HeaderMap; +use rustfs_heal::heal::{ + outcome::{HealObjectDisposition, HealTraversalCoverage}, + storage::{ECStoreHealStorage, HealObjectOptions as ObjectOptions, HealPutObjReader as PutObjReader}, + task::{HealOptions, HealPriority, HealRequest, HealTask, HealType}, +}; +use rustfs_heal_contracts::heal_channel::{DriveState, HealScanMode}; +use std::{sync::Arc, time::Duration}; +use tokio::io::AsyncReadExt as _; + +mod storage_api; +use storage_api::integration::{ + DiskAPI as _, DiskError, DiskOption, Endpoint, NamespaceLocking as _, ObjectIO as _, ReadOptions, new_disk, +}; + +fn deep_heal_task(storage: &Arc, bucket: &str, object: &str, dry_run: bool) -> HealTask { + HealTask::from_request( + HealRequest::new( + HealType::Object { + bucket: bucket.to_owned(), + object: object.to_owned(), + version_id: None, + }, + HealOptions { + scan_mode: HealScanMode::Deep, + dry_run, + ..Default::default() + }, + HealPriority::Normal, + ), + storage.clone(), + ) +} + +#[test] +fn deep_heal_rebuilds_truncated_xlmeta_with_authoritative_outcomes() { + // Like the torn-minority regression, this real storage scenario composes + // deep async futures that exceed libtest's default Linux thread stack. + std::thread::Builder::new() + .name("deep-heal-truncated-xlmeta".to_string()) + .stack_size(32 * 1024 * 1024) + .spawn(|| { + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("deep heal test runtime should build"); + runtime.block_on(deep_heal_truncated_xlmeta_scenario()); + }) + .expect("deep heal test thread should spawn") + .join() + .expect("deep heal test thread should finish"); +} + +async fn deep_heal_truncated_xlmeta_scenario() { + let temp = tempfile::tempdir().expect("create caller-owned disks"); + let env = rustfs_test_utils::TestECStoreEnv::builder() + .disk_count(16) + .base_dir(temp.path()) + .build() + .await; + let storage = Arc::new(ECStoreHealStorage::new(env.ecstore.clone())); + let bucket = "truncated-xlmeta"; + env.make_bucket(bucket, false).await; + let payload = vec![0x7b; 4 * 1024 * 1024]; + let mut endpoint = Endpoint::try_from(env.disk_paths[0].to_str().expect("UTF-8 disk path")).expect("target endpoint"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(0); + let disk = new_disk(&endpoint, &DiskOption::default()) + .await + .expect("open target disk inspector"); + let read_options = ReadOptions { + read_data: true, + ..Default::default() + }; + + for damage in ["length-prefix", "metadata-body", "crc-tail", "versioned"] { + let versioned = damage == "versioned"; + let bucket = if versioned { "truncated-xlmeta-versioned" } else { bucket }; + if versioned { + env.make_bucket(bucket, true).await; + } + let object = format!("{damage}/object.bin"); + let mut reader = PutObjReader::from_vec(payload.clone()); + env.ecstore + .put_object( + bucket, + &object, + &mut reader, + &ObjectOptions { + versioned, + ..Default::default() + }, + ) + .await + .expect("write source object"); + // The namespace fence waits for the detached PUT publication owner; + // no GET/HEAD can repair the target before the explicit heal. + let lock = env + .ecstore + .new_ns_lock(bucket, &object) + .await + .expect("fixture namespace lock"); + let settled = lock + .get_write_lock(Duration::from_secs(30)) + .await + .expect("PUT publication must finish"); + let original_info = disk + .read_version("", bucket, &object, "", &read_options) + .await + .expect("read original metadata"); + assert_eq!(original_info.version_id.is_some(), versioned); + assert_eq!((original_info.erasure.data_blocks, original_info.erasure.parity_blocks), (12, 4)); + let data_dir = original_info.data_dir.expect("non-inline object has a data directory"); + let target_dir = env.disk_paths[0].join(bucket).join(&object); + let target_meta = target_dir.join("xl.meta"); + let original_part = tokio::fs::read(target_dir.join(data_dir.to_string()).join("part.1")) + .await + .expect("read original shard"); + let mut original_metadata = Vec::new(); + for path in &env.disk_paths { + original_metadata.push( + tokio::fs::read(path.join(bucket).join(&object).join("xl.meta")) + .await + .expect("snapshot every member"), + ); + } + let original = &original_metadata[0]; + let metadata_len = usize::try_from(u32::from_be_bytes(original[9..13].try_into().expect("bin32 length"))) + .expect("metadata length fits usize"); + assert_eq!(original.len(), 13 + metadata_len + 5, "fixture must exclude inline data"); + let cut = match damage { + "length-prefix" => 12, + "metadata-body" | "versioned" => 13 + metadata_len / 2, + "crc-tail" => original.len() - 1, + _ => unreachable!("fixed damage matrix"), + }; + tokio::fs::write(&target_meta, &original[..cut]) + .await + .expect("inject exact truncation"); + drop(settled); + let read_error = disk + .read_version("", bucket, &object, "", &read_options) + .await + .expect_err("target must be unreadable before heal"); + + let dry_run = deep_heal_task(&storage, bucket, &object, true); + dry_run.execute().await.expect("dry-run traversal completes"); + assert_eq!(tokio::fs::read(&target_meta).await.expect("read dry-run target"), original[..cut]); + assert_eq!(dry_run.get_outcome().await.objects[0].disposition, HealObjectDisposition::DryRunObserved); + + let task = deep_heal_task(&storage, bucket, &object, false); + task.execute().await.expect("deep heal traversal completes"); + let outcome = task.get_outcome().await; + assert_eq!(outcome.coverage, HealTraversalCoverage::Complete); + assert_eq!(outcome.counters.processed, 1); + assert_eq!(outcome.counters.healed, 1, "{damage}: {outcome:?}"); + assert_eq!(outcome.counters.unknown, 0); + assert_eq!(outcome.counters.skipped, 0); + assert_eq!(outcome.counters.failed, 0); + assert_eq!(outcome.counters.attempt_failures, 0); + assert_eq!(outcome.objects[0].disposition, HealObjectDisposition::Repaired); + assert_eq!(read_error, DiskError::FileCorrupt); + let results = task.get_result_items().await; + assert_eq!(results.len(), 1); + assert_eq!(results[0].before.drives.len(), 16); + assert_eq!(results[0].after.drives.len(), 16); + assert_eq!(results[0].before.drives[0].state, DriveState::Corrupt.to_string()); + assert!( + results[0] + .after + .drives + .iter() + .all(|drive| drive.state == DriveState::Ok.to_string()) + ); + + let repaired = disk + .read_version("", bucket, &object, "", &read_options) + .await + .expect("physical metadata must decode after heal"); + assert_eq!(repaired.version_id, original_info.version_id); + assert_eq!(repaired.mod_time, original_info.mod_time); + assert_eq!(repaired.data_dir, original_info.data_dir); + assert_eq!(repaired.parts, original_info.parts); + assert_eq!(repaired.erasure, original_info.erasure); + assert_eq!( + tokio::fs::read(target_dir.join(data_dir.to_string()).join("part.1")) + .await + .expect("read rebuilt shard"), + original_part + ); + for (index, path) in env.disk_paths.iter().enumerate().skip(1) { + assert_eq!( + tokio::fs::read(path.join(bucket).join(&object).join("xl.meta")) + .await + .expect("read healthy member"), + original_metadata[index] + ); + } + + let healed_bytes = tokio::fs::read(&target_meta).await.expect("read committed metadata"); + let repeat = deep_heal_task(&storage, bucket, &object, false); + repeat.execute().await.expect("repeated heal completes"); + let repeat_outcome = repeat.get_outcome().await; + assert_eq!(repeat_outcome.counters.healed, 0); + assert_eq!(repeat_outcome.counters.unchanged, 1); + assert_eq!(repeat_outcome.counters.unknown, 0); + assert_eq!(tokio::fs::read(&target_meta).await.expect("read repeated-heal target"), healed_bytes); + let mut reader = env + .ecstore + .get_object_reader(bucket, &object, None, HeaderMap::new(), &ObjectOptions::default()) + .await + .expect("read healed object"); + let mut actual = Vec::new(); + reader.read_to_end(&mut actual).await.expect("read complete healed body"); + drop(reader); + assert_eq!(actual, payload); + + // A future format is not evidence of corruption. Even with healthy + // quorum, the target must remain unknown and must not be overwritten. + let mut future = healed_bytes; + future[4..6].copy_from_slice(&2_u16.to_le_bytes()); + tokio::fs::write(&target_meta, &future) + .await + .expect("install unsupported major version"); + let unsupported = deep_heal_task(&storage, bucket, &object, false); + unsupported + .execute() + .await + .expect("unsupported member permits traversal completion"); + let unknown = unsupported.get_outcome().await; + assert_eq!( + ( + unknown.counters.processed, + unknown.counters.healed, + unknown.counters.skipped, + unknown.counters.unknown, + unknown.counters.failed + ), + (1, 0, 1, 1, 0) + ); + assert_eq!(unknown.objects[0].disposition, HealObjectDisposition::Unknown); + assert_eq!(tokio::fs::read(&target_meta).await.expect("read unsupported metadata"), future); + + if damage == "metadata-body" { + for damaged in [4, 5] { + for (index, path) in env.disk_paths.iter().enumerate() { + let bytes = if index < damaged { + &original_metadata[index][..cut] + } else { + &original_metadata[index] + }; + tokio::fs::write(path.join(bucket).join(&object).join("xl.meta"), bytes) + .await + .expect("install exact-quorum fixture"); + } + let quorum_task = deep_heal_task(&storage, bucket, &object, false); + let result = quorum_task.execute().await; + let outcome = quorum_task.get_outcome().await; + if damaged == 4 { + result.expect("twelve authoritative members must repair four damaged copies"); + assert_eq!(outcome.counters.healed, 1); + assert_eq!(outcome.counters.unknown, 0); + assert!( + quorum_task.get_result_items().await[0] + .after + .drives + .iter() + .all(|drive| drive.state == DriveState::Ok.to_string()) + ); + } else { + assert_eq!(outcome.counters.healed, 0, "eleven healthy members must not authorize repair"); + } + for (index, path) in env.disk_paths.iter().enumerate() { + let actual = tokio::fs::read(path.join(bucket).join(&object).join("xl.meta")) + .await + .expect("inspect quorum fixture"); + if index < damaged && damaged == 4 { + let mut endpoint = + Endpoint::try_from(path.to_str().expect("UTF-8 member path")).expect("member endpoint"); + endpoint.set_pool_index(0); + endpoint.set_set_index(0); + endpoint.set_disk_index(index); + let member = new_disk(&endpoint, &DiskOption::default()) + .await + .expect("open repaired member"); + let info = member + .read_version("", bucket, &object, "", &read_options) + .await + .expect("every repaired member must decode"); + assert_eq!(info.data_dir, original_info.data_dir); + assert_eq!(info.version_id, original_info.version_id); + assert_eq!(info.parts, original_info.parts); + assert_eq!(info.erasure.index, original_info.erasure.distribution[index]); + assert_ne!(actual, original_metadata[index][..cut]); + } else { + let expected = if index < damaged { + &original_metadata[index][..cut] + } else { + &original_metadata[index] + }; + assert_eq!(actual, expected, "quorum-minus-one must preserve every copy"); + } + } + } + } + } +} diff --git a/crates/heal/tests/mrf_partial_write_test.rs b/crates/heal/tests/mrf_partial_write_test.rs index 0f23e9b7e..24b444030 100644 --- a/crates/heal/tests/mrf_partial_write_test.rs +++ b/crates/heal/tests/mrf_partial_write_test.rs @@ -27,7 +27,9 @@ use tokio::io::AsyncReadExt; mod storage_api; use storage_api::endpoint_index::{EndpointServerPools, Endpoints, init_local_disks}; -use storage_api::integration::{DiskAPI, DiskStore, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, ReadOptions}; +use storage_api::integration::{ + DiskAPI, DiskStore, ObjectIO, ObjectOperations, ObjectOptions, PutObjReader, RUSTFS_META_BUCKET, ReadOptions, +}; const SNAPSHOT_LIMIT: usize = 64 * 1024 * 1024; @@ -47,7 +49,10 @@ async fn partial_write_persistence_failure_is_reported_and_retained_for_retry() .await .expect("initial responsibility must commit"); assert!(snapshot_contains("old.bin").await); - for path in &env.disk_paths { + // Linux pins the disk root with a directory descriptor, so renaming that + // root does not interrupt I/O. Block the metadata volume below it instead. + let metadata_roots: Vec<_> = env.disk_paths.iter().map(|path| path.join(RUSTFS_META_BUCKET)).collect(); + for path in &metadata_roots { tokio::fs::rename(path, path.with_extension("offline")) .await .expect("detach journal disk"); @@ -60,7 +65,7 @@ async fn partial_write_persistence_failure_is_reported_and_retained_for_retry() Err(MrfDurableAdmissionError::Persistence), "failed checkpoint publication must not be acknowledged as durable success" ); - for path in &env.disk_paths { + for path in &metadata_roots { tokio::fs::remove_file(path).await.expect("remove journal fault"); tokio::fs::rename(path.with_extension("offline"), path) .await diff --git a/crates/heal/tests/storage_api.rs b/crates/heal/tests/storage_api.rs index 62bf11772..9c0fdf956 100644 --- a/crates/heal/tests/storage_api.rs +++ b/crates/heal/tests/storage_api.rs @@ -20,7 +20,9 @@ pub(crate) mod endpoint_index { // still touch directly. #[allow(unused_imports)] pub(crate) mod integration { - pub(crate) use rustfs_ecstore::api::disk::{DiskAPI, DiskStore, ReadOptions}; + pub(crate) use rustfs_ecstore::api::disk::{ + DiskAPI, DiskError, DiskOption, DiskStore, Endpoint, RUSTFS_META_BUCKET, ReadOptions, new_disk, + }; pub(crate) use rustfs_ecstore::api::object::{ObjectOptions, PutObjReader}; pub(crate) use rustfs_ecstore::api::storage::ECStore; pub(crate) use rustfs_storage_api::BucketOperations; diff --git a/crates/scanner/Cargo.toml b/crates/scanner/Cargo.toml index 635fe3423..6e92392e4 100644 --- a/crates/scanner/Cargo.toml +++ b/crates/scanner/Cargo.toml @@ -107,6 +107,7 @@ hex-simd.workspace = true [dev-dependencies] rustfs-heal.workspace = true +rustfs-test-utils.workspace = true tracing-subscriber = { workspace = true, features = ["json", "env-filter", "time"] } serial_test = { workspace = true } temp-env = { workspace = true, features = ["async_closure"] } diff --git a/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs b/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs index 142008cef..f792b2361 100644 --- a/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs +++ b/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs @@ -219,6 +219,7 @@ async fn mrf_ownership_cancelled_batch_restores_sync_without_per_item_clones() { struct NoticeStorage { calls: std::sync::Mutex>, retry_started: tokio::sync::Notify, + bucket_incarnation_id: Uuid, } #[async_trait::async_trait] @@ -235,6 +236,9 @@ impl HealStorageAPI for NoticeStorage { ..Default::default() })) } + async fn mrf_bucket_incarnation_id(&self, _: &str) -> rustfs_heal::Result> { + Ok(Some(self.bucket_incarnation_id)) + } async fn list_buckets(&self) -> rustfs_heal::Result> { Ok(Vec::new()) } @@ -315,13 +319,23 @@ async fn mrf_ownership_manager_completion_preserves_scanner_pending() { } // The production ingress channel is a process singleton; isolation keeps // its receiver and lease generations independent from other scanner tests. + // Partial writes require a committed journal and a complete bucket identity + // before the real consumer may dispatch them. + let journal_root = tempfile::tempdir().expect("MRF journal fixture"); + let _journal_env = rustfs_test_utils::TestECStoreEnv::builder() + .base_dir(journal_root.path()) + .build() + .await; let (mut scanner, temp_dir) = build_test_scanner().await; let _guard = TestGuard::new(u64::MAX, usize::MAX, &mut scanner, temp_dir); let bucket = format!("mrf-ownership-{}", Uuid::new_v4()); scanner.new_cache.info.name = bucket.clone(); scanner.update_cache.info.name = bucket.clone(); scanner.heal_object_select = 1; - let storage = Arc::new(NoticeStorage::default()); + let storage = Arc::new(NoticeStorage { + bucket_incarnation_id: Uuid::new_v4(), + ..Default::default() + }); let manager = Arc::new(HealManager::new( storage.clone(), Some(HealConfig {