diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index a43f591f6..95c66a23c 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -516,7 +516,6 @@ async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> c struct MrfRuntime { queue: MrfQueue, - retained_replay_intents: Vec, config: MrfConsumerConfig, new_since_flush: usize, /// True while the in-memory pending set has changed since the last @@ -536,7 +535,7 @@ impl MrfRuntime { fn snapshot(&self) -> (Vec, Vec) { let mut authoritative = Vec::new(); let mut legacy = Vec::new(); - for intent in self.retained_replay_intents.iter().chain(self.queue.intents()) { + for intent in self.queue.intents() { let scoped_identity = !matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption) && intent.scope.is_some(); if !encode_intent(intent, &mut authoritative) { @@ -674,11 +673,10 @@ pub async fn replay_journal_once(manager: &Arc) -> usize { struct ReplayOutcome { replayed: usize, journal_on_disk: bool, - retained_replay_intents: Vec, } -fn replay_must_retain_journal(rearm_incomplete: bool, pending_depth: usize, retained_replay_depth: usize) -> bool { - rearm_incomplete || pending_depth > 0 || retained_replay_depth > 0 +fn replay_must_retain_journal(rearm_incomplete: bool, pending_depth: usize) -> bool { + rearm_incomplete || pending_depth > 0 } /// Shared replay core: read + decode + re-arm, then drain what fits. The @@ -700,7 +698,6 @@ async fn replay_into( return ReplayOutcome { replayed: 0, journal_on_disk: false, - retained_replay_intents: Vec::new(), }; } }, @@ -739,7 +736,6 @@ async fn replay_into( // Drain the replayed intents immediately; whatever the manager refuses // stays armed in `queue` for the consumer's retry loop. - let mut retained_replay_intents = Vec::new(); if backoff_until.is_none() { while let Some(mut intent) = queue.pop_front() { if !matches!( @@ -752,9 +748,7 @@ async fn replay_into( break; } match submit_mrf_heal_request(manager, &intent).await { - Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => { - retained_replay_intents.push(intent); - } + Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {} Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => { intent.attempts = intent.attempts.saturating_add(1); if intent.attempts < MRF_MAX_ATTEMPTS { @@ -785,7 +779,7 @@ async fn replay_into( } } } - let journal_on_disk = if replay_must_retain_journal(rearm_incomplete, queue.depth(), retained_replay_intents.len()) { + let journal_on_disk = if replay_must_retain_journal(rearm_incomplete, queue.depth()) { true } else { !delete_journals().await @@ -793,7 +787,6 @@ async fn replay_into( ReplayOutcome { replayed, journal_on_disk, - retained_replay_intents, } } @@ -803,7 +796,6 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive let config = MrfConsumerConfig::default(); let mut runtime = MrfRuntime { queue: MrfQueue::new(config.queue_capacity, config.journal_max_bytes), - retained_replay_intents: Vec::new(), config: config.clone(), new_since_flush: 0, dirty: false, @@ -815,7 +807,6 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive // on disk whenever any replayed intent still needs a successor snapshot. let replay = replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await; runtime.journal_on_disk = replay.journal_on_disk; - runtime.retained_replay_intents = replay.retained_replay_intents; // Anything still pending (e.g. the manager was full and backoff armed) // must be re-persisted by the next flush before replay can delete the // startup anchor. @@ -833,7 +824,7 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive // provably current AND idle (a dirty or pending state // gets one last persist attempt, matching the shutdown // retry the unconditional flush used to provide). - if runtime.dirty || runtime.queue.depth() > 0 || !runtime.retained_replay_intents.is_empty() { + if runtime.dirty || runtime.queue.depth() > 0 { runtime.flush().await; } tracing::info!( @@ -862,7 +853,6 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive match tick_action( runtime.dirty, runtime.queue.depth(), - runtime.retained_replay_intents.len(), runtime.journal_on_disk, ) { TickAction::Flush => { @@ -877,8 +867,8 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive runtime.dispatch(manager.as_ref()).await; } TickAction::DeleteJournal => { - // Only remove a stale journal after every replayed - // intent has a durable successor proof. + // All replayed intents have either been accepted, + // merged, or replaced by a pending successor snapshot. if delete_journals().await { runtime.journal_on_disk = false; gauge!("rustfs_heal_mrf_journal_bytes").set(0.0); @@ -907,13 +897,11 @@ enum TickAction { Idle, } -fn tick_action(dirty: bool, depth: usize, retained_replay_depth: usize, journal_on_disk: bool) -> TickAction { +fn tick_action(dirty: bool, depth: usize, journal_on_disk: bool) -> TickAction { if dirty { TickAction::Flush } else if depth > 0 { TickAction::Retry - } else if retained_replay_depth > 0 { - TickAction::Idle } else if journal_on_disk { TickAction::DeleteJournal } else { @@ -946,42 +934,34 @@ mod tests { // Dirty dominates: a changed pending set flushes even when idle // otherwise. - assert!(matches!(tick_action(true, 0, 0, false), Flush)); - assert!(matches!(tick_action(true, 3, 0, true), Flush)); + assert!(matches!(tick_action(true, 0, false), Flush)); + assert!(matches!(tick_action(true, 3, true), Flush)); // Clean backlog: no rewrite, but keep draining so an expired // admission backoff retries on time. - assert!(matches!(tick_action(false, 1, 0, false), Retry)); - assert!(matches!(tick_action(false, 2, 0, true), Retry)); - - // Replayed records accepted by the manager are still restart anchors - // until a durable successor proof can tombstone them. - assert!(matches!(tick_action(false, 0, 1, true), Idle)); + assert!(matches!(tick_action(false, 1, false), Retry)); + assert!(matches!(tick_action(false, 2, true), Retry)); // Quiescent with a stale journal file on disk: remove it. - assert!(matches!(tick_action(false, 0, 0, true), DeleteJournal)); + assert!(matches!(tick_action(false, 0, true), DeleteJournal)); // Fully quiescent: nothing to do. - assert!(matches!(tick_action(false, 0, 0, false), Idle)); + assert!(matches!(tick_action(false, 0, false), Idle)); } #[test] fn replay_cleanup_retains_journal_for_unarmed_or_refused_records() { assert!( - replay_must_retain_journal(true, 0, 0), + replay_must_retain_journal(true, 0), "a rejected replay record still needs its disk anchor" ); assert!( - replay_must_retain_journal(false, 1, 0), + replay_must_retain_journal(false, 1), "a Full admission retry must keep the startup journal until the next snapshot" ); assert!( - replay_must_retain_journal(false, 0, 1), - "an accepted replay record still needs a durable successor before cleanup" - ); - assert!( - !replay_must_retain_journal(false, 0, 0), - "only a fully consumed replay snapshot with no retained anchors may be deleted" + !replay_must_retain_journal(false, 0), + "only a fully consumed replay snapshot may be deleted" ); } @@ -1019,37 +999,6 @@ mod tests { rustfs_common::mrf_channel::release_mrf_intent(&replay); } - #[test] - fn retained_replay_anchor_remains_in_successor_snapshot() { - let retained = intent("accepted-replay", "object", 0); - let mut runtime = MrfRuntime { - queue: MrfQueue::new(8, 8192), - retained_replay_intents: vec![retained.clone()], - config: MrfConsumerConfig::default(), - new_since_flush: 0, - dirty: false, - journal_on_disk: true, - backoff_until: None, - }; - assert_eq!( - runtime.queue.try_push_typed(intent("new-pending", "object", 0)), - MrfQueuePushResult::Enqueued - ); - - let (authoritative, legacy) = runtime.snapshot(); - let (decoded, truncated) = decode_journal(&authoritative); - let (legacy_decoded, legacy_truncated) = decode_journal(&legacy); - - assert_eq!(truncated, 0); - assert_eq!(legacy_truncated, 0); - assert_eq!(decoded.len(), 2); - assert_eq!(legacy_decoded.len(), 2); - assert!( - decoded.iter().any(|intent| intent.bucket == retained.bucket), - "accepted replay anchor must remain crash-replayable" - ); - } - #[test] fn replay_can_arm_more_records_than_live_queue_budget() { let mut queue = MrfQueue::new(1, intent("bucket", "object-0", 0).estimated_bytes());