mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-22 10:33:25 +00:00
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 <housemecn@gmail.com>
This commit is contained in:
Generated
+1
@@ -10752,6 +10752,7 @@ dependencies = [
|
||||
"rustfs-s3-types",
|
||||
"rustfs-scanner-metrics",
|
||||
"rustfs-storage-api",
|
||||
"rustfs-test-utils",
|
||||
"rustfs-utils",
|
||||
"s3s",
|
||||
"serde",
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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).
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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<ECStoreHealStorage>, 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");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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"] }
|
||||
|
||||
@@ -219,6 +219,7 @@ async fn mrf_ownership_cancelled_batch_restores_sync_without_per_item_clones() {
|
||||
struct NoticeStorage {
|
||||
calls: std::sync::Mutex<HashMap<String, u32>>,
|
||||
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<Option<Uuid>> {
|
||||
Ok(Some(self.bucket_incarnation_id))
|
||||
}
|
||||
async fn list_buckets(&self) -> rustfs_heal::Result<Vec<BucketInfo>> {
|
||||
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 {
|
||||
|
||||
Reference in New Issue
Block a user