From 0686277ee4c100f8cb3d1b4a7645ee1126434af6 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 7 Sep 2026 23:21:47 +0800 Subject: [PATCH] fix(heal): retain MRF replay anchors until successor proof (#7413) Keep replayed MRF records crash-replayable after manager admission until a later durable successor proof can tombstone them. Queue-full and transient replay submission failures now also preserve the old journal anchor instead of allowing cleanup to erase the only recovery source. Co-authored-by: zhi22915 --- crates/heal/src/heal/mrf_queue.rs | 118 ++++++++++++++++++++++++------ 1 file changed, 94 insertions(+), 24 deletions(-) diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 5d6e86562..698d2d585 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -516,6 +516,7 @@ 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 @@ -524,9 +525,8 @@ struct MrfRuntime { /// waiting out an admission backoff must not re-fsync every local disk /// twice a second. dirty: bool, - /// True while a journal snapshot exists on disk that no longer reflects - /// an all-consumed pending set; the next idle tick removes it (MinIO - /// deletes its `list.bin` after replay for the same reason). + /// True while a journal snapshot exists on disk that may still be needed + /// for replay or cleanup. journal_on_disk: bool, /// Earliest instant a full-admission retry may proceed. backoff_until: Option, @@ -536,7 +536,7 @@ impl MrfRuntime { fn snapshot(&self) -> (Vec, Vec) { let mut authoritative = Vec::new(); let mut legacy = Vec::new(); - for intent in self.queue.intents() { + for intent in self.retained_replay_intents.iter().chain(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,10 +674,11 @@ 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) -> bool { - rearm_incomplete || pending_depth > 0 +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 } /// Shared replay core: read + decode + re-arm, then drain what fits. The @@ -699,6 +700,7 @@ async fn replay_into( return ReplayOutcome { replayed: 0, journal_on_disk: false, + retained_replay_intents: Vec::new(), }; } }, @@ -738,23 +740,42 @@ 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() { match submit_mrf_heal_request(manager, &intent).await { - Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {} + Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => { + retained_replay_intents.push(intent); + } Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => { intent.attempts = intent.attempts.saturating_add(1); if intent.attempts < MRF_MAX_ATTEMPTS { queue.push_back(intent); *backoff_until = Some(tokio::time::Instant::now()); + } else { + rearm_incomplete = true; + counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1); + rustfs_common::mrf_channel::release_mrf_intent(&intent); + } + break; + } + Ok(HealAdmissionResult::Dropped(_)) => {} + Err(_) => { + intent.attempts = intent.attempts.saturating_add(1); + if intent.attempts < MRF_MAX_ATTEMPTS { + queue.push_back(intent); + *backoff_until = Some(tokio::time::Instant::now()); + } else { + rearm_incomplete = true; + counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1); + rustfs_common::mrf_channel::release_mrf_intent(&intent); } break; } - Ok(HealAdmissionResult::Dropped(_)) | Err(_) => {} } } } - let journal_on_disk = if replay_must_retain_journal(rearm_incomplete, queue.depth()) { + let journal_on_disk = if replay_must_retain_journal(rearm_incomplete, queue.depth(), retained_replay_intents.len()) { true } else { !delete_journals().await @@ -762,6 +783,7 @@ async fn replay_into( ReplayOutcome { replayed, journal_on_disk, + retained_replay_intents, } } @@ -771,6 +793,7 @@ 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, @@ -782,6 +805,7 @@ 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. @@ -799,7 +823,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 { + if runtime.dirty || runtime.queue.depth() > 0 || !runtime.retained_replay_intents.is_empty() { runtime.flush().await; } tracing::info!( @@ -825,7 +849,12 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive } } _ = flush_tick.tick() => { - match tick_action(runtime.dirty, runtime.queue.depth(), runtime.journal_on_disk) { + match tick_action( + runtime.dirty, + runtime.queue.depth(), + runtime.retained_replay_intents.len(), + runtime.journal_on_disk, + ) { TickAction::Flush => { runtime.flush().await; runtime.dispatch(manager.as_ref()).await; @@ -838,8 +867,8 @@ async fn run_mrf_consumer(manager: Arc, mut receiver: mpsc::Receive runtime.dispatch(manager.as_ref()).await; } TickAction::DeleteJournal => { - // All intents consumed: remove the journal so a restart - // replays nothing (mirrors MinIO's post-replay unlink). + // Only remove a stale journal after every replayed + // intent has a durable successor proof. if delete_journals().await { runtime.journal_on_disk = false; gauge!("rustfs_heal_mrf_journal_bytes").set(0.0); @@ -868,11 +897,13 @@ enum TickAction { Idle, } -fn tick_action(dirty: bool, depth: usize, journal_on_disk: bool) -> TickAction { +fn tick_action(dirty: bool, depth: usize, retained_replay_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 { @@ -905,34 +936,73 @@ mod tests { // Dirty dominates: a changed pending set flushes even when idle // otherwise. - assert!(matches!(tick_action(true, 0, false), Flush)); - assert!(matches!(tick_action(true, 3, true), Flush)); + assert!(matches!(tick_action(true, 0, 0, false), Flush)); + assert!(matches!(tick_action(true, 3, 0, true), Flush)); // Clean backlog: no rewrite, but keep draining so an expired // admission backoff retries on time. - assert!(matches!(tick_action(false, 1, false), Retry)); - assert!(matches!(tick_action(false, 2, true), Retry)); + 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)); // Quiescent with a stale journal file on disk: remove it. - assert!(matches!(tick_action(false, 0, true), DeleteJournal)); + assert!(matches!(tick_action(false, 0, 0, true), DeleteJournal)); // Fully quiescent: nothing to do. - assert!(matches!(tick_action(false, 0, false), Idle)); + assert!(matches!(tick_action(false, 0, 0, false), Idle)); } #[test] fn replay_cleanup_retains_journal_for_unarmed_or_refused_records() { assert!( - replay_must_retain_journal(true, 0), + replay_must_retain_journal(true, 0, 0), "a rejected replay record still needs its disk anchor" ); assert!( - replay_must_retain_journal(false, 1), + replay_must_retain_journal(false, 1, 0), "a Full admission retry must keep the startup journal until the next snapshot" ); assert!( - !replay_must_retain_journal(false, 0), - "only a fully consumed replay snapshot may be deleted" + 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" + ); + } + + #[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" ); }