From 88f8a7fb2122b4e4b5c6be1385d3acaf5639bd25 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 00:55:39 +0800 Subject: [PATCH] feat(common): add durable MRF proof matching (#7429) Add a fail-closed MRF repair proof adapter that only consumes anchors when the durable anchor and verified proof share the same full identity, ingress lease, and bucket incarnation. Legacy replay intents without leases cannot become dischargeable anchors, so the current retained journal behavior remains unchanged until a durable writer and producer proof source are connected. Co-authored-by: zhi22915 --- crates/common/src/mrf_channel.rs | 169 +++++++++++++++++++++++++++++++ 1 file changed, 169 insertions(+) diff --git a/crates/common/src/mrf_channel.rs b/crates/common/src/mrf_channel.rs index 5104f2010..a4f1b78be 100644 --- a/crates/common/src/mrf_channel.rs +++ b/crates/common/src/mrf_channel.rs @@ -90,6 +90,61 @@ pub struct MrfIntent { pub attempts: u8, } +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct MrfDurableRepairAnchor { + pub kind: MrfKind, + pub bucket: Arc, + pub object: Arc, + pub version_id: Option<[u8; 16]>, + pub scope: Option, + pub lease: MrfIngressLease, + pub bucket_incarnation_id: Uuid, +} + +impl MrfDurableRepairAnchor { + /// Build a dischargeable anchor only when the caller supplies the storage + /// incarnation and the original ingress lease. Legacy replay records lack + /// both pieces and therefore remain fail-closed. + pub fn from_intent(intent: &MrfIntent, bucket_incarnation_id: Uuid) -> Option { + if bucket_incarnation_id.is_nil() { + return None; + } + let lease = intent.lease?; + let (version_id, scope) = canonical_identity(intent.kind, intent.version_id, intent.scope); + Some(Self { + kind: intent.kind, + bucket: intent.bucket.clone(), + object: intent.object.clone(), + version_id, + scope, + lease, + bucket_incarnation_id, + }) + } + + pub fn is_proven_by(&self, event: &MrfVerifiedRepairEvent) -> bool { + let Some(lease) = event.lease else { + return false; + }; + self.kind == event.kind + && self.bucket == event.bucket + && self.object == event.object + && self.version_id == event.version_id + && self.scope == event.scope + && self.lease == lease + && self.bucket_incarnation_id == event.bucket_incarnation_id + } +} + +/// Consume only anchors proven by a complete verified-repair identity. The +/// caller remains responsible for persisting the resulting anchor set before +/// deleting older replay files. +pub fn consume_verified_mrf_repair_events(anchors: &mut Vec, events: &[MrfVerifiedRepairEvent]) -> usize { + let before = anchors.len(); + anchors.retain(|anchor| !events.iter().any(|event| anchor.is_proven_by(event))); + before.saturating_sub(anchors.len()) +} + #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] pub struct MrfScope { pub pool_index: u32, @@ -604,6 +659,120 @@ mod tests { assert_eq!(metadata_scope, None); } + #[test] + fn durable_repair_anchor_requires_lease_and_bucket_incarnation() { + let mut intent = MrfIntent { + bucket: Arc::from("durable-anchor-bucket"), + object: Arc::from("object"), + version_id: Some([0; 16]), + kind: MrfKind::PartialWrite, + scope: Some(MrfScope { + pool_index: 1, + set_index: 2, + }), + lease: None, + enqueued_at_ms: 0, + attempts: 0, + }; + assert!( + MrfDurableRepairAnchor::from_intent(&intent, Uuid::new_v4()).is_none(), + "legacy replay records without the ingress lease must remain anchored" + ); + + intent.lease = Some(MrfIngressLease::new(7)); + assert!( + MrfDurableRepairAnchor::from_intent(&intent, Uuid::nil()).is_none(), + "nil bucket incarnation cannot prove durable successor ownership" + ); + + let anchor = MrfDurableRepairAnchor::from_intent(&intent, Uuid::new_v4()) + .expect("complete identity should create a durable repair anchor"); + assert_eq!(anchor.version_id, None, "nil UUID is canonicalized before matching"); + assert_eq!( + anchor.scope, + Some(MrfScope { + pool_index: 1, + set_index: 2 + }) + ); + } + + #[test] + fn verified_repair_events_consume_only_exact_durable_anchors() { + let bucket = Arc::::from("proof-bucket"); + let object = Arc::::from("object"); + let incarnation = Uuid::new_v4(); + let lease = MrfIngressLease::new(11); + let anchor = MrfDurableRepairAnchor { + kind: MrfKind::PartialWrite, + bucket: bucket.clone(), + object: object.clone(), + version_id: Some([3; 16]), + scope: Some(MrfScope { + pool_index: 4, + set_index: 5, + }), + lease, + bucket_incarnation_id: incarnation, + }; + let event = MrfVerifiedRepairEvent { + kind: anchor.kind, + bucket, + object, + version_id: anchor.version_id, + scope: anchor.scope, + lease: Some(lease), + bucket_incarnation_id: incarnation, + disposition: MrfVerifiedRepairDisposition::Repaired, + }; + + for rejected in [ + MrfVerifiedRepairEvent { + lease: None, + ..event.clone() + }, + MrfVerifiedRepairEvent { + lease: Some(MrfIngressLease::new(12)), + ..event.clone() + }, + MrfVerifiedRepairEvent { + bucket_incarnation_id: Uuid::new_v4(), + ..event.clone() + }, + MrfVerifiedRepairEvent { + version_id: Some([4; 16]), + ..event.clone() + }, + MrfVerifiedRepairEvent { + scope: Some(MrfScope { + pool_index: 4, + set_index: 6, + }), + ..event.clone() + }, + MrfVerifiedRepairEvent { + kind: MrfKind::DecodeFailure, + ..event.clone() + }, + MrfVerifiedRepairEvent { + bucket: Arc::from("other-bucket"), + ..event.clone() + }, + MrfVerifiedRepairEvent { + object: Arc::from("other"), + ..event.clone() + }, + ] { + let mut retained = vec![anchor.clone()]; + assert_eq!(consume_verified_mrf_repair_events(&mut retained, &[rejected]), 0); + assert_eq!(retained, vec![anchor.clone()]); + } + + let mut retained = vec![anchor]; + assert_eq!(consume_verified_mrf_repair_events(&mut retained, &[event]), 1); + assert!(retained.is_empty()); + } + #[tokio::test] async fn try_send_delivers_and_respects_capacity() { let mut receiver = init_mrf_channel().expect("first initialization should succeed");