mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-23 04:39:04 +00:00
Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 82fb0a8843 | |||
| 307510749e | |||
| 5496e14960 | |||
| ec3b7a7dc6 |
@@ -23,32 +23,21 @@
|
||||
//! unconsumed intents is the consumer's job (see `rustfs-heal`
|
||||
//! `heal::mrf_queue`), mirroring MinIO's `.heal/mrf/list.bin`.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::collections::hash_map::RandomState;
|
||||
use std::hash::{BuildHasher, Hash};
|
||||
use std::sync::atomic::AtomicU64;
|
||||
use std::sync::atomic::AtomicUsize;
|
||||
use std::sync::{
|
||||
Arc, Mutex, OnceLock,
|
||||
Arc, OnceLock,
|
||||
atomic::{AtomicBool, Ordering},
|
||||
};
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::sync::mpsc;
|
||||
use uuid::Uuid;
|
||||
|
||||
/// Bounded capacity of the global MRF channel. Backpressure is resolved by
|
||||
/// dropping (and counting) intents, never by blocking the producer.
|
||||
const MRF_CHANNEL_CAPACITY: usize = 8192;
|
||||
const MRF_COALESCER_SHARDS: usize = 16;
|
||||
const MRF_COALESCER_MAX_KEYS: usize = 8192;
|
||||
const MRF_COALESCER_MAX_BYTES: usize = 16 * 1024 * 1024;
|
||||
const MRF_COALESCER_TTL: Duration = Duration::from_secs(60);
|
||||
const MRF_MAX_IDENTITY_COMPONENT: usize = 1024;
|
||||
|
||||
/// Why an intent was produced. Drives the heal priority mapping on the
|
||||
/// consumer side (DecodeFailure -> Urgent, MetadataCorruption -> High,
|
||||
/// PartialWrite -> Normal).
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum MrfKind {
|
||||
/// Erasure decode failed while serving a read (read path).
|
||||
DecodeFailure,
|
||||
@@ -78,52 +67,12 @@ pub struct MrfIntent {
|
||||
/// Version the intent targets, as raw UUID bytes.
|
||||
pub version_id: Option<[u8; 16]>,
|
||||
pub kind: MrfKind,
|
||||
/// Stable erasure-set scope when the producer has it. Kept optional so
|
||||
/// metadata corruption and legacy producers do not invent a scope.
|
||||
pub scope: Option<MrfScope>,
|
||||
/// Generation of the node-local ingress lease. It is not persisted in
|
||||
/// the journal; replayed records acquire a fresh lease when re-enqueued.
|
||||
pub lease: Option<MrfIngressLease>,
|
||||
pub enqueued_at_ms: u64,
|
||||
/// Times this intent has already been offered to the heal manager.
|
||||
/// Dropped by the consumer once it reaches `MRF_MAX_ATTEMPTS`.
|
||||
pub attempts: u8,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
|
||||
pub struct MrfScope {
|
||||
pub pool_index: u32,
|
||||
pub set_index: u32,
|
||||
}
|
||||
|
||||
/// Opaque generation used to release exactly the admission that created an
|
||||
/// ingress entry. A generation prevents a late terminal callback from
|
||||
/// deleting a newer retry for the same identity (ABA).
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
|
||||
pub struct MrfIngressLease(u64);
|
||||
|
||||
impl MrfIngressLease {
|
||||
const fn new(value: u64) -> Self {
|
||||
Self(value)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum MrfDropReason {
|
||||
Disabled,
|
||||
Uninitialized,
|
||||
Full,
|
||||
OversizedIdentity,
|
||||
CoalescerFull,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum MrfIngressResult {
|
||||
Enqueued,
|
||||
Coalesced,
|
||||
Dropped(MrfDropReason),
|
||||
}
|
||||
|
||||
/// Consumer-side retry ceiling before an intent is given up on.
|
||||
pub const MRF_MAX_ATTEMPTS: u8 = 3;
|
||||
|
||||
@@ -138,159 +87,6 @@ impl MrfIntent {
|
||||
|
||||
static GLOBAL_MRF_SENDER: OnceLock<mpsc::Sender<MrfIntent>> = OnceLock::new();
|
||||
|
||||
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
|
||||
struct MrfIdentityKey {
|
||||
kind: MrfKind,
|
||||
bucket: Arc<str>,
|
||||
object: Arc<str>,
|
||||
version_id: Option<[u8; 16]>,
|
||||
scope: Option<MrfScope>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct IngressEntry {
|
||||
lease: MrfIngressLease,
|
||||
expires_at: Instant,
|
||||
bytes: usize,
|
||||
}
|
||||
|
||||
type MrfCoalescerShard = Mutex<HashMap<MrfIdentityKey, IngressEntry>>;
|
||||
type MrfCoalescer = Box<[MrfCoalescerShard]>;
|
||||
|
||||
static MRF_COALESCER: OnceLock<MrfCoalescer> = OnceLock::new();
|
||||
static NEXT_MRF_LEASE: AtomicU64 = AtomicU64::new(1);
|
||||
static MRF_COALESCER_COUNT: AtomicUsize = AtomicUsize::new(0);
|
||||
static MRF_COALESCER_BYTES: AtomicUsize = AtomicUsize::new(0);
|
||||
static MRF_HASH_STATE: OnceLock<RandomState> = OnceLock::new();
|
||||
|
||||
fn coalescer() -> &'static [MrfCoalescerShard] {
|
||||
MRF_COALESCER.get_or_init(|| {
|
||||
(0..MRF_COALESCER_SHARDS)
|
||||
.map(|_| Mutex::new(HashMap::new()))
|
||||
.collect::<Vec<_>>()
|
||||
.into_boxed_slice()
|
||||
})
|
||||
}
|
||||
|
||||
fn key_shard(key: &MrfIdentityKey) -> usize {
|
||||
let hash = MRF_HASH_STATE.get_or_init(RandomState::new).hash_one(key);
|
||||
usize::try_from(hash).unwrap_or(0) % MRF_COALESCER_SHARDS
|
||||
}
|
||||
|
||||
fn canonical_version(version_id: Option<Uuid>) -> Option<[u8; 16]> {
|
||||
version_id
|
||||
.filter(|version| !version.is_nil())
|
||||
.map(|version| *version.as_bytes())
|
||||
}
|
||||
|
||||
fn canonical_identity(
|
||||
kind: MrfKind,
|
||||
version_id: Option<[u8; 16]>,
|
||||
scope: Option<MrfScope>,
|
||||
) -> (Option<[u8; 16]>, Option<MrfScope>) {
|
||||
let version_id = version_id.filter(|bytes| *bytes != [0; 16]);
|
||||
match kind {
|
||||
MrfKind::MetadataCorruption => (None, None),
|
||||
MrfKind::DecodeFailure | MrfKind::PartialWrite => (version_id, scope),
|
||||
}
|
||||
}
|
||||
|
||||
fn identity_estimated_bytes(key: &MrfIdentityKey) -> usize {
|
||||
64usize
|
||||
.saturating_add(key.bucket.len())
|
||||
.saturating_add(key.object.len())
|
||||
.saturating_add(key.version_id.map_or(0, |_| 16))
|
||||
.saturating_add(key.scope.map_or(0, |_| 8))
|
||||
}
|
||||
|
||||
fn reserve(counter: &AtomicUsize, limit: usize, amount: usize) -> bool {
|
||||
let mut current = counter.load(Ordering::Relaxed);
|
||||
loop {
|
||||
let Some(next) = current.checked_add(amount) else {
|
||||
return false;
|
||||
};
|
||||
if next > limit {
|
||||
return false;
|
||||
}
|
||||
match counter.compare_exchange_weak(current, next, Ordering::Relaxed, Ordering::Relaxed) {
|
||||
Ok(_) => return true,
|
||||
Err(observed) => current = observed,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn coalescer_admit(key: MrfIdentityKey) -> Result<MrfIngressLease, MrfIngressResult> {
|
||||
let shard = key_shard(&key);
|
||||
let mut entries = coalescer()[shard]
|
||||
.lock()
|
||||
.map_err(|_| MrfIngressResult::Dropped(MrfDropReason::CoalescerFull))?;
|
||||
let now = Instant::now();
|
||||
let before = entries.len();
|
||||
let mut expired_bytes = 0usize;
|
||||
entries.retain(|_, entry| {
|
||||
if entry.expires_at > now {
|
||||
true
|
||||
} else {
|
||||
expired_bytes = expired_bytes.saturating_add(entry.bytes);
|
||||
false
|
||||
}
|
||||
});
|
||||
let evicted = before.saturating_sub(entries.len());
|
||||
if evicted > 0 {
|
||||
MRF_COALESCER_COUNT.fetch_sub(evicted, Ordering::Relaxed);
|
||||
MRF_COALESCER_BYTES.fetch_sub(expired_bytes, Ordering::Relaxed);
|
||||
let evicted = u64::try_from(evicted).unwrap_or(u64::MAX);
|
||||
metrics::counter!("rustfs_heal_mrf_coalescer_expired_total").increment(evicted);
|
||||
metrics::counter!("rustfs_heal_mrf_coalescer_evictions_total").increment(evicted);
|
||||
}
|
||||
if entries.contains_key(&key) {
|
||||
metrics::counter!("rustfs_heal_mrf_coalesced_total").increment(1);
|
||||
return Err(MrfIngressResult::Coalesced);
|
||||
}
|
||||
let bytes = identity_estimated_bytes(&key);
|
||||
let count_reserved = reserve(&MRF_COALESCER_COUNT, MRF_COALESCER_MAX_KEYS, 1);
|
||||
let bytes_reserved = count_reserved && reserve(&MRF_COALESCER_BYTES, MRF_COALESCER_MAX_BYTES, bytes);
|
||||
if !count_reserved || !bytes_reserved {
|
||||
if count_reserved {
|
||||
MRF_COALESCER_COUNT.fetch_sub(1, Ordering::Relaxed);
|
||||
}
|
||||
metrics::counter!("rustfs_heal_mrf_dropped_total", "reason" => "coalescer_full").increment(1);
|
||||
return Err(MrfIngressResult::Dropped(MrfDropReason::CoalescerFull));
|
||||
}
|
||||
let lease = MrfIngressLease::new(NEXT_MRF_LEASE.fetch_add(1, Ordering::Relaxed));
|
||||
if entries
|
||||
.insert(
|
||||
key,
|
||||
IngressEntry {
|
||||
lease,
|
||||
expires_at: now + MRF_COALESCER_TTL,
|
||||
bytes,
|
||||
},
|
||||
)
|
||||
.is_some()
|
||||
{
|
||||
MRF_COALESCER_COUNT.fetch_sub(1, Ordering::Relaxed);
|
||||
MRF_COALESCER_BYTES.fetch_sub(bytes, Ordering::Relaxed);
|
||||
metrics::counter!("rustfs_heal_mrf_coalesced_total").increment(1);
|
||||
return Err(MrfIngressResult::Coalesced);
|
||||
}
|
||||
Ok(lease)
|
||||
}
|
||||
|
||||
fn coalescer_release(key: &MrfIdentityKey, lease: Option<MrfIngressLease>) {
|
||||
let Some(lease) = lease else {
|
||||
return;
|
||||
};
|
||||
if let Ok(mut entries) = coalescer()[key_shard(key)].lock() {
|
||||
let should_remove = entries.get(key).is_some_and(|entry| entry.lease == lease);
|
||||
if should_remove {
|
||||
let bytes = entries.remove(key).map(|entry| entry.bytes).unwrap_or(0);
|
||||
MRF_COALESCER_COUNT.fetch_sub(1, Ordering::Relaxed);
|
||||
MRF_COALESCER_BYTES.fetch_sub(bytes, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Delivery kill-switch, set from `RUSTFS_HEAL_MRF_ENABLE`. Producers check
|
||||
/// this before touching the channel so the disabled path stays allocation- and
|
||||
/// sync-free.
|
||||
@@ -326,90 +122,21 @@ pub fn init_mrf_channel() -> Result<mpsc::Receiver<MrfIntent>, &'static str> {
|
||||
/// This runs on IO error paths, so it stays synchronous and cheap: one
|
||||
/// bounded allocation for the two `Arc<str>` handles plus the channel slot.
|
||||
pub fn try_send_mrf_intent(kind: MrfKind, bucket: &str, object: &str, version_id: Option<Uuid>) -> bool {
|
||||
matches!(
|
||||
try_send_mrf_intent_typed(kind, bucket, object, version_id, None),
|
||||
MrfIngressResult::Enqueued
|
||||
)
|
||||
}
|
||||
|
||||
/// Typed ingress result. `Coalesced` means an equivalent in-flight channel
|
||||
/// intent already exists; it is not a second executable or durable admission.
|
||||
pub fn try_send_mrf_intent_typed(
|
||||
kind: MrfKind,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version_id: Option<Uuid>,
|
||||
scope: Option<MrfScope>,
|
||||
) -> MrfIngressResult {
|
||||
if !mrf_delivery_enabled() {
|
||||
return MrfIngressResult::Dropped(MrfDropReason::Disabled);
|
||||
return false;
|
||||
}
|
||||
let Some(sender) = GLOBAL_MRF_SENDER.get() else {
|
||||
return MrfIngressResult::Dropped(MrfDropReason::Uninitialized);
|
||||
};
|
||||
if bucket.len() > MRF_MAX_IDENTITY_COMPONENT || object.len() > MRF_MAX_IDENTITY_COMPONENT {
|
||||
return MrfIngressResult::Dropped(MrfDropReason::OversizedIdentity);
|
||||
}
|
||||
let (version_id, scope) = canonical_identity(kind, canonical_version(version_id), scope);
|
||||
let key = MrfIdentityKey {
|
||||
kind,
|
||||
bucket: Arc::from(bucket),
|
||||
object: Arc::from(object),
|
||||
version_id,
|
||||
scope,
|
||||
};
|
||||
let lease = match coalescer_admit(key.clone()) {
|
||||
Ok(lease) => lease,
|
||||
Err(result) => return result,
|
||||
return false;
|
||||
};
|
||||
let intent = MrfIntent {
|
||||
bucket: key.bucket.clone(),
|
||||
object: key.object.clone(),
|
||||
version_id: key.version_id,
|
||||
bucket: Arc::from(bucket),
|
||||
object: Arc::from(object),
|
||||
version_id: version_id.map(|vid| *vid.as_bytes()),
|
||||
kind,
|
||||
scope,
|
||||
lease: Some(lease),
|
||||
enqueued_at_ms: unix_now_ms(),
|
||||
attempts: 0,
|
||||
};
|
||||
match sender.try_send(intent) {
|
||||
Ok(()) => MrfIngressResult::Enqueued,
|
||||
Err(mpsc::error::TrySendError::Full(_)) => {
|
||||
coalescer_release(&key, Some(lease));
|
||||
metrics::counter!("rustfs_heal_mrf_dropped_total", "reason" => "channel_full").increment(1);
|
||||
MrfIngressResult::Dropped(MrfDropReason::Full)
|
||||
}
|
||||
Err(mpsc::error::TrySendError::Closed(_)) => {
|
||||
coalescer_release(&key, Some(lease));
|
||||
MrfIngressResult::Dropped(MrfDropReason::Uninitialized)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 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);
|
||||
}
|
||||
|
||||
pub fn release_mrf_identity(
|
||||
kind: MrfKind,
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
version_id: Option<[u8; 16]>,
|
||||
scope: Option<MrfScope>,
|
||||
lease: Option<MrfIngressLease>,
|
||||
) {
|
||||
let (version_id, scope) = canonical_identity(kind, version_id, scope);
|
||||
coalescer_release(
|
||||
&MrfIdentityKey {
|
||||
kind,
|
||||
bucket: Arc::from(bucket),
|
||||
object: Arc::from(object),
|
||||
version_id,
|
||||
scope,
|
||||
},
|
||||
lease,
|
||||
);
|
||||
sender.try_send(intent).is_ok()
|
||||
}
|
||||
|
||||
fn unix_now_ms() -> u64 {
|
||||
@@ -417,8 +144,7 @@ fn unix_now_ms() -> u64 {
|
||||
// failure would be a bug rather than something to handle here.
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.ok()
|
||||
.and_then(|d| u64::try_from(d.as_millis()).ok())
|
||||
.map(|d| d.as_millis() as u64)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
@@ -489,60 +215,12 @@ mod tests {
|
||||
object: Arc::from("object"),
|
||||
version_id: Some([0u8; 16]),
|
||||
kind: MrfKind::DecodeFailure,
|
||||
scope: None,
|
||||
lease: None,
|
||||
enqueued_at_ms: 0,
|
||||
attempts: 0,
|
||||
};
|
||||
assert!(intent.estimated_bytes() >= intent.bucket.len() + intent.object.len());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ingress_duplicate_identity_coalesces_and_releases_for_retry() {
|
||||
let key = MrfIdentityKey {
|
||||
kind: MrfKind::DecodeFailure,
|
||||
bucket: Arc::from("ingress-test-bucket"),
|
||||
object: Arc::from("ingress-test-object"),
|
||||
version_id: Some([9; 16]),
|
||||
scope: Some(MrfScope {
|
||||
pool_index: 3,
|
||||
set_index: 4,
|
||||
}),
|
||||
};
|
||||
let lease = coalescer_admit(key.clone()).expect("first identity should be admitted");
|
||||
for _ in 0..999 {
|
||||
assert_eq!(coalescer_admit(key.clone()), Err(MrfIngressResult::Coalesced));
|
||||
}
|
||||
coalescer_release(&key, Some(lease));
|
||||
let retry_lease = coalescer_admit(key.clone()).expect("released identity must admit a retry");
|
||||
coalescer_release(&key, Some(retry_lease));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn ingress_identity_preserves_kind_scope_and_version_boundaries() {
|
||||
let (nil_version, nil_scope) = canonical_identity(
|
||||
MrfKind::DecodeFailure,
|
||||
Some([0; 16]),
|
||||
Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
);
|
||||
assert_eq!(nil_version, None, "nil UUID is the unversioned identity");
|
||||
assert!(nil_scope.is_some());
|
||||
|
||||
let (metadata_version, metadata_scope) = canonical_identity(
|
||||
MrfKind::MetadataCorruption,
|
||||
Some([7; 16]),
|
||||
Some(MrfScope {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
}),
|
||||
);
|
||||
assert_eq!(metadata_version, None);
|
||||
assert_eq!(metadata_scope, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn try_send_delivers_and_respects_capacity() {
|
||||
let mut receiver = init_mrf_channel().expect("first initialization should succeed");
|
||||
@@ -552,7 +230,6 @@ mod tests {
|
||||
let intent = receiver.recv().await.expect("intent should arrive");
|
||||
assert_eq!(intent.kind, MrfKind::DecodeFailure);
|
||||
assert_eq!(intent.bucket.as_ref(), "b");
|
||||
release_mrf_intent(&intent);
|
||||
|
||||
// Disable delivery: producers become no-ops.
|
||||
set_mrf_delivery_enabled(false);
|
||||
@@ -562,8 +239,8 @@ mod tests {
|
||||
// Fill the bounded channel past capacity: excess intents are dropped,
|
||||
// never blocking.
|
||||
let mut accepted = 0;
|
||||
for index in 0..(MRF_CHANNEL_CAPACITY + 64) {
|
||||
if try_send_mrf_intent(MrfKind::PartialWrite, "b", &format!("o-{index}"), None) {
|
||||
for _ in 0..(MRF_CHANNEL_CAPACITY + 64) {
|
||||
if try_send_mrf_intent(MrfKind::PartialWrite, "b", "o", None) {
|
||||
accepted += 1;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -236,6 +236,15 @@ pub struct DataUsageInfo {
|
||||
/// without relying on synchronized clocks.
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub usage_snapshot_authoritative_baseline: Option<DataUsageSnapshotIdentity>,
|
||||
/// Per-set freshness for an observational aggregate. A set entry is
|
||||
/// never sufficient to make the aggregate authoritative; it only records
|
||||
/// which last-known-good generation contributed to the view.
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub usage_snapshot_set_states: Vec<DataUsageSnapshotSetState>,
|
||||
/// An observational view may contain only the sets that completed this
|
||||
/// cycle (or retained a compatible last-known-good cache).
|
||||
#[serde(default)]
|
||||
pub usage_snapshot_partial: bool,
|
||||
/// Deprecated kept here for backward compatibility reasons
|
||||
pub bucket_sizes: HashMap<String, u64>,
|
||||
/// Per-disk snapshot information when available
|
||||
@@ -252,6 +261,22 @@ pub struct DataUsageSnapshotIdentity {
|
||||
pub scanner_epoch: Option<u64>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
|
||||
pub struct DataUsageSnapshotSetState {
|
||||
pub pool_index: u64,
|
||||
pub set_index: u64,
|
||||
#[serde(default)]
|
||||
pub scanner_cycle: Option<u64>,
|
||||
#[serde(default)]
|
||||
pub scanner_epoch: Option<u64>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub scan_plan_digest: Option<[u8; 32]>,
|
||||
#[serde(default)]
|
||||
pub complete: bool,
|
||||
#[serde(default)]
|
||||
pub tombstone: bool,
|
||||
}
|
||||
|
||||
impl DataUsageInfo {
|
||||
pub fn snapshot_identity(&self) -> DataUsageSnapshotIdentity {
|
||||
DataUsageSnapshotIdentity {
|
||||
@@ -291,7 +316,7 @@ pub fn data_usage_snapshot_is_newer(candidate: &DataUsageInfo, baseline: &DataUs
|
||||
/// rollback delete/recreate fences the previous bucket incarnation too.
|
||||
pub fn observed_data_usage_is_newer(observed: &DataUsageInfo, authoritative: &DataUsageInfo) -> bool {
|
||||
observed.usage_snapshot_converged == Some(false)
|
||||
&& observed.is_complete_bucket_usage_snapshot()
|
||||
&& (observed.is_complete_bucket_usage_snapshot() || observed.is_valid_partial_snapshot())
|
||||
&& observed.usage_snapshot_authoritative_baseline.as_ref() == Some(&authoritative.snapshot_identity())
|
||||
&& data_usage_snapshot_is_newer(observed, authoritative)
|
||||
}
|
||||
@@ -1436,6 +1461,39 @@ impl DataUsageInfo {
|
||||
&& u64::try_from(self.buckets_usage.len()).ok() == Some(self.buckets_count)
|
||||
}
|
||||
|
||||
/// Validate provenance before an observational view can be selected for
|
||||
/// admin display. Partial data is accepted only with unique set states,
|
||||
/// a plan digest for every state, and at least one usable generation.
|
||||
pub fn is_valid_partial_snapshot(&self) -> bool {
|
||||
if !self.usage_snapshot_partial
|
||||
|| self.usage_snapshot_converged != Some(false)
|
||||
|| self.last_update.is_none()
|
||||
|| self.scanner_cycle.is_none()
|
||||
|| self.scanner_epoch.is_none()
|
||||
|| self.usage_snapshot_set_states.is_empty()
|
||||
|| u64::try_from(self.buckets_usage.len()).ok() != Some(self.buckets_count)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
|
||||
let mut previous = None;
|
||||
let mut plan_digest = None;
|
||||
let mut has_source = false;
|
||||
for state in &self.usage_snapshot_set_states {
|
||||
if state.scan_plan_digest.is_none()
|
||||
|| plan_digest.is_some_and(|digest| Some(digest) != state.scan_plan_digest)
|
||||
|| state.scanner_cycle.is_some() != state.scanner_epoch.is_some()
|
||||
|| previous.is_some_and(|(pool, set)| (pool, set) >= (state.pool_index, state.set_index))
|
||||
{
|
||||
return false;
|
||||
}
|
||||
previous = Some((state.pool_index, state.set_index));
|
||||
plan_digest = state.scan_plan_digest;
|
||||
has_source |= state.scanner_cycle.is_some() && !state.tombstone;
|
||||
}
|
||||
has_source
|
||||
}
|
||||
|
||||
/// Add object metadata to data usage statistics
|
||||
pub fn add_object(&mut self, object_path: &str, meta_object: &rustfs_filemeta::MetaObject) {
|
||||
// This method is kept for backward compatibility
|
||||
@@ -2263,6 +2321,55 @@ mod tests {
|
||||
assert!(!observed_data_usage_is_newer(&candidate(2, 9, Some(false), true), &authoritative));
|
||||
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(true), true), &authoritative));
|
||||
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(false), false), &authoritative));
|
||||
|
||||
let mut partial = candidate(2, 11, Some(false), false);
|
||||
partial.usage_snapshot_partial = true;
|
||||
partial.usage_snapshot_set_states = vec![DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(10),
|
||||
scanner_epoch: Some(2),
|
||||
scan_plan_digest: Some([1; 32]),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
}];
|
||||
assert!(observed_data_usage_is_newer(&partial, &authoritative));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mixed_topology_snapshot_is_rejected() {
|
||||
let mut partial = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(2)),
|
||||
scanner_cycle: Some(11),
|
||||
scanner_epoch: Some(2),
|
||||
buckets_count: 0,
|
||||
usage_snapshot_converged: Some(false),
|
||||
usage_snapshot_partial: true,
|
||||
usage_snapshot_set_states: vec![
|
||||
DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(11),
|
||||
scanner_epoch: Some(2),
|
||||
scan_plan_digest: Some([1; 32]),
|
||||
complete: true,
|
||||
tombstone: false,
|
||||
},
|
||||
DataUsageSnapshotSetState {
|
||||
pool_index: 1,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(10),
|
||||
scanner_epoch: Some(2),
|
||||
scan_plan_digest: Some([2; 32]),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
},
|
||||
],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(!partial.is_valid_partial_snapshot());
|
||||
partial.usage_snapshot_set_states[1].scan_plan_digest = Some([1; 32]);
|
||||
assert!(partial.is_valid_partial_snapshot());
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -73,6 +73,16 @@ struct CachedBucketUsage {
|
||||
// mutation. A strictly later generation is required before the mutation
|
||||
// evidence can be discarded.
|
||||
pending_scanner_position: Option<(u64, u64)>,
|
||||
// Deletes are visible to admin immediately, but quota admission keeps
|
||||
// them pending until a complete scanner generation reconciles the set.
|
||||
// This marker intentionally remains process-local: the delete request
|
||||
// updates this overlay before the scanner writes a durable snapshot. If
|
||||
// the process restarts first, loading the persisted complete snapshot
|
||||
// restores the pre-reconciliation (larger) baseline, which is
|
||||
// conservative for quota admission. A persisted post-delete snapshot is
|
||||
// necessarily a complete scanner reconciliation and therefore creates a
|
||||
// fresh cache entry with no pending hold.
|
||||
pending_negative_delta: u64,
|
||||
}
|
||||
|
||||
type UsageMemoryCache = Arc<RwLock<HashMap<String, CachedBucketUsage>>>;
|
||||
@@ -948,7 +958,12 @@ async fn load_observed_data_usage_snapshot(store: Arc<ECStore>) -> Option<DataUs
|
||||
};
|
||||
|
||||
match parse_usage_snapshot(&data) {
|
||||
Ok(info) if info.usage_snapshot_converged == Some(false) && info.is_complete_bucket_usage_snapshot() => Some(info),
|
||||
Ok(info)
|
||||
if info.usage_snapshot_converged == Some(false)
|
||||
&& (info.is_complete_bucket_usage_snapshot() || info.is_valid_partial_snapshot()) =>
|
||||
{
|
||||
Some(info)
|
||||
}
|
||||
Ok(_) => {
|
||||
error!(
|
||||
event = "data_usage_snapshot_load_failed",
|
||||
@@ -993,7 +1008,7 @@ async fn load_admin_data_usage_from_backend(store: Arc<ECStore>) -> Result<DataU
|
||||
}
|
||||
|
||||
fn discard_incomplete_bucket_usage(data_usage_info: &mut DataUsageInfo) {
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() {
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
|
||||
data_usage_info.usage_snapshot_complete = false;
|
||||
data_usage_info.buckets_usage.clear();
|
||||
data_usage_info.bucket_sizes.clear();
|
||||
@@ -1643,6 +1658,7 @@ fn cached_bucket_usage_from_backend(usage: BucketUsageInfo, updated_at: SystemTi
|
||||
dirty: false,
|
||||
stale_snapshot_pending: false,
|
||||
pending_scanner_position: None,
|
||||
pending_negative_delta: 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1656,6 +1672,7 @@ fn cached_bucket_usage_now(usage: BucketUsageInfo) -> CachedBucketUsage {
|
||||
dirty: false,
|
||||
stale_snapshot_pending: false,
|
||||
pending_scanner_position: None,
|
||||
pending_negative_delta: 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1808,6 +1825,7 @@ pub async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64,
|
||||
.or_insert_with(|| cached_bucket_usage_now(BucketUsageInfo::default()));
|
||||
|
||||
entry.usage.size = entry.usage.size.saturating_sub(deleted_size);
|
||||
entry.pending_negative_delta = entry.pending_negative_delta.saturating_add(deleted_size);
|
||||
if removed_current_object {
|
||||
entry.usage.objects_count = entry.usage.objects_count.saturating_sub(1);
|
||||
entry.usage.versions_count = entry.usage.versions_count.saturating_sub(1);
|
||||
@@ -1863,7 +1881,7 @@ pub async fn get_bucket_usage_memory(bucket: &str) -> Option<u64> {
|
||||
cache
|
||||
.get(bucket)
|
||||
.filter(|cached| cached.authoritative)
|
||||
.map(|cached| cached.usage.size)
|
||||
.map(|cached| cached.usage.size.saturating_add(cached.pending_negative_delta))
|
||||
}
|
||||
|
||||
async fn update_usage_cache_if_needed() {
|
||||
@@ -2943,6 +2961,45 @@ mod tests {
|
||||
assert_eq!(selected.usage_snapshot_converged, Some(true));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn persisted_authoritative_stalls_but_memory_overlay_remains_visible() {
|
||||
let authoritative = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH),
|
||||
scanner_epoch: Some(4),
|
||||
scanner_cycle: Some(10),
|
||||
usage_snapshot_complete: true,
|
||||
..Default::default()
|
||||
};
|
||||
let mut partial = authoritative.clone();
|
||||
partial.last_update = Some(SystemTime::UNIX_EPOCH + Duration::from_secs(1));
|
||||
partial.scanner_cycle = Some(11);
|
||||
partial.usage_snapshot_complete = false;
|
||||
partial.usage_snapshot_partial = true;
|
||||
partial.usage_snapshot_converged = Some(false);
|
||||
partial.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
|
||||
partial.usage_snapshot_set_states = vec![rustfs_data_usage::DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(10),
|
||||
scanner_epoch: Some(4),
|
||||
scan_plan_digest: Some([1; 32]),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
}];
|
||||
partial.buckets_usage.insert(
|
||||
"bucket".to_string(),
|
||||
BucketUsageInfo {
|
||||
size: 100,
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
partial.buckets_count = 1;
|
||||
|
||||
let (selected, _) = select_admin_data_usage_snapshot(authoritative, true, Some(partial));
|
||||
assert!(selected.usage_snapshot_partial);
|
||||
assert_eq!(selected.buckets_usage.get("bucket").map(|usage| usage.size), Some(100));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn authoritative_save_cleanup_removes_observed_snapshot_best_effort() {
|
||||
let store = UsageCasStore::default();
|
||||
@@ -4665,6 +4722,55 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn partial_usage_is_observational_not_authoritative_for_quota() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let mut partial = data_usage_info_for_test("bucket-a", 10, 100, SystemTime::now());
|
||||
partial.usage_snapshot_complete = false;
|
||||
partial.usage_snapshot_partial = true;
|
||||
replace_bucket_usage_memory_from_info(&partial).await;
|
||||
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::now());
|
||||
replace_bucket_usage_memory_from_info(&baseline).await;
|
||||
record_bucket_object_write_memory("bucket-a", None, 25).await;
|
||||
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(125));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn negative_delta_waits_for_set_reconciliation() {
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
|
||||
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::UNIX_EPOCH + Duration::from_secs(100));
|
||||
replace_bucket_usage_memory_from_info(&baseline).await;
|
||||
record_bucket_object_delete_memory("bucket-a", 25, true).await;
|
||||
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
|
||||
|
||||
// Simulate a process restart: the request-path overlay is gone, but
|
||||
// the persisted authoritative snapshot is still the pre-reconciliation
|
||||
// baseline. Quota must remain conservative until a complete scanner
|
||||
// result proves the delete.
|
||||
clear_usage_memory_cache_for_test().await;
|
||||
replace_bucket_usage_memory_from_info(&baseline).await;
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
|
||||
|
||||
let reconciled = data_usage_info_for_test("bucket-a", 0, 75, SystemTime::UNIX_EPOCH + Duration::from_secs(101));
|
||||
replace_bucket_usage_memory_from_info(&reconciled).await;
|
||||
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(75));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial]
|
||||
async fn memory_overlay_counts_versioned_overwrite_as_new_version() {
|
||||
|
||||
@@ -1412,11 +1412,8 @@ pub(in crate::set_disk) async fn submit_read_repair_heal_with_submitter(
|
||||
|
||||
// Reservation won: this sighting owns the repair records for the object,
|
||||
// including the durable journal intent when the caller asked for one.
|
||||
if let Some((kind, version_uuid)) = mrf_intent
|
||||
&& let (Ok(pool_index), Ok(set_index)) = (u32::try_from(pool_index), u32::try_from(set_index))
|
||||
{
|
||||
let scope = rustfs_common::mrf_channel::MrfScope { pool_index, set_index };
|
||||
let _ = rustfs_common::mrf_channel::try_send_mrf_intent_typed(kind, bucket, object, version_uuid, Some(scope));
|
||||
if let Some((kind, version_uuid)) = mrf_intent {
|
||||
rustfs_common::mrf_channel::try_send_mrf_intent(kind, bucket, object, version_uuid);
|
||||
}
|
||||
|
||||
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
|
||||
|
||||
@@ -6650,23 +6650,12 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
|
||||
async fn add_partial(&self, bucket: &str, object: &str, version_id: &str) -> Result<()> {
|
||||
// MRF journal intent: partial-write recovery must survive a restart
|
||||
// (HS-01); the heal request below remains the in-memory fast path.
|
||||
let version_uuid = if version_id.is_empty() {
|
||||
Some(None)
|
||||
} else {
|
||||
uuid::Uuid::try_parse(version_id).ok().map(Some)
|
||||
};
|
||||
if let Some(version_uuid) = version_uuid
|
||||
&& let (Ok(pool_index), Ok(set_index)) = (u32::try_from(self.pool_index), u32::try_from(self.set_index))
|
||||
{
|
||||
let scope = rustfs_common::mrf_channel::MrfScope { pool_index, set_index };
|
||||
let _ = rustfs_common::mrf_channel::try_send_mrf_intent_typed(
|
||||
rustfs_common::mrf_channel::MrfKind::PartialWrite,
|
||||
bucket,
|
||||
object,
|
||||
version_uuid,
|
||||
Some(scope),
|
||||
);
|
||||
}
|
||||
rustfs_common::mrf_channel::try_send_mrf_intent(
|
||||
rustfs_common::mrf_channel::MrfKind::PartialWrite,
|
||||
bucket,
|
||||
object,
|
||||
uuid::Uuid::try_parse(version_id).ok(),
|
||||
);
|
||||
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
|
||||
bucket.to_string(),
|
||||
Some(object.to_string()),
|
||||
|
||||
@@ -119,9 +119,6 @@ struct MrfRepairNoticeTarget {
|
||||
bucket: Arc<str>,
|
||||
object: Arc<str>,
|
||||
version_id: Option<[u8; 16]>,
|
||||
kind: rustfs_common::mrf_channel::MrfKind,
|
||||
scope: Option<rustfs_common::mrf_channel::MrfScope>,
|
||||
lease: Option<rustfs_common::mrf_channel::MrfIngressLease>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
@@ -892,19 +889,7 @@ impl HealManager {
|
||||
}
|
||||
|
||||
fn remove_mrf_repair_notice_targets_for_task(&self, task_id: &str) {
|
||||
let targets = lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).remove(task_id);
|
||||
if let Some(targets) = targets {
|
||||
for target in targets {
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
&target.bucket,
|
||||
&target.object,
|
||||
target.version_id,
|
||||
target.scope,
|
||||
target.lease,
|
||||
);
|
||||
}
|
||||
}
|
||||
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).remove(task_id);
|
||||
}
|
||||
|
||||
fn insert_mrf_repair_notice_target(
|
||||
@@ -1294,20 +1279,7 @@ impl HealManager {
|
||||
}
|
||||
self.task_aliases.lock().await.clear();
|
||||
self.retrying_heals.lock().await.clear();
|
||||
let mrf_targets = {
|
||||
let mut registry = lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets);
|
||||
registry.drain().flat_map(|(_, targets)| targets).collect::<Vec<_>>()
|
||||
};
|
||||
for target in mrf_targets {
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
&target.bucket,
|
||||
&target.object,
|
||||
target.version_id,
|
||||
target.scope,
|
||||
target.lease,
|
||||
);
|
||||
}
|
||||
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).clear();
|
||||
crate::set_heal_queue_length(0);
|
||||
|
||||
// update state
|
||||
@@ -1339,32 +1311,12 @@ impl HealManager {
|
||||
.await
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn submit_mrf_heal_request_with_receipt(
|
||||
&self,
|
||||
request: HealRequest,
|
||||
bucket: Arc<str>,
|
||||
object: Arc<str>,
|
||||
version_id: Option<[u8; 16]>,
|
||||
) -> Result<HealAdmissionReceipt> {
|
||||
let kind = match &request.heal_type {
|
||||
HealType::Metadata { .. } => rustfs_common::mrf_channel::MrfKind::MetadataCorruption,
|
||||
HealType::ECDecode { .. } => rustfs_common::mrf_channel::MrfKind::DecodeFailure,
|
||||
_ => rustfs_common::mrf_channel::MrfKind::PartialWrite,
|
||||
};
|
||||
self.submit_mrf_heal_request_with_receipt_and_identity(request, bucket, object, version_id, kind, None, None)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn submit_mrf_heal_request_with_receipt_and_identity(
|
||||
&self,
|
||||
request: HealRequest,
|
||||
bucket: Arc<str>,
|
||||
object: Arc<str>,
|
||||
version_id: Option<[u8; 16]>,
|
||||
kind: rustfs_common::mrf_channel::MrfKind,
|
||||
scope: Option<rustfs_common::mrf_channel::MrfScope>,
|
||||
lease: Option<rustfs_common::mrf_channel::MrfIngressLease>,
|
||||
) -> Result<HealAdmissionReceipt> {
|
||||
self.submit_heal_request_with_receipt_alias_and_mrf_notice(
|
||||
request,
|
||||
@@ -1373,9 +1325,6 @@ impl HealManager {
|
||||
bucket,
|
||||
object,
|
||||
version_id,
|
||||
kind,
|
||||
scope,
|
||||
lease,
|
||||
}),
|
||||
)
|
||||
.await
|
||||
@@ -1590,7 +1539,7 @@ impl HealManager {
|
||||
Self::insert_mrf_repair_notice_target(&mut targets, &task_id, target);
|
||||
}
|
||||
if let Some(displaced_task_id) = &displaced_task_id {
|
||||
self.remove_mrf_repair_notice_targets_for_task(displaced_task_id);
|
||||
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).remove(displaced_task_id);
|
||||
}
|
||||
drop(retrying_heals);
|
||||
drop(queue);
|
||||
|
||||
@@ -506,18 +506,7 @@ impl HealManager {
|
||||
&displaced_terminal,
|
||||
)
|
||||
.await;
|
||||
if let Some(targets) = lock_mrf_repair_notice_targets(&mrf_repair_notice_targets).remove(&displaced_task_id) {
|
||||
for target in targets {
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
&target.bucket,
|
||||
&target.object,
|
||||
target.version_id,
|
||||
target.scope,
|
||||
target.lease,
|
||||
);
|
||||
}
|
||||
}
|
||||
lock_mrf_repair_notice_targets(&mrf_repair_notice_targets).remove(&displaced_task_id);
|
||||
}
|
||||
if matches!(admission, HealAdmissionResult::Accepted) {
|
||||
if should_notify {
|
||||
|
||||
@@ -299,11 +299,7 @@ impl PriorityHealQueue {
|
||||
|
||||
/// Create a deduplication key from a heal request
|
||||
pub(super) fn make_dedup_key(request: &HealRequest) -> String {
|
||||
let base = Self::make_dedup_key_for_type(&request.heal_type);
|
||||
match (&request.heal_type, request.options.set_key()) {
|
||||
(HealType::Object { .. } | HealType::ECDecode { .. }, Some(scope)) => format!("{base}:scope:{scope}"),
|
||||
_ => base,
|
||||
}
|
||||
Self::make_dedup_key_for_type(&request.heal_type)
|
||||
}
|
||||
|
||||
pub(super) fn make_dedup_key_for_type(heal_type: &HealType) -> String {
|
||||
|
||||
@@ -349,8 +349,6 @@ impl HealManager {
|
||||
let notice_targets = take_mrf_repair_notice_targets(&mrf_repair_notice_targets_clone, &task_id);
|
||||
if successful_completion {
|
||||
emit_mrf_repaired_events(notice_targets);
|
||||
} else {
|
||||
release_mrf_repair_notice_targets(notice_targets);
|
||||
}
|
||||
task_aliases_clone
|
||||
.lock()
|
||||
@@ -640,19 +638,7 @@ pub(super) fn running_heal_set_counts(active_heals: &HashMap<String, Arc<HealTas
|
||||
}
|
||||
|
||||
fn remove_mrf_repair_notice_targets(registry: &Arc<StdMutex<HashMap<String, Vec<MrfRepairNoticeTarget>>>>, task_id: &str) {
|
||||
let targets = lock_mrf_repair_notice_targets(registry).remove(task_id);
|
||||
if let Some(targets) = targets {
|
||||
for target in targets {
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
&target.bucket,
|
||||
&target.object,
|
||||
target.version_id,
|
||||
target.scope,
|
||||
target.lease,
|
||||
);
|
||||
}
|
||||
}
|
||||
lock_mrf_repair_notice_targets(registry).remove(task_id);
|
||||
}
|
||||
|
||||
fn take_mrf_repair_notice_targets(
|
||||
@@ -685,27 +671,6 @@ fn move_mrf_repair_notice_targets(
|
||||
fn emit_mrf_repaired_events(targets: Vec<MrfRepairNoticeTarget>) {
|
||||
for target in targets {
|
||||
rustfs_common::mrf_channel::note_mrf_repaired(&target.bucket, &target.object, target.version_id);
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
&target.bucket,
|
||||
&target.object,
|
||||
target.version_id,
|
||||
target.scope,
|
||||
target.lease,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
fn release_mrf_repair_notice_targets(targets: Vec<MrfRepairNoticeTarget>) {
|
||||
for target in targets {
|
||||
rustfs_common::mrf_channel::release_mrf_identity(
|
||||
target.kind,
|
||||
&target.bucket,
|
||||
&target.object,
|
||||
target.version_id,
|
||||
target.scope,
|
||||
target.lease,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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::{HashSet, VecDeque};
|
||||
use std::collections::VecDeque;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::sync::mpsc;
|
||||
@@ -48,27 +48,15 @@ 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 {
|
||||
@@ -113,90 +101,40 @@ 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,
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
/// Returns `false` (after counting) when either ceiling would be crossed.
|
||||
pub(crate) fn try_push(&mut self, intent: MrfIntent) -> bool {
|
||||
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 MrfQueuePushResult::Rejected;
|
||||
return false;
|
||||
}
|
||||
self.bytes += cost;
|
||||
self.pending_keys.insert(key);
|
||||
self.pending.push_back(intent);
|
||||
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)
|
||||
true
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
@@ -219,24 +157,10 @@ impl MrfQueue {
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Append one encoded record to `out`.
|
||||
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]);
|
||||
pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) {
|
||||
let start = out.len();
|
||||
out.push(MRF_JOURNAL_FORMAT);
|
||||
out.push(if scope.is_some() {
|
||||
MRF_JOURNAL_VERSION_SCOPED
|
||||
} else {
|
||||
MRF_JOURNAL_VERSION
|
||||
});
|
||||
out.push(MRF_JOURNAL_VERSION);
|
||||
out.push(match intent.kind {
|
||||
rustfs_common::mrf_channel::MrfKind::DecodeFailure => 1,
|
||||
rustfs_common::mrf_channel::MrfKind::MetadataCorruption => 2,
|
||||
@@ -244,36 +168,27 @@ pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) -> bool {
|
||||
});
|
||||
out.push(intent.attempts);
|
||||
out.extend_from_slice(&intent.enqueued_at_ms.to_le_bytes());
|
||||
match version_id {
|
||||
match intent.version_id {
|
||||
Some(bytes) => {
|
||||
out.push(1);
|
||||
out.extend_from_slice(&bytes);
|
||||
}
|
||||
None => out.push(0),
|
||||
}
|
||||
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.len() as u32).to_le_bytes());
|
||||
out.extend_from_slice(&(intent.object.len() as u32).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..]);
|
||||
let Ok(checksum) = u32::try_from(hasher.finalize()) else {
|
||||
out.truncate(start);
|
||||
return false;
|
||||
};
|
||||
out.extend_from_slice(&checksum.to_le_bytes());
|
||||
true
|
||||
out.extend_from_slice(&(hasher.finalize() as u32).to_le_bytes());
|
||||
}
|
||||
|
||||
fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
|
||||
if data.len() < MRF_RECORD_FIXED_HEAD + 8 {
|
||||
return None;
|
||||
}
|
||||
if data[0] != MRF_JOURNAL_FORMAT || !matches!(data[1], MRF_JOURNAL_VERSION | MRF_JOURNAL_VERSION_SCOPED) {
|
||||
if data[0] != MRF_JOURNAL_FORMAT || data[1] != MRF_JOURNAL_VERSION {
|
||||
return None;
|
||||
}
|
||||
let kind = match data[2] {
|
||||
@@ -283,38 +198,24 @@ 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().ok()?);
|
||||
let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().expect("slice length checked"));
|
||||
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().ok()?;
|
||||
let bytes: [u8; 16] = data[cursor..cursor + 16].try_into().expect("slice length checked");
|
||||
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 = 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;
|
||||
}
|
||||
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;
|
||||
cursor += 8;
|
||||
let body_end = cursor.checked_add(bucket_len)?.checked_add(object_len)?;
|
||||
let record_end = body_end.checked_add(4)?;
|
||||
@@ -323,7 +224,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 u32::try_from(hasher.finalize()).ok()? != u32::from_le_bytes(data[body_end..record_end].try_into().ok()?) {
|
||||
if (hasher.finalize() as u32) != u32::from_le_bytes(data[body_end..record_end].try_into().expect("slice length checked")) {
|
||||
return None;
|
||||
}
|
||||
let bucket = std::sync::Arc::from(std::str::from_utf8(&data[cursor..cursor + bucket_len]).ok()?);
|
||||
@@ -334,12 +235,6 @@ 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,
|
||||
},
|
||||
@@ -374,9 +269,9 @@ async fn journal_disks() -> Vec<DiskStore> {
|
||||
map.values().flatten().cloned().collect()
|
||||
}
|
||||
|
||||
async fn read_journal(path: &str) -> Option<Vec<u8>> {
|
||||
async fn read_journal() -> Option<Vec<u8>> {
|
||||
for disk in journal_disks().await {
|
||||
match disk.read_all(super::RUSTFS_META_BUCKET, path).await {
|
||||
match disk.read_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH).await {
|
||||
Ok(bytes) => return Some(bytes.to_vec()),
|
||||
Err(_) => continue,
|
||||
}
|
||||
@@ -387,51 +282,35 @@ async fn read_journal(path: &str) -> 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(path: &str, data: &[u8]) -> bool {
|
||||
async fn write_journal(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, path, payload.clone()).await {
|
||||
match disk
|
||||
.write_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_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(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
|
||||
async fn delete_journal() {
|
||||
for disk in journal_disks().await {
|
||||
let _ = disk
|
||||
.delete(
|
||||
super::RUSTFS_META_BUCKET,
|
||||
path,
|
||||
MRF_JOURNAL_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) {
|
||||
@@ -452,10 +331,7 @@ 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
|
||||
.filter(|bytes| *bytes != [0; 16])
|
||||
.map(|bytes| Uuid::from_bytes(bytes).to_string());
|
||||
let version_id = intent.version_id.map(|bytes| Uuid::from_bytes(bytes).to_string());
|
||||
let (heal_type, priority) = match intent.kind {
|
||||
rustfs_common::mrf_channel::MrfKind::DecodeFailure => (
|
||||
HealType::ECDecode {
|
||||
@@ -475,28 +351,18 @@ pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest {
|
||||
HealPriority::Normal,
|
||||
),
|
||||
};
|
||||
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);
|
||||
let mut request = HealRequest::new(heal_type, HealOptions::default(), 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_and_identity(
|
||||
.submit_mrf_heal_request_with_receipt(
|
||||
build_heal_request(intent),
|
||||
intent.bucket.clone(),
|
||||
intent.object.clone(),
|
||||
intent.version_id,
|
||||
intent.kind,
|
||||
intent.scope,
|
||||
intent.lease,
|
||||
)
|
||||
.await?;
|
||||
Ok(receipt.result)
|
||||
@@ -521,38 +387,16 @@ struct MrfRuntime {
|
||||
}
|
||||
|
||||
impl MrfRuntime {
|
||||
fn snapshot(&self) -> (Vec<u8>, Vec<u8>) {
|
||||
let mut authoritative = Vec::new();
|
||||
let mut legacy = Vec::new();
|
||||
fn snapshot(&self) -> Vec<u8> {
|
||||
let mut buf = Vec::new();
|
||||
for intent in self.queue.intents() {
|
||||
let scoped_identity =
|
||||
!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption) && intent.scope.is_some();
|
||||
if !encode_intent(intent, &mut authoritative) {
|
||||
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);
|
||||
}
|
||||
encode_intent(intent, &mut buf);
|
||||
}
|
||||
(authoritative, legacy)
|
||||
buf
|
||||
}
|
||||
|
||||
async fn flush(&mut self) {
|
||||
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;
|
||||
let persisted = write_journal(&self.snapshot()).await;
|
||||
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
|
||||
@@ -560,7 +404,7 @@ impl MrfRuntime {
|
||||
if persisted {
|
||||
self.dirty = false;
|
||||
}
|
||||
self.journal_on_disk |= authoritative_persisted || legacy_persisted;
|
||||
self.journal_on_disk = true;
|
||||
}
|
||||
|
||||
/// Drain pending intents into the heal manager until it is full, the
|
||||
@@ -586,7 +430,6 @@ 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);
|
||||
@@ -595,13 +438,11 @@ 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);
|
||||
@@ -610,8 +451,8 @@ impl MrfRuntime {
|
||||
}
|
||||
}
|
||||
}
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(self.queue.depth()));
|
||||
gauge!("rustfs_heal_mrf_queue_bytes").set(metric_f64(self.queue.bytes()));
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(self.queue.depth() as f64);
|
||||
gauge!("rustfs_heal_mrf_queue_bytes").set(self.queue.bytes() as f64);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -655,12 +496,7 @@ 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.replayed
|
||||
}
|
||||
|
||||
struct ReplayOutcome {
|
||||
replayed: usize,
|
||||
journal_on_disk: bool,
|
||||
replay_into(manager, &mut queue, &mut backoff_until).await
|
||||
}
|
||||
|
||||
/// Shared replay core: read + decode + re-arm + delete, then drain what fits.
|
||||
@@ -668,25 +504,11 @@ async fn replay_into(
|
||||
manager: &Arc<HealManager>,
|
||||
queue: &mut MrfQueue,
|
||||
backoff_until: &mut Option<tokio::time::Instant>,
|
||||
) -> 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,
|
||||
};
|
||||
}
|
||||
},
|
||||
) -> usize {
|
||||
let Some(data) = read_journal().await else {
|
||||
return 0;
|
||||
};
|
||||
let mut intents = Vec::new();
|
||||
let (decoded, truncated) = decode_journal(&data);
|
||||
intents.extend(decoded);
|
||||
let (intents, truncated) = decode_journal(&data);
|
||||
if truncated > 0 {
|
||||
tracing::warn!(
|
||||
target: "rustfs::heal::mrf",
|
||||
@@ -694,15 +516,12 @@ async fn replay_into(
|
||||
"MRF journal had a torn tail; truncated records were discarded"
|
||||
);
|
||||
}
|
||||
counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(intents.len()).unwrap_or(u64::MAX));
|
||||
counter!("rustfs_heal_mrf_replayed_total").increment(intents.len() as u64);
|
||||
let replayed = intents.len();
|
||||
for intent in intents {
|
||||
let result = queue.try_push_typed(intent.clone());
|
||||
if !matches!(result, MrfQueuePushResult::Enqueued) {
|
||||
rustfs_common::mrf_channel::release_mrf_intent(&intent);
|
||||
}
|
||||
queue.try_push(intent);
|
||||
}
|
||||
let journal_on_disk = !delete_journals().await;
|
||||
delete_journal().await;
|
||||
|
||||
// Drain the replayed intents immediately; whatever the manager refuses
|
||||
// stays armed in `queue` for the consumer's retry loop.
|
||||
@@ -722,10 +541,7 @@ async fn replay_into(
|
||||
}
|
||||
}
|
||||
}
|
||||
ReplayOutcome {
|
||||
replayed,
|
||||
journal_on_disk,
|
||||
}
|
||||
replayed
|
||||
}
|
||||
|
||||
/// Replay the journal, then keep draining the channel into the heal manager
|
||||
@@ -743,8 +559,7 @@ 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.
|
||||
let replay = replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
|
||||
runtime.journal_on_disk = replay.journal_on_disk;
|
||||
replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
|
||||
// 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.
|
||||
@@ -772,14 +587,9 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
|
||||
return;
|
||||
}
|
||||
for intent in batch.drain(..) {
|
||||
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);
|
||||
}
|
||||
if runtime.queue.try_push(intent) {
|
||||
runtime.new_since_flush += 1;
|
||||
runtime.dirty = true;
|
||||
}
|
||||
}
|
||||
runtime.dispatch(manager.as_ref()).await;
|
||||
@@ -803,14 +613,13 @@ 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).
|
||||
if delete_journals().await {
|
||||
runtime.journal_on_disk = false;
|
||||
gauge!("rustfs_heal_mrf_journal_bytes").set(0.0);
|
||||
}
|
||||
delete_journal().await;
|
||||
runtime.journal_on_disk = false;
|
||||
gauge!("rustfs_heal_mrf_journal_bytes").set(0.0);
|
||||
}
|
||||
TickAction::Idle => {}
|
||||
}
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(runtime.queue.depth()));
|
||||
gauge!("rustfs_heal_mrf_queue_depth").set(runtime.queue.depth() as f64);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -855,8 +664,6 @@ 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,
|
||||
}
|
||||
@@ -887,101 +694,17 @@ 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", "o2", 0)));
|
||||
assert!(!queue.try_push(intent("b", "o3", 0)), "count ceiling must drop");
|
||||
assert!(queue.try_push(intent("b", "o", 0)));
|
||||
assert!(!queue.try_push(intent("b", "o", 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", "object2", 0)),
|
||||
!tiny.try_push(intent("bucket", "object", 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![
|
||||
@@ -992,8 +715,6 @@ mod tests {
|
||||
object: StdArc::from("object/c"),
|
||||
version_id: None,
|
||||
kind: MrfKind::MetadataCorruption,
|
||||
scope: None,
|
||||
lease: None,
|
||||
enqueued_at_ms: 5,
|
||||
attempts: 1,
|
||||
},
|
||||
@@ -1045,8 +766,6 @@ mod tests {
|
||||
object: StdArc::from("o"),
|
||||
version_id: None,
|
||||
kind: MrfKind::MetadataCorruption,
|
||||
scope: None,
|
||||
lease: None,
|
||||
enqueued_at_ms: 0,
|
||||
attempts: 0,
|
||||
});
|
||||
@@ -1058,8 +777,6 @@ mod tests {
|
||||
object: StdArc::from("o"),
|
||||
version_id: None,
|
||||
kind: MrfKind::PartialWrite,
|
||||
scope: None,
|
||||
lease: None,
|
||||
enqueued_at_ms: 0,
|
||||
attempts: 0,
|
||||
});
|
||||
|
||||
@@ -35,7 +35,6 @@ use storage_api::endpoint_index::{Endpoint, EndpointServerPools, Endpoints, Pool
|
||||
|
||||
const META_BUCKET: &str = ".rustfs.sys";
|
||||
const JOURNAL_REL: &str = "buckets/.heal/mrf/journal.bin";
|
||||
const SCOPED_JOURNAL_REL: &str = "buckets/.heal/mrf/journal-scoped.bin";
|
||||
|
||||
async fn heal_env() -> (Vec<std::path::PathBuf>, Arc<dyn HealStorageAPI>) {
|
||||
let env = rustfs_test_utils::TestECStoreEnv::builder()
|
||||
@@ -80,18 +79,14 @@ fn journal_record(kind: u8, bucket: &str, object: &str, version: Option<[u8; 16]
|
||||
body
|
||||
}
|
||||
|
||||
fn write_journal_path_to_disks(disk_paths: &[std::path::PathBuf], relative_path: &str, data: &[u8]) {
|
||||
fn write_journal_to_disks(disk_paths: &[std::path::PathBuf], data: &[u8]) {
|
||||
for path in disk_paths {
|
||||
let journal = path.join(META_BUCKET).join(relative_path);
|
||||
let journal = path.join(META_BUCKET).join(JOURNAL_REL);
|
||||
std::fs::create_dir_all(journal.parent().expect("journal parent")).expect("create journal dir");
|
||||
std::fs::write(&journal, data).expect("write journal fixture");
|
||||
}
|
||||
}
|
||||
|
||||
fn write_journal_to_disks(disk_paths: &[std::path::PathBuf], data: &[u8]) {
|
||||
write_journal_path_to_disks(disk_paths, JOURNAL_REL, data);
|
||||
}
|
||||
|
||||
async fn wait_until<F, Fut>(deadline: Duration, mut probe: F) -> bool
|
||||
where
|
||||
F: FnMut() -> Fut,
|
||||
@@ -192,73 +187,8 @@ async fn journal_replay_arms_intents_and_deletes_the_file() {
|
||||
.all(|path| !Path::new(path).join(META_BUCKET).join(JOURNAL_REL).exists()),
|
||||
"the journal file must be removed after a successful replay"
|
||||
);
|
||||
assert!(
|
||||
disk_paths
|
||||
.iter()
|
||||
.all(|path| !Path::new(path).join(META_BUCKET).join(SCOPED_JOURNAL_REL).exists()),
|
||||
"the authoritative journal file must also be removed after replay"
|
||||
);
|
||||
|
||||
let snapshot = manager.operations_snapshot().await;
|
||||
assert_eq!(snapshot.queued_by_priority.urgent, 1, "the decode-failure record must replay as Urgent");
|
||||
assert!(snapshot.queued_by_priority.normal >= 1, "the partial-write record must replay as Normal");
|
||||
}
|
||||
|
||||
/// A canonical snapshot and its compatibility mirror may differ after a
|
||||
/// partial flush. Replay must choose the complete canonical epoch instead of
|
||||
/// combining records that never coexisted in memory.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
#[serial]
|
||||
async fn authoritative_journal_is_not_merged_with_legacy_mirror() {
|
||||
let (disk_paths, storage) = heal_env().await;
|
||||
let mut endpoints: Vec<Endpoint> = disk_paths
|
||||
.iter()
|
||||
.map(|p| Endpoint::try_from(p.to_string_lossy().as_ref()).expect("endpoint from disk path"))
|
||||
.collect();
|
||||
for (i, endpoint) in endpoints.iter_mut().enumerate() {
|
||||
endpoint.set_pool_index(0);
|
||||
endpoint.set_set_index(0);
|
||||
endpoint.set_disk_index(i);
|
||||
}
|
||||
let pool = PoolEndpoints {
|
||||
legacy: false,
|
||||
set_count: 1,
|
||||
drives_per_set: endpoints.len(),
|
||||
endpoints: Endpoints::from(endpoints),
|
||||
cmd_line: "mrf-authoritative-test".to_string(),
|
||||
platform: String::new(),
|
||||
};
|
||||
init_local_disks(EndpointServerPools::from(vec![pool]))
|
||||
.await
|
||||
.expect("local disks should register");
|
||||
|
||||
let authoritative = journal_record(1, "authoritative-bucket", "authoritative-object", None, 0);
|
||||
let legacy = journal_record(1, "legacy-bucket", "legacy-object", None, 0);
|
||||
write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &authoritative);
|
||||
write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &legacy);
|
||||
|
||||
let manager = make_manager(storage);
|
||||
let replayed = mrf_queue::replay_journal_once(&manager).await;
|
||||
assert_eq!(replayed, 1, "only the authoritative snapshot epoch may replay");
|
||||
|
||||
let snapshot = manager.operations_snapshot().await;
|
||||
assert_eq!(snapshot.queued_by_source.mrf, 1);
|
||||
assert!(
|
||||
disk_paths.iter().all(|path| {
|
||||
!Path::new(path).join(META_BUCKET).join(JOURNAL_REL).exists()
|
||||
&& !Path::new(path).join(META_BUCKET).join(SCOPED_JOURNAL_REL).exists()
|
||||
}),
|
||||
"replay cleanup must remove both journal paths"
|
||||
);
|
||||
|
||||
// A scoped-only snapshot is valid during a rollout where no legacy
|
||||
// compatibility mirror was written. Missing legacy files must not leave
|
||||
// the runtime in a permanent cleanup-retry state.
|
||||
let scoped_only = journal_record(1, "scoped-only-bucket", "scoped-only-object", None, 0);
|
||||
write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &scoped_only);
|
||||
assert_eq!(mrf_queue::replay_journal_once(&manager).await, 1);
|
||||
assert!(disk_paths.iter().all(|path| {
|
||||
!Path::new(path).join(META_BUCKET).join(JOURNAL_REL).exists()
|
||||
&& !Path::new(path).join(META_BUCKET).join(SCOPED_JOURNAL_REL).exists()
|
||||
}));
|
||||
}
|
||||
|
||||
@@ -28,8 +28,9 @@ use rustfs_common::heal_channel::HealScanMode;
|
||||
use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS;
|
||||
pub use rustfs_data_usage::{
|
||||
AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME,
|
||||
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry,
|
||||
PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path, prefix_usage_in_cache,
|
||||
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME,
|
||||
PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path,
|
||||
prefix_usage_in_cache,
|
||||
};
|
||||
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
|
||||
use tokio::time::{Duration, Instant, sleep, timeout};
|
||||
@@ -344,6 +345,18 @@ pub struct DataUsageCacheInfo {
|
||||
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
|
||||
#[serde(default)]
|
||||
pub cache_key_format: u16,
|
||||
/// Whether the entries retained while a set scan was incomplete come
|
||||
/// from a prior complete set snapshot. This is observational input only.
|
||||
#[serde(default)]
|
||||
pub lkg_snapshot_complete: bool,
|
||||
#[serde(default)]
|
||||
pub lkg_next_cycle: Option<u64>,
|
||||
#[serde(default)]
|
||||
pub lkg_last_update: Option<SystemTime>,
|
||||
#[serde(default)]
|
||||
pub lkg_leader_epoch: Option<u64>,
|
||||
#[serde(default)]
|
||||
pub lkg_scan_plan_digest: Option<DataUsageScanPlanDigest>,
|
||||
}
|
||||
|
||||
impl Serialize for DataUsageCacheInfo {
|
||||
@@ -353,7 +366,7 @@ impl Serialize for DataUsageCacheInfo {
|
||||
{
|
||||
// Keep this metadata map-encoded so older readers can ignore fields
|
||||
// appended by newer scanner versions during rolling upgrades.
|
||||
let mut state = serializer.serialize_map(Some(16))?;
|
||||
let mut state = serializer.serialize_map(Some(21))?;
|
||||
state.serialize_entry("name", &self.name)?;
|
||||
state.serialize_entry("next_cycle", &self.next_cycle)?;
|
||||
state.serialize_entry("leader_epoch", &self.leader_epoch)?;
|
||||
@@ -370,6 +383,11 @@ impl Serialize for DataUsageCacheInfo {
|
||||
state.serialize_entry("snapshot_complete", &self.snapshot_complete)?;
|
||||
state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?;
|
||||
state.serialize_entry("cache_key_format", &self.cache_key_format)?;
|
||||
state.serialize_entry("lkg_snapshot_complete", &self.lkg_snapshot_complete)?;
|
||||
state.serialize_entry("lkg_next_cycle", &self.lkg_next_cycle)?;
|
||||
state.serialize_entry("lkg_last_update", &self.lkg_last_update)?;
|
||||
state.serialize_entry("lkg_leader_epoch", &self.lkg_leader_epoch)?;
|
||||
state.serialize_entry("lkg_scan_plan_digest", &self.lkg_scan_plan_digest)?;
|
||||
state.end()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2274,9 +2274,8 @@ async fn final_data_usage_publication_defer_reason(
|
||||
}
|
||||
}
|
||||
ScannerCycleStatus::Deferred(reason) => Some(reason),
|
||||
// Incomplete cycles do not publish a usage snapshot. Keep the
|
||||
// decision permissive so existing partial-cycle handling remains
|
||||
// unchanged if a future scanner path emits a bookkeeping update.
|
||||
// Incomplete cycles may publish a non-authoritative observational
|
||||
// snapshot when at least one set has a usable current/LKG view.
|
||||
ScannerCycleStatus::Incomplete => None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -198,7 +198,7 @@ where
|
||||
data_usage_info.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
|
||||
}
|
||||
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() {
|
||||
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
|
||||
error!(
|
||||
target: "rustfs::scanner",
|
||||
event = EVENT_SCANNER_PERSIST_STATE,
|
||||
|
||||
@@ -1375,14 +1375,13 @@ impl FolderScanner {
|
||||
// Single-flight (backlog#1894 axis A) — the
|
||||
// recording mode and its guarantees are pinned by
|
||||
// corrupt_metadata_recording below.
|
||||
let mrf_result = rustfs_common::mrf_channel::try_send_mrf_intent_typed(
|
||||
let mrf_accepted = rustfs_common::mrf_channel::try_send_mrf_intent(
|
||||
rustfs_common::mrf_channel::MrfKind::MetadataCorruption,
|
||||
&item.bucket,
|
||||
&object,
|
||||
None,
|
||||
None,
|
||||
);
|
||||
match corrupt_metadata_recording(mrf_result) {
|
||||
match corrupt_metadata_recording(mrf_accepted) {
|
||||
CorruptMetadataRecording::LedgerOnly => {
|
||||
// Recorded as Full (retry-later): admission
|
||||
// for this target happens in the MRF
|
||||
|
||||
@@ -50,12 +50,11 @@ pub(super) enum CorruptMetadataRecording {
|
||||
ImmediateAndLedger,
|
||||
}
|
||||
|
||||
pub(super) fn corrupt_metadata_recording(result: rustfs_common::mrf_channel::MrfIngressResult) -> CorruptMetadataRecording {
|
||||
match result {
|
||||
rustfs_common::mrf_channel::MrfIngressResult::Enqueued | rustfs_common::mrf_channel::MrfIngressResult::Coalesced => {
|
||||
CorruptMetadataRecording::LedgerOnly
|
||||
}
|
||||
rustfs_common::mrf_channel::MrfIngressResult::Dropped(_) => CorruptMetadataRecording::ImmediateAndLedger,
|
||||
pub(super) fn corrupt_metadata_recording(mrf_accepted: bool) -> CorruptMetadataRecording {
|
||||
if mrf_accepted {
|
||||
CorruptMetadataRecording::LedgerOnly
|
||||
} else {
|
||||
CorruptMetadataRecording::ImmediateAndLedger
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -48,20 +48,8 @@ fn scanner_alert_wire_names_match_canonical_event_names() {
|
||||
/// the backstop survives regardless of delivery.
|
||||
#[test]
|
||||
fn corrupt_metadata_recording_maps_delivery_to_backstop() {
|
||||
assert_eq!(
|
||||
corrupt_metadata_recording(rustfs_common::mrf_channel::MrfIngressResult::Enqueued),
|
||||
CorruptMetadataRecording::LedgerOnly
|
||||
);
|
||||
assert_eq!(
|
||||
corrupt_metadata_recording(rustfs_common::mrf_channel::MrfIngressResult::Coalesced),
|
||||
CorruptMetadataRecording::LedgerOnly
|
||||
);
|
||||
assert_eq!(
|
||||
corrupt_metadata_recording(rustfs_common::mrf_channel::MrfIngressResult::Dropped(
|
||||
rustfs_common::mrf_channel::MrfDropReason::Full
|
||||
)),
|
||||
CorruptMetadataRecording::ImmediateAndLedger
|
||||
);
|
||||
assert_eq!(corrupt_metadata_recording(true), CorruptMetadataRecording::LedgerOnly);
|
||||
assert_eq!(corrupt_metadata_recording(false), CorruptMetadataRecording::ImmediateAndLedger);
|
||||
}
|
||||
|
||||
fn cooldown_map_len() -> usize {
|
||||
|
||||
@@ -18,8 +18,8 @@ use crate::scanner_folder::{ScannerItem, scan_data_folder};
|
||||
use crate::sleeper::SCANNER_SLEEPER;
|
||||
use crate::{
|
||||
DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT, DataUsageCache, DataUsageCacheInfo, DataUsageCachePrepareOutcome,
|
||||
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, ScannerError, SizeSummary,
|
||||
TierStats,
|
||||
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState,
|
||||
ScannerError, SizeSummary, TierStats,
|
||||
};
|
||||
use futures::future::join_all;
|
||||
use metrics::counter;
|
||||
@@ -278,6 +278,17 @@ async fn publish_usage_snapshot(
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn publish_observational_snapshot(
|
||||
updates: &mpsc::Sender<DataUsageInfo>,
|
||||
mut data_usage_info: DataUsageInfo,
|
||||
) -> Result<bool> {
|
||||
data_usage_info.usage_snapshot_complete = false;
|
||||
data_usage_info.usage_snapshot_partial = true;
|
||||
data_usage_info.usage_snapshot_converged = Some(false);
|
||||
send_data_usage_update(updates, data_usage_info).await?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
enum ScannerCycleActivityStatus {
|
||||
Unchanged,
|
||||
|
||||
@@ -188,7 +188,7 @@ pub(super) fn completed_data_usage_info(
|
||||
}
|
||||
|
||||
let mut total = DataUsageEntry::default();
|
||||
let mut buckets_usage = HashMap::with_capacity(all_buckets.len());
|
||||
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
|
||||
for bucket in all_buckets {
|
||||
let mut merged = DataUsageEntry::default();
|
||||
for result in results {
|
||||
@@ -200,10 +200,14 @@ pub(super) fn completed_data_usage_info(
|
||||
if !total.checked_merge(&merged) {
|
||||
return None;
|
||||
}
|
||||
buckets_usage.insert(bucket.clone(), checked_bucket_usage_info(&merged)?);
|
||||
bucket_entries.insert(bucket.clone(), merged);
|
||||
}
|
||||
|
||||
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
|
||||
let buckets_usage = bucket_entries
|
||||
.iter()
|
||||
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
|
||||
.collect::<Option<HashMap<_, _>>>()?;
|
||||
let bucket_sizes = buckets_usage
|
||||
.iter()
|
||||
.map(|(bucket, usage)| (bucket.clone(), usage.size))
|
||||
@@ -225,6 +229,145 @@ pub(super) fn completed_data_usage_info(
|
||||
Some((data_usage_info, merged_last_update))
|
||||
}
|
||||
|
||||
/// Build a non-authoritative view from the set snapshots that completed this
|
||||
/// cycle plus compatible per-set last-known-good caches. The caller must
|
||||
/// persist this result only on the observational object; a missing set is
|
||||
/// intentionally represented by an incomplete state and is never treated as
|
||||
/// an empty set.
|
||||
pub(super) fn observational_data_usage_info(
|
||||
results: &[DataUsageCache],
|
||||
expected_sources: &HashSet<DataUsageCacheSource>,
|
||||
all_buckets: &[String],
|
||||
expected_plan_digest: DataUsageScanPlanDigest,
|
||||
scanner_cycle: u64,
|
||||
leader_epoch: u64,
|
||||
) -> Option<(DataUsageInfo, SystemTime)> {
|
||||
let mut by_source = HashMap::with_capacity(results.len());
|
||||
for result in results {
|
||||
let source = result.info.source?;
|
||||
if !expected_sources.contains(&source) || by_source.insert(source, result).is_some() {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
|
||||
let mut usable = Vec::new();
|
||||
let mut set_states = Vec::with_capacity(expected_sources.len());
|
||||
let mut sources = expected_sources.iter().copied().collect::<Vec<_>>();
|
||||
sources.sort_by_key(|source| (source.pool_index, source.set_index));
|
||||
for source in sources {
|
||||
let result = by_source.get(&source).copied();
|
||||
let current = result.filter(|result| {
|
||||
result.info.snapshot_complete
|
||||
&& result.info.next_cycle == scanner_cycle
|
||||
&& result.info.leader_epoch == leader_epoch
|
||||
&& result.info.scan_plan_digest == Some(expected_plan_digest)
|
||||
});
|
||||
let lkg = result.filter(|result| {
|
||||
!result.info.snapshot_complete
|
||||
&& result.info.lkg_snapshot_complete
|
||||
&& result.info.lkg_scan_plan_digest == Some(expected_plan_digest)
|
||||
&& result.info.lkg_leader_epoch.is_some_and(|epoch| {
|
||||
epoch < leader_epoch
|
||||
|| (epoch == leader_epoch && result.info.lkg_next_cycle.is_some_and(|cycle| cycle <= scanner_cycle))
|
||||
})
|
||||
});
|
||||
let current_snapshot = current.is_some();
|
||||
let selected = current.or(lkg);
|
||||
if let Some(selected) = selected {
|
||||
let (cycle, epoch, digest, last_update, complete) = if current_snapshot {
|
||||
(
|
||||
Some(selected.info.next_cycle),
|
||||
Some(selected.info.leader_epoch),
|
||||
selected.info.scan_plan_digest.map(|digest| digest.0),
|
||||
selected.info.last_update,
|
||||
true,
|
||||
)
|
||||
} else {
|
||||
(
|
||||
selected.info.lkg_next_cycle,
|
||||
selected.info.lkg_leader_epoch,
|
||||
selected.info.lkg_scan_plan_digest.map(|digest| digest.0),
|
||||
selected.info.lkg_last_update,
|
||||
false,
|
||||
)
|
||||
};
|
||||
set_states.push(DataUsageSnapshotSetState {
|
||||
pool_index: u64::try_from(source.pool_index).ok()?,
|
||||
set_index: u64::try_from(source.set_index).ok()?,
|
||||
scanner_cycle: cycle,
|
||||
scanner_epoch: epoch,
|
||||
scan_plan_digest: digest,
|
||||
complete,
|
||||
tombstone: false,
|
||||
});
|
||||
usable.push((selected, last_update));
|
||||
} else {
|
||||
set_states.push(DataUsageSnapshotSetState {
|
||||
pool_index: u64::try_from(source.pool_index).ok()?,
|
||||
set_index: u64::try_from(source.set_index).ok()?,
|
||||
scanner_cycle: None,
|
||||
scanner_epoch: None,
|
||||
scan_plan_digest: Some(expected_plan_digest.0),
|
||||
complete: false,
|
||||
tombstone: false,
|
||||
});
|
||||
}
|
||||
}
|
||||
if usable.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let mut total = DataUsageEntry::default();
|
||||
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
|
||||
let mut merged_last_update = None;
|
||||
for (result, last_update) in usable {
|
||||
if let Some(update) = last_update {
|
||||
merged_last_update = Some(merged_last_update.map_or(update, |current: SystemTime| current.max(update)));
|
||||
}
|
||||
for bucket in all_buckets {
|
||||
let Some(entry) = result.checked_flatten(bucket) else {
|
||||
continue;
|
||||
};
|
||||
let bucket_entry = bucket_entries.entry(bucket.clone()).or_insert_with(DataUsageEntry::default);
|
||||
if !bucket_entry.checked_merge(&entry) {
|
||||
return None;
|
||||
}
|
||||
if !total.checked_merge(&entry) {
|
||||
return None;
|
||||
}
|
||||
}
|
||||
}
|
||||
let merged_last_update = merged_last_update?;
|
||||
let buckets_usage = bucket_entries
|
||||
.iter()
|
||||
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
|
||||
.collect::<Option<HashMap<_, _>>>()?;
|
||||
Some((
|
||||
DataUsageInfo {
|
||||
last_update: Some(merged_last_update),
|
||||
scanner_cycle: Some(scanner_cycle),
|
||||
scanner_epoch: Some(leader_epoch),
|
||||
objects_total_count: u64::try_from(total.objects).ok()?,
|
||||
versions_total_count: u64::try_from(total.versions).ok()?,
|
||||
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
|
||||
objects_total_size: u64::try_from(total.size).ok()?,
|
||||
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
|
||||
buckets_count: u64::try_from(buckets_usage.len()).ok()?,
|
||||
bucket_sizes: buckets_usage
|
||||
.iter()
|
||||
.map(|(bucket, usage)| (bucket.clone(), usage.size))
|
||||
.collect(),
|
||||
buckets_usage,
|
||||
usage_snapshot_complete: false,
|
||||
usage_snapshot_partial: true,
|
||||
usage_snapshot_converged: Some(false),
|
||||
usage_snapshot_set_states: set_states,
|
||||
..Default::default()
|
||||
},
|
||||
merged_last_update,
|
||||
))
|
||||
}
|
||||
|
||||
pub(super) async fn send_cache_root_entry_info(
|
||||
bucket_result_tx: &mpsc::Sender<DataUsageEntryInfo>,
|
||||
cache: &DataUsageCache,
|
||||
|
||||
@@ -40,6 +40,21 @@ impl ScannerIOCache for SetDisks {
|
||||
let set_label = self.set_index.to_string();
|
||||
|
||||
let source = DataUsageCacheSource::new(self.pool_index, self.set_index);
|
||||
let mut old_cache = DataUsageCache::default();
|
||||
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
|
||||
warn!(
|
||||
target: "rustfs::scanner::io",
|
||||
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_IO,
|
||||
pool = self.pool_index,
|
||||
set = self.set_index,
|
||||
cache_name = DATA_USAGE_CACHE_NAME,
|
||||
state = "old_cache_load_failed",
|
||||
error = %e,
|
||||
"Scanner old data usage cache load failed; rebuilding from bucket caches"
|
||||
);
|
||||
}
|
||||
if buckets.is_empty() {
|
||||
let now = SystemTime::now();
|
||||
let mut cache = DataUsageCache {
|
||||
@@ -80,6 +95,24 @@ impl ScannerIOCache for SetDisks {
|
||||
"Scanner set state found no online disks"
|
||||
);
|
||||
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
|
||||
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
|
||||
let mut incomplete_scope = lkg.clone().unwrap_or_default();
|
||||
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
|
||||
incomplete_scope.info.next_cycle = want_cycle;
|
||||
incomplete_scope.info.last_update = None;
|
||||
incomplete_scope.info.leader_epoch = leader_epoch;
|
||||
incomplete_scope.info.source = Some(source);
|
||||
incomplete_scope.info.snapshot_complete = false;
|
||||
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
|
||||
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
|
||||
if let Some(lkg) = lkg {
|
||||
incomplete_scope.info.lkg_snapshot_complete = true;
|
||||
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
|
||||
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
|
||||
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
|
||||
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
|
||||
}
|
||||
let _ = updates.send(incomplete_scope).await;
|
||||
return Ok(());
|
||||
}
|
||||
// Preserve the original set topology across capability filtering. During
|
||||
@@ -162,6 +195,24 @@ impl ScannerIOCache for SetDisks {
|
||||
"Scanner set state found no usable namespace scanner disks"
|
||||
);
|
||||
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
|
||||
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
|
||||
let mut incomplete_scope = lkg.clone().unwrap_or_default();
|
||||
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
|
||||
incomplete_scope.info.next_cycle = want_cycle;
|
||||
incomplete_scope.info.last_update = None;
|
||||
incomplete_scope.info.leader_epoch = leader_epoch;
|
||||
incomplete_scope.info.source = Some(source);
|
||||
incomplete_scope.info.snapshot_complete = false;
|
||||
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
|
||||
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
|
||||
if let Some(lkg) = lkg {
|
||||
incomplete_scope.info.lkg_snapshot_complete = true;
|
||||
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
|
||||
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
|
||||
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
|
||||
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
|
||||
}
|
||||
let _ = updates.send(incomplete_scope).await;
|
||||
return Ok(());
|
||||
}
|
||||
let set_disk_inventory = Arc::new(scanner_set_disk_inventory(self.as_ref()).await);
|
||||
@@ -203,22 +254,15 @@ impl ScannerIOCache for SetDisks {
|
||||
record_disk_bucket_scans_active(0, &pool_label, &set_label);
|
||||
let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone());
|
||||
|
||||
let mut old_cache = DataUsageCache::default();
|
||||
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
|
||||
warn!(
|
||||
target: "rustfs::scanner::io",
|
||||
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
|
||||
component = LOG_COMPONENT_SCANNER,
|
||||
subsystem = LOG_SUBSYSTEM_IO,
|
||||
pool = self.pool_index,
|
||||
set = self.set_index,
|
||||
cache_name = DATA_USAGE_CACHE_NAME,
|
||||
state = "old_cache_load_failed",
|
||||
error = %e,
|
||||
"Scanner old data usage cache load failed; rebuilding from bucket caches"
|
||||
);
|
||||
}
|
||||
match old_cache.prepare_for_scan(
|
||||
let old_lkg = old_cache.info.snapshot_complete.then(|| {
|
||||
(
|
||||
old_cache.info.next_cycle,
|
||||
old_cache.info.last_update,
|
||||
old_cache.info.leader_epoch,
|
||||
old_cache.info.scan_plan_digest,
|
||||
)
|
||||
});
|
||||
let prepare_outcome = match old_cache.prepare_for_scan(
|
||||
DATA_USAGE_ROOT,
|
||||
want_cycle,
|
||||
leader_epoch,
|
||||
@@ -259,7 +303,16 @@ impl ScannerIOCache for SetDisks {
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
DataUsageCachePrepareOutcome::Reused | DataUsageCachePrepareOutcome::Reset => {}
|
||||
outcome => outcome,
|
||||
};
|
||||
if matches!(prepare_outcome, DataUsageCachePrepareOutcome::Reused)
|
||||
&& let Some((cycle, last_update, epoch, digest)) = old_lkg
|
||||
{
|
||||
old_cache.info.lkg_snapshot_complete = true;
|
||||
old_cache.info.lkg_next_cycle = Some(cycle);
|
||||
old_cache.info.lkg_last_update = last_update;
|
||||
old_cache.info.lkg_leader_epoch = Some(epoch);
|
||||
old_cache.info.lkg_scan_plan_digest = digest;
|
||||
}
|
||||
|
||||
let mut cache = DataUsageCache {
|
||||
@@ -1099,23 +1152,29 @@ impl ScannerIOCache for SetDisks {
|
||||
cache.info.next_cycle = want_cycle;
|
||||
cache.info.last_update.get_or_insert_with(SystemTime::now);
|
||||
cache.info.snapshot_complete = true;
|
||||
cache.info.lkg_snapshot_complete = false;
|
||||
cache.info.lkg_next_cycle = None;
|
||||
cache.info.lkg_last_update = None;
|
||||
cache.info.lkg_leader_epoch = None;
|
||||
cache.info.lkg_scan_plan_digest = None;
|
||||
cache.clone()
|
||||
};
|
||||
let _ = persist_and_publish_cache_snapshot(self.clone(), &updates, cache_snapshot, cache_cycle_floor.as_ref()).await;
|
||||
} else {
|
||||
let incomplete_scope = DataUsageCache {
|
||||
info: DataUsageCacheInfo {
|
||||
name: DATA_USAGE_ROOT.to_string(),
|
||||
next_cycle: want_cycle,
|
||||
leader_epoch,
|
||||
source: Some(source),
|
||||
snapshot_complete: false,
|
||||
scan_plan_digest: Some(scan_plan_digest),
|
||||
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
|
||||
..Default::default()
|
||||
},
|
||||
cache: HashMap::new(),
|
||||
};
|
||||
let mut incomplete_scope = cache_mutex.lock().await.clone();
|
||||
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
|
||||
incomplete_scope.info.next_cycle = want_cycle;
|
||||
incomplete_scope.info.last_update = None;
|
||||
incomplete_scope.info.leader_epoch = leader_epoch;
|
||||
incomplete_scope.info.source = Some(source);
|
||||
incomplete_scope.info.snapshot_complete = false;
|
||||
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
|
||||
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
|
||||
incomplete_scope.info.lkg_snapshot_complete = old_cache.info.lkg_snapshot_complete;
|
||||
incomplete_scope.info.lkg_next_cycle = old_cache.info.lkg_next_cycle;
|
||||
incomplete_scope.info.lkg_last_update = old_cache.info.lkg_last_update;
|
||||
incomplete_scope.info.lkg_leader_epoch = old_cache.info.lkg_leader_epoch;
|
||||
incomplete_scope.info.lkg_scan_plan_digest = old_cache.info.lkg_scan_plan_digest;
|
||||
if let Err(e) = updates.send(incomplete_scope).await {
|
||||
error!(
|
||||
target: "rustfs::scanner::io",
|
||||
|
||||
@@ -234,6 +234,7 @@ impl ScannerIOCycle for ECStore {
|
||||
let active_set_scans_clone = active_set_scans.clone();
|
||||
|
||||
let (tx, mut rx) = mpsc::channel::<DataUsageCache>(1);
|
||||
let failed_scope_tx = tx.clone();
|
||||
|
||||
// Spawn task to receive and store results
|
||||
let receiver_fut = tokio::spawn(async move {
|
||||
@@ -314,6 +315,21 @@ impl ScannerIOCycle for ECStore {
|
||||
state = "set_scan_failed",
|
||||
"Scanner set scan failed; continuing cycle"
|
||||
);
|
||||
let _ = failed_scope_tx
|
||||
.send(DataUsageCache {
|
||||
info: DataUsageCacheInfo {
|
||||
name: DATA_USAGE_ROOT.to_string(),
|
||||
next_cycle: want_cycle_clone,
|
||||
leader_epoch,
|
||||
source: Some(source),
|
||||
snapshot_complete: false,
|
||||
scan_plan_digest: Some(scan_plan_digest),
|
||||
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
|
||||
..Default::default()
|
||||
},
|
||||
cache: HashMap::new(),
|
||||
})
|
||||
.await;
|
||||
let mut first_err = first_err_mutex_clone.lock().await;
|
||||
record_set_scan_failure(&mut first_err, e);
|
||||
}
|
||||
@@ -370,6 +386,19 @@ impl ScannerIOCycle for ECStore {
|
||||
budget_elapsed,
|
||||
ctx.is_cancelled(),
|
||||
);
|
||||
let observational_usage = completed_usage
|
||||
.is_none()
|
||||
.then(|| {
|
||||
observational_data_usage_info(
|
||||
&results,
|
||||
&expected_sources,
|
||||
&all_bucket_names,
|
||||
scan_plan_digest,
|
||||
want_cycle,
|
||||
leader_epoch,
|
||||
)
|
||||
})
|
||||
.flatten();
|
||||
let structurally_complete_snapshot = result.is_ok() && completed_all_sets && completed_usage.is_some();
|
||||
let cycle_status = classify_nsscanner_cycle(
|
||||
structurally_complete_snapshot,
|
||||
@@ -381,6 +410,10 @@ impl ScannerIOCycle for ECStore {
|
||||
);
|
||||
if let Some((data_usage_info, _)) = completed_usage {
|
||||
publish_usage_snapshot(&updates, cycle_status, data_usage_info).await?;
|
||||
} else if !ctx.is_cancelled()
|
||||
&& let Some((data_usage_info, _)) = observational_usage
|
||||
{
|
||||
publish_observational_snapshot(&updates, data_usage_info).await?;
|
||||
}
|
||||
let dirty_usage_clear = should_clear_dirty_usage_snapshot(
|
||||
result.is_ok(),
|
||||
|
||||
@@ -105,6 +105,160 @@ fn completed_data_usage_info_for_test(
|
||||
completed_data_usage_info(results, &expected_sources, all_buckets, true, budget_elapsed, cancelled)
|
||||
}
|
||||
|
||||
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
|
||||
let mut cache = completed_root_cache(bucket, objects, 10, source);
|
||||
cache.info.snapshot_complete = false;
|
||||
cache.info.next_cycle = 8;
|
||||
cache.info.leader_epoch = 3;
|
||||
cache.info.lkg_snapshot_complete = true;
|
||||
cache.info.lkg_next_cycle = Some(7);
|
||||
cache.info.lkg_last_update = cache.info.last_update;
|
||||
cache.info.lkg_leader_epoch = Some(3);
|
||||
cache.info.lkg_scan_plan_digest = Some(TEST_PLAN_DIGEST);
|
||||
cache
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn partial_usage_is_observational_not_authoritative_for_quota() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let current_source = DataUsageCacheSource::new(0, 0);
|
||||
let stalled_source = DataUsageCacheSource::new(1, 0);
|
||||
let mut current = completed_root_cache("bucket", 2, 20, current_source);
|
||||
current.info.next_cycle = 8;
|
||||
current.info.leader_epoch = 3;
|
||||
let stalled = lkg_root_cache("bucket", 1, stalled_source);
|
||||
let expected = HashSet::from([current_source, stalled_source]);
|
||||
|
||||
assert!(
|
||||
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, true, false, false).is_none()
|
||||
);
|
||||
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("a completed set should produce an observational view");
|
||||
assert!(observed.usage_snapshot_partial);
|
||||
assert!(!observed.usage_snapshot_complete);
|
||||
assert_eq!(observed.usage_snapshot_converged, Some(false));
|
||||
assert_eq!(observed.usage_snapshot_set_states.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lkg_scope_does_not_count_as_current_cycle_completion() {
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut lkg = lkg_root_cache("bucket", 1, source);
|
||||
lkg.info.last_update = None;
|
||||
let expected = HashSet::from([source]);
|
||||
assert!(!scanner_results_form_complete_snapshot(&[lkg], &expected));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut current = completed_root_cache("bucket", 3, 20, source);
|
||||
current.info.next_cycle = 8;
|
||||
current.info.leader_epoch = 3;
|
||||
let expected = HashSet::from([source]);
|
||||
let (observed, _) = observational_data_usage_info(&[current], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("complete set data is a valid observational baseline");
|
||||
assert_eq!(observed.objects_total_size, 30);
|
||||
assert_eq!(observed.usage_snapshot_set_states[0].complete, true);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn negative_delta_waits_for_set_reconciliation() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut stalled = lkg_root_cache("bucket", 4, source);
|
||||
stalled.info.lkg_scan_plan_digest = Some(DataUsageScanPlanDigest([9; 32]));
|
||||
let expected = HashSet::from([source]);
|
||||
assert!(observational_data_usage_info(&[stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn set_membership_add_remove_uses_generation_and_tombstone() {
|
||||
let state = DataUsageSnapshotSetState {
|
||||
pool_index: 1,
|
||||
set_index: 2,
|
||||
scanner_cycle: Some(9),
|
||||
scanner_epoch: Some(4),
|
||||
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
|
||||
complete: false,
|
||||
tombstone: true,
|
||||
};
|
||||
let encoded = serde_json::to_vec(&state).expect("set state should serialize");
|
||||
let decoded: DataUsageSnapshotSetState = serde_json::from_slice(&encoded).expect("set state should deserialize");
|
||||
assert_eq!(decoded, state);
|
||||
|
||||
let snapshot = DataUsageInfo {
|
||||
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
|
||||
scanner_cycle: Some(9),
|
||||
scanner_epoch: Some(4),
|
||||
buckets_count: 0,
|
||||
usage_snapshot_converged: Some(false),
|
||||
usage_snapshot_partial: true,
|
||||
usage_snapshot_set_states: vec![
|
||||
DataUsageSnapshotSetState {
|
||||
pool_index: 0,
|
||||
set_index: 0,
|
||||
scanner_cycle: Some(9),
|
||||
scanner_epoch: Some(4),
|
||||
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
|
||||
complete: true,
|
||||
tombstone: false,
|
||||
},
|
||||
state,
|
||||
],
|
||||
..Default::default()
|
||||
};
|
||||
assert!(snapshot.is_valid_partial_snapshot());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn old_set_completion_cannot_overwrite_new_aggregate() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut old = completed_root_cache("bucket", 1, 20, source);
|
||||
old.info.next_cycle = 7;
|
||||
old.info.leader_epoch = 2;
|
||||
let expected = HashSet::from([source]);
|
||||
assert!(observational_data_usage_info(&[old], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_aggregate_survives_restart_and_leader_failover() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let source = DataUsageCacheSource::new(0, 0);
|
||||
let mut lkg = lkg_root_cache("bucket", 5, source);
|
||||
lkg.info.lkg_leader_epoch = Some(4);
|
||||
lkg.info.lkg_next_cycle = Some(9);
|
||||
let expected = HashSet::from([source]);
|
||||
let (observed, _) = observational_data_usage_info(&[lkg], &expected, &all_buckets, TEST_PLAN_DIGEST, 10, 5)
|
||||
.expect("compatible LKG should survive a leader change");
|
||||
assert_eq!(observed.usage_snapshot_set_states[0].scanner_epoch, Some(4));
|
||||
assert_eq!(observed.objects_total_size, 50);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn usage_aggregate_cost_is_linear_in_set_count() {
|
||||
let all_buckets = vec!["bucket".to_string()];
|
||||
let mut results = Vec::new();
|
||||
let mut expected = HashSet::new();
|
||||
for index in 0..32 {
|
||||
let source = DataUsageCacheSource::new(index, 0);
|
||||
expected.insert(source);
|
||||
let mut cache = completed_root_cache("bucket", 1, 20, source);
|
||||
cache.info.next_cycle = 8;
|
||||
cache.info.leader_epoch = 3;
|
||||
results.push(cache);
|
||||
}
|
||||
let (observed, _) = observational_data_usage_info(&results, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("all set snapshots should aggregate");
|
||||
assert_eq!(observed.objects_total_count, 32);
|
||||
let reversed = results.iter().rev().cloned().collect::<Vec<_>>();
|
||||
let (reversed_observed, _) = observational_data_usage_info(&reversed, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
|
||||
.expect("reordered set snapshots should aggregate");
|
||||
assert_eq!(observed.usage_snapshot_set_states, reversed_observed.usage_snapshot_set_states);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completed_data_usage_info_publishes_tier_stats_across_sets() {
|
||||
let all_buckets = vec!["bucket-a".to_string(), "bucket-b".to_string()];
|
||||
|
||||
Reference in New Issue
Block a user