diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index 784969051..e70d4f353 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -1095,6 +1095,14 @@ pub(in crate::set_disk) struct ReadRepairHealSubmission<'a> { pub(in crate::set_disk) set_index: usize, pub(in crate::set_disk) part_number: Option, pub(in crate::set_disk) reason: &'static str, + /// Durable MRF journal intent to file alongside the read-repair request + /// (backlog#1894 axis A): the intent kind plus its native `Uuid` + /// version id (the submission's string form stays display-only). Bound + /// to the reservation — the intent is only delivered when this sighting + /// wins the dedup TTL, so a burst of reads failing on the same object + /// books exactly one journal record instead of one per retry. `None` + /// keeps the historical no-intent behavior. + pub(in crate::set_disk) mrf_intent: Option<(rustfs_common::mrf_channel::MrfKind, Option)>, } pub(in crate::set_disk) fn send_read_repair_heal_request( @@ -1126,6 +1134,7 @@ pub(in crate::set_disk) async fn submit_read_repair_heal( set_index, part_number, reason, + mrf_intent: None, }, send_read_repair_heal_request, ) @@ -1144,6 +1153,7 @@ pub(in crate::set_disk) async fn submit_read_repair_heal_with_submitter( set_index, part_number, reason, + mrf_intent, } = submission; let Some(dedup_key) = reserve_read_repair_heal(bucket, object, version_id, pool_index, set_index).await else { @@ -1155,6 +1165,12 @@ pub(in crate::set_disk) async fn submit_read_repair_heal_with_submitter( return; }; + // Reservation won: this sighting owns the repair records for the object, + // including the durable journal intent when the caller asked for one. + if let Some((kind, version_uuid)) = mrf_intent { + rustfs_common::mrf_channel::try_send_mrf_intent(kind, bucket, object, version_uuid); + } + let mut request = rustfs_common::heal_channel::create_heal_request_with_options( bucket.to_string(), Some(object.to_string()), @@ -8710,6 +8726,42 @@ mod tests { assert_eq!(responses[0].error, Error::ErasureReadQuorum.to_string()); } + #[tokio::test] + #[serial_test::serial] + async fn mrf_intent_is_filed_once_per_read_repair_reservation() { + // Serial: owns the process-global MRF channel for this test binary + // (same key as the other channel-owning tests above). + let bucket = format!("mrf-intent-bucket-{}", Uuid::new_v4()); + let object = format!("object-{}", Uuid::new_v4()); + let mut receiver = rustfs_common::mrf_channel::init_mrf_channel().expect("first channel init in this binary"); + rustfs_common::mrf_channel::set_mrf_delivery_enabled(true); + + fn intent_submission<'a>(bucket: &'a str, object: &'a str) -> ReadRepairHealSubmission<'a> { + ReadRepairHealSubmission { + bucket, + object, + version_id: None, + pool_index: 9, + set_index: 9, + part_number: Some(1), + reason: "decode_error", + mrf_intent: Some((rustfs_common::mrf_channel::MrfKind::DecodeFailure, None)), + } + } + + // First sighting wins the reservation: the journal intent is filed + // synchronously before the admission task is spawned. + submit_read_repair_heal_with_submitter(intent_submission(&bucket, &object), accepted_read_repair_submitter).await; + let first = receiver.try_recv().expect("first sighting must file exactly one MRF intent"); + assert_eq!(*first.bucket, bucket); + assert_eq!(*first.object, object); + + // Second sighting within the dedup TTL is a duplicate: no request, no + // second journal record. + submit_read_repair_heal_with_submitter(intent_submission(&bucket, &object), accepted_read_repair_submitter).await; + assert!(receiver.try_recv().is_err(), "duplicate sighting must not file another MRF intent"); + } + #[tokio::test] async fn reserve_read_repair_heal_dedupes_by_object_version_and_set() { let object = format!("object-{}", Uuid::new_v4()); @@ -8818,6 +8870,7 @@ mod tests { set_index: 2, part_number: Some(1), reason: "test", + mrf_intent: None, }, failed_read_repair_submitter, ) @@ -8846,6 +8899,7 @@ mod tests { set_index: 3, part_number: Some(2), reason: "test", + mrf_intent: None, }, dropped_read_repair_submitter, ) @@ -8874,6 +8928,7 @@ mod tests { set_index: 4, part_number: None, reason: "test", + mrf_intent: None, }, accepted_read_repair_submitter, ) diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 07bbcf5f9..2b556b143 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -1077,23 +1077,23 @@ impl SetDisks { "Recoverable decode error triggered read repair" ); let version_id = fi.version_id.as_ref().map(ToString::to_string); - // MRF journal intent: keeps a durable Urgent ECDecode - // request alive across restarts even when the in-memory - // read-repair request is dropped or lost (HS-01). - rustfs_common::mrf_channel::try_send_mrf_intent( - rustfs_common::mrf_channel::MrfKind::DecodeFailure, - bucket, - object, - fi.version_id, - ); - submit_read_repair_heal( - bucket, - object, - version_id.as_deref(), - pool_index, - set_index, - Some(part_number), - "decode_error", + // Single-flight (backlog#1894 axis A): the durable + // MRF intent (Urgent ECDecode across restarts, HS-01) + // is bound to the read-repair reservation, so only the + // first sighting within the dedup TTL books a journal + // record instead of one per retried read. + submit_read_repair_heal_with_submitter( + ReadRepairHealSubmission { + bucket, + object, + version_id: version_id.as_deref(), + pool_index, + set_index, + part_number: Some(part_number), + reason: "decode_error", + mrf_intent: Some((rustfs_common::mrf_channel::MrfKind::DecodeFailure, fi.version_id)), + }, + send_read_repair_heal_request, ) .await; has_err = false; @@ -2577,6 +2577,7 @@ mod metadata_cache_tests { set_index: 0, part_number: Some(1), reason: "missing_shards", + mrf_intent: None, }, slow_read_repair_submitter, ) @@ -2611,6 +2612,7 @@ mod metadata_cache_tests { set_index: 0, part_number: Some(1), reason: "missing_shards", + mrf_intent: None, }, dropped_read_repair_submitter, ) @@ -2647,6 +2649,7 @@ mod metadata_cache_tests { set_index: 0, part_number: Some(1), reason: "missing_shards", + mrf_intent: None, }, capture_read_repair_submitter, ) diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 8a7e1e3b5..362a32966 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -25,9 +25,12 @@ //! set, rewritten on a group-commit cadence (every flush interval or flush //! threshold new intents). A rewrite is atomic at the record level only — a //! torn tail simply truncates during replay because every record carries its -//! own CRC32. Losing the last flush window (≤500 ms) is acceptable: replayed -//! duplicates are merged by the manager's dedup key, and read-repair remains -//! the safety net. +//! own CRC32. Losing the last flush window (≤500 ms) is acceptable because +//! every producer keeps its own safety net: read-repair re-detects on the +//! next failing read, and the scanner's corrupt-metadata branch leaves a +//! pending-ledger entry behind even when its MRF intent is accepted +//! (backlog#1894 axis A), so a lost intent is retried by the ledger rather +//! than waiting for the failed-object TTL to re-scan the path. use super::{DiskStore, HealDiskExt as _, local_disk_map_read}; use crate::heal::manager::HealManager; diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 70db02ba4..fa4a48c9f 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -645,6 +645,32 @@ enum GetSizeFailureAction { HealMetadata { object: String }, } +/// How the corrupt-metadata branch records the repair after attempting an +/// MRF intent (backlog#1894 axis A). +#[derive(Debug, PartialEq, Eq)] +enum CorruptMetadataRecording { + /// Intent accepted: the MRF consumer owns the repair (High Metadata + /// heal, durable after the journal's group-commit flush), so the + /// immediate heal request is skipped — the manager would otherwise book + /// two tasks for one target. A pending-ledger entry stays behind as the + /// backstop for what the journal cannot cover on its own (a crash inside + /// the flush window, or the consumer exhausting its admission attempts); + /// the repaired-notice fanout (axis B) drops the entry once the repair + /// lands. + LedgerOnly, + /// Intent rejected (feature disabled, channel uninitialized, or full): + /// the historical immediate heal request plus the ledger entry. + ImmediateAndLedger, +} + +fn corrupt_metadata_recording(mrf_accepted: bool) -> CorruptMetadataRecording { + if mrf_accepted { + CorruptMetadataRecording::LedgerOnly + } else { + CorruptMetadataRecording::ImmediateAndLedger + } +} + fn build_bucket_heal_request(bucket: String, priority: HealChannelPriority) -> HealChannelRequest { HealChannelRequest { bucket, @@ -2478,29 +2504,46 @@ impl FolderScanner { } if let GetSizeFailureAction::HealMetadata { object } = failure_action { - // MRF journal intent: durable High-priority Metadata - // heal across restarts (HS-01); the scanner heal - // request below stays as the immediate path. - rustfs_common::mrf_channel::try_send_mrf_intent( + // Single-flight (backlog#1894 axis A) — the + // recording mode and its guarantees are pinned by + // corrupt_metadata_recording below. + let mrf_accepted = rustfs_common::mrf_channel::try_send_mrf_intent( rustfs_common::mrf_channel::MrfKind::MetadataCorruption, &item.bucket, &object, None, ); - self.send_required_scanner_heal_request( - PendingScannerHealKind::Object, - item.bucket.clone(), - Some(object.clone()), - None, - build_object_heal_request( - item.bucket.clone(), - object.clone(), - None, - self.scan_mode, - HealChannelPriority::High, - ), - ) - .await?; + match corrupt_metadata_recording(mrf_accepted) { + CorruptMetadataRecording::LedgerOnly => { + // Recorded as Full (retry-later): admission + // for this target happens in the MRF + // consumer, not in the manager's queue here. + self.update_pending_scanner_heal_after_admission( + PendingScannerHealKind::Object, + &item.bucket, + Some(&object), + None, + self.scan_mode, + HealAdmissionResult::Full, + ); + } + CorruptMetadataRecording::ImmediateAndLedger => { + self.send_required_scanner_heal_request( + PendingScannerHealKind::Object, + item.bucket.clone(), + Some(object.clone()), + None, + build_object_heal_request( + item.bucket.clone(), + object.clone(), + None, + self.scan_mode, + HealChannelPriority::High, + ), + ) + .await?; + } + } } timer.sleep().await; @@ -3395,6 +3438,17 @@ mod tests { assert_eq!(EVENT_SCANNER_BIG_PREFIX, EventName::ScannerBigPrefix.to_string()); } + /// Single-flight decision for the corrupt-metadata branch (backlog#1894 + /// axis A): an accepted MRF intent must drop the immediate heal request + /// (the consumer files one; the manager would double-book) while a + /// rejected one must keep it — in both cases a ledger entry remains, so + /// the backstop survives regardless of delivery. + #[test] + fn corrupt_metadata_recording_maps_delivery_to_backstop() { + assert_eq!(corrupt_metadata_recording(true), CorruptMetadataRecording::LedgerOnly); + assert_eq!(corrupt_metadata_recording(false), CorruptMetadataRecording::ImmediateAndLedger); + } + fn cooldown_map_len() -> usize { SCANNER_ALERT_EMISSION_COOLDOWN .lock()