feat(heal): publish verified MRF repair events

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-08 00:00:28 +08:00
parent 7c85c72fd1
commit 1f65bd828d
3 changed files with 349 additions and 9 deletions
+86
View File
@@ -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<str>,
pub object: Arc<str>,
pub version_id: Option<[u8; 16]>,
pub scope: Option<MrfScope>,
pub lease: Option<MrfIngressLease>,
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<std::sync::Mutex<std::collections::VecDeque<MrfRepairedEvent>>> = OnceLock::new();
static MRF_VERIFIED_REPAIR_EVENTS: OnceLock<std::sync::Mutex<std::collections::VecDeque<MrfVerifiedRepairEvent>>> =
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<MrfRepairedEvent> {
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<MrfVerifiedRepairEvent> {
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);
}
}
+78 -9
View File
@@ -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(&notice_targets);
}
publish_verified_mrf_repair_events(&notice_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<MrfRepairNoticeTarget>) {
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<MrfRepairNoticeTarget>) {
}
}
pub(super) fn mrf_verified_repair_event_for_target(
target: &MrfRepairNoticeTarget,
outcome: &crate::heal::outcome::HealObjectOutcome,
) -> Option<rustfs_common::mrf_channel::MrfVerifiedRepairEvent> {
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<String> {
match &task.heal_type {
HealType::ErasureSet { set_disk_id, .. } => Some(set_disk_id.clone()),
+185
View File
@@ -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::<str>::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();