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");