mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 20:43:04 +00:00
fix(heal): park healthy legacy MRF intents
Co-Authored-By: heihutu <heihutu@gmail.com> Co-Authored-By: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
@@ -163,6 +163,15 @@ pub struct MrfDurableRepairAnchor {
|
||||
pub bucket_incarnation_id: Uuid,
|
||||
}
|
||||
|
||||
/// A completed deep check proved that a legacy object is present and healthy,
|
||||
/// but the format has no independent identity commitment that can discharge a
|
||||
/// durable partial-write responsibility. Keep the journal entry and pause
|
||||
/// automatic retries until an operator has repaired the source of truth.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct MrfUnverifiedLegacyEvent {
|
||||
pub anchor: MrfDurableRepairAnchor,
|
||||
}
|
||||
|
||||
impl MrfDurableRepairAnchor {
|
||||
/// Build a dischargeable anchor only when the caller supplies the storage
|
||||
/// incarnation and the original ingress lease. Legacy replay records lack
|
||||
@@ -782,6 +791,8 @@ const MRF_REPAIRED_EVENT_CAP: usize = 4096;
|
||||
static MRF_REPAIRED_EVENTS: OnceLock<std::sync::Mutex<std::collections::VecDeque<MrfRepairedEvent>>> = OnceLock::new();
|
||||
static MRF_VERIFIED_REPAIR_EVENTS: OnceLock<std::sync::Mutex<std::collections::VecDeque<MrfVerifiedRepairEvent>>> =
|
||||
OnceLock::new();
|
||||
static MRF_UNVERIFIED_LEGACY_EVENTS: OnceLock<std::sync::Mutex<std::collections::VecDeque<MrfUnverifiedLegacyEvent>>> =
|
||||
OnceLock::new();
|
||||
|
||||
/// Record a legacy notification for compatibility. This is not an
|
||||
/// acknowledgement of storage verification or durable repair completion.
|
||||
@@ -859,6 +870,44 @@ pub fn take_mrf_verified_repair_events_for(bucket: &str) -> Vec<MrfVerifiedRepai
|
||||
taken
|
||||
}
|
||||
|
||||
/// Record a bounded, exact-identity notice that a legacy durable partial-write
|
||||
/// was checked but cannot be proven from its on-disk format. This is not a
|
||||
/// repair proof and must never release the durable responsibility.
|
||||
pub fn note_mrf_unverified_legacy(event: MrfUnverifiedLegacyEvent) {
|
||||
let registry = MRF_UNVERIFIED_LEGACY_EVENTS.get_or_init(|| std::sync::Mutex::new(std::collections::VecDeque::new()));
|
||||
let Ok(mut events) = registry.lock() else {
|
||||
return;
|
||||
};
|
||||
if events.contains(&event) {
|
||||
return;
|
||||
}
|
||||
if events.len() >= MRF_REPAIRED_EVENT_CAP {
|
||||
events.pop_front();
|
||||
}
|
||||
events.push_back(event);
|
||||
}
|
||||
|
||||
/// Take held-legacy notices for one bucket. The durable ledger still verifies
|
||||
/// the full anchor before changing the journal entry to its held state.
|
||||
pub fn take_mrf_unverified_legacy_events_for(bucket: &str) -> Vec<MrfUnverifiedLegacyEvent> {
|
||||
let Some(registry) = MRF_UNVERIFIED_LEGACY_EVENTS.get() else {
|
||||
return Vec::new();
|
||||
};
|
||||
let Ok(mut events) = registry.lock() else {
|
||||
return Vec::new();
|
||||
};
|
||||
let mut taken = Vec::new();
|
||||
events.retain(|event| {
|
||||
if event.anchor.bucket.as_ref() == bucket {
|
||||
taken.push(event.clone());
|
||||
false
|
||||
} else {
|
||||
true
|
||||
}
|
||||
});
|
||||
taken
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -1233,6 +1282,53 @@ mod tests {
|
||||
assert!(take_mrf_verified_repair_events_for("verified-bucket-a").is_empty());
|
||||
assert_eq!(take_mrf_verified_repair_events_for("verified-bucket-b").len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unverified_legacy_events_are_bucket_scoped_and_bounded() {
|
||||
let bucket = Arc::<str>::from("unverified-legacy-event-bucket");
|
||||
let _ = take_mrf_unverified_legacy_events_for(bucket.as_ref());
|
||||
let event = MrfUnverifiedLegacyEvent {
|
||||
anchor: MrfDurableRepairAnchor {
|
||||
kind: MrfKind::PartialWrite,
|
||||
bucket: bucket.clone(),
|
||||
object: Arc::from("unverified-object"),
|
||||
version_id: None,
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
delete_marker_purge: None,
|
||||
lease: MrfIngressLease::new(912),
|
||||
bucket_incarnation_id: Uuid::new_v4(),
|
||||
},
|
||||
};
|
||||
note_mrf_unverified_legacy(event.clone());
|
||||
note_mrf_unverified_legacy(event.clone());
|
||||
assert_eq!(take_mrf_unverified_legacy_events_for(bucket.as_ref()), vec![event]);
|
||||
assert!(take_mrf_unverified_legacy_events_for(bucket.as_ref()).is_empty());
|
||||
|
||||
let cap_bucket = Arc::<str>::from("unverified-legacy-event-cap-bucket");
|
||||
for index in 0..MRF_REPAIRED_EVENT_CAP + 4 {
|
||||
note_mrf_unverified_legacy(MrfUnverifiedLegacyEvent {
|
||||
anchor: MrfDurableRepairAnchor {
|
||||
kind: MrfKind::PartialWrite,
|
||||
bucket: cap_bucket.clone(),
|
||||
object: Arc::from(format!("object-{index}")),
|
||||
version_id: None,
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
delete_marker_purge: None,
|
||||
lease: MrfIngressLease::new(u64::try_from(index).expect("test index fits")),
|
||||
bucket_incarnation_id: Uuid::new_v4(),
|
||||
},
|
||||
});
|
||||
}
|
||||
let events = take_mrf_unverified_legacy_events_for(cap_bucket.as_ref());
|
||||
assert_eq!(events.len(), MRF_REPAIRED_EVENT_CAP);
|
||||
assert_eq!(events.first().expect("bounded ring is non-empty").anchor.object.as_ref(), "object-4");
|
||||
}
|
||||
}
|
||||
#[tokio::test]
|
||||
async fn partial_write_durable_admission_waits_for_persistence_and_renews_lease() {
|
||||
|
||||
@@ -1166,7 +1166,7 @@ impl SetDisks {
|
||||
|
||||
if !latest_meta.deleted && !latest_meta.is_remote() && !protected {
|
||||
result.detail =
|
||||
"Legacy object uses standard repair; independent object identity remains unverified".to_owned();
|
||||
rustfs_heal_contracts::heal_channel::LEGACY_OBJECT_IDENTITY_UNVERIFIED_DETAIL.to_owned();
|
||||
}
|
||||
|
||||
if disks_to_heal_count == 0 {
|
||||
|
||||
@@ -24,6 +24,12 @@ pub const HEAL_DELETE_DANGLING: bool = true;
|
||||
pub const RUSTFS_RESERVED_BUCKET: &str = "rustfs";
|
||||
pub const RUSTFS_RESERVED_BUCKET_PATH: &str = "/rustfs";
|
||||
|
||||
/// Detail attached to a completed deep heal when a healthy legacy object has
|
||||
/// no independent identity commitment. Durable MRF handling uses this exact
|
||||
/// reason to pause proofless retries without treating the object as repaired.
|
||||
pub const LEGACY_OBJECT_IDENTITY_UNVERIFIED_DETAIL: &str =
|
||||
"Legacy object uses standard repair; independent object identity remains unverified";
|
||||
|
||||
#[derive(Clone, Copy, Debug, Serialize, Deserialize)]
|
||||
pub enum HealItemType {
|
||||
Metadata,
|
||||
|
||||
@@ -174,6 +174,7 @@ pub(super) struct MrfRepairNoticeTarget {
|
||||
pub(super) scope: Option<rustfs_common::mrf_channel::MrfScope>,
|
||||
pub(super) delete_marker_purge: Option<rustfs_common::mrf_channel::MrfDeleteMarkerPurgeIdentity>,
|
||||
pub(super) lease: Option<rustfs_common::mrf_channel::MrfIngressLease>,
|
||||
pub(super) durable_anchor: Option<rustfs_common::mrf_channel::MrfDurableRepairAnchor>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
@@ -1808,6 +1809,7 @@ impl HealManager {
|
||||
scope: None,
|
||||
delete_marker_purge: None,
|
||||
lease: None,
|
||||
durable_anchor: None,
|
||||
},
|
||||
)
|
||||
.await
|
||||
|
||||
@@ -438,6 +438,7 @@ impl HealManager {
|
||||
release_mrf_repair_notice_targets(¬ice_targets);
|
||||
}
|
||||
publish_verified_mrf_repair_events(¬ice_targets, &completed_status_for_verified_events);
|
||||
publish_unverified_legacy_mrf_events(¬ice_targets, &completed_status_for_verified_events);
|
||||
// update statistics
|
||||
let mut stats = statistics_clone.write().await;
|
||||
match completed_status {
|
||||
@@ -880,6 +881,66 @@ pub(super) fn publish_verified_mrf_repair_events(targets: &[MrfRepairNoticeTarge
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn unverified_legacy_mrf_event_for_target(
|
||||
target: &MrfRepairNoticeTarget,
|
||||
outcome: &crate::heal::outcome::HealObjectOutcome,
|
||||
) -> Option<rustfs_common::mrf_channel::MrfUnverifiedLegacyEvent> {
|
||||
use crate::heal::outcome::{HealObjectDisposition, HealObjectKind};
|
||||
use rustfs_common::mrf_channel::MrfKind;
|
||||
|
||||
let anchor = target.durable_anchor.as_ref()?;
|
||||
if target.kind != MrfKind::PartialWrite
|
||||
|| outcome.disposition != HealObjectDisposition::Unknown
|
||||
|| outcome.detail.as_deref() != Some(rustfs_heal_contracts::heal_channel::LEGACY_OBJECT_IDENTITY_UNVERIFIED_DETAIL)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
let version_id = target.version_id.filter(|bytes| *bytes != [0; 16]);
|
||||
let expected_version = version_id.map(|bytes| uuid::Uuid::from_bytes(bytes).to_string());
|
||||
let expected_pool = target.scope.and_then(|scope| usize::try_from(scope.pool_index).ok());
|
||||
let expected_set = target.scope.and_then(|scope| usize::try_from(scope.set_index).ok());
|
||||
if anchor.kind != target.kind
|
||||
|| anchor.bucket != target.bucket
|
||||
|| anchor.object != target.object
|
||||
|| anchor.version_id != version_id
|
||||
|| anchor.scope != target.scope
|
||||
|| Some(anchor.lease) != target.lease
|
||||
|| outcome.identity.kind != HealObjectKind::Object
|
||||
|| outcome.identity.bucket != target.bucket.as_ref()
|
||||
|| outcome.identity.object != target.object.as_ref()
|
||||
|| outcome.identity.version_id != expected_version
|
||||
|| outcome.identity.pool_index != expected_pool
|
||||
|| outcome.identity.set_index != expected_set
|
||||
|| outcome.identity.bucket_incarnation_id != Some(anchor.bucket_incarnation_id)
|
||||
{
|
||||
return None;
|
||||
}
|
||||
Some(rustfs_common::mrf_channel::MrfUnverifiedLegacyEvent { anchor: anchor.clone() })
|
||||
}
|
||||
|
||||
pub(super) fn publish_unverified_legacy_mrf_events(targets: &[MrfRepairNoticeTarget], completed: &CompletedHealStatus) {
|
||||
use crate::heal::outcome::HealExecutionOutcome;
|
||||
|
||||
if !matches!(completed.status, HealTaskStatus::Completed | HealTaskStatus::Failed { .. }) {
|
||||
return;
|
||||
}
|
||||
let Some(outcome) = completed.outcome.as_ref() else {
|
||||
return;
|
||||
};
|
||||
if outcome.execution != HealExecutionOutcome::Completed {
|
||||
return;
|
||||
}
|
||||
for target in targets {
|
||||
if let Some(event) = outcome
|
||||
.objects
|
||||
.iter()
|
||||
.find_map(|object| unverified_legacy_mrf_event_for_target(target, object))
|
||||
{
|
||||
rustfs_common::mrf_channel::note_mrf_unverified_legacy(event);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn heal_request_set_key_for_task(task: &HealTask) -> Option<String> {
|
||||
match &task.heal_type {
|
||||
HealType::ErasureSet { set_disk_id, .. } => Some(set_disk_id.clone()),
|
||||
|
||||
@@ -1228,6 +1228,7 @@ fn mrf_verified_repair_event_requires_positive_exact_identity() {
|
||||
}),
|
||||
delete_marker_purge: None,
|
||||
lease: None,
|
||||
durable_anchor: None,
|
||||
};
|
||||
let matching = HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
@@ -1335,6 +1336,102 @@ fn mrf_verified_repair_event_requires_positive_exact_identity() {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unverified_legacy_notice_requires_the_exact_partial_write_target() {
|
||||
use crate::heal::outcome::{HealObjectIdentity, HealObjectKind, HealObjectOutcome};
|
||||
use rustfs_common::mrf_channel::{
|
||||
MrfDurableRepairAnchor, MrfIngressResult, MrfIntent, MrfKind, MrfScope, try_rearm_mrf_replay_intent,
|
||||
};
|
||||
|
||||
let bucket = Arc::<str>::from("legacy-held-bucket");
|
||||
let object = Arc::<str>::from("legacy-held-object");
|
||||
let scope = MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
};
|
||||
let mut intent = MrfIntent {
|
||||
bucket: bucket.clone(),
|
||||
object: object.clone(),
|
||||
version_id: None,
|
||||
kind: MrfKind::PartialWrite,
|
||||
delete_marker_purge: None,
|
||||
scope: Some(scope),
|
||||
lease: None,
|
||||
enqueued_at_ms: 1,
|
||||
attempts: 0,
|
||||
};
|
||||
assert_eq!(try_rearm_mrf_replay_intent(&mut intent), MrfIngressResult::Enqueued);
|
||||
let lease = intent.lease.expect("replayed intent should own a generation");
|
||||
let anchor = MrfDurableRepairAnchor {
|
||||
kind: MrfKind::PartialWrite,
|
||||
bucket: bucket.clone(),
|
||||
object: object.clone(),
|
||||
version_id: None,
|
||||
scope: Some(scope),
|
||||
delete_marker_purge: None,
|
||||
lease,
|
||||
bucket_incarnation_id: uuid::Uuid::new_v4(),
|
||||
};
|
||||
let target = MrfRepairNoticeTarget {
|
||||
bucket: bucket.clone(),
|
||||
object: object.clone(),
|
||||
version_id: None,
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: Some(scope),
|
||||
delete_marker_purge: None,
|
||||
lease: Some(lease),
|
||||
durable_anchor: Some(anchor.clone()),
|
||||
};
|
||||
let incarnation = anchor.bucket_incarnation_id;
|
||||
let outcome = HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
kind: HealObjectKind::Object,
|
||||
bucket: bucket.to_string(),
|
||||
object: object.to_string(),
|
||||
version_id: None,
|
||||
bucket_incarnation_id: Some(incarnation),
|
||||
pool_index: Some(1),
|
||||
set_index: Some(2),
|
||||
},
|
||||
disposition: HealObjectDisposition::Unknown,
|
||||
detail: Some(rustfs_heal_contracts::heal_channel::LEGACY_OBJECT_IDENTITY_UNVERIFIED_DETAIL.to_string()),
|
||||
};
|
||||
|
||||
let event = super::scheduler::unverified_legacy_mrf_event_for_target(&target, &outcome)
|
||||
.expect("exact legacy partial-write outcome should park only its runtime retry");
|
||||
assert_eq!(event.anchor, anchor);
|
||||
|
||||
let wrong_identity = HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
object: "other-object".to_string(),
|
||||
..outcome.identity.clone()
|
||||
},
|
||||
..outcome.clone()
|
||||
};
|
||||
assert!(super::scheduler::unverified_legacy_mrf_event_for_target(&target, &wrong_identity).is_none());
|
||||
let missing_incarnation = HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
bucket_incarnation_id: None,
|
||||
..outcome.identity.clone()
|
||||
},
|
||||
..outcome.clone()
|
||||
};
|
||||
assert!(super::scheduler::unverified_legacy_mrf_event_for_target(&target, &missing_incarnation).is_none());
|
||||
let wrong_incarnation = HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
bucket_incarnation_id: Some(uuid::Uuid::new_v4()),
|
||||
..outcome.identity.clone()
|
||||
},
|
||||
..outcome.clone()
|
||||
};
|
||||
assert!(super::scheduler::unverified_legacy_mrf_event_for_target(&target, &wrong_incarnation).is_none());
|
||||
let wrong_reason = HealObjectOutcome {
|
||||
detail: Some("object was readable".to_string()),
|
||||
..outcome
|
||||
};
|
||||
assert!(super::scheduler::unverified_legacy_mrf_event_for_target(&target, &wrong_reason).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completed_mrf_notice_publishes_only_verified_positive_events() {
|
||||
use crate::heal::outcome::{HealObjectIdentity, HealObjectKind, HealObjectOutcome, HealTaskOutcome};
|
||||
@@ -1355,6 +1452,7 @@ fn completed_mrf_notice_publishes_only_verified_positive_events() {
|
||||
}),
|
||||
delete_marker_purge: None,
|
||||
lease: None,
|
||||
durable_anchor: None,
|
||||
};
|
||||
let mismatch_target = MrfRepairNoticeTarget {
|
||||
object: Arc::from("object-b"),
|
||||
|
||||
@@ -20,7 +20,9 @@
|
||||
//! cannot discharge them; only an exact storage-verified proof may do so.
|
||||
//! Both paths share the existing committed snapshot format and legacy mirrors.
|
||||
//! Replay retains partial-write intents for live retries when a member is
|
||||
//! still offline at startup. A lost proof causes another repair, not deletion.
|
||||
//! still offline at startup. Healthy legacy objects without identity proof are
|
||||
//! held in memory after one check; their unchanged journal records are retried
|
||||
//! on process restart. A lost proof causes another repair, not deletion.
|
||||
|
||||
use super::{DiskStore, HealDiskExt as _, local_disk_map_read};
|
||||
use crate::heal::manager::{HealManager, MrfRepairNoticeTarget};
|
||||
@@ -584,7 +586,11 @@ pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest {
|
||||
request
|
||||
}
|
||||
|
||||
async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> crate::Result<HealAdmissionResult> {
|
||||
async fn submit_mrf_heal_request(
|
||||
manager: &HealManager,
|
||||
intent: &MrfIntent,
|
||||
durable_anchor: Option<MrfDurableRepairAnchor>,
|
||||
) -> crate::Result<HealAdmissionResult> {
|
||||
let receipt = manager
|
||||
.submit_mrf_heal_request_with_receipt_and_identity(
|
||||
build_heal_request(intent),
|
||||
@@ -596,6 +602,7 @@ async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> c
|
||||
scope: intent.scope,
|
||||
delete_marker_purge: intent.delete_marker_purge.as_ref().map(MrfDeleteMarkerPurge::identity),
|
||||
lease: intent.lease,
|
||||
durable_anchor,
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
@@ -763,7 +770,7 @@ impl MrfRuntime {
|
||||
// attempts counter) changes the encoded snapshot; mark it dirty
|
||||
// either way.
|
||||
self.dirty = true;
|
||||
match submit_mrf_heal_request(manager, &intent).await {
|
||||
match submit_mrf_heal_request(manager, &intent, None).await {
|
||||
// Accepted intents leave the pending set; the next flush persists the
|
||||
// smaller snapshot. This is not a durable successor receipt and
|
||||
// does not discharge the producer's existing retry hints.
|
||||
@@ -796,8 +803,7 @@ impl MrfRuntime {
|
||||
}
|
||||
}
|
||||
}
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(self.queue.depth() + self.partial_writes.depth()));
|
||||
gauge!("rustfs_heal_mrf_queue_bytes").set(metric_f64(self.queue.bytes() + self.partial_writes.bytes()));
|
||||
self.publish_metrics();
|
||||
}
|
||||
|
||||
fn retained_replay_journal(&self) -> bool {
|
||||
@@ -856,12 +862,38 @@ impl MrfRuntime {
|
||||
let mut buckets: Vec<Arc<str>> = anchors.iter().map(|anchor| anchor.bucket.clone()).collect();
|
||||
buckets.sort_unstable();
|
||||
buckets.dedup();
|
||||
for bucket in buckets {
|
||||
for bucket in &buckets {
|
||||
rustfs_common::mrf_channel::consume_recorded_verified_mrf_repair_events_for(bucket.as_ref(), &mut anchors);
|
||||
}
|
||||
let remaining: HashSet<_> = anchors.into_iter().collect();
|
||||
self.durable_replay_anchors.retain(|anchor| remaining.contains(anchor));
|
||||
self.dirty |= self.partial_writes.retain_unproven(&remaining);
|
||||
for bucket in buckets {
|
||||
for event in rustfs_common::mrf_channel::take_mrf_unverified_legacy_events_for(bucket.as_ref()) {
|
||||
// Parking is an in-memory dispatch decision. The unchanged
|
||||
// durable journal remains the recovery/retry authority.
|
||||
self.partial_writes.park_unverified_legacy(&event.anchor);
|
||||
}
|
||||
}
|
||||
self.publish_metrics();
|
||||
}
|
||||
|
||||
fn publish_metrics(&self) {
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(self.queue.depth() + self.partial_writes.depth()));
|
||||
gauge!("rustfs_heal_mrf_queue_bytes").set(metric_f64(self.queue.bytes() + self.partial_writes.bytes()));
|
||||
let held = self.partial_writes.unverified_legacy_count();
|
||||
gauge!("rustfs_heal_mrf_unverified_legacy").set(metric_f64(held));
|
||||
let now_ms = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
|
||||
.unwrap_or_default();
|
||||
let oldest_age_seconds = self
|
||||
.partial_writes
|
||||
.oldest_unverified_legacy_enqueued_at_ms()
|
||||
.map(|enqueued_at_ms| now_ms.saturating_sub(enqueued_at_ms) / 1_000)
|
||||
.unwrap_or_default();
|
||||
gauge!("rustfs_heal_mrf_unverified_legacy_oldest_age_seconds")
|
||||
.set(metric_f64(usize::try_from(oldest_age_seconds).unwrap_or(usize::MAX)));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1081,16 +1113,16 @@ async fn replay_into(
|
||||
break;
|
||||
}
|
||||
if intent.kind.is_durable() {
|
||||
// Preserve the executable record as well as its proof anchor:
|
||||
// a target that is still offline during replay needs live retries.
|
||||
// Adopt durable replay into its single checkpointed owner. The
|
||||
// runtime dispatches it after publishing the successor, avoiding
|
||||
// a duplicate task in the startup-replay and steady-state paths.
|
||||
if let Some(anchor) = manager.durable_mrf_repair_anchor(&intent).await {
|
||||
durable_replay_anchors.push(anchor);
|
||||
durable_replay_anchors.push(anchor.clone());
|
||||
}
|
||||
let _ = submit_mrf_heal_request(manager, &intent).await;
|
||||
partial_writes.push(intent);
|
||||
continue;
|
||||
}
|
||||
match submit_mrf_heal_request(manager, &intent).await {
|
||||
match submit_mrf_heal_request(manager, &intent, None).await {
|
||||
Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {
|
||||
if let Some(anchor) = manager.durable_mrf_repair_anchor(&intent).await {
|
||||
durable_replay_anchors.push(anchor);
|
||||
@@ -1196,6 +1228,7 @@ async fn run_mrf_consumer(
|
||||
if runtime.dirty {
|
||||
runtime.flush().await;
|
||||
}
|
||||
runtime.publish_metrics();
|
||||
|
||||
let mut flush_tick = tokio::time::interval(runtime.config.flush_interval);
|
||||
flush_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
|
||||
@@ -1301,7 +1334,7 @@ async fn run_mrf_consumer(
|
||||
}
|
||||
TickAction::Idle => {}
|
||||
}
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(runtime.queue.depth() + runtime.partial_writes.depth()));
|
||||
runtime.publish_metrics();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,6 +21,11 @@ struct Responsibility {
|
||||
intent: MrfIntent,
|
||||
anchor: Option<MrfDurableRepairAnchor>,
|
||||
persisted: bool,
|
||||
/// A completed check identified a healthy legacy object that cannot
|
||||
/// discharge its durable intent without an independent payload proof.
|
||||
/// This is process-local only: the unchanged journal rechecks it on restart.
|
||||
unverified_legacy: bool,
|
||||
retry_queued: bool,
|
||||
next_attempt: Instant,
|
||||
}
|
||||
|
||||
@@ -49,6 +54,8 @@ impl PartialWrites {
|
||||
}
|
||||
let key = queue_key(&intent);
|
||||
let previous = self.entries.get(&key);
|
||||
let previous_was_held = previous.is_some_and(|entry| entry.unverified_legacy);
|
||||
let mut retry_queued = previous.is_some_and(|entry| entry.retry_queued);
|
||||
let old_cost = previous.map_or(0, |entry| Self::cost(&entry.intent));
|
||||
let next_bytes = self.bytes.saturating_sub(old_cost).saturating_add(Self::cost(&intent));
|
||||
if (previous.is_none() && self.entries.len() >= capacity) || next_bytes > byte_budget {
|
||||
@@ -57,8 +64,9 @@ impl PartialWrites {
|
||||
if previous.is_some_and(|entry| entry.intent.lease == intent.lease) {
|
||||
return Ok(());
|
||||
}
|
||||
if previous.is_none() {
|
||||
if previous.is_none() || (previous_was_held && !retry_queued) {
|
||||
self.retry_order.push_back(key.clone());
|
||||
retry_queued = true;
|
||||
}
|
||||
// Replacing the generation preserves the logical repair obligation,
|
||||
// but requires a new checkpoint and proof before it can be released.
|
||||
@@ -68,6 +76,8 @@ impl PartialWrites {
|
||||
intent,
|
||||
anchor: None,
|
||||
persisted: false,
|
||||
unverified_legacy: false,
|
||||
retry_queued,
|
||||
next_attempt: Instant::now(),
|
||||
},
|
||||
);
|
||||
@@ -91,6 +101,37 @@ impl PartialWrites {
|
||||
self.entries.values().filter_map(|entry| entry.anchor.as_ref())
|
||||
}
|
||||
|
||||
pub(super) fn park_unverified_legacy(&mut self, anchor: &MrfDurableRepairAnchor) -> bool {
|
||||
let key = MrfQueueKey {
|
||||
kind: anchor.kind,
|
||||
bucket: anchor.bucket.clone(),
|
||||
object: anchor.object.clone(),
|
||||
version_id: anchor.version_id,
|
||||
scope: anchor.scope,
|
||||
delete_marker_purge: anchor.delete_marker_purge,
|
||||
};
|
||||
let Some(entry) = self.entries.get_mut(&key) else {
|
||||
return false;
|
||||
};
|
||||
if entry.anchor.as_ref() != Some(anchor) || entry.unverified_legacy {
|
||||
return false;
|
||||
}
|
||||
entry.unverified_legacy = true;
|
||||
true
|
||||
}
|
||||
|
||||
pub(super) fn unverified_legacy_count(&self) -> usize {
|
||||
self.entries.values().filter(|entry| entry.unverified_legacy).count()
|
||||
}
|
||||
|
||||
pub(super) fn oldest_unverified_legacy_enqueued_at_ms(&self) -> Option<u64> {
|
||||
self.entries
|
||||
.values()
|
||||
.filter(|entry| entry.unverified_legacy)
|
||||
.map(|entry| entry.intent.enqueued_at_ms)
|
||||
.min()
|
||||
}
|
||||
|
||||
pub(super) fn mark_persisted(&mut self) {
|
||||
for entry in self.entries.values_mut() {
|
||||
entry.persisted = true;
|
||||
@@ -107,10 +148,12 @@ impl PartialWrites {
|
||||
if entry.anchor.is_none() {
|
||||
entry.anchor = manager.durable_mrf_repair_anchor(&entry.intent).await;
|
||||
}
|
||||
if entry.anchor.is_some() {
|
||||
if let Some(anchor) = entry.anchor.clone() {
|
||||
// Every outcome retains responsibility until a verified proof.
|
||||
// A full manager or an offline target only postpones another try.
|
||||
let _ = submit_mrf_heal_request(manager, &entry.intent).await;
|
||||
// A healthy legacy object without an independent identity proof
|
||||
// is parked in memory after one check; its disk journal remains
|
||||
// unchanged and startup replay checks it again.
|
||||
let _ = submit_mrf_heal_request(manager, &entry.intent, Some(anchor)).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -126,10 +169,17 @@ impl PartialWrites {
|
||||
let Some(key) = self.retry_order.pop_front() else {
|
||||
break;
|
||||
};
|
||||
let due = self
|
||||
.entries
|
||||
.get(&key)
|
||||
.is_some_and(|entry| entry.persisted && entry.next_attempt <= now);
|
||||
let Some(entry) = self.entries.get_mut(&key) else {
|
||||
continue;
|
||||
};
|
||||
entry.retry_queued = false;
|
||||
if entry.unverified_legacy {
|
||||
// Held obligations remain in the durable responsibility map,
|
||||
// but leave the hot retry index until a new generation arrives.
|
||||
continue;
|
||||
}
|
||||
let due = entry.persisted && entry.next_attempt <= now;
|
||||
entry.retry_queued = true;
|
||||
self.retry_order.push_back(key.clone());
|
||||
if due {
|
||||
ready.push(key);
|
||||
@@ -161,6 +211,7 @@ mod tests {
|
||||
use super::*;
|
||||
use rustfs_common::mrf_channel::{MrfKind, MrfScope};
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use uuid::Uuid;
|
||||
|
||||
fn intent(object: &str) -> MrfIntent {
|
||||
@@ -176,7 +227,7 @@ mod tests {
|
||||
}),
|
||||
lease: None,
|
||||
enqueued_at_ms: 1,
|
||||
attempts: u8::MAX,
|
||||
attempts: 0,
|
||||
};
|
||||
assert_eq!(try_rearm_mrf_replay_intent(&mut intent), MrfIngressResult::Enqueued);
|
||||
intent
|
||||
@@ -217,6 +268,41 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unverified_legacy_responsibility_is_held_until_restart_without_being_released() {
|
||||
let original = intent("legacy");
|
||||
let mut writes = PartialWrites::default();
|
||||
writes
|
||||
.admit(original.clone(), 1, 8192)
|
||||
.expect("durable intent should be retained");
|
||||
writes.mark_persisted();
|
||||
let anchor = MrfDurableRepairAnchor::from_intent(&original, Uuid::new_v4()).expect("exact durable anchor");
|
||||
writes.entries.get_mut(&queue_key(&original)).expect("resident intent").anchor = Some(anchor.clone());
|
||||
|
||||
assert!(writes.park_unverified_legacy(&anchor));
|
||||
assert!(!writes.park_unverified_legacy(&anchor), "duplicate notices are idempotent");
|
||||
assert_eq!(writes.unverified_legacy_count(), 1);
|
||||
assert_eq!(writes.depth(), 1, "holding the intent must preserve responsibility");
|
||||
assert!(writes.ready_keys(Instant::now() + Duration::from_secs(60), 1).is_empty());
|
||||
assert!(writes.retry_order.is_empty(), "held intents leave the retry index after one pass");
|
||||
|
||||
let replacement = intent("legacy");
|
||||
writes
|
||||
.admit(replacement, 1, 8192)
|
||||
.expect("a new generation should become retryable");
|
||||
writes.mark_persisted();
|
||||
assert_eq!(writes.unverified_legacy_count(), 0);
|
||||
assert_eq!(writes.ready_keys(Instant::now() + Duration::from_secs(60), 1).len(), 1);
|
||||
|
||||
let mut restarted = PartialWrites::default();
|
||||
restarted
|
||||
.admit(original, 1, 8192)
|
||||
.expect("the unchanged journal re-arms the same responsibility after restart");
|
||||
restarted.mark_persisted();
|
||||
assert_eq!(restarted.unverified_legacy_count(), 0);
|
||||
assert_eq!(restarted.ready_keys(Instant::now() + Duration::from_secs(60), 1).len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn partial_write_retention_adopts_replay_anchor_before_generation_replacement() {
|
||||
use super::super::{MrfQueue, MrfRuntime};
|
||||
|
||||
@@ -93,6 +93,12 @@ impl HealTask {
|
||||
"Heal object stage entered"
|
||||
);
|
||||
self.check_control_flags().await?;
|
||||
if self.source == HealRequestSource::Mrf {
|
||||
// Durable partial writes always use the incarnation-fenced storage
|
||||
// path, including for present objects. A normal heal can inspect a
|
||||
// recreated bucket after the replay anchor was captured.
|
||||
return self.heal_mrf_partial_write_object(bucket, object, version_id).await;
|
||||
}
|
||||
let mut object_exists = match self.await_with_control(self.storage.object_exists(bucket, object)).await {
|
||||
Ok(exists) => exists,
|
||||
Err(err @ Error::TransientSkip { .. }) => {
|
||||
@@ -429,10 +435,6 @@ impl HealTask {
|
||||
|
||||
/// Recreate missing object (for EC decode scenarios)
|
||||
async fn recreate_missing_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
|
||||
if self.source == HealRequestSource::Mrf {
|
||||
return self.recreate_missing_mrf_object(bucket, object, version_id).await;
|
||||
}
|
||||
|
||||
debug!(
|
||||
target: "rustfs::heal::task",
|
||||
event = EVENT_HEAL_OBJECT_STAGE,
|
||||
@@ -534,8 +536,8 @@ impl HealTask {
|
||||
}
|
||||
}
|
||||
|
||||
/// Durable MRF responsibilities may complete only with an exact storage proof.
|
||||
async fn recreate_missing_mrf_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
|
||||
/// Durable MRF partial writes may complete only with an incarnation-fenced storage proof.
|
||||
async fn heal_mrf_partial_write_object(&self, bucket: &str, object: &str, version_id: Option<&str>) -> Result<()> {
|
||||
let heal_opts = HealOpts {
|
||||
recursive: false,
|
||||
dry_run: self.options.dry_run,
|
||||
@@ -577,7 +579,30 @@ impl HealTask {
|
||||
storage_result.receipt.as_ref().map(|receipt| &receipt.disposition),
|
||||
Some(HealObjectDisposition::AuthoritativelyAbsent)
|
||||
);
|
||||
if !self.record_verified_storage_receipt(expected, storage_result.receipt).await {
|
||||
let receipt_missing = storage_result.receipt.is_none();
|
||||
if !self
|
||||
.record_verified_storage_receipt(expected.clone(), storage_result.receipt)
|
||||
.await
|
||||
{
|
||||
let ok_drive_state = DriveState::Ok.to_string();
|
||||
let healthy_legacy_without_repair = receipt_missing
|
||||
&& storage_result.item.detail == rustfs_heal_contracts::heal_channel::LEGACY_OBJECT_IDENTITY_UNVERIFIED_DETAIL
|
||||
&& storage_result.item.drives_healed() == Some(0)
|
||||
&& !storage_result.item.after.drives.is_empty()
|
||||
&& storage_result
|
||||
.item
|
||||
.after
|
||||
.drives
|
||||
.iter()
|
||||
.all(|drive| drive.state == ok_drive_state);
|
||||
if healthy_legacy_without_repair {
|
||||
self.outcome.write().await.record(HealObjectOutcome {
|
||||
identity: expected,
|
||||
disposition: HealObjectDisposition::Unknown,
|
||||
detail: Some(storage_result.item.detail.clone()),
|
||||
});
|
||||
}
|
||||
self.record_result_item(storage_result.item).await;
|
||||
return Err(Error::TaskExecutionFailed {
|
||||
message: format!("Missing exact storage proof for durable MRF repair {bucket}/{object}"),
|
||||
});
|
||||
|
||||
@@ -588,6 +588,15 @@ async fn partial_write_sigkill_replay_scenario(protected: bool) {
|
||||
.await,
|
||||
"legacy replay attempts must finish"
|
||||
);
|
||||
// Cover two MRF admission backoff periods so a missed hold notice
|
||||
// cannot pass merely because its task finished between polls.
|
||||
let retry_window = tokio::time::Instant::now() + Duration::from_secs(12);
|
||||
while tokio::time::Instant::now() < retry_window {
|
||||
let snapshot = manager.operations_snapshot().await;
|
||||
assert_eq!(snapshot.queue_length, 0, "an unverified legacy result must not refill the manager queue");
|
||||
assert_eq!(snapshot.active_tasks, 0, "a held intent must not stay active");
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
}
|
||||
assert!(snapshot_contains("crash.bin").await, "unverified legacy responsibility must remain");
|
||||
}
|
||||
manager.stop().await.expect("restarted manager should stop");
|
||||
|
||||
@@ -284,6 +284,16 @@ Reports scanner-driven bitrot state together with heal queue execution state. `h
|
||||
|
||||
Use this route when `metrics.source_work` shows `heal` or `bitrot` queued or missed work. Scanner-originated object checks should appear under `scanner/low`, manual admin heal under `admin/high`. If scanner work grows but admin work remains blocked, treat that as heal queue pressure rather than scanner pacing pressure.
|
||||
|
||||
Durable partial-write responsibilities also have node-local MRF metrics:
|
||||
|
||||
| Metric | Meaning |
|
||||
|---|---|
|
||||
| `rustfs_heal_mrf_queue_depth` | In-memory MRF work plus retained durable responsibilities, including held legacy entries. |
|
||||
| `rustfs_heal_mrf_unverified_legacy` | Durable `PartialWrite` entries whose completed deep check found healthy legacy data without an independent payload identity proof. The journal entry remains intact; the current process pauses automatic re-dispatch for that entry. |
|
||||
| `rustfs_heal_mrf_unverified_legacy_oldest_age_seconds` | Age of the oldest such retained responsibility on this node. |
|
||||
|
||||
These gauges are emitted per node through the configured OTLP metrics exporter (`RUSTFS_OBS_ENDPOINT` or `RUSTFS_OBS_METRIC_ENDPOINT`). The background-heal status endpoint reports execution tasks; it does not include durable MRF responsibilities. A process restart replays unchanged journal records and checks them again; it does not delete or certify an entry, and it replays all other durable intents on that node as well. For an unversioned object, restore the expected content from a trusted canonical source with protected shard integrity enabled, then restart the node that owns the journal so its intent can obtain an exact receipt. For a versioned object, a protected rewrite creates a new version and does not resolve a responsibility for the old version; preserve the source and only retire that exact version when the intended data is backed up and an authoritative absence proof is appropriate. A successful receipt may discharge only the matching responsibility; an unverified result remains retained and becomes held again. If no trusted copy or expected checksum is available, preserve the held responsibility and investigate its source of truth before retrying. The protected-copy migration in the [shard-integrity audit workflow](shard-integrity-audit.md) writes a new key and preserves its source; completing that migration alone does not prove or discharge an intent for the original key. Never remove MRF journal files manually.
|
||||
|
||||
## Heal runtime controls
|
||||
|
||||
Heal knobs are environment-only and read by `HealConfig::default` (`crates/heal/src/heal/manager.rs`), the MRF queue (`crates/heal/src/heal/mrf_queue.rs`), or the erasure-set healer (`crates/heal/src/heal/erasure_healer.rs`). The admin `heal` config subsystem accepts only `bitrot_cycle` (`HEAL_KEYS`), which is documented in the scanner table above. Constants live in `crates/config/src/constants/heal.rs` unless another file is named.
|
||||
|
||||
Reference in New Issue
Block a user