From c8fc50ada247cd53d99ecb1d80c93197b342b504 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 7 Sep 2026 19:55:29 +0800 Subject: [PATCH] feat(heal): accept verified object repair receipts (#7397) Introduce an object heal receipt wrapper so storage owners can report a positive repaired or verified result without changing legacy heal consumers. Task-level object heal now records only receipts that match the requested object identity, version, pool, set, and carry a bucket incarnation. Legacy and mismatched receipts still fall back to unknown outcome accounting. Co-authored-by: zhi22915 --- crates/heal/src/heal/outcome.rs | 67 +++++++++++++++ crates/heal/src/heal/storage.rs | 74 +++++++++++++++- crates/heal/src/heal/task.rs | 22 ++++- crates/heal/src/heal/task/heal_object.rs | 10 ++- crates/heal/src/heal/task/tests.rs | 103 +++++++++++++++++++++++ 5 files changed, 272 insertions(+), 4 deletions(-) diff --git a/crates/heal/src/heal/outcome.rs b/crates/heal/src/heal/outcome.rs index a283b0117..c8b48da6d 100644 --- a/crates/heal/src/heal/outcome.rs +++ b/crates/heal/src/heal/outcome.rs @@ -91,6 +91,29 @@ pub struct HealObjectOutcome { pub detail: Option, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct HealObjectReceipt { + pub identity: HealObjectIdentity, + pub disposition: HealObjectDisposition, +} + +impl HealObjectReceipt { + pub(crate) fn verified_for(&self, expected: &HealObjectIdentity) -> bool { + matches!( + self.disposition, + HealObjectDisposition::Repaired + | HealObjectDisposition::VerifiedHealthy + | HealObjectDisposition::AuthoritativelyAbsent + ) && self.identity.kind == expected.kind + && self.identity.bucket == expected.bucket + && self.identity.object == expected.object + && self.identity.version_id == expected.version_id + && self.identity.pool_index == expected.pool_index + && self.identity.set_index == expected.set_index + && self.identity.bucket_incarnation_id.is_some() + } +} + impl HealObjectOutcome { fn retained_bytes(&self) -> usize { size_of::() @@ -470,4 +493,48 @@ mod canonical_outcome_tests { assert_eq!(outcome.counters.processed, u64::MAX); assert_eq!(outcome.coverage, HealTraversalCoverage::Partial); } + + #[test] + fn positive_receipt_requires_exact_identity_and_bucket_incarnation() { + let expected = item(HealObjectDisposition::Unknown).identity; + let mut receipt = HealObjectReceipt { + identity: expected.clone(), + disposition: HealObjectDisposition::Repaired, + }; + + assert!( + !receipt.verified_for(&expected), + "a positive storage receipt without bucket incarnation must remain untrusted" + ); + + let incarnation = Uuid::new_v4(); + receipt.identity.bucket_incarnation_id = Some(incarnation); + assert!(receipt.verified_for(&expected)); + + receipt.identity.version_id = Some("older-version".to_string()); + assert!( + !receipt.verified_for(&expected), + "a storage receipt for a different object/version tuple must not clear the requested responsibility" + ); + + receipt.identity = HealObjectIdentity { + bucket_incarnation_id: Some(incarnation), + pool_index: Some(1), + ..expected.clone() + }; + assert!( + !receipt.verified_for(&expected), + "a storage receipt for a different erasure location must not clear the requested responsibility" + ); + + receipt.identity = HealObjectIdentity { + bucket_incarnation_id: Some(incarnation), + ..expected + }; + receipt.disposition = HealObjectDisposition::Unknown; + assert!( + !receipt.verified_for(&receipt.identity), + "legacy success without a positive disposition remains unknown" + ); + } } diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index b6c5731c9..c4554b327 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -15,12 +15,13 @@ use crate::{Error, Result}; use async_trait::async_trait; use base64_simd::URL_SAFE_NO_PAD; -use rustfs_heal_contracts::heal_channel::{HealOpts, HealScanMode}; +use rustfs_heal_contracts::heal_channel::{DriveState, HealOpts, HealScanMode}; use rustfs_madmin::heal_commands::HealResultItem; use serde::{Deserialize, Serialize}; use std::sync::Arc; use tracing::{debug, error, warn}; +use super::outcome::{HealObjectDisposition, HealObjectIdentity, HealObjectKind, HealObjectReceipt}; use super::progress::stable_generation; use super::storage_api::owner::{EcstoreHealLifecycleExpiryContext, ecstore_load_admin_data_usage_from_backend_cached}; use super::storage_api::storage::{ @@ -67,6 +68,23 @@ impl HealLifecycleExpiryContext { } } +#[derive(Debug, Default)] +pub struct HealStorageObjectResult { + pub item: HealResultItem, + pub error: Option, + pub receipt: Option, +} + +impl From<(HealResultItem, Option)> for HealStorageObjectResult { + fn from((item, error): (HealResultItem, Option)) -> Self { + Self { + item, + error, + receipt: None, + } + } +} + const LOG_COMPONENT_HEAL: &str = "heal"; const LOG_SUBSYSTEM_STORAGE: &str = "storage"; const EVENT_HEAL_STORAGE_OBJECT_IO: &str = "heal_storage_object_io"; @@ -374,6 +392,16 @@ pub trait HealStorageAPI: Send + Sync { opts: &HealOpts, ) -> Result<(HealResultItem, Option)>; + async fn heal_object_with_receipt( + &self, + bucket: &str, + object: &str, + version_id: Option<&str>, + opts: &HealOpts, + ) -> Result { + self.heal_object(bucket, object, version_id, opts).await.map(Into::into) + } + /// Heal bucket using ecstore async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result; @@ -1062,6 +1090,50 @@ impl HealStorageAPI for ECStoreHealStorage { } } + async fn heal_object_with_receipt( + &self, + bucket: &str, + object: &str, + version_id: Option<&str>, + opts: &HealOpts, + ) -> Result { + let (item, error) = self.heal_object(bucket, object, version_id, opts).await?; + let receipt = if error.is_none() && !opts.dry_run { + let ok_drive_state = DriveState::Ok.to_string(); + let all_after_drives_ok = item.after.drives.iter().all(|drive| drive.state == ok_drive_state); + match ( + self.ecstore.bucket_incarnation_id(bucket).await, + item.drives_reported(), + item.drives_healed(), + all_after_drives_ok, + ) { + (Ok(bucket_incarnation_id), Some(_), Some(drives_healed), true) => { + let disposition = if drives_healed > 0 { + HealObjectDisposition::Repaired + } else { + HealObjectDisposition::VerifiedHealthy + }; + Some(HealObjectReceipt { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: bucket.to_string(), + object: object.to_string(), + version_id: version_id.map(ToOwned::to_owned), + bucket_incarnation_id: Some(bucket_incarnation_id), + pool_index: opts.pool, + set_index: opts.set, + }, + disposition, + }) + } + _ => None, + } + } else { + None + }; + Ok(HealStorageObjectResult { item, error, receipt }) + } + async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result { debug!( target: "rustfs::heal::storage", diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index f1707be24..bb2c659b1 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -17,7 +17,7 @@ use crate::heal::{ erasure_healer::target_outcomes_complete, outcome::{ HealAbortReason, HealDeferredReason, HealFailureClass, HealObjectDisposition, HealObjectIdentity, HealObjectKind, - HealObjectOutcome, HealTaskOutcome, + HealObjectOutcome, HealObjectReceipt, HealTaskOutcome, }, progress::HealProgress, resume::{ @@ -592,6 +592,26 @@ impl HealTask { Some(self.outcome_identity(bucket, object, version, self.options.pool_index, self.options.set_index)) } + pub(super) async fn record_verified_storage_receipt( + &self, + expected: HealObjectIdentity, + receipt: Option, + ) -> bool { + let Some(receipt) = receipt else { + return false; + }; + if !receipt.verified_for(&expected) { + return false; + } + let mut outcome = self.outcome.write().await; + outcome.record(HealObjectOutcome { + identity: receipt.identity, + disposition: receipt.disposition, + detail: None, + }); + true + } + async fn record_deferred_object(&self, reason: HealDeferredReason) { if let Some(identity) = self.single_object_identity() { let mut outcome = self.outcome.write().await; diff --git a/crates/heal/src/heal/task/heal_object.rs b/crates/heal/src/heal/task/heal_object.rs index eebb4a95d..9e2c859e7 100644 --- a/crates/heal/src/heal/task/heal_object.rs +++ b/crates/heal/src/heal/task/heal_object.rs @@ -163,7 +163,7 @@ impl HealTask { set: self.options.set_index, }; - let heal_fut = self.storage.heal_object(bucket, object, version_id, &heal_opts); + let heal_fut = self.storage.heal_object_with_receipt(bucket, object, version_id, &heal_opts); let heal_result = if self.source == HealRequestSource::ReadRepair { let result = heal_fut.await; if self.cancel_token.is_cancelled() { @@ -176,7 +176,9 @@ impl HealTask { }; match heal_result { - Ok((result, error)) => { + Ok(storage_result) => { + let result = storage_result.item; + let error = storage_result.error; if let Some(e) = error { if self.skip_dangling_delete_grace_error(bucket, object, &e).await { return Ok(()); @@ -264,6 +266,10 @@ impl HealTask { let mut progress = self.progress.write().await; progress.update_object_progress(1, 1, 0, 0, object_size); } + let expected_identity = + self.outcome_identity(bucket, object, version_id, self.options.pool_index, self.options.set_index); + self.record_verified_storage_receipt(expected_identity, storage_result.receipt) + .await; self.record_result_item(result).await; Ok(()) } diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 23118c92b..30ac040ba 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -14,6 +14,7 @@ use super::super::{DiskOption, DiskStore, Endpoint, new_disk}; use super::*; +use crate::heal::storage::HealStorageObjectResult; mod deferred_retry; @@ -1048,6 +1049,7 @@ struct MockStorage { object_exists_by_name: Mutex>, heal_object_outcome: Mutex>, heal_object_outcomes: Mutex>>, + heal_object_receipts: Mutex>>, format_no_heal_required: Mutex, format_error: Mutex>, global_format_calls: Mutex, @@ -1151,6 +1153,90 @@ async fn execute_emits_heal_trace_task_state() { assert_eq!(trace_attr_string(&completed, "state").as_deref(), Some("completed")); } +fn object_receipt(object: &str, version_id: Option<&str>, disposition: HealObjectDisposition) -> HealObjectReceipt { + HealObjectReceipt { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: "bucket-a".to_string(), + object: object.to_string(), + version_id: version_id.map(ToOwned::to_owned), + bucket_incarnation_id: Some(Uuid::new_v4()), + pool_index: None, + set_index: None, + }, + disposition, + } +} + +#[tokio::test] +async fn object_heal_records_matching_positive_storage_receipt() { + let storage = Arc::new(MockStorage { + heal_object_receipts: Mutex::new(HashMap::from([( + "object-a".to_string(), + VecDeque::from([object_receipt("object-a", Some("version-a"), HealObjectDisposition::Repaired)]), + )])), + ..Default::default() + }); + let task = HealTask::from_request( + HealRequest::object("bucket-a".to_string(), "object-a".to_string(), Some("version-a".to_string())), + storage, + ); + + task.execute().await.expect("mock object heal should complete"); + + let outcome = task.get_outcome().await; + assert_eq!(outcome.counters.healed, 1); + assert_eq!(outcome.counters.unknown, 0); + let object = outcome.objects.front().expect("positive receipt should be recorded"); + assert_eq!(object.identity.object, "object-a"); + assert_eq!(object.identity.version_id.as_deref(), Some("version-a")); + assert!(object.identity.bucket_incarnation_id.is_some()); + assert_eq!(object.disposition, HealObjectDisposition::Repaired); +} + +#[tokio::test] +async fn object_heal_rejects_mismatched_or_legacy_storage_receipts() { + let storage = Arc::new(MockStorage { + heal_object_receipts: Mutex::new(HashMap::from([( + "object-a".to_string(), + VecDeque::from([object_receipt( + "object-a", + Some("old-version"), + HealObjectDisposition::Repaired, + )]), + )])), + ..Default::default() + }); + let task = HealTask::from_request( + HealRequest::object("bucket-a".to_string(), "object-a".to_string(), Some("version-a".to_string())), + storage, + ); + + task.execute() + .await + .expect("a mismatched receipt must not fail the legacy heal result"); + let outcome = task.get_outcome().await; + assert_eq!(outcome.counters.healed, 0); + assert_eq!(outcome.counters.unknown, 1); + assert_eq!( + outcome + .objects + .front() + .expect("legacy fallback should be recorded") + .disposition, + HealObjectDisposition::Unknown + ); + + let legacy = HealTask::from_request( + HealRequest::object("bucket-a".to_string(), "object-b".to_string(), None), + Arc::new(MockStorage::default()), + ); + legacy.execute().await.expect("legacy mock object heal should complete"); + let legacy_outcome = legacy.get_outcome().await; + assert_eq!(legacy_outcome.counters.healed, 0); + assert_eq!(legacy_outcome.counters.unknown, 1); +} + async fn recv_trace_task_state(trace: &mut TraceSubscription, task_id: &str, state: &str) -> TraceEvent { for _ in 0..32 { let event = tokio::time::timeout(Duration::from_secs(1), trace.recv()) @@ -1408,6 +1494,23 @@ impl HealStorageAPI for MockStorage { )) } + async fn heal_object_with_receipt( + &self, + bucket: &str, + object: &str, + version_id: Option<&str>, + opts: &HealOpts, + ) -> Result { + let (item, error) = self.heal_object(bucket, object, version_id, opts).await?; + let receipt = self + .heal_object_receipts + .lock() + .unwrap() + .get_mut(object) + .and_then(VecDeque::pop_front); + Ok(HealStorageObjectResult { item, error, receipt }) + } + async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result { self.bucket_heal_calls.lock().unwrap().push(bucket.to_string()); self.bucket_heal_opts.lock().unwrap().push(*opts);