From 002b0000557347ba85f67ad0e56a417dc6a93a70 Mon Sep 17 00:00:00 2001 From: cxymds Date: Tue, 15 Sep 2026 08:16:20 +0800 Subject: [PATCH] fix(heal): report verified normal shard repairs (#7890) --- crates/ecstore/src/set_disk/ops/heal.rs | 11 ++ crates/heal/src/heal/storage.rs | 25 ++- .../heal/tests/shard_identity_receipt_test.rs | 159 ++++++++++++++++++ crates/madmin/src/heal_commands.rs | 7 + 4 files changed, 199 insertions(+), 3 deletions(-) diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index fba1c496c..e187c4139 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -1732,6 +1732,17 @@ impl SetDisks { return Ok((result, Some(error))); } + result.repair_verified = protected + && !latest_meta.deleted + && !latest_meta.is_remote() + && !read_repair_uses_shared_lock + && result.drives_healed().is_some_and(|healed| healed > 0) + && result + .after + .drives + .iter() + .all(|drive| drive.state == DriveState::Ok.to_string()); + // The object is healthy here; sweep any data dirs left behind // by pre-#3510 unversioned overwrites, which the dangling paths // above never touch (issues #3231, #3191). Best effort — a diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index 6558a86c0..4566949e1 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -105,7 +105,7 @@ fn verified_object_receipt( item: &HealResultItem, bucket_incarnation_id: Uuid, ) -> Option { - if opts.dry_run || !item.integrity_verified { + if opts.dry_run || (!item.integrity_verified && !item.repair_verified) { return None; } let resolved_version = Uuid::from_bytes(item.resolved_version_id?); @@ -116,6 +116,9 @@ fn verified_object_receipt( } item.drives_reported()?; let drives_healed = item.drives_healed()?; + if item.repair_verified && drives_healed == 0 { + return None; + } let ok_drive_state = DriveState::Ok.to_string(); if !item.after.drives.iter().all(|drive| drive.state == ok_drive_state) { return None; @@ -132,8 +135,10 @@ fn verified_object_receipt( }, disposition: if drives_healed > 0 { HealObjectDisposition::Repaired - } else { + } else if item.integrity_verified { HealObjectDisposition::VerifiedHealthy + } else { + return None; }, }) } @@ -682,7 +687,7 @@ impl ECStoreHealStorage { } else { None } - } else if error.is_none() && !opts.dry_run && item.integrity_verified { + } else if error.is_none() && !opts.dry_run && (item.integrity_verified || item.repair_verified) { let bucket_incarnation_id = match expected { Some(expected) => Some(expected), None => self.ecstore.bucket_incarnation_id(bucket).await.ok(), @@ -1848,6 +1853,20 @@ mod tests { "the exact version cannot certify unverified shard integrity" ); assert!(verified_object_receipt("bucket", "object", None, &options, &item, incarnation).is_none()); + item.repair_verified = true; + item.resolved_version_id = Some([0; 16]); + assert_eq!( + verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation) + .expect("a protected repaired shard should carry a repair receipt") + .disposition, + HealObjectDisposition::Repaired + ); + item.before.drives[0].state = "ok".to_string(); + assert!( + verified_object_receipt("bucket", "object", Some(&null), &options, &item, incarnation).is_none(), + "repair proof without a repaired drive must remain unresolved" + ); + item.repair_verified = false; item.integrity_verified = true; item.resolved_version_id = None; assert!( diff --git a/crates/heal/tests/shard_identity_receipt_test.rs b/crates/heal/tests/shard_identity_receipt_test.rs index a2278b324..6243582e1 100644 --- a/crates/heal/tests/shard_identity_receipt_test.rs +++ b/crates/heal/tests/shard_identity_receipt_test.rs @@ -379,3 +379,162 @@ fn admin_pool_set_repairs_five_shards_and_retains_exact_terminal_after_restart() .join() .expect("heal fixture result"); } + +#[test] +#[serial] +fn admin_normal_protected_repairs_report_repaired_for_data_and_parity() { + std::thread::Builder::new() + .stack_size(8 * 1024 * 1024) + .spawn(|| { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("heal fixture runtime") + .block_on(async { + use rustfs_heal::heal::{ + manager::{HealConfig, HealManager}, + task::{HealOptions, HealPriority, HealRequest, HealTaskStatus, HealType}, + }; + use rustfs_heal_contracts::heal_channel::HealRequestSource; + use std::sync::Arc; + use std::time::Duration; + use storage_api::integration::WriteCompletion; + + let root = tempfile::tempdir().expect("normal repair fixture"); + let env = TestECStoreEnv::builder() + .base_dir(root.path()) + .prefix("admin_normal_repair_receipts") + .build() + .await; + let bucket = "admin-normal-repair"; + env.make_bucket(bucket, false).await; + let set = env.ecstore.pools[0].get_disks(0); + let disks = set.disks.read().await.iter().flatten().cloned().collect::>(); + let mut faults = Vec::new(); + for (object, want_data, byte) in [("missing-data", true, 0x31_u8), ("missing-parity", false, 0x72_u8)] { + let body = vec![byte; 1024 * 1024 + 37]; + env.ecstore + .put_object( + bucket, + object, + &mut PutObjReader::from_vec(body.clone()), + &ObjectOptions { + write_completion: WriteCompletion::TailDrained, + shard_integrity_write_mode: Some(ShardIntegrityWriteMode::Protected), + ..Default::default() + }, + ) + .await + .expect("commit protected fixture"); + let mut metadata = Vec::new(); + for disk in &disks { + metadata.push( + disk.read_version("", bucket, object, "", &ReadOptions::default()) + .await + .expect("fixture metadata"), + ); + } + let target_slot = metadata + .iter() + .position(|part| (part.erasure.index <= part.erasure.data_blocks) == want_data) + .expect("fixture must expose both data and parity shards"); + let target = &metadata[target_slot]; + let path = env.disk_paths[target_slot] + .join(bucket) + .join(object) + .join(target.data_dir.expect("external shard").to_string()) + .join("part.1"); + let original = tokio::fs::read(&path).await.expect("original shard bytes"); + tokio::fs::remove_file(&path).await.expect("inject missing shard"); + faults.push((object.to_owned(), path, original, body)); + } + + let storage = Arc::new(ECStoreHealStorage::new(env.ecstore.clone())); + let manager = HealManager::new( + storage.clone(), + Some(HealConfig { + enable_auto_heal: false, + ..Default::default() + }), + ); + manager.start().await.expect("start heal manager"); + let mut request = HealRequest::new( + HealType::ErasureSet { + buckets: Vec::new(), + set_disk_id: "pool_0_set_0".to_owned(), + }, + HealOptions { + recursive: true, + scan_mode: HealScanMode::Normal, + pool_index: Some(0), + set_index: Some(0), + timeout: Some(Duration::from_secs(60)), + ..Default::default() + }, + HealPriority::High, + ); + request.source = HealRequestSource::Admin; + let token = request.id.clone(); + manager + .submit_heal_request(request) + .await + .expect("admit normal all-buckets heal"); + let terminal = tokio::time::timeout(Duration::from_secs(60), async { + loop { + let report = manager.get_task_report(&token).await.expect("same-token report"); + if matches!( + report.status, + HealTaskStatus::Completed | HealTaskStatus::Failed { .. } | HealTaskStatus::Timeout + ) { + break report; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + }) + .await + .expect("terminal deadline"); + assert_eq!(terminal.status, HealTaskStatus::Completed); + let outcome = terminal.outcome.as_deref().expect("canonical terminal"); + assert_eq!( + ( + outcome.counters.processed, + outcome.counters.healed, + outcome.counters.unchanged, + outcome.counters.skipped, + outcome.counters.failed, + outcome.counters.unknown, + ), + (2, 2, 0, 0, 0, 0), + "normal protected repairs must be canonical: {outcome:?}" + ); + assert_eq!(outcome.objects.len(), 2); + for (object, path, original, body) in &faults { + assert_eq!(&tokio::fs::read(path).await.expect("rebuilt physical shard"), original); + let receipt = outcome + .objects + .iter() + .find(|receipt| receipt.identity.bucket == bucket && receipt.identity.object == *object) + .expect("one receipt per repaired object"); + assert_eq!(receipt.disposition, HealObjectDisposition::Repaired); + assert_eq!((receipt.identity.pool_index, receipt.identity.set_index), (Some(0), Some(0))); + assert!(receipt.identity.version_id.is_some()); + assert_eq!( + receipt.identity.bucket_incarnation_id, + Some(storage.admit_bucket_incarnation(bucket).await.expect("original bucket")) + ); + let mut reader = env + .ecstore + .get_object_reader(bucket, object, None, Default::default(), &ObjectOptions::default()) + .await + .expect("read repaired object"); + let mut actual = Vec::new(); + reader.stream.read_to_end(&mut actual).await.expect("read repaired body"); + assert_eq!(&actual, body); + } + manager.stop().await.expect("stop heal manager"); + }); + }) + .expect("heal fixture thread") + .join() + .expect("heal fixture result"); +} diff --git a/crates/madmin/src/heal_commands.rs b/crates/madmin/src/heal_commands.rs index b2ac27930..8c65f2ef8 100644 --- a/crates/madmin/src/heal_commands.rs +++ b/crates/madmin/src/heal_commands.rs @@ -40,6 +40,11 @@ pub struct HealResultItem { /// or old-peer results default to unproven and cannot advance heal receipts. #[serde(skip)] pub integrity_verified: bool, + /// Storage-owner proof that a protected missing shard was reconstructed and + /// committed. This is intentionally in-process only and does not certify a + /// healthy or unchanged object. + #[serde(skip)] + pub repair_verified: bool, #[serde(rename = "resultId")] pub result_index: usize, #[serde(rename = "type")] @@ -113,6 +118,8 @@ mod tests { let mut wire = serde_json::to_value(&item).expect("heal result should serialize"); assert!(wire.get("resolved_version_id").is_none()); assert!(wire.get("resolvedVersionId").is_none()); + assert!(wire.get("repair_verified").is_none()); + assert!(wire.get("repairVerified").is_none()); wire["resolved_version_id"] = serde_json::json!(vec![0; 16]); let decoded: HealResultItem = serde_json::from_value(wire).expect("legacy wire shape should remain readable"); assert_eq!(decoded.resolved_version_id, None, "wire input cannot supply owner proof");