mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-21 18:13:34 +00:00
fix(heal): inspect stale members during deep scans (#8009)
* fix(heal): inspect stale members during deep scans * test(heal): prove convergence outcomes
This commit is contained in:
@@ -978,12 +978,25 @@ impl SetDisks {
|
||||
result.parity_blocks = result.disk_count - read_quorum as usize;
|
||||
result.data_blocks = read_quorum as usize;
|
||||
|
||||
let ((mut online_disks, quorum_mod_time, quorum_etag), disk_len) = {
|
||||
let ((quorum_disks, quorum_mod_time, quorum_etag), disk_len) = {
|
||||
let disks = self.disks.read().await;
|
||||
let disk_len = disks.len();
|
||||
(Self::list_online_disks(&disks, &parts_metadata, &errs, read_quorum as usize), disk_len)
|
||||
};
|
||||
|
||||
// A deep heal must inspect every reachable member, not only the
|
||||
// quorum that agreed on the canonical metadata. A returning disk
|
||||
// can legitimately carry an older xl.meta (or an explicit-version
|
||||
// absence) while still being the exact stale target that needs
|
||||
// repair. Restricting the repair set to quorum members silently
|
||||
// turns that stale member into `verified_healthy`/`unchanged`
|
||||
// (backlog#2612, #2596). Normal scans retain the quorum fast path.
|
||||
let mut online_disks = if opts.scan_mode == HealScanMode::Deep {
|
||||
self.get_disks_internal().await
|
||||
} else {
|
||||
quorum_disks
|
||||
};
|
||||
|
||||
trace!(
|
||||
event = EVENT_SET_DISK_HEAL,
|
||||
component = LOG_COMPONENT_ECSTORE,
|
||||
@@ -4087,6 +4100,89 @@ mod heal_result_report_tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn deep_heal_rebuilds_near_tail_truncated_part() {
|
||||
let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await;
|
||||
let bucket = "deep-heal-near-tail-truncation";
|
||||
let object = "object.bin";
|
||||
for disk in &disks {
|
||||
disk.make_volume(bucket).await.expect("bucket volume should be created");
|
||||
}
|
||||
|
||||
let expected_payload = vec![0x6d; 1024 * 1024];
|
||||
set.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(expected_payload.clone()),
|
||||
&ObjectOptions {
|
||||
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");
|
||||
|
||||
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 (result, error) = set
|
||||
.heal_object(
|
||||
bucket,
|
||||
object,
|
||||
"",
|
||||
&HealOpts {
|
||||
no_lock: true,
|
||||
scan_mode: HealScanMode::Deep,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("deep heal should finish after a one-byte tail truncation");
|
||||
|
||||
assert!(error.is_none(), "deep heal should recover the truncated shard: {error:?}");
|
||||
assert_eq!(result.drives_healed(), Some(1));
|
||||
assert_eq!(result.before.drives[1].state, DriveState::Corrupt.to_string());
|
||||
assert_eq!(
|
||||
tokio::fs::metadata(&truncated_part)
|
||||
.await
|
||||
.expect("repaired shard should exist")
|
||||
.len(),
|
||||
original_len,
|
||||
"deep heal must restore the complete shard length"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn deep_heal_inline_bitrot_rebuilds_physical_data_and_parity_shards() {
|
||||
use crate::storage_api_contracts::range::HTTPRangeSpec;
|
||||
|
||||
@@ -36,6 +36,7 @@ use rustfs_heal::heal::{
|
||||
};
|
||||
use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode};
|
||||
use serial_test::serial;
|
||||
use sha2::{Digest, Sha256};
|
||||
use std::{
|
||||
path::{Path, PathBuf},
|
||||
sync::Arc,
|
||||
@@ -722,6 +723,91 @@ mod serial_tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
#[serial]
|
||||
async fn test_recursive_deep_heal_reports_stale_delete_marker_repair() {
|
||||
let (disk_paths, ecstore, storage) = heal_env().await;
|
||||
let manager = HealManager::new(
|
||||
storage.clone(),
|
||||
Some(HealConfig {
|
||||
heal_interval: Duration::from_millis(1),
|
||||
..Default::default()
|
||||
}),
|
||||
);
|
||||
manager.start().await.expect("heal manager should start");
|
||||
|
||||
let bucket = "b5-stale-delete-marker-recursive";
|
||||
let object = "obj.bin";
|
||||
create_versioned_bucket(&ecstore, bucket).await;
|
||||
put_unversioned(&ecstore, bucket, "control/object.bin", &versioned_test_data(42)).await;
|
||||
let historical_data = versioned_test_data(43);
|
||||
let historical = put_versioned(&ecstore, bucket, object, &historical_data).await;
|
||||
let target = &disk_paths[1];
|
||||
let target_meta = xl_meta_path(&object_dir(target, bucket, object));
|
||||
let stale_meta = std::fs::read(&target_meta).expect("target xl.meta must exist before marker creation");
|
||||
let stale_digest = Sha256::digest(&stale_meta);
|
||||
let marker = put_delete_marker(&ecstore, bucket, object).await;
|
||||
std::fs::write(&target_meta, &stale_meta).expect("restore stale target xl.meta");
|
||||
|
||||
let request = HealRequest::new(
|
||||
HealType::Bucket {
|
||||
bucket: bucket.to_string(),
|
||||
},
|
||||
HealOptions {
|
||||
recursive: true,
|
||||
scan_mode: HealScanMode::Deep,
|
||||
pool_index: Some(0),
|
||||
set_index: Some(0),
|
||||
..Default::default()
|
||||
},
|
||||
HealPriority::Normal,
|
||||
);
|
||||
let task_id = request.id.clone();
|
||||
assert!(
|
||||
manager
|
||||
.submit_heal_request(request)
|
||||
.await
|
||||
.expect("submit recursive heal")
|
||||
.is_admitted()
|
||||
);
|
||||
wait_for_task(&manager, &task_id, Duration::from_secs(60)).await;
|
||||
|
||||
let report = manager.get_task_report(&task_id).await.expect("completed task report");
|
||||
let outcome = report.outcome.expect("recursive task must have an outcome");
|
||||
assert_eq!(outcome.counters.failed, 0);
|
||||
assert_eq!(outcome.counters.unknown, 0);
|
||||
assert_eq!(outcome.counters.processed, 3);
|
||||
assert_eq!(outcome.counters.healed, 1, "delete marker repair must be counted as healed");
|
||||
assert_eq!(outcome.counters.unchanged, 2);
|
||||
|
||||
let repaired_meta = std::fs::read(&target_meta).expect("repaired target xl.meta must exist");
|
||||
assert_ne!(
|
||||
Sha256::digest(&repaired_meta),
|
||||
stale_digest,
|
||||
"deep heal must replace stale rejoined metadata"
|
||||
);
|
||||
|
||||
let marker_info = physical_version(target, bucket, object, &marker);
|
||||
assert!(
|
||||
marker_info.deleted && marker_info.is_latest,
|
||||
"target metadata must converge to delete-marker latest"
|
||||
);
|
||||
let historical_info = physical_version(target, bucket, object, &historical);
|
||||
assert!(!historical_info.is_latest, "historical version must no longer be marked latest");
|
||||
assert_eq!(
|
||||
read_version(&ecstore, bucket, object, &historical).await,
|
||||
historical_data,
|
||||
"historical version must remain readable after marker convergence"
|
||||
);
|
||||
let marker_receipt = outcome
|
||||
.objects
|
||||
.iter()
|
||||
.find(|item| item.identity.version_id.as_deref() == Some(marker.as_str()))
|
||||
.expect("deep heal must emit a receipt for the repaired delete marker");
|
||||
assert_eq!(marker_receipt.disposition, HealObjectDisposition::Repaired);
|
||||
manager.stop().await.expect("heal manager should stop");
|
||||
}
|
||||
|
||||
/// A missing xl.meta must produce an exact marker repair receipt, while a
|
||||
/// healthy replay must prove metadata health without claiming payload integrity.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
|
||||
Reference in New Issue
Block a user