diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 698d2d585..34d6072e1 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(), }; } }, @@ -740,13 +737,10 @@ 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) => { - 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 { @@ -775,7 +769,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 @@ -783,7 +777,6 @@ async fn replay_into( ReplayOutcome { replayed, journal_on_disk, - retained_replay_intents, } } @@ -793,7 +786,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, @@ -805,7 +797,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. @@ -823,7 +814,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!( @@ -852,7 +843,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 => { @@ -867,8 +857,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); @@ -897,13 +887,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 { @@ -936,73 +924,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" - ); - } - - #[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" + !replay_must_retain_journal(false, 0), + "only a fully consumed replay snapshot may be deleted" ); } diff --git a/crates/log-analyzer/src/rules/seed/disk.rs b/crates/log-analyzer/src/rules/seed/disk.rs index 8cd1f031a..c90e44646 100644 --- a/crates/log-analyzer/src/rules/seed/disk.rs +++ b/crates/log-analyzer/src/rules/seed/disk.rs @@ -68,14 +68,17 @@ pub(super) fn rules() -> Vec { ) }, Rule { - anchors: strings(["reporting peer disks offline after consecutive storage_info failures"]), + anchors: strings(["Storage inventory probe failed; current drive health is unknown"]), ..base( "peer-disks-offline", P2Degraded, "disk", - "peer 磁盘被整体判定离线", - contains("reporting peer disks offline after consecutive storage_info failures"), - "对某 peer 连续 storage_info 失败,判定其磁盘整体离线。", + "peer 存储清单探测失败", + any([ + contains("Storage inventory probe failed; current drive health is unknown"), + contains("reporting peer disks offline after consecutive storage_info failures"), + ]), + "某 peer 的 storage_info 探测失败,当前磁盘健康状态未知。", "检查该 peer 节点存活与 RPC 端口可达。", ) }, diff --git a/crates/log-analyzer/src/rules/seed/tests.rs b/crates/log-analyzer/src/rules/seed/tests.rs index 442546854..88bf61a6a 100644 --- a/crates/log-analyzer/src/rules/seed/tests.rs +++ b/crates/log-analyzer/src/rules/seed/tests.rs @@ -110,7 +110,7 @@ fn every_rule_has_a_positive_sample() { ("remote-peer-faulty", msg("Remote peer health check failed for node2: marking as faulty")), ( "peer-disks-offline", - msg("reporting peer disks offline after consecutive storage_info failures"), + msg("Storage inventory probe failed; current drive health is unknown"), ), ("drive-faulty-error", msg("remote drive is faulty")), ( @@ -318,6 +318,10 @@ fn smoke_samples_hit_exact_rule_sets() { &["disk-marked-faulty"], ); exact(&msg("erasure write quorum (required=8, achieved=5)"), &["ec-write-quorum"]); + exact( + &msg("reporting peer disks offline after consecutive storage_info failures"), + &["peer-disks-offline"], + ); exact( &Sample { message: "Metacache listing quorum failed", diff --git a/rustfs/src/error.rs b/rustfs/src/error.rs index 773f7b869..be2a86187 100644 --- a/rustfs/src/error.rs +++ b/rustfs/src/error.rs @@ -496,9 +496,9 @@ impl From for ApiError { let message = if matches!(&err, StorageError::QuotaExceeded { .. }) { err.to_string() - } else if matches!(&err, StorageError::MaxVersionsExceeded) { - ApiError::error_code_to_message(&code) - } else if code == S3ErrorCode::InternalError && matches!(&err, StorageError::Io(_)) { + } else if matches!(&err, StorageError::MaxVersionsExceeded) + || (code == S3ErrorCode::InternalError && matches!(&err, StorageError::Io(_))) + { ApiError::error_code_to_message(&code) } else if code == S3ErrorCode::InternalError { err.to_string()