From 1f65bd828dc416879cc3415c204aca34feccd63e Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 00:00:28 +0800 Subject: [PATCH] feat(heal): publish verified MRF repair events Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- crates/common/src/mrf_channel.rs | 86 ++++++++++ crates/heal/src/heal/manager/scheduler.rs | 87 ++++++++-- crates/heal/src/heal/manager/tests.rs | 185 ++++++++++++++++++++++ 3 files changed, 349 insertions(+), 9 deletions(-) diff --git a/crates/common/src/mrf_channel.rs b/crates/common/src/mrf_channel.rs index 83e1663ea..5104f2010 100644 --- a/crates/common/src/mrf_channel.rs +++ b/crates/common/src/mrf_channel.rs @@ -432,12 +432,33 @@ pub struct MrfRepairedEvent { pub version_id: Option<[u8; 16]>, } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum MrfVerifiedRepairDisposition { + Repaired, + VerifiedHealthy, + AuthoritativelyAbsent, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MrfVerifiedRepairEvent { + pub kind: MrfKind, + pub bucket: Arc, + pub object: Arc, + pub version_id: Option<[u8; 16]>, + pub scope: Option, + pub lease: Option, + pub bucket_incarnation_id: Uuid, + pub disposition: MrfVerifiedRepairDisposition, +} + /// Bound on the repaired-event backlog. Notices are best-effort hints; when /// the ring is full the oldest are dropped and the affected ledger entries /// simply expire through their own attempts/age limits. const MRF_REPAIRED_EVENT_CAP: usize = 4096; static MRF_REPAIRED_EVENTS: OnceLock>> = OnceLock::new(); +static MRF_VERIFIED_REPAIR_EVENTS: OnceLock>> = + OnceLock::new(); /// Record a legacy notification for compatibility. This is not an /// acknowledgement of storage verification or durable repair completion. @@ -478,6 +499,43 @@ pub fn take_mrf_repaired_events_for(bucket: &str) -> Vec { taken } +/// Record a storage-owned MRF completion proof. Unlike the legacy repaired +/// event, this identity is complete enough for future durable ledgers to make +/// an exact responsibility decision. +pub fn note_mrf_verified_repair(event: MrfVerifiedRepairEvent) { + let registry = MRF_VERIFIED_REPAIR_EVENTS.get_or_init(|| std::sync::Mutex::new(std::collections::VecDeque::new())); + let Ok(mut events) = registry.lock() else { + return; + }; + if events.len() >= MRF_REPAIRED_EVENT_CAP { + events.pop_front(); + } + events.push_back(event); +} + +/// Take verified repair events recorded for `bucket`, leaving other buckets' +/// proofs in place. Consumers still have to match kind, object, version, scope +/// lease and incarnation before discharging durable responsibility. +pub fn take_mrf_verified_repair_events_for(bucket: &str) -> Vec { + let Some(registry) = MRF_VERIFIED_REPAIR_EVENTS.get() else { + return Vec::new(); + }; + let Ok(mut events) = registry.lock() else { + return Vec::new(); + }; + let mut taken = Vec::new(); + let mut retained = std::collections::VecDeque::with_capacity(events.len()); + while let Some(event) = events.pop_front() { + if event.bucket.as_ref() == bucket { + taken.push(event); + } else { + retained.push_back(event); + } + } + *events = retained; + taken +} + #[cfg(test)] mod tests { use super::*; @@ -610,4 +668,32 @@ mod tests { assert_eq!(flooded.len(), MRF_REPAIRED_EVENT_CAP); assert_eq!(flooded[0].object.as_ref(), "object-9", "the oldest notices past the cap are dropped"); } + + #[test] + fn verified_repair_events_preserve_full_identity_and_bucket_scope() { + let bucket_incarnation_id = Uuid::new_v4(); + let event = MrfVerifiedRepairEvent { + kind: MrfKind::PartialWrite, + bucket: Arc::from("verified-bucket-a"), + object: Arc::from("object-a"), + version_id: Some([4u8; 16]), + scope: Some(MrfScope { + pool_index: 2, + set_index: 3, + }), + lease: Some(MrfIngressLease::new(42)), + bucket_incarnation_id, + disposition: MrfVerifiedRepairDisposition::Repaired, + }; + note_mrf_verified_repair(event.clone()); + note_mrf_verified_repair(MrfVerifiedRepairEvent { + bucket: Arc::from("verified-bucket-b"), + ..event.clone() + }); + + let taken = take_mrf_verified_repair_events_for("verified-bucket-a"); + assert_eq!(taken, vec![event]); + assert!(take_mrf_verified_repair_events_for("verified-bucket-a").is_empty()); + assert_eq!(take_mrf_verified_repair_events_for("verified-bucket-b").len(), 1); + } } diff --git a/crates/heal/src/heal/manager/scheduler.rs b/crates/heal/src/heal/manager/scheduler.rs index 6c5fc70a6..c9ea5a8a7 100644 --- a/crates/heal/src/heal/manager/scheduler.rs +++ b/crates/heal/src/heal/manager/scheduler.rs @@ -313,6 +313,7 @@ impl HealManager { completed_status_entry.outcome = Some(Arc::new(task.get_outcome().await)); } let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. }); + let completed_status_for_verified_events = completed_status_entry.clone(); // Keep retry ownership continuous: status snapshots acquire // these locks in the same active -> retrying order. let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) = @@ -358,6 +359,15 @@ impl HealManager { tests::pause_completed_retention_handoff(&task_id).await; if completed_task.is_some() { + let notice_targets = if terminal_completion { + take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id) + } else { + Vec::new() + }; + if terminal_completion { + release_mrf_repair_notice_targets(¬ice_targets); + } + publish_verified_mrf_repair_events(¬ice_targets, &completed_status_for_verified_events); // update statistics let mut stats = statistics_clone.write().await; match completed_status { @@ -372,14 +382,6 @@ impl HealManager { } stats.update_running_tasks(usize_to_u64_saturated(active_count)); drop(stats); - if terminal_completion { - let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id); - // Neither task status nor the diagnostic outcome - // window supplies a storage-owned repair receipt. - // Release only the ingress lease for rediscovery; - // preserve the producer's existing retry hints. - release_mrf_repair_notice_targets(notice_targets); - } } if let (Some((retry_request, retry_delay, retry_error)), Some(retry_cancel_token)) = @@ -705,7 +707,7 @@ fn move_mrf_repair_notice_targets( } } -fn release_mrf_repair_notice_targets(targets: Vec) { +fn release_mrf_repair_notice_targets(targets: &[MrfRepairNoticeTarget]) { for target in targets { rustfs_common::mrf_channel::release_mrf_identity( target.kind, @@ -718,6 +720,73 @@ fn release_mrf_repair_notice_targets(targets: Vec) { } } +pub(super) fn mrf_verified_repair_event_for_target( + target: &MrfRepairNoticeTarget, + outcome: &crate::heal::outcome::HealObjectOutcome, +) -> Option { + use crate::heal::outcome::{HealObjectDisposition, HealObjectKind}; + use rustfs_common::mrf_channel::{MrfKind, MrfVerifiedRepairDisposition}; + + let disposition = match outcome.disposition { + HealObjectDisposition::Repaired => MrfVerifiedRepairDisposition::Repaired, + HealObjectDisposition::VerifiedHealthy => MrfVerifiedRepairDisposition::VerifiedHealthy, + HealObjectDisposition::AuthoritativelyAbsent => MrfVerifiedRepairDisposition::AuthoritativelyAbsent, + _ => return None, + }; + if target.kind != MrfKind::PartialWrite { + return None; + } + let expected_kind = HealObjectKind::Object; + if outcome.identity.kind != expected_kind + || outcome.identity.bucket.as_str() != target.bucket.as_ref() + || outcome.identity.object.as_str() != target.object.as_ref() + { + return None; + } + let version_id = target.version_id.filter(|bytes| *bytes != [0; 16]); + let expected_version = version_id.map(|bytes| uuid::Uuid::from_bytes(bytes).to_string()); + if outcome.identity.version_id != expected_version { + return None; + } + let expected_pool = target.scope.and_then(|scope| usize::try_from(scope.pool_index).ok()); + let expected_set = target.scope.and_then(|scope| usize::try_from(scope.set_index).ok()); + if outcome.identity.pool_index != expected_pool || outcome.identity.set_index != expected_set { + return None; + } + let bucket_incarnation_id = outcome.identity.bucket_incarnation_id?; + Some(rustfs_common::mrf_channel::MrfVerifiedRepairEvent { + kind: target.kind, + bucket: target.bucket.clone(), + object: target.object.clone(), + version_id, + scope: target.scope, + lease: target.lease, + bucket_incarnation_id, + disposition, + }) +} + +pub(super) fn publish_verified_mrf_repair_events(targets: &[MrfRepairNoticeTarget], completed: &CompletedHealStatus) { + if completed.status != HealTaskStatus::Completed { + return; + } + let Some(outcome) = completed.outcome.as_ref() else { + return; + }; + if outcome.execution != crate::heal::outcome::HealExecutionOutcome::Completed { + return; + } + for target in targets { + if let Some(event) = outcome + .objects + .iter() + .find_map(|object| mrf_verified_repair_event_for_target(target, object)) + { + rustfs_common::mrf_channel::note_mrf_verified_repair(event); + } + } +} + pub(super) fn heal_request_set_key_for_task(task: &HealTask) -> Option { match &task.heal_type { HealType::ErasureSet { set_disk_id, .. } => Some(set_disk_id.clone()), diff --git a/crates/heal/src/heal/manager/tests.rs b/crates/heal/src/heal/manager/tests.rs index 12495bc2c..bbf102c72 100644 --- a/crates/heal/src/heal/manager/tests.rs +++ b/crates/heal/src/heal/manager/tests.rs @@ -1094,6 +1094,191 @@ fn queued_request_id_for_dedup_key_tracks_the_representative() { assert!(queue.queued_request_id_for_dedup_key(&first_key).is_none()); } +#[test] +fn mrf_verified_repair_event_requires_positive_exact_identity() { + use crate::heal::outcome::{HealObjectIdentity, HealObjectKind, HealObjectOutcome}; + use rustfs_common::mrf_channel::{MrfKind, MrfScope, MrfVerifiedRepairDisposition}; + + let version = uuid::Uuid::new_v4(); + let incarnation = uuid::Uuid::new_v4(); + let target = MrfRepairNoticeTarget { + bucket: Arc::from("bucket"), + object: Arc::from("object"), + version_id: Some(*version.as_bytes()), + kind: MrfKind::PartialWrite, + scope: Some(MrfScope { + pool_index: 1, + set_index: 2, + }), + lease: None, + }; + let matching = HealObjectOutcome { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: "bucket".to_string(), + object: "object".to_string(), + version_id: Some(version.to_string()), + bucket_incarnation_id: Some(incarnation), + pool_index: Some(1), + set_index: Some(2), + }, + disposition: HealObjectDisposition::Repaired, + detail: None, + }; + + let event = mrf_verified_repair_event_for_target(&target, &matching).expect("matching positive receipt should publish"); + assert_eq!(event.kind, MrfKind::PartialWrite); + assert_eq!(event.bucket.as_ref(), "bucket"); + assert_eq!(event.object.as_ref(), "object"); + assert_eq!(event.version_id, Some(*version.as_bytes())); + assert_eq!( + event.scope, + Some(MrfScope { + pool_index: 1, + set_index: 2 + }) + ); + assert_eq!(event.lease, None); + assert_eq!(event.bucket_incarnation_id, incarnation); + assert_eq!(event.disposition, MrfVerifiedRepairDisposition::Repaired); + + assert!( + mrf_verified_repair_event_for_target( + &MrfRepairNoticeTarget { + kind: MrfKind::DecodeFailure, + ..target.clone() + }, + &matching + ) + .is_none(), + "only receipt-producing partial-write object heals can publish verified events today" + ); + + for rejected in [ + HealObjectOutcome { + disposition: HealObjectDisposition::Unknown, + ..matching.clone() + }, + HealObjectOutcome { + identity: HealObjectIdentity { + bucket_incarnation_id: None, + ..matching.identity.clone() + }, + ..matching.clone() + }, + HealObjectOutcome { + identity: HealObjectIdentity { + object: "other".to_string(), + ..matching.identity.clone() + }, + ..matching.clone() + }, + HealObjectOutcome { + identity: HealObjectIdentity { + pool_index: Some(3), + ..matching.identity.clone() + }, + ..matching + }, + ] { + assert!( + mrf_verified_repair_event_for_target(&target, &rejected).is_none(), + "legacy, incomplete or mismatched outcomes must not discharge MRF responsibility" + ); + } +} + +#[test] +fn completed_mrf_notice_publishes_only_verified_positive_events() { + use crate::heal::outcome::{HealObjectIdentity, HealObjectKind, HealObjectOutcome, HealTaskOutcome}; + use rustfs_common::mrf_channel::{MrfKind, MrfScope, take_mrf_verified_repair_events_for}; + + let bucket = Arc::::from("verified-mrf-completed-bucket"); + let _ = take_mrf_verified_repair_events_for(bucket.as_ref()); + let version = uuid::Uuid::new_v4(); + let incarnation = uuid::Uuid::new_v4(); + let matching_target = MrfRepairNoticeTarget { + bucket: bucket.clone(), + object: Arc::from("object-a"), + version_id: Some(*version.as_bytes()), + kind: MrfKind::PartialWrite, + scope: Some(MrfScope { + pool_index: 1, + set_index: 2, + }), + lease: None, + }; + let mismatch_target = MrfRepairNoticeTarget { + object: Arc::from("object-b"), + ..matching_target.clone() + }; + let mut outcome = HealTaskOutcome::default(); + outcome.execution = HealExecutionOutcome::Completed; + outcome.coverage = crate::heal::outcome::HealTraversalCoverage::Complete; + outcome.record(HealObjectOutcome { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: bucket.to_string(), + object: "object-a".to_string(), + version_id: Some(version.to_string()), + bucket_incarnation_id: Some(incarnation), + pool_index: Some(1), + set_index: Some(2), + }, + disposition: HealObjectDisposition::VerifiedHealthy, + detail: None, + }); + outcome.record(HealObjectOutcome { + identity: HealObjectIdentity { + kind: HealObjectKind::Object, + bucket: bucket.to_string(), + object: "object-b".to_string(), + version_id: Some(version.to_string()), + bucket_incarnation_id: None, + pool_index: Some(1), + set_index: Some(2), + }, + disposition: HealObjectDisposition::Repaired, + detail: None, + }); + let completed = CompletedHealStatus { + outcome: Some(Arc::new(outcome)), + ..completed_retention_fixture(SystemTime::now()) + }; + + publish_verified_mrf_repair_events(&[matching_target.clone(), mismatch_target], &completed); + + let events = take_mrf_verified_repair_events_for(bucket.as_ref()); + assert_eq!(events.len(), 1); + assert_eq!(events[0].object.as_ref(), "object-a"); + assert_eq!(events[0].lease, None); + assert_eq!(events[0].bucket_incarnation_id, incarnation); + + let failed = CompletedHealStatus { + status: HealTaskStatus::Failed { + error: "terminal failure".to_string(), + }, + ..completed.clone() + }; + publish_verified_mrf_repair_events(std::slice::from_ref(&matching_target), &failed); + assert!( + take_mrf_verified_repair_events_for(bucket.as_ref()).is_empty(), + "failed terminal tasks must not publish a verified repair event" + ); + + let mut completed_with_errors_outcome = completed.outcome.as_ref().expect("completed outcome").as_ref().clone(); + completed_with_errors_outcome.execution = crate::heal::outcome::HealExecutionOutcome::CompletedWithErrors; + let completed_with_errors = CompletedHealStatus { + outcome: Some(Arc::new(completed_with_errors_outcome)), + ..completed + }; + publish_verified_mrf_repair_events(std::slice::from_ref(&matching_target), &completed_with_errors); + assert!( + take_mrf_verified_repair_events_for(bucket.as_ref()).is_empty(), + "non-success canonical outcomes must not publish a verified repair event" + ); +} + #[test] fn test_priority_queue_ordering() { let mut queue = PriorityHealQueue::new();