mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 20:46:11 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 61536636ec | |||
| 7937c4c0ad | |||
| 7ab93d72d2 |
@@ -1 +1 @@
|
||||
sha256=95c8adc016bbc0df9fb2afa24a108bcdf6567ec4d0518725a6cae301593ab556
|
||||
sha256=0fe8408874ccec3620262a9812d67920ddd72dc9edf0e36e0d0aed3f8bad026e
|
||||
|
||||
@@ -90,61 +90,6 @@ pub struct MrfIntent {
|
||||
pub attempts: u8,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub struct MrfDurableRepairAnchor {
|
||||
pub kind: MrfKind,
|
||||
pub bucket: Arc<str>,
|
||||
pub object: Arc<str>,
|
||||
pub version_id: Option<[u8; 16]>,
|
||||
pub scope: Option<MrfScope>,
|
||||
pub lease: MrfIngressLease,
|
||||
pub bucket_incarnation_id: Uuid,
|
||||
}
|
||||
|
||||
impl MrfDurableRepairAnchor {
|
||||
/// Build a dischargeable anchor only when the caller supplies the storage
|
||||
/// incarnation and the original ingress lease. Legacy replay records lack
|
||||
/// both pieces and therefore remain fail-closed.
|
||||
pub fn from_intent(intent: &MrfIntent, bucket_incarnation_id: Uuid) -> Option<Self> {
|
||||
if bucket_incarnation_id.is_nil() {
|
||||
return None;
|
||||
}
|
||||
let lease = intent.lease?;
|
||||
let (version_id, scope) = canonical_identity(intent.kind, intent.version_id, intent.scope);
|
||||
Some(Self {
|
||||
kind: intent.kind,
|
||||
bucket: intent.bucket.clone(),
|
||||
object: intent.object.clone(),
|
||||
version_id,
|
||||
scope,
|
||||
lease,
|
||||
bucket_incarnation_id,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn is_proven_by(&self, event: &MrfVerifiedRepairEvent) -> bool {
|
||||
let Some(lease) = event.lease else {
|
||||
return false;
|
||||
};
|
||||
self.kind == event.kind
|
||||
&& self.bucket == event.bucket
|
||||
&& self.object == event.object
|
||||
&& self.version_id == event.version_id
|
||||
&& self.scope == event.scope
|
||||
&& self.lease == lease
|
||||
&& self.bucket_incarnation_id == event.bucket_incarnation_id
|
||||
}
|
||||
}
|
||||
|
||||
/// Consume only anchors proven by a complete verified-repair identity. The
|
||||
/// caller remains responsible for persisting the resulting anchor set before
|
||||
/// deleting older replay files.
|
||||
pub fn consume_verified_mrf_repair_events(anchors: &mut Vec<MrfDurableRepairAnchor>, events: &[MrfVerifiedRepairEvent]) -> usize {
|
||||
let before = anchors.len();
|
||||
anchors.retain(|anchor| !events.iter().any(|event| anchor.is_proven_by(event)));
|
||||
before.saturating_sub(anchors.len())
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
|
||||
pub struct MrfScope {
|
||||
pub pool_index: u32,
|
||||
@@ -441,27 +386,6 @@ pub fn try_send_mrf_intent_typed(
|
||||
}
|
||||
}
|
||||
|
||||
/// Acquire a fresh process-local lease for one durable replay record.
|
||||
///
|
||||
/// Journal records deliberately do not persist leases. A replay consumer must
|
||||
/// call this before submitting the record so a later verified repair event can
|
||||
/// identify the exact replay admission. Replay does not reserve the live
|
||||
/// producer coalescer key: the replay queue first deduplicates legacy records,
|
||||
/// then manager admission owns task-level deduplication with live producers.
|
||||
pub fn try_rearm_mrf_replay_intent(intent: &mut MrfIntent) -> MrfIngressResult {
|
||||
if intent.lease.is_some() {
|
||||
return MrfIngressResult::Enqueued;
|
||||
}
|
||||
if intent.bucket.len() > MRF_MAX_IDENTITY_COMPONENT || intent.object.len() > MRF_MAX_IDENTITY_COMPONENT {
|
||||
return MrfIngressResult::Dropped(MrfDropReason::OversizedIdentity);
|
||||
}
|
||||
let (version_id, scope) = canonical_identity(intent.kind, intent.version_id, intent.scope);
|
||||
intent.version_id = version_id;
|
||||
intent.scope = scope;
|
||||
intent.lease = Some(MrfIngressLease::new(NEXT_MRF_LEASE.fetch_add(1, Ordering::Relaxed)));
|
||||
MrfIngressResult::Enqueued
|
||||
}
|
||||
|
||||
/// Release the ingress key once the consumer owns the intent.
|
||||
pub fn release_mrf_intent(intent: &MrfIntent) {
|
||||
release_mrf_identity(intent.kind, &intent.bucket, &intent.object, intent.version_id, intent.scope, intent.lease);
|
||||
@@ -508,33 +432,12 @@ pub struct MrfRepairedEvent {
|
||||
pub version_id: Option<[u8; 16]>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum MrfVerifiedRepairDisposition {
|
||||
Repaired,
|
||||
VerifiedHealthy,
|
||||
AuthoritativelyAbsent,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct MrfVerifiedRepairEvent {
|
||||
pub kind: MrfKind,
|
||||
pub bucket: Arc<str>,
|
||||
pub object: Arc<str>,
|
||||
pub version_id: Option<[u8; 16]>,
|
||||
pub scope: Option<MrfScope>,
|
||||
pub lease: Option<MrfIngressLease>,
|
||||
pub bucket_incarnation_id: Uuid,
|
||||
pub disposition: MrfVerifiedRepairDisposition,
|
||||
}
|
||||
|
||||
/// Bound on the repaired-event backlog. Notices are best-effort hints; when
|
||||
/// the ring is full the oldest are dropped and the affected ledger entries
|
||||
/// simply expire through their own attempts/age limits.
|
||||
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();
|
||||
|
||||
/// Record a legacy notification for compatibility. This is not an
|
||||
/// acknowledgement of storage verification or durable repair completion.
|
||||
@@ -575,43 +478,6 @@ pub fn take_mrf_repaired_events_for(bucket: &str) -> Vec<MrfRepairedEvent> {
|
||||
taken
|
||||
}
|
||||
|
||||
/// Record a storage-owned MRF completion proof. Unlike the legacy repaired
|
||||
/// event, this identity is complete enough for future durable ledgers to make
|
||||
/// an exact responsibility decision.
|
||||
pub fn note_mrf_verified_repair(event: MrfVerifiedRepairEvent) {
|
||||
let registry = MRF_VERIFIED_REPAIR_EVENTS.get_or_init(|| std::sync::Mutex::new(std::collections::VecDeque::new()));
|
||||
let Ok(mut events) = registry.lock() else {
|
||||
return;
|
||||
};
|
||||
if events.len() >= MRF_REPAIRED_EVENT_CAP {
|
||||
events.pop_front();
|
||||
}
|
||||
events.push_back(event);
|
||||
}
|
||||
|
||||
/// Take verified repair events recorded for `bucket`, leaving other buckets'
|
||||
/// proofs in place. Consumers still have to match kind, object, version, scope
|
||||
/// lease and incarnation before discharging durable responsibility.
|
||||
pub fn take_mrf_verified_repair_events_for(bucket: &str) -> Vec<MrfVerifiedRepairEvent> {
|
||||
let Some(registry) = MRF_VERIFIED_REPAIR_EVENTS.get() else {
|
||||
return Vec::new();
|
||||
};
|
||||
let Ok(mut events) = registry.lock() else {
|
||||
return Vec::new();
|
||||
};
|
||||
let mut taken = Vec::new();
|
||||
let mut retained = std::collections::VecDeque::with_capacity(events.len());
|
||||
while let Some(event) = events.pop_front() {
|
||||
if event.bucket.as_ref() == bucket {
|
||||
taken.push(event);
|
||||
} else {
|
||||
retained.push_back(event);
|
||||
}
|
||||
}
|
||||
*events = retained;
|
||||
taken
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -680,147 +546,6 @@ mod tests {
|
||||
assert_eq!(metadata_scope, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn durable_repair_anchor_requires_lease_and_bucket_incarnation() {
|
||||
let mut intent = MrfIntent {
|
||||
bucket: Arc::from("durable-anchor-bucket"),
|
||||
object: Arc::from("object"),
|
||||
version_id: Some([0; 16]),
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
lease: None,
|
||||
enqueued_at_ms: 0,
|
||||
attempts: 0,
|
||||
};
|
||||
assert!(
|
||||
MrfDurableRepairAnchor::from_intent(&intent, Uuid::new_v4()).is_none(),
|
||||
"legacy replay records without the ingress lease must remain anchored"
|
||||
);
|
||||
|
||||
intent.lease = Some(MrfIngressLease::new(7));
|
||||
assert!(
|
||||
MrfDurableRepairAnchor::from_intent(&intent, Uuid::nil()).is_none(),
|
||||
"nil bucket incarnation cannot prove durable successor ownership"
|
||||
);
|
||||
|
||||
let anchor = MrfDurableRepairAnchor::from_intent(&intent, Uuid::new_v4())
|
||||
.expect("complete identity should create a durable repair anchor");
|
||||
assert_eq!(anchor.version_id, None, "nil UUID is canonicalized before matching");
|
||||
assert_eq!(
|
||||
anchor.scope,
|
||||
Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn durable_replay_rearm_assigns_a_fresh_dischargeable_lease() {
|
||||
let unique = Uuid::new_v4();
|
||||
let mut intent = MrfIntent {
|
||||
bucket: Arc::from(format!("replay-{unique}")),
|
||||
object: Arc::from("object"),
|
||||
version_id: Some([0; 16]),
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 2,
|
||||
set_index: 3,
|
||||
}),
|
||||
lease: None,
|
||||
enqueued_at_ms: 1,
|
||||
attempts: 0,
|
||||
};
|
||||
|
||||
assert_eq!(try_rearm_mrf_replay_intent(&mut intent), MrfIngressResult::Enqueued);
|
||||
assert_eq!(intent.version_id, None, "nil versions remain canonical during replay");
|
||||
assert!(intent.lease.is_some(), "replay admission must carry a fresh lease");
|
||||
assert!(
|
||||
MrfDurableRepairAnchor::from_intent(&intent, Uuid::new_v4()).is_some(),
|
||||
"a rearmed replay record can participate in exact durable proof matching"
|
||||
);
|
||||
release_mrf_intent(&intent);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn verified_repair_events_consume_only_exact_durable_anchors() {
|
||||
let bucket = Arc::<str>::from("proof-bucket");
|
||||
let object = Arc::<str>::from("object");
|
||||
let incarnation = Uuid::new_v4();
|
||||
let lease = MrfIngressLease::new(11);
|
||||
let anchor = MrfDurableRepairAnchor {
|
||||
kind: MrfKind::PartialWrite,
|
||||
bucket: bucket.clone(),
|
||||
object: object.clone(),
|
||||
version_id: Some([3; 16]),
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 4,
|
||||
set_index: 5,
|
||||
}),
|
||||
lease,
|
||||
bucket_incarnation_id: incarnation,
|
||||
};
|
||||
let event = MrfVerifiedRepairEvent {
|
||||
kind: anchor.kind,
|
||||
bucket,
|
||||
object,
|
||||
version_id: anchor.version_id,
|
||||
scope: anchor.scope,
|
||||
lease: Some(lease),
|
||||
bucket_incarnation_id: incarnation,
|
||||
disposition: MrfVerifiedRepairDisposition::Repaired,
|
||||
};
|
||||
|
||||
for rejected in [
|
||||
MrfVerifiedRepairEvent {
|
||||
lease: None,
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
lease: Some(MrfIngressLease::new(12)),
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
bucket_incarnation_id: Uuid::new_v4(),
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
version_id: Some([4; 16]),
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 4,
|
||||
set_index: 6,
|
||||
}),
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
kind: MrfKind::DecodeFailure,
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
bucket: Arc::from("other-bucket"),
|
||||
..event.clone()
|
||||
},
|
||||
MrfVerifiedRepairEvent {
|
||||
object: Arc::from("other"),
|
||||
..event.clone()
|
||||
},
|
||||
] {
|
||||
let mut retained = vec![anchor.clone()];
|
||||
assert_eq!(consume_verified_mrf_repair_events(&mut retained, &[rejected]), 0);
|
||||
assert_eq!(retained, vec![anchor.clone()]);
|
||||
}
|
||||
|
||||
let mut retained = vec![anchor];
|
||||
assert_eq!(consume_verified_mrf_repair_events(&mut retained, &[event]), 1);
|
||||
assert!(retained.is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn try_send_delivers_and_respects_capacity() {
|
||||
let mut receiver = init_mrf_channel().expect("first initialization should succeed");
|
||||
@@ -885,32 +610,4 @@ mod tests {
|
||||
assert_eq!(flooded.len(), MRF_REPAIRED_EVENT_CAP);
|
||||
assert_eq!(flooded[0].object.as_ref(), "object-9", "the oldest notices past the cap are dropped");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn verified_repair_events_preserve_full_identity_and_bucket_scope() {
|
||||
let bucket_incarnation_id = Uuid::new_v4();
|
||||
let event = MrfVerifiedRepairEvent {
|
||||
kind: MrfKind::PartialWrite,
|
||||
bucket: Arc::from("verified-bucket-a"),
|
||||
object: Arc::from("object-a"),
|
||||
version_id: Some([4u8; 16]),
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 2,
|
||||
set_index: 3,
|
||||
}),
|
||||
lease: Some(MrfIngressLease::new(42)),
|
||||
bucket_incarnation_id,
|
||||
disposition: MrfVerifiedRepairDisposition::Repaired,
|
||||
};
|
||||
note_mrf_verified_repair(event.clone());
|
||||
note_mrf_verified_repair(MrfVerifiedRepairEvent {
|
||||
bucket: Arc::from("verified-bucket-b"),
|
||||
..event.clone()
|
||||
});
|
||||
|
||||
let taken = take_mrf_verified_repair_events_for("verified-bucket-a");
|
||||
assert_eq!(taken, vec![event]);
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -313,7 +313,6 @@ impl HealManager {
|
||||
completed_status_entry.outcome = Some(Arc::new(task.get_outcome().await));
|
||||
}
|
||||
let terminal_completion = !matches!(completed_status, HealTaskStatus::Retrying { .. });
|
||||
let completed_status_for_verified_events = completed_status_entry.clone();
|
||||
// Keep retry ownership continuous: status snapshots acquire
|
||||
// these locks in the same active -> retrying order.
|
||||
let mut retrying_heals_guard = if let (Some((request, _, error)), Some(cancel_token)) =
|
||||
@@ -359,15 +358,6 @@ impl HealManager {
|
||||
tests::pause_completed_retention_handoff(&task_id).await;
|
||||
|
||||
if completed_task.is_some() {
|
||||
let notice_targets = if terminal_completion {
|
||||
take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id)
|
||||
} else {
|
||||
Vec::new()
|
||||
};
|
||||
if terminal_completion {
|
||||
release_mrf_repair_notice_targets(¬ice_targets);
|
||||
}
|
||||
publish_verified_mrf_repair_events(¬ice_targets, &completed_status_for_verified_events);
|
||||
// update statistics
|
||||
let mut stats = statistics_clone.write().await;
|
||||
match completed_status {
|
||||
@@ -382,6 +372,14 @@ impl HealManager {
|
||||
}
|
||||
stats.update_running_tasks(usize_to_u64_saturated(active_count));
|
||||
drop(stats);
|
||||
if terminal_completion {
|
||||
let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id);
|
||||
// Neither task status nor the diagnostic outcome
|
||||
// window supplies a storage-owned repair receipt.
|
||||
// Release only the ingress lease for rediscovery;
|
||||
// preserve the producer's existing retry hints.
|
||||
release_mrf_repair_notice_targets(notice_targets);
|
||||
}
|
||||
}
|
||||
|
||||
if let (Some((retry_request, retry_delay, retry_error)), Some(retry_cancel_token)) =
|
||||
@@ -707,7 +705,7 @@ fn move_mrf_repair_notice_targets(
|
||||
}
|
||||
}
|
||||
|
||||
fn release_mrf_repair_notice_targets(targets: &[MrfRepairNoticeTarget]) {
|
||||
fn release_mrf_repair_notice_targets(targets: Vec<MrfRepairNoticeTarget>) {
|
||||
for target in targets {
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
@@ -720,73 +718,6 @@ fn release_mrf_repair_notice_targets(targets: &[MrfRepairNoticeTarget]) {
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn mrf_verified_repair_event_for_target(
|
||||
target: &MrfRepairNoticeTarget,
|
||||
outcome: &crate::heal::outcome::HealObjectOutcome,
|
||||
) -> Option<rustfs_common::mrf_channel::MrfVerifiedRepairEvent> {
|
||||
use crate::heal::outcome::{HealObjectDisposition, HealObjectKind};
|
||||
use rustfs_common::mrf_channel::{MrfKind, MrfVerifiedRepairDisposition};
|
||||
|
||||
let disposition = match outcome.disposition {
|
||||
HealObjectDisposition::Repaired => MrfVerifiedRepairDisposition::Repaired,
|
||||
HealObjectDisposition::VerifiedHealthy => MrfVerifiedRepairDisposition::VerifiedHealthy,
|
||||
HealObjectDisposition::AuthoritativelyAbsent => MrfVerifiedRepairDisposition::AuthoritativelyAbsent,
|
||||
_ => return None,
|
||||
};
|
||||
if target.kind != MrfKind::PartialWrite {
|
||||
return None;
|
||||
}
|
||||
let expected_kind = HealObjectKind::Object;
|
||||
if outcome.identity.kind != expected_kind
|
||||
|| outcome.identity.bucket.as_str() != target.bucket.as_ref()
|
||||
|| outcome.identity.object.as_str() != target.object.as_ref()
|
||||
{
|
||||
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());
|
||||
if outcome.identity.version_id != expected_version {
|
||||
return None;
|
||||
}
|
||||
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 outcome.identity.pool_index != expected_pool || outcome.identity.set_index != expected_set {
|
||||
return None;
|
||||
}
|
||||
let bucket_incarnation_id = outcome.identity.bucket_incarnation_id?;
|
||||
Some(rustfs_common::mrf_channel::MrfVerifiedRepairEvent {
|
||||
kind: target.kind,
|
||||
bucket: target.bucket.clone(),
|
||||
object: target.object.clone(),
|
||||
version_id,
|
||||
scope: target.scope,
|
||||
lease: target.lease,
|
||||
bucket_incarnation_id,
|
||||
disposition,
|
||||
})
|
||||
}
|
||||
|
||||
pub(super) fn publish_verified_mrf_repair_events(targets: &[MrfRepairNoticeTarget], completed: &CompletedHealStatus) {
|
||||
if completed.status != HealTaskStatus::Completed {
|
||||
return;
|
||||
}
|
||||
let Some(outcome) = completed.outcome.as_ref() else {
|
||||
return;
|
||||
};
|
||||
if outcome.execution != crate::heal::outcome::HealExecutionOutcome::Completed {
|
||||
return;
|
||||
}
|
||||
for target in targets {
|
||||
if let Some(event) = outcome
|
||||
.objects
|
||||
.iter()
|
||||
.find_map(|object| mrf_verified_repair_event_for_target(target, object))
|
||||
{
|
||||
rustfs_common::mrf_channel::note_mrf_verified_repair(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()),
|
||||
|
||||
@@ -1094,191 +1094,6 @@ fn queued_request_id_for_dedup_key_tracks_the_representative() {
|
||||
assert!(queue.queued_request_id_for_dedup_key(&first_key).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mrf_verified_repair_event_requires_positive_exact_identity() {
|
||||
use crate::heal::outcome::{HealObjectIdentity, HealObjectKind, HealObjectOutcome};
|
||||
use rustfs_common::mrf_channel::{MrfKind, MrfScope, MrfVerifiedRepairDisposition};
|
||||
|
||||
let version = uuid::Uuid::new_v4();
|
||||
let incarnation = uuid::Uuid::new_v4();
|
||||
let target = MrfRepairNoticeTarget {
|
||||
bucket: Arc::from("bucket"),
|
||||
object: Arc::from("object"),
|
||||
version_id: Some(*version.as_bytes()),
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
lease: None,
|
||||
};
|
||||
let matching = HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
kind: HealObjectKind::Object,
|
||||
bucket: "bucket".to_string(),
|
||||
object: "object".to_string(),
|
||||
version_id: Some(version.to_string()),
|
||||
bucket_incarnation_id: Some(incarnation),
|
||||
pool_index: Some(1),
|
||||
set_index: Some(2),
|
||||
},
|
||||
disposition: HealObjectDisposition::Repaired,
|
||||
detail: None,
|
||||
};
|
||||
|
||||
let event = mrf_verified_repair_event_for_target(&target, &matching).expect("matching positive receipt should publish");
|
||||
assert_eq!(event.kind, MrfKind::PartialWrite);
|
||||
assert_eq!(event.bucket.as_ref(), "bucket");
|
||||
assert_eq!(event.object.as_ref(), "object");
|
||||
assert_eq!(event.version_id, Some(*version.as_bytes()));
|
||||
assert_eq!(
|
||||
event.scope,
|
||||
Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2
|
||||
})
|
||||
);
|
||||
assert_eq!(event.lease, None);
|
||||
assert_eq!(event.bucket_incarnation_id, incarnation);
|
||||
assert_eq!(event.disposition, MrfVerifiedRepairDisposition::Repaired);
|
||||
|
||||
assert!(
|
||||
mrf_verified_repair_event_for_target(
|
||||
&MrfRepairNoticeTarget {
|
||||
kind: MrfKind::DecodeFailure,
|
||||
..target.clone()
|
||||
},
|
||||
&matching
|
||||
)
|
||||
.is_none(),
|
||||
"only receipt-producing partial-write object heals can publish verified events today"
|
||||
);
|
||||
|
||||
for rejected in [
|
||||
HealObjectOutcome {
|
||||
disposition: HealObjectDisposition::Unknown,
|
||||
..matching.clone()
|
||||
},
|
||||
HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
bucket_incarnation_id: None,
|
||||
..matching.identity.clone()
|
||||
},
|
||||
..matching.clone()
|
||||
},
|
||||
HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
object: "other".to_string(),
|
||||
..matching.identity.clone()
|
||||
},
|
||||
..matching.clone()
|
||||
},
|
||||
HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
pool_index: Some(3),
|
||||
..matching.identity.clone()
|
||||
},
|
||||
..matching
|
||||
},
|
||||
] {
|
||||
assert!(
|
||||
mrf_verified_repair_event_for_target(&target, &rejected).is_none(),
|
||||
"legacy, incomplete or mismatched outcomes must not discharge MRF responsibility"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completed_mrf_notice_publishes_only_verified_positive_events() {
|
||||
use crate::heal::outcome::{HealObjectIdentity, HealObjectKind, HealObjectOutcome, HealTaskOutcome};
|
||||
use rustfs_common::mrf_channel::{MrfKind, MrfScope, take_mrf_verified_repair_events_for};
|
||||
|
||||
let bucket = Arc::<str>::from("verified-mrf-completed-bucket");
|
||||
let _ = take_mrf_verified_repair_events_for(bucket.as_ref());
|
||||
let version = uuid::Uuid::new_v4();
|
||||
let incarnation = uuid::Uuid::new_v4();
|
||||
let matching_target = MrfRepairNoticeTarget {
|
||||
bucket: bucket.clone(),
|
||||
object: Arc::from("object-a"),
|
||||
version_id: Some(*version.as_bytes()),
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
lease: None,
|
||||
};
|
||||
let mismatch_target = MrfRepairNoticeTarget {
|
||||
object: Arc::from("object-b"),
|
||||
..matching_target.clone()
|
||||
};
|
||||
let mut outcome = HealTaskOutcome::default();
|
||||
outcome.execution = HealExecutionOutcome::Completed;
|
||||
outcome.coverage = crate::heal::outcome::HealTraversalCoverage::Complete;
|
||||
outcome.record(HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
kind: HealObjectKind::Object,
|
||||
bucket: bucket.to_string(),
|
||||
object: "object-a".to_string(),
|
||||
version_id: Some(version.to_string()),
|
||||
bucket_incarnation_id: Some(incarnation),
|
||||
pool_index: Some(1),
|
||||
set_index: Some(2),
|
||||
},
|
||||
disposition: HealObjectDisposition::VerifiedHealthy,
|
||||
detail: None,
|
||||
});
|
||||
outcome.record(HealObjectOutcome {
|
||||
identity: HealObjectIdentity {
|
||||
kind: HealObjectKind::Object,
|
||||
bucket: bucket.to_string(),
|
||||
object: "object-b".to_string(),
|
||||
version_id: Some(version.to_string()),
|
||||
bucket_incarnation_id: None,
|
||||
pool_index: Some(1),
|
||||
set_index: Some(2),
|
||||
},
|
||||
disposition: HealObjectDisposition::Repaired,
|
||||
detail: None,
|
||||
});
|
||||
let completed = CompletedHealStatus {
|
||||
outcome: Some(Arc::new(outcome)),
|
||||
..completed_retention_fixture(SystemTime::now())
|
||||
};
|
||||
|
||||
publish_verified_mrf_repair_events(&[matching_target.clone(), mismatch_target], &completed);
|
||||
|
||||
let events = take_mrf_verified_repair_events_for(bucket.as_ref());
|
||||
assert_eq!(events.len(), 1);
|
||||
assert_eq!(events[0].object.as_ref(), "object-a");
|
||||
assert_eq!(events[0].lease, None);
|
||||
assert_eq!(events[0].bucket_incarnation_id, incarnation);
|
||||
|
||||
let failed = CompletedHealStatus {
|
||||
status: HealTaskStatus::Failed {
|
||||
error: "terminal failure".to_string(),
|
||||
},
|
||||
..completed.clone()
|
||||
};
|
||||
publish_verified_mrf_repair_events(std::slice::from_ref(&matching_target), &failed);
|
||||
assert!(
|
||||
take_mrf_verified_repair_events_for(bucket.as_ref()).is_empty(),
|
||||
"failed terminal tasks must not publish a verified repair event"
|
||||
);
|
||||
|
||||
let mut completed_with_errors_outcome = completed.outcome.as_ref().expect("completed outcome").as_ref().clone();
|
||||
completed_with_errors_outcome.execution = crate::heal::outcome::HealExecutionOutcome::CompletedWithErrors;
|
||||
let completed_with_errors = CompletedHealStatus {
|
||||
outcome: Some(Arc::new(completed_with_errors_outcome)),
|
||||
..completed
|
||||
};
|
||||
publish_verified_mrf_repair_events(std::slice::from_ref(&matching_target), &completed_with_errors);
|
||||
assert!(
|
||||
take_mrf_verified_repair_events_for(bucket.as_ref()).is_empty(),
|
||||
"non-success canonical outcomes must not publish a verified repair event"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_priority_queue_ordering() {
|
||||
let mut queue = PriorityHealQueue::new();
|
||||
|
||||
@@ -36,7 +36,7 @@
|
||||
use super::{DiskStore, HealDiskExt as _, local_disk_map_read};
|
||||
use crate::heal::manager::{HealManager, MrfRepairNoticeTarget};
|
||||
use metrics::{counter, gauge};
|
||||
use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIngressResult, MrfIntent};
|
||||
use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIntent};
|
||||
use rustfs_heal_contracts::heal_channel::{HealAdmissionDropReason, HealAdmissionResult};
|
||||
use std::collections::{HashSet, VecDeque};
|
||||
use std::sync::Arc;
|
||||
@@ -702,9 +702,9 @@ async fn replay_into(
|
||||
}
|
||||
},
|
||||
};
|
||||
let mut intents = Vec::new();
|
||||
let (decoded, truncated) = decode_journal(&data);
|
||||
let replayed = decoded.len();
|
||||
let intents = decoded;
|
||||
intents.extend(decoded);
|
||||
if truncated > 0 {
|
||||
tracing::warn!(
|
||||
target: "rustfs::heal::mrf",
|
||||
@@ -712,7 +712,8 @@ async fn replay_into(
|
||||
"MRF journal had a torn tail; truncated records were discarded"
|
||||
);
|
||||
}
|
||||
counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(replayed).unwrap_or(u64::MAX));
|
||||
counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(intents.len()).unwrap_or(u64::MAX));
|
||||
let replayed = intents.len();
|
||||
let replay_bytes = intents
|
||||
.iter()
|
||||
.fold(0usize, |total, intent| total.saturating_add(intent.estimated_bytes()));
|
||||
@@ -738,15 +739,6 @@ async fn replay_into(
|
||||
// stays armed in `queue` for the consumer's retry loop.
|
||||
if backoff_until.is_none() {
|
||||
while let Some(mut intent) = queue.pop_front() {
|
||||
if !matches!(
|
||||
rustfs_common::mrf_channel::try_rearm_mrf_replay_intent(&mut intent),
|
||||
MrfIngressResult::Enqueued
|
||||
) {
|
||||
queue.push_back(intent);
|
||||
rearm_incomplete = true;
|
||||
*backoff_until = Some(tokio::time::Instant::now());
|
||||
break;
|
||||
}
|
||||
match submit_mrf_heal_request(manager, &intent).await {
|
||||
Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {}
|
||||
Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => {
|
||||
@@ -761,9 +753,7 @@ async fn replay_into(
|
||||
}
|
||||
break;
|
||||
}
|
||||
Ok(HealAdmissionResult::Dropped(_)) => {
|
||||
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
||||
}
|
||||
Ok(HealAdmissionResult::Dropped(_)) => {}
|
||||
Err(_) => {
|
||||
intent.attempts = intent.attempts.saturating_add(1);
|
||||
if intent.attempts < MRF_MAX_ATTEMPTS {
|
||||
@@ -965,40 +955,6 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn durable_replay_acquires_a_fresh_lease_before_manager_admission() {
|
||||
let unique = uuid::Uuid::new_v4();
|
||||
let original = intent(&format!("replay-{unique}"), "object", 0);
|
||||
assert!(original.lease.is_none(), "legacy journal records do not persist process leases");
|
||||
let mut queue = MrfQueue::new(2, usize::MAX);
|
||||
assert_eq!(queue.try_push_typed(original.clone()), MrfQueuePushResult::Enqueued);
|
||||
assert_eq!(
|
||||
queue.try_push_typed(original),
|
||||
MrfQueuePushResult::Coalesced,
|
||||
"legacy duplicates are one durable responsibility before a lease is assigned"
|
||||
);
|
||||
let mut replay = queue.pop_front().expect("one deduplicated replay record");
|
||||
assert_eq!(
|
||||
rustfs_common::mrf_channel::try_rearm_mrf_replay_intent(&mut replay),
|
||||
MrfIngressResult::Enqueued
|
||||
);
|
||||
assert!(replay.lease.is_some(), "manager admission must receive the replay lease");
|
||||
assert!(
|
||||
rustfs_common::mrf_channel::MrfDurableRepairAnchor::from_intent(&replay, uuid::Uuid::new_v4()).is_some(),
|
||||
"the replay identity must be usable by the durable proof consumer"
|
||||
);
|
||||
let mut encoded = Vec::new();
|
||||
assert!(encode_intent(&replay, &mut encoded));
|
||||
let (decoded, truncated) = decode_journal(&encoded);
|
||||
assert_eq!(truncated, 0);
|
||||
assert_eq!(decoded.len(), 1);
|
||||
assert!(
|
||||
decoded[0].lease.is_none(),
|
||||
"process-local leases must not enter the durable journal format"
|
||||
);
|
||||
rustfs_common::mrf_channel::release_mrf_intent(&replay);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replay_can_arm_more_records_than_live_queue_budget() {
|
||||
let mut queue = MrfQueue::new(1, intent("bucket", "object-0", 0).estimated_bytes());
|
||||
|
||||
Reference in New Issue
Block a user