mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-08 04:58:12 +00:00
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 <qiuzgang@gmail.com>
This commit is contained in:
@@ -90,6 +90,61 @@ pub struct MrfIntent {
|
||||
pub attempts: u8,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub struct MrfDurableRepairAnchor {
|
||||
pub kind: MrfKind,
|
||||
pub bucket: Arc<str>,
|
||||
pub object: Arc<str>,
|
||||
pub version_id: Option<[u8; 16]>,
|
||||
pub scope: Option<MrfScope>,
|
||||
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<Self> {
|
||||
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<MrfDurableRepairAnchor>, 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::<str>::from("proof-bucket");
|
||||
let object = Arc::<str>::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");
|
||||
|
||||
Reference in New Issue
Block a user