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 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-07 23:21:47 +08:00
committed by GitHub
parent 7373a5902e
commit 0686277ee4
+94 -24
View File
@@ -516,6 +516,7 @@ async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> c
struct MrfRuntime {
queue: MrfQueue,
retained_replay_intents: Vec<MrfIntent>,
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<tokio::time::Instant>,
@@ -536,7 +536,7 @@ impl MrfRuntime {
fn snapshot(&self) -> (Vec<u8>, Vec<u8>) {
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<HealManager>) -> usize {
struct ReplayOutcome {
replayed: usize,
journal_on_disk: bool,
retained_replay_intents: Vec<MrfIntent>,
}
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<HealManager>, 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<HealManager>, 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<HealManager>, 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<HealManager>, 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<HealManager>, 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"
);
}