|
|
|
@@ -37,7 +37,7 @@ use crate::heal::manager::HealManager;
|
|
|
|
|
use metrics::{counter, gauge};
|
|
|
|
|
use rustfs_common::heal_channel::{HealAdmissionDropReason, HealAdmissionResult};
|
|
|
|
|
use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIntent};
|
|
|
|
|
use std::collections::VecDeque;
|
|
|
|
|
use std::collections::{HashSet, VecDeque};
|
|
|
|
|
use std::sync::Arc;
|
|
|
|
|
use std::time::Duration;
|
|
|
|
|
use tokio::sync::mpsc;
|
|
|
|
@@ -48,15 +48,27 @@ use crate::heal::task::{HealOptions, HealPriority, HealRequest, HealType};
|
|
|
|
|
/// Journal location inside the metadata bucket, following the resume-state
|
|
|
|
|
/// layout.
|
|
|
|
|
pub(crate) const MRF_JOURNAL_PATH: &str = "buckets/.heal/mrf/journal.bin";
|
|
|
|
|
/// The scoped path is the authoritative snapshot for new readers and carries
|
|
|
|
|
/// both v1 and v2 records. The legacy path is only a v1 compatibility mirror;
|
|
|
|
|
/// older readers ignore the authoritative path, while new readers never merge
|
|
|
|
|
/// the two files. This prevents a partial two-file flush from fabricating a
|
|
|
|
|
/// mixed epoch.
|
|
|
|
|
pub(crate) const MRF_SCOPED_JOURNAL_PATH: &str = "buckets/.heal/mrf/journal-scoped.bin";
|
|
|
|
|
|
|
|
|
|
/// Record format tag.
|
|
|
|
|
const MRF_JOURNAL_FORMAT: u8 = 1;
|
|
|
|
|
/// Record layout version.
|
|
|
|
|
const MRF_JOURNAL_VERSION: u8 = 1;
|
|
|
|
|
const MRF_JOURNAL_VERSION_SCOPED: u8 = 2;
|
|
|
|
|
|
|
|
|
|
/// Fixed header size: format, version, kind, attempts, enqueued_at_ms,
|
|
|
|
|
/// has_version flag.
|
|
|
|
|
const MRF_RECORD_FIXED_HEAD: usize = 1 + 1 + 1 + 1 + 8 + 1;
|
|
|
|
|
const MRF_MAX_IDENTITY_COMPONENT: usize = 1024;
|
|
|
|
|
|
|
|
|
|
fn metric_f64(value: usize) -> f64 {
|
|
|
|
|
f64::from(u32::try_from(value).unwrap_or(u32::MAX))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Debug, Clone)]
|
|
|
|
|
pub(crate) struct MrfConsumerConfig {
|
|
|
|
@@ -101,40 +113,90 @@ impl Default for MrfConsumerConfig {
|
|
|
|
|
/// incoming intent (never a resident one) and counts the loss.
|
|
|
|
|
pub(crate) struct MrfQueue {
|
|
|
|
|
pending: VecDeque<MrfIntent>,
|
|
|
|
|
pending_keys: HashSet<MrfQueueKey>,
|
|
|
|
|
bytes: usize,
|
|
|
|
|
capacity: usize,
|
|
|
|
|
byte_budget: usize,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
|
|
|
|
struct MrfQueueKey {
|
|
|
|
|
kind: rustfs_common::mrf_channel::MrfKind,
|
|
|
|
|
bucket: Arc<str>,
|
|
|
|
|
object: Arc<str>,
|
|
|
|
|
version_id: Option<[u8; 16]>,
|
|
|
|
|
scope: Option<rustfs_common::mrf_channel::MrfScope>,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn queue_key(intent: &MrfIntent) -> MrfQueueKey {
|
|
|
|
|
let version_id = intent.version_id.filter(|bytes| *bytes != [0; 16]);
|
|
|
|
|
let scope = (!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption))
|
|
|
|
|
.then_some(intent.scope)
|
|
|
|
|
.flatten();
|
|
|
|
|
MrfQueueKey {
|
|
|
|
|
kind: intent.kind,
|
|
|
|
|
bucket: intent.bucket.clone(),
|
|
|
|
|
object: intent.object.clone(),
|
|
|
|
|
version_id,
|
|
|
|
|
scope,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
|
|
|
|
pub(crate) enum MrfQueuePushResult {
|
|
|
|
|
Enqueued,
|
|
|
|
|
Coalesced,
|
|
|
|
|
Rejected,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl MrfQueue {
|
|
|
|
|
pub(crate) fn new(capacity: usize, byte_budget: usize) -> Self {
|
|
|
|
|
Self {
|
|
|
|
|
pending: VecDeque::new(),
|
|
|
|
|
pending_keys: HashSet::new(),
|
|
|
|
|
bytes: 0,
|
|
|
|
|
capacity,
|
|
|
|
|
byte_budget,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Returns `false` (after counting) when either ceiling would be crossed.
|
|
|
|
|
pub(crate) fn try_push(&mut self, intent: MrfIntent) -> bool {
|
|
|
|
|
pub(crate) fn try_push_typed(&mut self, intent: MrfIntent) -> MrfQueuePushResult {
|
|
|
|
|
if intent.bucket.len() > MRF_MAX_IDENTITY_COMPONENT || intent.object.len() > MRF_MAX_IDENTITY_COMPONENT {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "identity_oversized").increment(1);
|
|
|
|
|
return MrfQueuePushResult::Rejected;
|
|
|
|
|
}
|
|
|
|
|
let key = queue_key(&intent);
|
|
|
|
|
if self.pending_keys.contains(&key) {
|
|
|
|
|
counter!("rustfs_heal_mrf_coalesced_total", "layer" => "queue").increment(1);
|
|
|
|
|
return MrfQueuePushResult::Coalesced;
|
|
|
|
|
}
|
|
|
|
|
let cost = intent.estimated_bytes();
|
|
|
|
|
if self.pending.len() >= self.capacity || self.bytes + cost > self.byte_budget {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "queue_overflow").increment(1);
|
|
|
|
|
return false;
|
|
|
|
|
return MrfQueuePushResult::Rejected;
|
|
|
|
|
}
|
|
|
|
|
self.bytes += cost;
|
|
|
|
|
self.pending_keys.insert(key);
|
|
|
|
|
self.pending.push_back(intent);
|
|
|
|
|
true
|
|
|
|
|
MrfQueuePushResult::Enqueued
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Bool compatibility adapter: only a newly executable queue item is
|
|
|
|
|
/// reported as accepted; a coalesced duplicate is not durable admission.
|
|
|
|
|
#[cfg(test)]
|
|
|
|
|
pub(crate) fn try_push(&mut self, intent: MrfIntent) -> bool {
|
|
|
|
|
matches!(self.try_push_typed(intent), MrfQueuePushResult::Enqueued)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub(crate) fn pop_front(&mut self) -> Option<MrfIntent> {
|
|
|
|
|
let intent = self.pending.pop_front()?;
|
|
|
|
|
self.pending_keys.remove(&queue_key(&intent));
|
|
|
|
|
self.bytes = self.bytes.saturating_sub(intent.estimated_bytes());
|
|
|
|
|
Some(intent)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub(crate) fn push_back(&mut self, intent: MrfIntent) {
|
|
|
|
|
self.pending_keys.insert(queue_key(&intent));
|
|
|
|
|
self.bytes += intent.estimated_bytes();
|
|
|
|
|
self.pending.push_back(intent);
|
|
|
|
|
}
|
|
|
|
@@ -157,10 +219,24 @@ impl MrfQueue {
|
|
|
|
|
// ---------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
/// Append one encoded record to `out`.
|
|
|
|
|
pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) {
|
|
|
|
|
pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) -> bool {
|
|
|
|
|
let Ok(bucket_len) = u32::try_from(intent.bucket.len()) else {
|
|
|
|
|
return false;
|
|
|
|
|
};
|
|
|
|
|
let Ok(object_len) = u32::try_from(intent.object.len()) else {
|
|
|
|
|
return false;
|
|
|
|
|
};
|
|
|
|
|
let scope = (!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption))
|
|
|
|
|
.then_some(intent.scope)
|
|
|
|
|
.flatten();
|
|
|
|
|
let version_id = intent.version_id.filter(|bytes| *bytes != [0; 16]);
|
|
|
|
|
let start = out.len();
|
|
|
|
|
out.push(MRF_JOURNAL_FORMAT);
|
|
|
|
|
out.push(MRF_JOURNAL_VERSION);
|
|
|
|
|
out.push(if scope.is_some() {
|
|
|
|
|
MRF_JOURNAL_VERSION_SCOPED
|
|
|
|
|
} else {
|
|
|
|
|
MRF_JOURNAL_VERSION
|
|
|
|
|
});
|
|
|
|
|
out.push(match intent.kind {
|
|
|
|
|
rustfs_common::mrf_channel::MrfKind::DecodeFailure => 1,
|
|
|
|
|
rustfs_common::mrf_channel::MrfKind::MetadataCorruption => 2,
|
|
|
|
@@ -168,27 +244,36 @@ pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) {
|
|
|
|
|
});
|
|
|
|
|
out.push(intent.attempts);
|
|
|
|
|
out.extend_from_slice(&intent.enqueued_at_ms.to_le_bytes());
|
|
|
|
|
match intent.version_id {
|
|
|
|
|
match version_id {
|
|
|
|
|
Some(bytes) => {
|
|
|
|
|
out.push(1);
|
|
|
|
|
out.extend_from_slice(&bytes);
|
|
|
|
|
}
|
|
|
|
|
None => out.push(0),
|
|
|
|
|
}
|
|
|
|
|
out.extend_from_slice(&(intent.bucket.len() as u32).to_le_bytes());
|
|
|
|
|
out.extend_from_slice(&(intent.object.len() as u32).to_le_bytes());
|
|
|
|
|
if let Some(scope) = scope {
|
|
|
|
|
out.extend_from_slice(&scope.pool_index.to_le_bytes());
|
|
|
|
|
out.extend_from_slice(&scope.set_index.to_le_bytes());
|
|
|
|
|
}
|
|
|
|
|
out.extend_from_slice(&bucket_len.to_le_bytes());
|
|
|
|
|
out.extend_from_slice(&object_len.to_le_bytes());
|
|
|
|
|
out.extend_from_slice(intent.bucket.as_bytes());
|
|
|
|
|
out.extend_from_slice(intent.object.as_bytes());
|
|
|
|
|
let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc);
|
|
|
|
|
hasher.update(&out[start..]);
|
|
|
|
|
out.extend_from_slice(&(hasher.finalize() as u32).to_le_bytes());
|
|
|
|
|
let Ok(checksum) = u32::try_from(hasher.finalize()) else {
|
|
|
|
|
out.truncate(start);
|
|
|
|
|
return false;
|
|
|
|
|
};
|
|
|
|
|
out.extend_from_slice(&checksum.to_le_bytes());
|
|
|
|
|
true
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
|
|
|
|
|
if data.len() < MRF_RECORD_FIXED_HEAD + 8 {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
if data[0] != MRF_JOURNAL_FORMAT || data[1] != MRF_JOURNAL_VERSION {
|
|
|
|
|
if data[0] != MRF_JOURNAL_FORMAT || !matches!(data[1], MRF_JOURNAL_VERSION | MRF_JOURNAL_VERSION_SCOPED) {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
let kind = match data[2] {
|
|
|
|
@@ -198,24 +283,38 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
|
|
|
|
|
_ => return None,
|
|
|
|
|
};
|
|
|
|
|
let attempts = data[3];
|
|
|
|
|
let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().expect("slice length checked"));
|
|
|
|
|
let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().ok()?);
|
|
|
|
|
let has_version = data[12] != 0;
|
|
|
|
|
let mut cursor = MRF_RECORD_FIXED_HEAD;
|
|
|
|
|
let version_id = if has_version {
|
|
|
|
|
if data.len() < cursor + 16 {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
let bytes: [u8; 16] = data[cursor..cursor + 16].try_into().expect("slice length checked");
|
|
|
|
|
let bytes: [u8; 16] = data[cursor..cursor + 16].try_into().ok()?;
|
|
|
|
|
cursor += 16;
|
|
|
|
|
Some(bytes)
|
|
|
|
|
} else {
|
|
|
|
|
None
|
|
|
|
|
};
|
|
|
|
|
let scope = if data[1] == MRF_JOURNAL_VERSION_SCOPED {
|
|
|
|
|
if data.len() < cursor + 8 {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
let pool_index = u32::from_le_bytes(data[cursor..cursor + 4].try_into().ok()?);
|
|
|
|
|
let set_index = u32::from_le_bytes(data[cursor + 4..cursor + 8].try_into().ok()?);
|
|
|
|
|
cursor += 8;
|
|
|
|
|
Some(rustfs_common::mrf_channel::MrfScope { pool_index, set_index })
|
|
|
|
|
} else {
|
|
|
|
|
None
|
|
|
|
|
};
|
|
|
|
|
if data.len() < cursor + 8 {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
let bucket_len = u32::from_le_bytes(data[cursor..cursor + 4].try_into().expect("slice length checked")) as usize;
|
|
|
|
|
let object_len = u32::from_le_bytes(data[cursor + 4..cursor + 8].try_into().expect("slice length checked")) as usize;
|
|
|
|
|
let bucket_len = usize::try_from(u32::from_le_bytes(data[cursor..cursor + 4].try_into().ok()?)).ok()?;
|
|
|
|
|
let object_len = usize::try_from(u32::from_le_bytes(data[cursor + 4..cursor + 8].try_into().ok()?)).ok()?;
|
|
|
|
|
if bucket_len > MRF_MAX_IDENTITY_COMPONENT || object_len > MRF_MAX_IDENTITY_COMPONENT {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
cursor += 8;
|
|
|
|
|
let body_end = cursor.checked_add(bucket_len)?.checked_add(object_len)?;
|
|
|
|
|
let record_end = body_end.checked_add(4)?;
|
|
|
|
@@ -224,7 +323,7 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
|
|
|
|
|
}
|
|
|
|
|
let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc);
|
|
|
|
|
hasher.update(&data[..body_end]);
|
|
|
|
|
if (hasher.finalize() as u32) != u32::from_le_bytes(data[body_end..record_end].try_into().expect("slice length checked")) {
|
|
|
|
|
if u32::try_from(hasher.finalize()).ok()? != u32::from_le_bytes(data[body_end..record_end].try_into().ok()?) {
|
|
|
|
|
return None;
|
|
|
|
|
}
|
|
|
|
|
let bucket = std::sync::Arc::from(std::str::from_utf8(&data[cursor..cursor + bucket_len]).ok()?);
|
|
|
|
@@ -235,6 +334,12 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
|
|
|
|
|
object,
|
|
|
|
|
version_id,
|
|
|
|
|
kind,
|
|
|
|
|
scope: if matches!(kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption) {
|
|
|
|
|
None
|
|
|
|
|
} else {
|
|
|
|
|
scope
|
|
|
|
|
},
|
|
|
|
|
lease: None,
|
|
|
|
|
enqueued_at_ms,
|
|
|
|
|
attempts,
|
|
|
|
|
},
|
|
|
|
@@ -269,9 +374,9 @@ async fn journal_disks() -> Vec<DiskStore> {
|
|
|
|
|
map.values().flatten().cloned().collect()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async fn read_journal() -> Option<Vec<u8>> {
|
|
|
|
|
async fn read_journal(path: &str) -> Option<Vec<u8>> {
|
|
|
|
|
for disk in journal_disks().await {
|
|
|
|
|
match disk.read_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH).await {
|
|
|
|
|
match disk.read_all(super::RUSTFS_META_BUCKET, path).await {
|
|
|
|
|
Ok(bytes) => return Some(bytes.to_vec()),
|
|
|
|
|
Err(_) => continue,
|
|
|
|
|
}
|
|
|
|
@@ -282,35 +387,51 @@ async fn read_journal() -> Option<Vec<u8>> {
|
|
|
|
|
/// Write the snapshot to every local disk; returns true when at least one
|
|
|
|
|
/// disk accepted it, so a total write failure keeps the runtime dirty and
|
|
|
|
|
/// the next tick retries the persist.
|
|
|
|
|
async fn write_journal(data: &[u8]) -> bool {
|
|
|
|
|
async fn write_journal(path: &str, data: &[u8]) -> bool {
|
|
|
|
|
let payload = bytes::Bytes::copy_from_slice(data);
|
|
|
|
|
let mut any_persisted = false;
|
|
|
|
|
for disk in journal_disks().await {
|
|
|
|
|
match disk
|
|
|
|
|
.write_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH, payload.clone())
|
|
|
|
|
.await
|
|
|
|
|
{
|
|
|
|
|
match disk.write_all(super::RUSTFS_META_BUCKET, path, payload.clone()).await {
|
|
|
|
|
Ok(()) => any_persisted = true,
|
|
|
|
|
Err(err) => warn_mrf_journal_write(&err),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if !data.is_empty() {
|
|
|
|
|
counter!("rustfs_heal_mrf_journal_fsync_total").increment(1);
|
|
|
|
|
}
|
|
|
|
|
gauge!("rustfs_heal_mrf_journal_bytes").set(data.len() as f64);
|
|
|
|
|
any_persisted
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async fn delete_journal() {
|
|
|
|
|
for disk in journal_disks().await {
|
|
|
|
|
let _ = disk
|
|
|
|
|
async fn delete_journal(path: &str) -> bool {
|
|
|
|
|
let disks = journal_disks().await;
|
|
|
|
|
if disks.is_empty() {
|
|
|
|
|
counter!("rustfs_heal_mrf_journal_delete_failures_total").increment(1);
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
let mut all_deleted = true;
|
|
|
|
|
for disk in disks {
|
|
|
|
|
let result = disk
|
|
|
|
|
.delete(
|
|
|
|
|
super::RUSTFS_META_BUCKET,
|
|
|
|
|
MRF_JOURNAL_PATH,
|
|
|
|
|
path,
|
|
|
|
|
crate::heal::storage_api::owner::EcstoreDeleteOptions::default(),
|
|
|
|
|
)
|
|
|
|
|
.await;
|
|
|
|
|
if let Err(err) = result {
|
|
|
|
|
// Delete is idempotent: a compatibility mirror that was never
|
|
|
|
|
// written (or was already removed) is clean, not a retry state.
|
|
|
|
|
if !matches!(err, super::DiskError::FileNotFound | super::DiskError::VolumeNotFound) {
|
|
|
|
|
all_deleted = false;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if !all_deleted {
|
|
|
|
|
counter!("rustfs_heal_mrf_journal_delete_failures_total").increment(1);
|
|
|
|
|
}
|
|
|
|
|
all_deleted
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async fn delete_journals() -> bool {
|
|
|
|
|
let authoritative_deleted = delete_journal(MRF_SCOPED_JOURNAL_PATH).await;
|
|
|
|
|
let legacy_deleted = delete_journal(MRF_JOURNAL_PATH).await;
|
|
|
|
|
authoritative_deleted && legacy_deleted
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn warn_mrf_journal_write(err: &super::DiskError) {
|
|
|
|
@@ -331,7 +452,10 @@ fn warn_mrf_journal_write(err: &super::DiskError) {
|
|
|
|
|
pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest {
|
|
|
|
|
let bucket = intent.bucket.to_string();
|
|
|
|
|
let object = intent.object.to_string();
|
|
|
|
|
let version_id = intent.version_id.map(|bytes| Uuid::from_bytes(bytes).to_string());
|
|
|
|
|
let version_id = intent
|
|
|
|
|
.version_id
|
|
|
|
|
.filter(|bytes| *bytes != [0; 16])
|
|
|
|
|
.map(|bytes| Uuid::from_bytes(bytes).to_string());
|
|
|
|
|
let (heal_type, priority) = match intent.kind {
|
|
|
|
|
rustfs_common::mrf_channel::MrfKind::DecodeFailure => (
|
|
|
|
|
HealType::ECDecode {
|
|
|
|
@@ -351,18 +475,28 @@ pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest {
|
|
|
|
|
HealPriority::Normal,
|
|
|
|
|
),
|
|
|
|
|
};
|
|
|
|
|
let mut request = HealRequest::new(heal_type, HealOptions::default(), priority);
|
|
|
|
|
let mut options = HealOptions::default();
|
|
|
|
|
if !matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption)
|
|
|
|
|
&& let Some(scope) = intent.scope
|
|
|
|
|
{
|
|
|
|
|
options.pool_index = usize::try_from(scope.pool_index).ok();
|
|
|
|
|
options.set_index = usize::try_from(scope.set_index).ok();
|
|
|
|
|
}
|
|
|
|
|
let mut request = HealRequest::new(heal_type, options, priority);
|
|
|
|
|
request.source = rustfs_common::heal_channel::HealRequestSource::Mrf;
|
|
|
|
|
request
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> crate::Result<HealAdmissionResult> {
|
|
|
|
|
let receipt = manager
|
|
|
|
|
.submit_mrf_heal_request_with_receipt(
|
|
|
|
|
.submit_mrf_heal_request_with_receipt_and_identity(
|
|
|
|
|
build_heal_request(intent),
|
|
|
|
|
intent.bucket.clone(),
|
|
|
|
|
intent.object.clone(),
|
|
|
|
|
intent.version_id,
|
|
|
|
|
intent.kind,
|
|
|
|
|
intent.scope,
|
|
|
|
|
intent.lease,
|
|
|
|
|
)
|
|
|
|
|
.await?;
|
|
|
|
|
Ok(receipt.result)
|
|
|
|
@@ -387,16 +521,38 @@ struct MrfRuntime {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl MrfRuntime {
|
|
|
|
|
fn snapshot(&self) -> Vec<u8> {
|
|
|
|
|
let mut buf = Vec::new();
|
|
|
|
|
fn snapshot(&self) -> (Vec<u8>, Vec<u8>) {
|
|
|
|
|
let mut authoritative = Vec::new();
|
|
|
|
|
let mut legacy = Vec::new();
|
|
|
|
|
for intent in self.queue.intents() {
|
|
|
|
|
encode_intent(intent, &mut buf);
|
|
|
|
|
let scoped_identity =
|
|
|
|
|
!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption) && intent.scope.is_some();
|
|
|
|
|
if !encode_intent(intent, &mut authoritative) {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "journal_identity_oversized").increment(1);
|
|
|
|
|
}
|
|
|
|
|
if !scoped_identity && !encode_intent(intent, &mut legacy) {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "journal_identity_oversized").increment(1);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
buf
|
|
|
|
|
(authoritative, legacy)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async fn flush(&mut self) {
|
|
|
|
|
let persisted = write_journal(&self.snapshot()).await;
|
|
|
|
|
let (authoritative, legacy) = self.snapshot();
|
|
|
|
|
let authoritative_persisted = write_journal(MRF_SCOPED_JOURNAL_PATH, &authoritative).await;
|
|
|
|
|
if !authoritative.is_empty() {
|
|
|
|
|
counter!("rustfs_heal_mrf_journal_fsync_total").increment(1);
|
|
|
|
|
}
|
|
|
|
|
gauge!("rustfs_heal_mrf_journal_bytes").set(metric_f64(authoritative.len()));
|
|
|
|
|
// Publish the compatibility mirror only after the authoritative
|
|
|
|
|
// snapshot has reached at least one disk. This ordering prevents an
|
|
|
|
|
// old reader from observing a newer epoch that a new reader cannot
|
|
|
|
|
// see when the canonical write is unavailable.
|
|
|
|
|
let legacy_persisted = authoritative_persisted && write_journal(MRF_JOURNAL_PATH, &legacy).await;
|
|
|
|
|
// Keep dirty until both the authoritative snapshot and its
|
|
|
|
|
// compatibility mirror have been accepted; otherwise a one-sided
|
|
|
|
|
// failure would never retry the missing file.
|
|
|
|
|
let persisted = authoritative_persisted && legacy_persisted;
|
|
|
|
|
self.new_since_flush = 0;
|
|
|
|
|
// Keep the dirty flag when every disk write failed: a clean backlog
|
|
|
|
|
// would otherwise never rewrite, losing the periodic persist retry a
|
|
|
|
@@ -404,7 +560,7 @@ impl MrfRuntime {
|
|
|
|
|
if persisted {
|
|
|
|
|
self.dirty = false;
|
|
|
|
|
}
|
|
|
|
|
self.journal_on_disk = true;
|
|
|
|
|
self.journal_on_disk |= authoritative_persisted || legacy_persisted;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Drain pending intents into the heal manager until it is full, the
|
|
|
|
@@ -430,6 +586,7 @@ impl MrfRuntime {
|
|
|
|
|
intent.attempts = intent.attempts.saturating_add(1);
|
|
|
|
|
if intent.attempts >= MRF_MAX_ATTEMPTS {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1);
|
|
|
|
|
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
self.queue.push_back(intent);
|
|
|
|
@@ -438,11 +595,13 @@ impl MrfRuntime {
|
|
|
|
|
}
|
|
|
|
|
Ok(HealAdmissionResult::Dropped(_)) => {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "admission_policy").increment(1);
|
|
|
|
|
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
|
|
|
|
}
|
|
|
|
|
Err(_) => {
|
|
|
|
|
intent.attempts = intent.attempts.saturating_add(1);
|
|
|
|
|
if intent.attempts >= MRF_MAX_ATTEMPTS {
|
|
|
|
|
counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1);
|
|
|
|
|
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
|
|
|
|
continue;
|
|
|
|
|
}
|
|
|
|
|
self.queue.push_back(intent);
|
|
|
|
@@ -451,8 +610,8 @@ impl MrfRuntime {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
gauge!("rustfs_heal_mrf_queue_depth").set(self.queue.depth() as f64);
|
|
|
|
|
gauge!("rustfs_heal_mrf_queue_bytes").set(self.queue.bytes() as f64);
|
|
|
|
|
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(self.queue.depth()));
|
|
|
|
|
gauge!("rustfs_heal_mrf_queue_bytes").set(metric_f64(self.queue.bytes()));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -496,7 +655,12 @@ pub async fn replay_journal_once(manager: &Arc<HealManager>) -> usize {
|
|
|
|
|
let config = MrfConsumerConfig::default();
|
|
|
|
|
let mut queue = MrfQueue::new(config.queue_capacity, config.journal_max_bytes);
|
|
|
|
|
let mut backoff_until: Option<tokio::time::Instant> = None;
|
|
|
|
|
replay_into(manager, &mut queue, &mut backoff_until).await
|
|
|
|
|
replay_into(manager, &mut queue, &mut backoff_until).await.replayed
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
struct ReplayOutcome {
|
|
|
|
|
replayed: usize,
|
|
|
|
|
journal_on_disk: bool,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Shared replay core: read + decode + re-arm + delete, then drain what fits.
|
|
|
|
@@ -504,11 +668,25 @@ async fn replay_into(
|
|
|
|
|
manager: &Arc<HealManager>,
|
|
|
|
|
queue: &mut MrfQueue,
|
|
|
|
|
backoff_until: &mut Option<tokio::time::Instant>,
|
|
|
|
|
) -> usize {
|
|
|
|
|
let Some(data) = read_journal().await else {
|
|
|
|
|
return 0;
|
|
|
|
|
) -> ReplayOutcome {
|
|
|
|
|
// The scoped file is a complete authoritative snapshot. Fall back to the
|
|
|
|
|
// legacy mirror only when the authoritative path is unavailable; merging
|
|
|
|
|
// both files could combine records from different flush epochs.
|
|
|
|
|
let data = match read_journal(MRF_SCOPED_JOURNAL_PATH).await {
|
|
|
|
|
Some(data) => data,
|
|
|
|
|
None => match read_journal(MRF_JOURNAL_PATH).await {
|
|
|
|
|
Some(data) => data,
|
|
|
|
|
None => {
|
|
|
|
|
return ReplayOutcome {
|
|
|
|
|
replayed: 0,
|
|
|
|
|
journal_on_disk: false,
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
};
|
|
|
|
|
let (intents, truncated) = decode_journal(&data);
|
|
|
|
|
let mut intents = Vec::new();
|
|
|
|
|
let (decoded, truncated) = decode_journal(&data);
|
|
|
|
|
intents.extend(decoded);
|
|
|
|
|
if truncated > 0 {
|
|
|
|
|
tracing::warn!(
|
|
|
|
|
target: "rustfs::heal::mrf",
|
|
|
|
@@ -516,12 +694,15 @@ async fn replay_into(
|
|
|
|
|
"MRF journal had a torn tail; truncated records were discarded"
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
counter!("rustfs_heal_mrf_replayed_total").increment(intents.len() as u64);
|
|
|
|
|
counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(intents.len()).unwrap_or(u64::MAX));
|
|
|
|
|
let replayed = intents.len();
|
|
|
|
|
for intent in intents {
|
|
|
|
|
queue.try_push(intent);
|
|
|
|
|
let result = queue.try_push_typed(intent.clone());
|
|
|
|
|
if !matches!(result, MrfQueuePushResult::Enqueued) {
|
|
|
|
|
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
delete_journal().await;
|
|
|
|
|
let journal_on_disk = !delete_journals().await;
|
|
|
|
|
|
|
|
|
|
// Drain the replayed intents immediately; whatever the manager refuses
|
|
|
|
|
// stays armed in `queue` for the consumer's retry loop.
|
|
|
|
@@ -541,7 +722,10 @@ async fn replay_into(
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
replayed
|
|
|
|
|
ReplayOutcome {
|
|
|
|
|
replayed,
|
|
|
|
|
journal_on_disk,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Replay the journal, then keep draining the channel into the heal manager
|
|
|
|
@@ -559,7 +743,8 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
|
|
|
|
|
|
|
|
|
// Replay: read the journal, re-arm intents (duplicates are merged by the
|
|
|
|
|
// manager's dedup key), then drop the file so the next flush starts clean.
|
|
|
|
|
replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
|
|
|
|
|
let replay = replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
|
|
|
|
|
runtime.journal_on_disk = replay.journal_on_disk;
|
|
|
|
|
// The replay deleted the journal file; anything still pending (e.g. the
|
|
|
|
|
// manager was full and backoff armed) must be re-persisted by the next
|
|
|
|
|
// flush or a crash before it would lose those intents.
|
|
|
|
@@ -587,9 +772,14 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
for intent in batch.drain(..) {
|
|
|
|
|
if runtime.queue.try_push(intent) {
|
|
|
|
|
runtime.new_since_flush += 1;
|
|
|
|
|
runtime.dirty = true;
|
|
|
|
|
match runtime.queue.try_push_typed(intent.clone()) {
|
|
|
|
|
MrfQueuePushResult::Enqueued => {
|
|
|
|
|
runtime.new_since_flush += 1;
|
|
|
|
|
runtime.dirty = true;
|
|
|
|
|
}
|
|
|
|
|
MrfQueuePushResult::Coalesced | MrfQueuePushResult::Rejected => {
|
|
|
|
|
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
runtime.dispatch(manager.as_ref()).await;
|
|
|
|
@@ -613,13 +803,14 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
|
|
|
|
TickAction::DeleteJournal => {
|
|
|
|
|
// All intents consumed: remove the journal so a restart
|
|
|
|
|
// replays nothing (mirrors MinIO's post-replay unlink).
|
|
|
|
|
delete_journal().await;
|
|
|
|
|
runtime.journal_on_disk = false;
|
|
|
|
|
gauge!("rustfs_heal_mrf_journal_bytes").set(0.0);
|
|
|
|
|
if delete_journals().await {
|
|
|
|
|
runtime.journal_on_disk = false;
|
|
|
|
|
gauge!("rustfs_heal_mrf_journal_bytes").set(0.0);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
TickAction::Idle => {}
|
|
|
|
|
}
|
|
|
|
|
gauge!("rustfs_heal_mrf_queue_depth").set(runtime.queue.depth() as f64);
|
|
|
|
|
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(runtime.queue.depth()));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@@ -664,6 +855,8 @@ mod tests {
|
|
|
|
|
object: StdArc::from(object),
|
|
|
|
|
version_id: Some([7u8; 16]),
|
|
|
|
|
kind: MrfKind::DecodeFailure,
|
|
|
|
|
scope: None,
|
|
|
|
|
lease: None,
|
|
|
|
|
enqueued_at_ms: 1_700_000_000_000,
|
|
|
|
|
attempts,
|
|
|
|
|
}
|
|
|
|
@@ -694,17 +887,101 @@ mod tests {
|
|
|
|
|
fn queue_enforces_count_and_byte_ceilings() {
|
|
|
|
|
let mut queue = MrfQueue::new(2, usize::MAX);
|
|
|
|
|
assert!(queue.try_push(intent("b", "o", 0)));
|
|
|
|
|
assert!(queue.try_push(intent("b", "o", 0)));
|
|
|
|
|
assert!(!queue.try_push(intent("b", "o", 0)), "count ceiling must drop");
|
|
|
|
|
assert!(queue.try_push(intent("b", "o2", 0)));
|
|
|
|
|
assert!(!queue.try_push(intent("b", "o3", 0)), "count ceiling must drop");
|
|
|
|
|
|
|
|
|
|
let mut tiny = MrfQueue::new(usize::MAX, intent("bucket", "object", 0).estimated_bytes());
|
|
|
|
|
assert!(tiny.try_push(intent("bucket", "object", 0)));
|
|
|
|
|
assert!(
|
|
|
|
|
!tiny.try_push(intent("bucket", "object", 0)),
|
|
|
|
|
!tiny.try_push(intent("bucket", "object2", 0)),
|
|
|
|
|
"byte budget must drop before the second intent fits"
|
|
|
|
|
);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn duplicate_mrf_intents_coalesce_to_one_execution() {
|
|
|
|
|
let mut queue = MrfQueue::new(1000, usize::MAX);
|
|
|
|
|
let mut enqueued = 0;
|
|
|
|
|
let mut coalesced = 0;
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
enqueued += 1;
|
|
|
|
|
for _ in 0..999 {
|
|
|
|
|
match queue.try_push_typed(intent("bucket", "object", 0)) {
|
|
|
|
|
MrfQueuePushResult::Coalesced => coalesced += 1,
|
|
|
|
|
other => panic!("duplicate intent was not coalesced: {other:?}"),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
assert_eq!(enqueued, 1);
|
|
|
|
|
assert_eq!(coalesced, 999);
|
|
|
|
|
assert_eq!(queue.depth(), 1);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn mrf_dedupe_does_not_merge_adjacent_version_pool_or_kind() {
|
|
|
|
|
let mut queue = MrfQueue::new(8, usize::MAX);
|
|
|
|
|
let mut first = intent("bucket", "object", 0);
|
|
|
|
|
first.kind = MrfKind::PartialWrite;
|
|
|
|
|
first.scope = Some(rustfs_common::mrf_channel::MrfScope {
|
|
|
|
|
pool_index: 1,
|
|
|
|
|
set_index: 1,
|
|
|
|
|
});
|
|
|
|
|
assert!(queue.try_push(first.clone()));
|
|
|
|
|
first.version_id = Some([8u8; 16]);
|
|
|
|
|
assert!(queue.try_push(first));
|
|
|
|
|
let mut other_scope = intent("bucket", "object", 0);
|
|
|
|
|
other_scope.kind = MrfKind::PartialWrite;
|
|
|
|
|
other_scope.scope = Some(rustfs_common::mrf_channel::MrfScope {
|
|
|
|
|
pool_index: 2,
|
|
|
|
|
set_index: 1,
|
|
|
|
|
});
|
|
|
|
|
assert!(queue.try_push(other_scope));
|
|
|
|
|
let mut other_kind = intent("bucket", "object", 0);
|
|
|
|
|
other_kind.kind = MrfKind::DecodeFailure;
|
|
|
|
|
other_kind.scope = None;
|
|
|
|
|
assert!(queue.try_push(other_kind));
|
|
|
|
|
assert_eq!(queue.depth(), 4);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn mrf_dedupe_full_returns_rejected_with_durable_pending() {
|
|
|
|
|
let mut queue = MrfQueue::new(1, usize::MAX);
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "other", 0)), MrfQueuePushResult::Rejected);
|
|
|
|
|
assert_eq!(queue.depth(), 1);
|
|
|
|
|
let mut snapshot = Vec::new();
|
|
|
|
|
assert!(encode_intent(queue.intents().next().expect("resident intent"), &mut snapshot));
|
|
|
|
|
assert!(!snapshot.is_empty(), "the resident intent remains journalable after rejection");
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn mrf_dedupe_failure_releases_key_for_retry() {
|
|
|
|
|
let mut queue = MrfQueue::new(1, usize::MAX);
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
let _failed = queue.pop_front().expect("queued intent");
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "object", 1)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
assert_eq!(queue.depth(), 1);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn mrf_dedupe_key_and_map_are_bounded() {
|
|
|
|
|
let mut queue = MrfQueue::new(2, usize::MAX);
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "other", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
assert_eq!(queue.pending_keys.len(), 2);
|
|
|
|
|
assert_eq!(queue.try_push_typed(intent("bucket", "third", 0)), MrfQueuePushResult::Rejected);
|
|
|
|
|
assert_eq!(queue.depth(), 2);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn cross_node_duplicate_execution_remains_idempotent() {
|
|
|
|
|
// Node-local ingress maps intentionally do not merge across nodes;
|
|
|
|
|
// the manager's existing identity key absorbs the duplicate later.
|
|
|
|
|
let mut node_a = MrfQueue::new(8, usize::MAX);
|
|
|
|
|
let mut node_b = MrfQueue::new(8, usize::MAX);
|
|
|
|
|
assert_eq!(node_a.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
assert_eq!(node_b.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[test]
|
|
|
|
|
fn journal_roundtrip_preserves_intents() {
|
|
|
|
|
let intents = vec![
|
|
|
|
@@ -715,6 +992,8 @@ mod tests {
|
|
|
|
|
object: StdArc::from("object/c"),
|
|
|
|
|
version_id: None,
|
|
|
|
|
kind: MrfKind::MetadataCorruption,
|
|
|
|
|
scope: None,
|
|
|
|
|
lease: None,
|
|
|
|
|
enqueued_at_ms: 5,
|
|
|
|
|
attempts: 1,
|
|
|
|
|
},
|
|
|
|
@@ -766,6 +1045,8 @@ mod tests {
|
|
|
|
|
object: StdArc::from("o"),
|
|
|
|
|
version_id: None,
|
|
|
|
|
kind: MrfKind::MetadataCorruption,
|
|
|
|
|
scope: None,
|
|
|
|
|
lease: None,
|
|
|
|
|
enqueued_at_ms: 0,
|
|
|
|
|
attempts: 0,
|
|
|
|
|
});
|
|
|
|
@@ -777,6 +1058,8 @@ mod tests {
|
|
|
|
|
object: StdArc::from("o"),
|
|
|
|
|
version_id: None,
|
|
|
|
|
kind: MrfKind::PartialWrite,
|
|
|
|
|
scope: None,
|
|
|
|
|
lease: None,
|
|
|
|
|
enqueued_at_ms: 0,
|
|
|
|
|
attempts: 0,
|
|
|
|
|
});
|
|
|
|
|