diff --git a/crates/common/src/mrf_channel.rs b/crates/common/src/mrf_channel.rs index 83e1663ea..7f782bebf 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, @@ -386,6 +441,27 @@ pub fn try_send_mrf_intent_typed( } } +/// Acquire a fresh process-local lease for one durable replay record. +/// +/// Journal records deliberately do not persist leases. A replay consumer must +/// call this before submitting the record so a later verified repair event can +/// identify the exact replay admission. Replay does not reserve the live +/// producer coalescer key: the replay queue first deduplicates legacy records, +/// then manager admission owns task-level deduplication with live producers. +pub fn try_rearm_mrf_replay_intent(intent: &mut MrfIntent) -> MrfIngressResult { + if intent.lease.is_some() { + return MrfIngressResult::Enqueued; + } + if intent.bucket.len() > MRF_MAX_IDENTITY_COMPONENT || intent.object.len() > MRF_MAX_IDENTITY_COMPONENT { + return MrfIngressResult::Dropped(MrfDropReason::OversizedIdentity); + } + let (version_id, scope) = canonical_identity(intent.kind, intent.version_id, intent.scope); + intent.version_id = version_id; + intent.scope = scope; + intent.lease = Some(MrfIngressLease::new(NEXT_MRF_LEASE.fetch_add(1, Ordering::Relaxed))); + MrfIngressResult::Enqueued +} + /// Release the ingress key once the consumer owns the intent. pub fn release_mrf_intent(intent: &MrfIntent) { release_mrf_identity(intent.kind, &intent.bucket, &intent.object, intent.version_id, intent.scope, intent.lease); @@ -432,12 +508,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 +575,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::*; @@ -546,6 +680,147 @@ 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 durable_replay_rearm_assigns_a_fresh_dischargeable_lease() { + let unique = Uuid::new_v4(); + let mut intent = MrfIntent { + bucket: Arc::from(format!("replay-{unique}")), + object: Arc::from("object"), + version_id: Some([0; 16]), + kind: MrfKind::PartialWrite, + scope: Some(MrfScope { + pool_index: 2, + set_index: 3, + }), + lease: None, + enqueued_at_ms: 1, + attempts: 0, + }; + + assert_eq!(try_rearm_mrf_replay_intent(&mut intent), MrfIngressResult::Enqueued); + assert_eq!(intent.version_id, None, "nil versions remain canonical during replay"); + assert!(intent.lease.is_some(), "replay admission must carry a fresh lease"); + assert!( + MrfDurableRepairAnchor::from_intent(&intent, Uuid::new_v4()).is_some(), + "a rearmed replay record can participate in exact durable proof matching" + ); + release_mrf_intent(&intent); + } + + #[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"); @@ -610,4 +885,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 c1431472a..16e71eb79 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(); diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 34d6072e1..95c66a23c 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -36,7 +36,7 @@ use super::{DiskStore, HealDiskExt as _, local_disk_map_read}; use crate::heal::manager::{HealManager, MrfRepairNoticeTarget}; use metrics::{counter, gauge}; -use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIntent}; +use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIngressResult, MrfIntent}; use rustfs_heal_contracts::heal_channel::{HealAdmissionDropReason, HealAdmissionResult}; use std::collections::{HashSet, VecDeque}; use std::sync::Arc; @@ -702,9 +702,9 @@ async fn replay_into( } }, }; - let mut intents = Vec::new(); let (decoded, truncated) = decode_journal(&data); - intents.extend(decoded); + let replayed = decoded.len(); + let intents = decoded; if truncated > 0 { tracing::warn!( target: "rustfs::heal::mrf", @@ -712,8 +712,7 @@ async fn replay_into( "MRF journal had a torn tail; truncated records were discarded" ); } - counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(intents.len()).unwrap_or(u64::MAX)); - let replayed = intents.len(); + counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(replayed).unwrap_or(u64::MAX)); let replay_bytes = intents .iter() .fold(0usize, |total, intent| total.saturating_add(intent.estimated_bytes())); @@ -739,6 +738,15 @@ async fn replay_into( // stays armed in `queue` for the consumer's retry loop. if backoff_until.is_none() { while let Some(mut intent) = queue.pop_front() { + if !matches!( + rustfs_common::mrf_channel::try_rearm_mrf_replay_intent(&mut intent), + MrfIngressResult::Enqueued + ) { + queue.push_back(intent); + rearm_incomplete = true; + *backoff_until = Some(tokio::time::Instant::now()); + break; + } match submit_mrf_heal_request(manager, &intent).await { Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {} Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => { @@ -753,7 +761,9 @@ async fn replay_into( } break; } - Ok(HealAdmissionResult::Dropped(_)) => {} + Ok(HealAdmissionResult::Dropped(_)) => { + rustfs_common::mrf_channel::release_mrf_intent(&intent); + } Err(_) => { intent.attempts = intent.attempts.saturating_add(1); if intent.attempts < MRF_MAX_ATTEMPTS { @@ -955,6 +965,40 @@ mod tests { ); } + #[test] + fn durable_replay_acquires_a_fresh_lease_before_manager_admission() { + let unique = uuid::Uuid::new_v4(); + let original = intent(&format!("replay-{unique}"), "object", 0); + assert!(original.lease.is_none(), "legacy journal records do not persist process leases"); + let mut queue = MrfQueue::new(2, usize::MAX); + assert_eq!(queue.try_push_typed(original.clone()), MrfQueuePushResult::Enqueued); + assert_eq!( + queue.try_push_typed(original), + MrfQueuePushResult::Coalesced, + "legacy duplicates are one durable responsibility before a lease is assigned" + ); + let mut replay = queue.pop_front().expect("one deduplicated replay record"); + assert_eq!( + rustfs_common::mrf_channel::try_rearm_mrf_replay_intent(&mut replay), + MrfIngressResult::Enqueued + ); + assert!(replay.lease.is_some(), "manager admission must receive the replay lease"); + assert!( + rustfs_common::mrf_channel::MrfDurableRepairAnchor::from_intent(&replay, uuid::Uuid::new_v4()).is_some(), + "the replay identity must be usable by the durable proof consumer" + ); + let mut encoded = Vec::new(); + assert!(encode_intent(&replay, &mut encoded)); + let (decoded, truncated) = decode_journal(&encoded); + assert_eq!(truncated, 0); + assert_eq!(decoded.len(), 1); + assert!( + decoded[0].lease.is_none(), + "process-local leases must not enter the durable journal format" + ); + rustfs_common::mrf_channel::release_mrf_intent(&replay); + } + #[test] fn replay_can_arm_more_records_than_live_queue_budget() { let mut queue = MrfQueue::new(1, intent("bucket", "object-0", 0).estimated_bytes());