Compare commits

..

1 Commits

Author SHA1 Message Date
马登山 f4f9b315a6 fix(heal): coalesce duplicate MRF intents 2026-08-23 09:05:45 +08:00
32 changed files with 1123 additions and 1650 deletions
+6 -3
View File
@@ -39,10 +39,11 @@ jobs:
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
steps:
- name: Checkout repository
- name: Checkout main branch
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -88,10 +89,11 @@ jobs:
# either casing.
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout repository
- name: Checkout main branch
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -176,10 +178,11 @@ jobs:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout repository
- name: Checkout main branch
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
+1 -2
View File
@@ -30,8 +30,7 @@ make build-docker BUILD_OS=ubuntu22.04
- Crate membership: `Cargo.toml` `[workspace].members`
- Architecture, layering, crate map: [ARCHITECTURE.md](ARCHITECTURE.md)
- Migration guardrails & readiness contracts: [docs/architecture/](docs/architecture/README.md)
- CI workflow steps: `.github/workflows/`; event, timeout, and required-status
matrix: [docs/testing/ci-gates.md](docs/testing/ci-gates.md)
- CI gates: `.github/workflows/ci.yml` (source of truth; never copy its steps into docs)
- Test-layer taxonomy, per-layer entry commands, serial/nextest rules, flake
policy: [docs/testing/README.md](docs/testing/README.md)
- Tier/ILM transition debugging (xl.meta inspection, versionId tracing):
-2
View File
@@ -70,8 +70,6 @@ make pre-pr
> For the full test-layer taxonomy (unit / ecstore black-box / e2e / s3s-e2e / S3 compatibility / chaos / fuzz / bench), each layer's entry command, the naming conventions the migration gate depends on, and the serial/nextest rules, see [docs/testing/README.md](docs/testing/README.md).
> For the event, timeout, required-status, and local reproduction matrix, see [docs/testing/ci-gates.md](docs/testing/ci-gates.md).
### 🔒 Automated Pre-commit Hooks
#### What `make pre-commit` and `make pre-pr` actually run
Generated
+27 -22
View File
@@ -1858,9 +1858,9 @@ dependencies = [
[[package]]
name = "cc"
version = "1.4.4"
version = "1.4.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ad534f4357a5264cce5019c989cf66a4f0dc4e0d1b1d15f8aacec0ff7360273"
checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d"
dependencies = [
"find-msvc-tools",
"jobserver",
@@ -2522,6 +2522,12 @@ dependencies = [
"subtle",
]
[[package]]
name = "cty"
version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b365fabc795046672053e29c954733ec3b05e4be654ab130fe8f1f94d7051f35"
[[package]]
name = "curve25519-dalek"
version = "4.1.3"
@@ -5982,6 +5988,15 @@ version = "0.2.16"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981"
[[package]]
name = "libmimalloc-sys"
version = "0.1.49"
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
dependencies = [
"cc",
"cty",
]
[[package]]
name = "libredox"
version = "0.1.20"
@@ -6382,6 +6397,14 @@ dependencies = [
"synstructure 0.13.2",
]
[[package]]
name = "mimalloc"
version = "0.1.52"
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
dependencies = [
"libmimalloc-sys",
]
[[package]]
name = "mime"
version = "0.3.17"
@@ -9139,11 +9162,13 @@ dependencies = [
"insta",
"jiff",
"libc",
"libmimalloc-sys",
"libsystemd",
"matchit 0.9.2",
"md-5 0.11.0",
"metrics",
"metrics-util",
"mimalloc",
"mime_guess",
"opentelemetry",
"opentelemetry_sdk",
@@ -9179,8 +9204,6 @@ dependencies = [
"rustfs-lock",
"rustfs-log-analyzer",
"rustfs-madmin",
"rustfs-mimalloc",
"rustfs-mimalloc-sys",
"rustfs-notify",
"rustfs-object-capacity",
"rustfs-object-data-cache",
@@ -9852,24 +9875,6 @@ dependencies = [
"tokio",
]
[[package]]
name = "rustfs-mimalloc"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a406f4aa07084301d485beec873af6dccc8e3f8762da244743df92038b1db1a6"
dependencies = [
"rustfs-mimalloc-sys",
]
[[package]]
name = "rustfs-mimalloc-sys"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c3051b819175f58445d4c369a72f0ab88149f3885ba8bea2aff3be01f53fe7cd"
dependencies = [
"cc",
]
[[package]]
name = "rustfs-notify"
version = "1.0.0-rc.3"
+2 -2
View File
@@ -350,8 +350,8 @@ russh-sftp = "2.4.0"
dav-server = "0.11.0"
# Performance Analysis and Memory Profiling
rustfs-mimalloc = { version = "0.5.0" }
rustfs-mimalloc-sys = { version = "0.5.0" }
mimalloc = { version = "0.1.52", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11" }
libmimalloc-sys = { version = "0.1.49", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11", features = ["extended"] }
hotpath = { version = "0.23.3", default-features = false }
# Snapshot testing for output format regression detection
insta = { version = "1.48" }
+333 -10
View File
@@ -23,21 +23,32 @@
//! 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, OnceLock,
Arc, Mutex, 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)]
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub enum MrfKind {
/// Erasure decode failed while serving a read (read path).
DecodeFailure,
@@ -67,12 +78,52 @@ 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;
@@ -87,6 +138,159 @@ 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.
@@ -122,21 +326,90 @@ 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 false;
return MrfIngressResult::Dropped(MrfDropReason::Disabled);
}
let Some(sender) = GLOBAL_MRF_SENDER.get() else {
return false;
return MrfIngressResult::Dropped(MrfDropReason::Uninitialized);
};
let intent = MrfIntent {
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: version_id.map(|vid| *vid.as_bytes()),
version_id,
scope,
};
let lease = match coalescer_admit(key.clone()) {
Ok(lease) => lease,
Err(result) => return result,
};
let intent = MrfIntent {
bucket: key.bucket.clone(),
object: key.object.clone(),
version_id: key.version_id,
kind,
scope,
lease: Some(lease),
enqueued_at_ms: unix_now_ms(),
attempts: 0,
};
sender.try_send(intent).is_ok()
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,
);
}
fn unix_now_ms() -> u64 {
@@ -144,7 +417,8 @@ 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)
.map(|d| d.as_millis() as u64)
.ok()
.and_then(|d| u64::try_from(d.as_millis()).ok())
.unwrap_or(0)
}
@@ -215,12 +489,60 @@ 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");
@@ -230,6 +552,7 @@ 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);
@@ -239,8 +562,8 @@ mod tests {
// Fill the bounded channel past capacity: excess intents are dropped,
// never blocking.
let mut accepted = 0;
for _ in 0..(MRF_CHANNEL_CAPACITY + 64) {
if try_send_mrf_intent(MrfKind::PartialWrite, "b", "o", None) {
for index in 0..(MRF_CHANNEL_CAPACITY + 64) {
if try_send_mrf_intent(MrfKind::PartialWrite, "b", &format!("o-{index}"), None) {
accepted += 1;
}
}
+1 -89
View File
@@ -307,38 +307,6 @@ pub struct DiskUsageStatus {
pub snapshot_exists: bool,
}
/// A bounded reconciliation record for an object whose logical size could not
/// be trusted at the scanner boundary. The scanner persists these records in
/// its cache; keeping the model here avoids a second, incompatible accounting
/// representation in storage-facing crates.
#[derive(Debug, Default, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct SizeReconciliationEntry {
/// Stable object/version identity key (not a metrics label).
pub key: String,
pub bucket: String,
pub object: String,
#[serde(default)]
pub version_id: Option<String>,
#[serde(default)]
pub generation: Option<String>,
/// Structured reason label; raw metadata values must never be stored here.
pub reason: String,
#[serde(default)]
pub physical_size: Option<u64>,
#[serde(default)]
pub first_seen: u64,
#[serde(default)]
pub attempts: u32,
}
/// Object scope refreshed by one scanner pass. Existing debts in this scope
/// are removed before the pass's unresolved records are inserted.
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct SizeReconciliationScope {
pub bucket: String,
pub object: String,
}
/// Size summary for a single object or group of objects
#[derive(Debug, Default, Clone)]
pub struct SizeSummary {
@@ -368,16 +336,6 @@ pub struct SizeSummary {
pub repl_target_stats: HashMap<String, ReplTargetSizeSummary>,
/// Per-tier accounting, keyed by storage class or remote tier name
pub tier_stats: HashMap<String, TierStats>,
/// Size-resolution debts observed while scanning this summary.
pub size_reconciliation: Vec<SizeReconciliationEntry>,
/// True when the per-object summary exceeded its bounded debt buffer.
/// Callers must retain prior ledger entries rather than treating the
/// partial list as a complete refresh.
pub size_reconciliation_truncated: bool,
/// Object scopes refreshed by this summary. They let the durable ledger
/// remove versions that resolved without allocating one key per healthy
/// version on the hot path.
pub reconciliation_scopes: Vec<SizeReconciliationScope>,
}
/// Replication target size summary
@@ -875,8 +833,7 @@ impl DataUsageEntry {
///
/// The canonical wire format is written by the hand-written map-encoded
/// `Serialize` on the scanner-side `DataUsageCacheInfo`
/// (`crates/scanner/src/data_usage_define.rs`), which carries the original 16
/// fields plus an optional reconciliation field.
/// (`crates/scanner/src/data_usage_define.rs`), which carries 16 fields.
/// This type decodes only the shared subset and is deliberately not
/// `Serialize`: a derived (array) encoding of this 6-field subset would
/// corrupt the cache for scanner readers, so no write path may exist here.
@@ -1820,51 +1777,6 @@ impl SizeSummary {
entry.pending_count = entry.pending_count.saturating_add(stats.pending_count);
entry.failed_count = entry.failed_count.saturating_add(stats.failed_count);
}
for entry in &other.size_reconciliation {
self.record_size_reconciliation(entry.clone());
}
self.size_reconciliation_truncated |= other.size_reconciliation_truncated;
for scope in &other.reconciliation_scopes {
self.record_reconciliation_scope(&scope.bucket, &scope.object);
}
}
/// Add one reconciliation debt, coalescing repeated observations in the
/// same object summary. The scanner cache applies its own larger bound.
pub fn record_size_reconciliation(&mut self, entry: SizeReconciliationEntry) {
const MAX_SUMMARY_RECONCILIATION_ENTRIES: usize = 1024;
if let Some(existing) = self.size_reconciliation.iter_mut().find(|value| value.key == entry.key) {
existing.reason = entry.reason;
existing.physical_size = entry.physical_size;
existing.generation = entry.generation;
existing.version_id = entry.version_id;
return;
}
if self.size_reconciliation.len() < MAX_SUMMARY_RECONCILIATION_ENTRIES {
self.size_reconciliation.push(entry);
} else {
self.size_reconciliation_truncated = true;
}
}
/// Mark one object scope as refreshed. Duplicate scopes are suppressed so
/// merging summaries remains bounded and deterministic.
pub fn record_reconciliation_scope(&mut self, bucket: &str, object: &str) {
if !self
.reconciliation_scopes
.iter()
.any(|scope| scope.bucket == bucket && scope.object == object)
{
if self.reconciliation_scopes.len() >= 1024 {
self.size_reconciliation_truncated = true;
return;
}
self.reconciliation_scopes.push(SizeReconciliationScope {
bucket: bucket.to_string(),
object: object.to_string(),
});
}
}
}
+6 -38
View File
@@ -784,24 +784,6 @@ pub(crate) fn create_deferred_bitrot_reader_with_stripe_handle(
///
/// # Returns
/// A Result containing the BitrotWriterWrapper or an error
/// Size hint handed to `DiskAPI::create_file` for a bitrot-wrapped shard.
///
/// A known length is grown by one checksum per shard so the on-disk file size
/// matches what the bitrot writer emits. A negative length is the
/// unknown-size sentinel (`HashReader::SIZE_PRESERVE_LAYER`, used by SSE and
/// compression) and must be preserved: `RemoteDisk::create_file` forwards it
/// in the `put_file_stream` query, and the receiver only treats `size > 0` as
/// a fixed body length when locating the authenticated trailer. Clamping it
/// to `0` would claim an empty body and misframe the stream. `0` stays `0`
/// because a genuinely empty object still means an empty body.
fn bitrot_create_file_size(length: i64, shard_size: usize, checksum_algo: &HashAlgorithm) -> i64 {
if length <= 0 {
return length;
}
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
}
pub async fn create_bitrot_writer(
is_inline_buffer: bool,
disk: Option<&DiskStore>,
@@ -814,7 +796,12 @@ pub async fn create_bitrot_writer(
let writer = if is_inline_buffer {
CustomWriter::new_inline_buffer()
} else if let Some(disk) = disk {
let length = bitrot_create_file_size(length, shard_size, &checksum_algo);
let length = if length > 0 {
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
} else {
0
};
let file = disk.create_file("", volume, path, length).await?;
#[cfg(feature = "hotpath")]
@@ -833,25 +820,6 @@ mod tests {
use rustfs_rio::ChunkReader;
use std::collections::VecDeque;
#[test]
fn bitrot_create_file_size_grows_known_length_by_checksums() {
// 10 bytes over 4-byte shards = 3 shards, each followed by a 32-byte hash.
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::HighwayHash256), 10 + 3 * 32);
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::None), 10);
}
#[test]
fn bitrot_create_file_size_keeps_empty_and_unknown_distinct() {
assert_eq!(bitrot_create_file_size(0, 4, &HashAlgorithm::HighwayHash256), 0);
// SSE/compression streams advertise SIZE_PRESERVE_LAYER (-1); the remote
// put_file_stream receiver relies on a non-positive size to parse the auth
// trailer from the stream tail, so the sentinel must survive untouched.
assert_eq!(
bitrot_create_file_size(rustfs_rio::HashReader::SIZE_PRESERVE_LAYER, 4, &HashAlgorithm::HighwayHash256),
rustfs_rio::HashReader::SIZE_PRESERVE_LAYER
);
}
struct TestChunkReader {
chunks: VecDeque<Bytes>,
}
@@ -1412,8 +1412,11 @@ 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 {
rustfs_common::mrf_channel::try_send_mrf_intent(kind, bucket, object, version_uuid);
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));
}
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
+31 -7
View File
@@ -2124,13 +2124,26 @@ impl SetDisks {
let put_object_size = known_put_object_storage_size(data.size());
let shard_file_size_raw = erasure.shard_file_size(put_object_size);
let is_inline_buffer = storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
let is_inline_buffer =
storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled();
let shard_file_size = shard_file_size_raw;
let shard_size = erasure.shard_size();
let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size);
let direct_inline_commit = matches!(write_path, SmallWritePath::Inline);
{
use std::io::Write;
let msg = format!(
"INLINE_DEBUG: bucket={} obj={} size={} shard_fs={} ds={} bs={} inline={} direct={} path={} iblock={} ver={}\n",
bucket, object, put_object_size, shard_file_size_raw, erasure.data_shards, fi.erasure.block_size,
is_inline_buffer, direct_inline_commit, write_path.metric_label(), storage_class_config.inline_block(), opts.versioned
);
if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open("/tmp/rustfs_inline_debug.log") {
let _ = f.write_all(msg.as_bytes());
}
let _ = std::io::stderr().write_all(msg.as_bytes());
}
rustfs_io_metrics::record_put_object_path(write_path.metric_label());
let writer_setup_stage_start = collect_stage_timing.then(Instant::now);
let (mut writers, errors) = if direct_inline_commit {
@@ -6637,12 +6650,23 @@ 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.
rustfs_common::mrf_channel::try_send_mrf_intent(
rustfs_common::mrf_channel::MrfKind::PartialWrite,
bucket,
object,
uuid::Uuid::try_parse(version_id).ok(),
);
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),
);
}
let mut request = rustfs_common::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
+1 -1
View File
@@ -3194,7 +3194,7 @@ impl ECStore {
// Default return value
let mut del_objects = vec![DeletedObject::default(); objects.len()];
let accounting = vec![None; objects.len()];
let mut accounting = vec![None; objects.len()];
let mut del_errs = Vec::with_capacity(objects.len());
for _ in 0..objects.len() {
@@ -271,7 +271,7 @@ pub(super) fn resolve_latest_object_info_candidates(
.filter(|candidate| latest_candidate_mod_time(candidate) == Some(latest_mod_time))
.collect::<Vec<_>>();
latest_candidates.sort_by_key(|candidate| std::cmp::Reverse(candidate.idx));
latest_candidates.sort_by(|left, right| right.idx.cmp(&left.idx));
let Some(winner) = latest_candidates.first() else {
return Err(Error::ErasureReadQuorum);
+54 -3
View File
@@ -119,6 +119,9 @@ 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)]
@@ -889,7 +892,19 @@ impl HealManager {
}
fn remove_mrf_repair_notice_targets_for_task(&self, task_id: &str) {
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).remove(task_id);
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,
);
}
}
}
fn insert_mrf_repair_notice_target(
@@ -1279,7 +1294,20 @@ impl HealManager {
}
self.task_aliases.lock().await.clear();
self.retrying_heals.lock().await.clear();
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).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,
);
}
crate::set_heal_queue_length(0);
// update state
@@ -1311,12 +1339,32 @@ 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,
@@ -1325,6 +1373,9 @@ impl HealManager {
bucket,
object,
version_id,
kind,
scope,
lease,
}),
)
.await
@@ -1539,7 +1590,7 @@ impl HealManager {
Self::insert_mrf_repair_notice_target(&mut targets, &task_id, target);
}
if let Some(displaced_task_id) = &displaced_task_id {
lock_mrf_repair_notice_targets(&self.mrf_repair_notice_targets).remove(displaced_task_id);
self.remove_mrf_repair_notice_targets_for_task(displaced_task_id);
}
drop(retrying_heals);
drop(queue);
+12 -1
View File
@@ -506,7 +506,18 @@ impl HealManager {
&displaced_terminal,
)
.await;
lock_mrf_repair_notice_targets(&mrf_repair_notice_targets).remove(&displaced_task_id);
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,
);
}
}
}
if matches!(admission, HealAdmissionResult::Accepted) {
if should_notify {
+5 -1
View File
@@ -299,7 +299,11 @@ impl PriorityHealQueue {
/// Create a deduplication key from a heal request
pub(super) fn make_dedup_key(request: &HealRequest) -> String {
Self::make_dedup_key_for_type(&request.heal_type)
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,
}
}
pub(super) fn make_dedup_key_for_type(heal_type: &HealType) -> String {
+36 -1
View File
@@ -349,6 +349,8 @@ 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()
@@ -638,7 +640,19 @@ 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) {
lock_mrf_repair_notice_targets(registry).remove(task_id);
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,
);
}
}
}
fn take_mrf_repair_notice_targets(
@@ -671,6 +685,27 @@ 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,
);
}
}
+346 -63
View File
@@ -37,7 +37,7 @@ use crate::heal::manager::HealManager;
use metrics::{counter, gauge};
use rustfs_common::heal_channel::{HealAdmissionDropReason, HealAdmissionResult};
use rustfs_common::mrf_channel::{MRF_MAX_ATTEMPTS, MrfIntent};
use std::collections::VecDeque;
use std::collections::{HashSet, VecDeque};
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::mpsc;
@@ -48,15 +48,27 @@ use crate::heal::task::{HealOptions, HealPriority, HealRequest, HealType};
/// Journal location inside the metadata bucket, following the resume-state
/// layout.
pub(crate) const MRF_JOURNAL_PATH: &str = "buckets/.heal/mrf/journal.bin";
/// The scoped path is the authoritative snapshot for new readers and carries
/// both v1 and v2 records. The legacy path is only a v1 compatibility mirror;
/// older readers ignore the authoritative path, while new readers never merge
/// the two files. This prevents a partial two-file flush from fabricating a
/// mixed epoch.
pub(crate) const MRF_SCOPED_JOURNAL_PATH: &str = "buckets/.heal/mrf/journal-scoped.bin";
/// Record format tag.
const MRF_JOURNAL_FORMAT: u8 = 1;
/// Record layout version.
const MRF_JOURNAL_VERSION: u8 = 1;
const MRF_JOURNAL_VERSION_SCOPED: u8 = 2;
/// Fixed header size: format, version, kind, attempts, enqueued_at_ms,
/// has_version flag.
const MRF_RECORD_FIXED_HEAD: usize = 1 + 1 + 1 + 1 + 8 + 1;
const MRF_MAX_IDENTITY_COMPONENT: usize = 1024;
fn metric_f64(value: usize) -> f64 {
f64::from(u32::try_from(value).unwrap_or(u32::MAX))
}
#[derive(Debug, Clone)]
pub(crate) struct MrfConsumerConfig {
@@ -101,40 +113,90 @@ impl Default for MrfConsumerConfig {
/// incoming intent (never a resident one) and counts the loss.
pub(crate) struct MrfQueue {
pending: VecDeque<MrfIntent>,
pending_keys: HashSet<MrfQueueKey>,
bytes: usize,
capacity: usize,
byte_budget: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct MrfQueueKey {
kind: rustfs_common::mrf_channel::MrfKind,
bucket: Arc<str>,
object: Arc<str>,
version_id: Option<[u8; 16]>,
scope: Option<rustfs_common::mrf_channel::MrfScope>,
}
fn queue_key(intent: &MrfIntent) -> MrfQueueKey {
let version_id = intent.version_id.filter(|bytes| *bytes != [0; 16]);
let scope = (!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption))
.then_some(intent.scope)
.flatten();
MrfQueueKey {
kind: intent.kind,
bucket: intent.bucket.clone(),
object: intent.object.clone(),
version_id,
scope,
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum MrfQueuePushResult {
Enqueued,
Coalesced,
Rejected,
}
impl MrfQueue {
pub(crate) fn new(capacity: usize, byte_budget: usize) -> Self {
Self {
pending: VecDeque::new(),
pending_keys: HashSet::new(),
bytes: 0,
capacity,
byte_budget,
}
}
/// Returns `false` (after counting) when either ceiling would be crossed.
pub(crate) fn try_push(&mut self, intent: MrfIntent) -> bool {
pub(crate) fn try_push_typed(&mut self, intent: MrfIntent) -> MrfQueuePushResult {
if intent.bucket.len() > MRF_MAX_IDENTITY_COMPONENT || intent.object.len() > MRF_MAX_IDENTITY_COMPONENT {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "identity_oversized").increment(1);
return MrfQueuePushResult::Rejected;
}
let key = queue_key(&intent);
if self.pending_keys.contains(&key) {
counter!("rustfs_heal_mrf_coalesced_total", "layer" => "queue").increment(1);
return MrfQueuePushResult::Coalesced;
}
let cost = intent.estimated_bytes();
if self.pending.len() >= self.capacity || self.bytes + cost > self.byte_budget {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "queue_overflow").increment(1);
return false;
return MrfQueuePushResult::Rejected;
}
self.bytes += cost;
self.pending_keys.insert(key);
self.pending.push_back(intent);
true
MrfQueuePushResult::Enqueued
}
/// Bool compatibility adapter: only a newly executable queue item is
/// reported as accepted; a coalesced duplicate is not durable admission.
#[cfg(test)]
pub(crate) fn try_push(&mut self, intent: MrfIntent) -> bool {
matches!(self.try_push_typed(intent), MrfQueuePushResult::Enqueued)
}
pub(crate) fn pop_front(&mut self) -> Option<MrfIntent> {
let intent = self.pending.pop_front()?;
self.pending_keys.remove(&queue_key(&intent));
self.bytes = self.bytes.saturating_sub(intent.estimated_bytes());
Some(intent)
}
pub(crate) fn push_back(&mut self, intent: MrfIntent) {
self.pending_keys.insert(queue_key(&intent));
self.bytes += intent.estimated_bytes();
self.pending.push_back(intent);
}
@@ -157,10 +219,24 @@ impl MrfQueue {
// ---------------------------------------------------------------------------
/// Append one encoded record to `out`.
pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) {
pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) -> bool {
let Ok(bucket_len) = u32::try_from(intent.bucket.len()) else {
return false;
};
let Ok(object_len) = u32::try_from(intent.object.len()) else {
return false;
};
let scope = (!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption))
.then_some(intent.scope)
.flatten();
let version_id = intent.version_id.filter(|bytes| *bytes != [0; 16]);
let start = out.len();
out.push(MRF_JOURNAL_FORMAT);
out.push(MRF_JOURNAL_VERSION);
out.push(if scope.is_some() {
MRF_JOURNAL_VERSION_SCOPED
} else {
MRF_JOURNAL_VERSION
});
out.push(match intent.kind {
rustfs_common::mrf_channel::MrfKind::DecodeFailure => 1,
rustfs_common::mrf_channel::MrfKind::MetadataCorruption => 2,
@@ -168,27 +244,36 @@ pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) {
});
out.push(intent.attempts);
out.extend_from_slice(&intent.enqueued_at_ms.to_le_bytes());
match intent.version_id {
match version_id {
Some(bytes) => {
out.push(1);
out.extend_from_slice(&bytes);
}
None => out.push(0),
}
out.extend_from_slice(&(intent.bucket.len() as u32).to_le_bytes());
out.extend_from_slice(&(intent.object.len() as u32).to_le_bytes());
if let Some(scope) = scope {
out.extend_from_slice(&scope.pool_index.to_le_bytes());
out.extend_from_slice(&scope.set_index.to_le_bytes());
}
out.extend_from_slice(&bucket_len.to_le_bytes());
out.extend_from_slice(&object_len.to_le_bytes());
out.extend_from_slice(intent.bucket.as_bytes());
out.extend_from_slice(intent.object.as_bytes());
let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc);
hasher.update(&out[start..]);
out.extend_from_slice(&(hasher.finalize() as u32).to_le_bytes());
let Ok(checksum) = u32::try_from(hasher.finalize()) else {
out.truncate(start);
return false;
};
out.extend_from_slice(&checksum.to_le_bytes());
true
}
fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
if data.len() < MRF_RECORD_FIXED_HEAD + 8 {
return None;
}
if data[0] != MRF_JOURNAL_FORMAT || data[1] != MRF_JOURNAL_VERSION {
if data[0] != MRF_JOURNAL_FORMAT || !matches!(data[1], MRF_JOURNAL_VERSION | MRF_JOURNAL_VERSION_SCOPED) {
return None;
}
let kind = match data[2] {
@@ -198,24 +283,38 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
_ => return None,
};
let attempts = data[3];
let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().expect("slice length checked"));
let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().ok()?);
let has_version = data[12] != 0;
let mut cursor = MRF_RECORD_FIXED_HEAD;
let version_id = if has_version {
if data.len() < cursor + 16 {
return None;
}
let bytes: [u8; 16] = data[cursor..cursor + 16].try_into().expect("slice length checked");
let bytes: [u8; 16] = data[cursor..cursor + 16].try_into().ok()?;
cursor += 16;
Some(bytes)
} else {
None
};
let scope = if data[1] == MRF_JOURNAL_VERSION_SCOPED {
if data.len() < cursor + 8 {
return None;
}
let pool_index = u32::from_le_bytes(data[cursor..cursor + 4].try_into().ok()?);
let set_index = u32::from_le_bytes(data[cursor + 4..cursor + 8].try_into().ok()?);
cursor += 8;
Some(rustfs_common::mrf_channel::MrfScope { pool_index, set_index })
} else {
None
};
if data.len() < cursor + 8 {
return None;
}
let bucket_len = u32::from_le_bytes(data[cursor..cursor + 4].try_into().expect("slice length checked")) as usize;
let object_len = u32::from_le_bytes(data[cursor + 4..cursor + 8].try_into().expect("slice length checked")) as usize;
let bucket_len = usize::try_from(u32::from_le_bytes(data[cursor..cursor + 4].try_into().ok()?)).ok()?;
let object_len = usize::try_from(u32::from_le_bytes(data[cursor + 4..cursor + 8].try_into().ok()?)).ok()?;
if bucket_len > MRF_MAX_IDENTITY_COMPONENT || object_len > MRF_MAX_IDENTITY_COMPONENT {
return None;
}
cursor += 8;
let body_end = cursor.checked_add(bucket_len)?.checked_add(object_len)?;
let record_end = body_end.checked_add(4)?;
@@ -224,7 +323,7 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
}
let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc);
hasher.update(&data[..body_end]);
if (hasher.finalize() as u32) != u32::from_le_bytes(data[body_end..record_end].try_into().expect("slice length checked")) {
if u32::try_from(hasher.finalize()).ok()? != u32::from_le_bytes(data[body_end..record_end].try_into().ok()?) {
return None;
}
let bucket = std::sync::Arc::from(std::str::from_utf8(&data[cursor..cursor + bucket_len]).ok()?);
@@ -235,6 +334,12 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> {
object,
version_id,
kind,
scope: if matches!(kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption) {
None
} else {
scope
},
lease: None,
enqueued_at_ms,
attempts,
},
@@ -269,9 +374,9 @@ async fn journal_disks() -> Vec<DiskStore> {
map.values().flatten().cloned().collect()
}
async fn read_journal() -> Option<Vec<u8>> {
async fn read_journal(path: &str) -> Option<Vec<u8>> {
for disk in journal_disks().await {
match disk.read_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH).await {
match disk.read_all(super::RUSTFS_META_BUCKET, path).await {
Ok(bytes) => return Some(bytes.to_vec()),
Err(_) => continue,
}
@@ -282,35 +387,51 @@ async fn read_journal() -> Option<Vec<u8>> {
/// Write the snapshot to every local disk; returns true when at least one
/// disk accepted it, so a total write failure keeps the runtime dirty and
/// the next tick retries the persist.
async fn write_journal(data: &[u8]) -> bool {
async fn write_journal(path: &str, data: &[u8]) -> bool {
let payload = bytes::Bytes::copy_from_slice(data);
let mut any_persisted = false;
for disk in journal_disks().await {
match disk
.write_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH, payload.clone())
.await
{
match disk.write_all(super::RUSTFS_META_BUCKET, path, payload.clone()).await {
Ok(()) => any_persisted = true,
Err(err) => warn_mrf_journal_write(&err),
}
}
if !data.is_empty() {
counter!("rustfs_heal_mrf_journal_fsync_total").increment(1);
}
gauge!("rustfs_heal_mrf_journal_bytes").set(data.len() as f64);
any_persisted
}
async fn delete_journal() {
for disk in journal_disks().await {
let _ = disk
async fn delete_journal(path: &str) -> bool {
let disks = journal_disks().await;
if disks.is_empty() {
counter!("rustfs_heal_mrf_journal_delete_failures_total").increment(1);
return false;
}
let mut all_deleted = true;
for disk in disks {
let result = disk
.delete(
super::RUSTFS_META_BUCKET,
MRF_JOURNAL_PATH,
path,
crate::heal::storage_api::owner::EcstoreDeleteOptions::default(),
)
.await;
if let Err(err) = result {
// Delete is idempotent: a compatibility mirror that was never
// written (or was already removed) is clean, not a retry state.
if !matches!(err, super::DiskError::FileNotFound | super::DiskError::VolumeNotFound) {
all_deleted = false;
}
}
}
if !all_deleted {
counter!("rustfs_heal_mrf_journal_delete_failures_total").increment(1);
}
all_deleted
}
async fn delete_journals() -> bool {
let authoritative_deleted = delete_journal(MRF_SCOPED_JOURNAL_PATH).await;
let legacy_deleted = delete_journal(MRF_JOURNAL_PATH).await;
authoritative_deleted && legacy_deleted
}
fn warn_mrf_journal_write(err: &super::DiskError) {
@@ -331,7 +452,10 @@ fn warn_mrf_journal_write(err: &super::DiskError) {
pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest {
let bucket = intent.bucket.to_string();
let object = intent.object.to_string();
let version_id = intent.version_id.map(|bytes| Uuid::from_bytes(bytes).to_string());
let version_id = intent
.version_id
.filter(|bytes| *bytes != [0; 16])
.map(|bytes| Uuid::from_bytes(bytes).to_string());
let (heal_type, priority) = match intent.kind {
rustfs_common::mrf_channel::MrfKind::DecodeFailure => (
HealType::ECDecode {
@@ -351,18 +475,28 @@ pub(crate) fn build_heal_request(intent: &MrfIntent) -> HealRequest {
HealPriority::Normal,
),
};
let mut request = HealRequest::new(heal_type, HealOptions::default(), priority);
let mut options = HealOptions::default();
if !matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption)
&& let Some(scope) = intent.scope
{
options.pool_index = usize::try_from(scope.pool_index).ok();
options.set_index = usize::try_from(scope.set_index).ok();
}
let mut request = HealRequest::new(heal_type, options, priority);
request.source = rustfs_common::heal_channel::HealRequestSource::Mrf;
request
}
async fn submit_mrf_heal_request(manager: &HealManager, intent: &MrfIntent) -> crate::Result<HealAdmissionResult> {
let receipt = manager
.submit_mrf_heal_request_with_receipt(
.submit_mrf_heal_request_with_receipt_and_identity(
build_heal_request(intent),
intent.bucket.clone(),
intent.object.clone(),
intent.version_id,
intent.kind,
intent.scope,
intent.lease,
)
.await?;
Ok(receipt.result)
@@ -387,16 +521,38 @@ struct MrfRuntime {
}
impl MrfRuntime {
fn snapshot(&self) -> Vec<u8> {
let mut buf = Vec::new();
fn snapshot(&self) -> (Vec<u8>, Vec<u8>) {
let mut authoritative = Vec::new();
let mut legacy = Vec::new();
for intent in self.queue.intents() {
encode_intent(intent, &mut buf);
let scoped_identity =
!matches!(intent.kind, rustfs_common::mrf_channel::MrfKind::MetadataCorruption) && intent.scope.is_some();
if !encode_intent(intent, &mut authoritative) {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "journal_identity_oversized").increment(1);
}
if !scoped_identity && !encode_intent(intent, &mut legacy) {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "journal_identity_oversized").increment(1);
}
}
buf
(authoritative, legacy)
}
async fn flush(&mut self) {
let persisted = write_journal(&self.snapshot()).await;
let (authoritative, legacy) = self.snapshot();
let authoritative_persisted = write_journal(MRF_SCOPED_JOURNAL_PATH, &authoritative).await;
if !authoritative.is_empty() {
counter!("rustfs_heal_mrf_journal_fsync_total").increment(1);
}
gauge!("rustfs_heal_mrf_journal_bytes").set(metric_f64(authoritative.len()));
// Publish the compatibility mirror only after the authoritative
// snapshot has reached at least one disk. This ordering prevents an
// old reader from observing a newer epoch that a new reader cannot
// see when the canonical write is unavailable.
let legacy_persisted = authoritative_persisted && write_journal(MRF_JOURNAL_PATH, &legacy).await;
// Keep dirty until both the authoritative snapshot and its
// compatibility mirror have been accepted; otherwise a one-sided
// failure would never retry the missing file.
let persisted = authoritative_persisted && legacy_persisted;
self.new_since_flush = 0;
// Keep the dirty flag when every disk write failed: a clean backlog
// would otherwise never rewrite, losing the periodic persist retry a
@@ -404,7 +560,7 @@ impl MrfRuntime {
if persisted {
self.dirty = false;
}
self.journal_on_disk = true;
self.journal_on_disk |= authoritative_persisted || legacy_persisted;
}
/// Drain pending intents into the heal manager until it is full, the
@@ -430,6 +586,7 @@ impl MrfRuntime {
intent.attempts = intent.attempts.saturating_add(1);
if intent.attempts >= MRF_MAX_ATTEMPTS {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1);
rustfs_common::mrf_channel::release_mrf_intent(&intent);
continue;
}
self.queue.push_back(intent);
@@ -438,11 +595,13 @@ impl MrfRuntime {
}
Ok(HealAdmissionResult::Dropped(_)) => {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "admission_policy").increment(1);
rustfs_common::mrf_channel::release_mrf_intent(&intent);
}
Err(_) => {
intent.attempts = intent.attempts.saturating_add(1);
if intent.attempts >= MRF_MAX_ATTEMPTS {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1);
rustfs_common::mrf_channel::release_mrf_intent(&intent);
continue;
}
self.queue.push_back(intent);
@@ -451,8 +610,8 @@ impl MrfRuntime {
}
}
}
gauge!("rustfs_heal_mrf_queue_depth").set(self.queue.depth() as f64);
gauge!("rustfs_heal_mrf_queue_bytes").set(self.queue.bytes() as f64);
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(self.queue.depth()));
gauge!("rustfs_heal_mrf_queue_bytes").set(metric_f64(self.queue.bytes()));
}
}
@@ -496,7 +655,12 @@ pub async fn replay_journal_once(manager: &Arc<HealManager>) -> usize {
let config = MrfConsumerConfig::default();
let mut queue = MrfQueue::new(config.queue_capacity, config.journal_max_bytes);
let mut backoff_until: Option<tokio::time::Instant> = None;
replay_into(manager, &mut queue, &mut backoff_until).await
replay_into(manager, &mut queue, &mut backoff_until).await.replayed
}
struct ReplayOutcome {
replayed: usize,
journal_on_disk: bool,
}
/// Shared replay core: read + decode + re-arm + delete, then drain what fits.
@@ -504,11 +668,25 @@ async fn replay_into(
manager: &Arc<HealManager>,
queue: &mut MrfQueue,
backoff_until: &mut Option<tokio::time::Instant>,
) -> usize {
let Some(data) = read_journal().await else {
return 0;
) -> ReplayOutcome {
// The scoped file is a complete authoritative snapshot. Fall back to the
// legacy mirror only when the authoritative path is unavailable; merging
// both files could combine records from different flush epochs.
let data = match read_journal(MRF_SCOPED_JOURNAL_PATH).await {
Some(data) => data,
None => match read_journal(MRF_JOURNAL_PATH).await {
Some(data) => data,
None => {
return ReplayOutcome {
replayed: 0,
journal_on_disk: false,
};
}
},
};
let (intents, truncated) = decode_journal(&data);
let mut intents = Vec::new();
let (decoded, truncated) = decode_journal(&data);
intents.extend(decoded);
if truncated > 0 {
tracing::warn!(
target: "rustfs::heal::mrf",
@@ -516,12 +694,15 @@ async fn replay_into(
"MRF journal had a torn tail; truncated records were discarded"
);
}
counter!("rustfs_heal_mrf_replayed_total").increment(intents.len() as u64);
counter!("rustfs_heal_mrf_replayed_total").increment(u64::try_from(intents.len()).unwrap_or(u64::MAX));
let replayed = intents.len();
for intent in intents {
queue.try_push(intent);
let result = queue.try_push_typed(intent.clone());
if !matches!(result, MrfQueuePushResult::Enqueued) {
rustfs_common::mrf_channel::release_mrf_intent(&intent);
}
}
delete_journal().await;
let journal_on_disk = !delete_journals().await;
// Drain the replayed intents immediately; whatever the manager refuses
// stays armed in `queue` for the consumer's retry loop.
@@ -541,7 +722,10 @@ async fn replay_into(
}
}
}
replayed
ReplayOutcome {
replayed,
journal_on_disk,
}
}
/// Replay the journal, then keep draining the channel into the heal manager
@@ -559,7 +743,8 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
// Replay: read the journal, re-arm intents (duplicates are merged by the
// manager's dedup key), then drop the file so the next flush starts clean.
replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
let replay = replay_into(&manager, &mut runtime.queue, &mut runtime.backoff_until).await;
runtime.journal_on_disk = replay.journal_on_disk;
// The replay deleted the journal file; anything still pending (e.g. the
// manager was full and backoff armed) must be re-persisted by the next
// flush or a crash before it would lose those intents.
@@ -587,9 +772,14 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
return;
}
for intent in batch.drain(..) {
if runtime.queue.try_push(intent) {
runtime.new_since_flush += 1;
runtime.dirty = true;
match runtime.queue.try_push_typed(intent.clone()) {
MrfQueuePushResult::Enqueued => {
runtime.new_since_flush += 1;
runtime.dirty = true;
}
MrfQueuePushResult::Coalesced | MrfQueuePushResult::Rejected => {
rustfs_common::mrf_channel::release_mrf_intent(&intent);
}
}
}
runtime.dispatch(manager.as_ref()).await;
@@ -613,13 +803,14 @@ async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receive
TickAction::DeleteJournal => {
// All intents consumed: remove the journal so a restart
// replays nothing (mirrors MinIO's post-replay unlink).
delete_journal().await;
runtime.journal_on_disk = false;
gauge!("rustfs_heal_mrf_journal_bytes").set(0.0);
if delete_journals().await {
runtime.journal_on_disk = false;
gauge!("rustfs_heal_mrf_journal_bytes").set(0.0);
}
}
TickAction::Idle => {}
}
gauge!("rustfs_heal_mrf_queue_depth").set(runtime.queue.depth() as f64);
gauge!("rustfs_heal_mrf_queue_depth").set(metric_f64(runtime.queue.depth()));
}
}
}
@@ -664,6 +855,8 @@ mod tests {
object: StdArc::from(object),
version_id: Some([7u8; 16]),
kind: MrfKind::DecodeFailure,
scope: None,
lease: None,
enqueued_at_ms: 1_700_000_000_000,
attempts,
}
@@ -694,17 +887,101 @@ mod tests {
fn queue_enforces_count_and_byte_ceilings() {
let mut queue = MrfQueue::new(2, usize::MAX);
assert!(queue.try_push(intent("b", "o", 0)));
assert!(queue.try_push(intent("b", "o", 0)));
assert!(!queue.try_push(intent("b", "o", 0)), "count ceiling must drop");
assert!(queue.try_push(intent("b", "o2", 0)));
assert!(!queue.try_push(intent("b", "o3", 0)), "count ceiling must drop");
let mut tiny = MrfQueue::new(usize::MAX, intent("bucket", "object", 0).estimated_bytes());
assert!(tiny.try_push(intent("bucket", "object", 0)));
assert!(
!tiny.try_push(intent("bucket", "object", 0)),
!tiny.try_push(intent("bucket", "object2", 0)),
"byte budget must drop before the second intent fits"
);
}
#[test]
fn duplicate_mrf_intents_coalesce_to_one_execution() {
let mut queue = MrfQueue::new(1000, usize::MAX);
let mut enqueued = 0;
let mut coalesced = 0;
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
enqueued += 1;
for _ in 0..999 {
match queue.try_push_typed(intent("bucket", "object", 0)) {
MrfQueuePushResult::Coalesced => coalesced += 1,
other => panic!("duplicate intent was not coalesced: {other:?}"),
}
}
assert_eq!(enqueued, 1);
assert_eq!(coalesced, 999);
assert_eq!(queue.depth(), 1);
}
#[test]
fn mrf_dedupe_does_not_merge_adjacent_version_pool_or_kind() {
let mut queue = MrfQueue::new(8, usize::MAX);
let mut first = intent("bucket", "object", 0);
first.kind = MrfKind::PartialWrite;
first.scope = Some(rustfs_common::mrf_channel::MrfScope {
pool_index: 1,
set_index: 1,
});
assert!(queue.try_push(first.clone()));
first.version_id = Some([8u8; 16]);
assert!(queue.try_push(first));
let mut other_scope = intent("bucket", "object", 0);
other_scope.kind = MrfKind::PartialWrite;
other_scope.scope = Some(rustfs_common::mrf_channel::MrfScope {
pool_index: 2,
set_index: 1,
});
assert!(queue.try_push(other_scope));
let mut other_kind = intent("bucket", "object", 0);
other_kind.kind = MrfKind::DecodeFailure;
other_kind.scope = None;
assert!(queue.try_push(other_kind));
assert_eq!(queue.depth(), 4);
}
#[test]
fn mrf_dedupe_full_returns_rejected_with_durable_pending() {
let mut queue = MrfQueue::new(1, usize::MAX);
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
assert_eq!(queue.try_push_typed(intent("bucket", "other", 0)), MrfQueuePushResult::Rejected);
assert_eq!(queue.depth(), 1);
let mut snapshot = Vec::new();
assert!(encode_intent(queue.intents().next().expect("resident intent"), &mut snapshot));
assert!(!snapshot.is_empty(), "the resident intent remains journalable after rejection");
}
#[test]
fn mrf_dedupe_failure_releases_key_for_retry() {
let mut queue = MrfQueue::new(1, usize::MAX);
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
let _failed = queue.pop_front().expect("queued intent");
assert_eq!(queue.try_push_typed(intent("bucket", "object", 1)), MrfQueuePushResult::Enqueued);
assert_eq!(queue.depth(), 1);
}
#[test]
fn mrf_dedupe_key_and_map_are_bounded() {
let mut queue = MrfQueue::new(2, usize::MAX);
assert_eq!(queue.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
assert_eq!(queue.try_push_typed(intent("bucket", "other", 0)), MrfQueuePushResult::Enqueued);
assert_eq!(queue.pending_keys.len(), 2);
assert_eq!(queue.try_push_typed(intent("bucket", "third", 0)), MrfQueuePushResult::Rejected);
assert_eq!(queue.depth(), 2);
}
#[test]
fn cross_node_duplicate_execution_remains_idempotent() {
// Node-local ingress maps intentionally do not merge across nodes;
// the manager's existing identity key absorbs the duplicate later.
let mut node_a = MrfQueue::new(8, usize::MAX);
let mut node_b = MrfQueue::new(8, usize::MAX);
assert_eq!(node_a.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
assert_eq!(node_b.try_push_typed(intent("bucket", "object", 0)), MrfQueuePushResult::Enqueued);
}
#[test]
fn journal_roundtrip_preserves_intents() {
let intents = vec![
@@ -715,6 +992,8 @@ mod tests {
object: StdArc::from("object/c"),
version_id: None,
kind: MrfKind::MetadataCorruption,
scope: None,
lease: None,
enqueued_at_ms: 5,
attempts: 1,
},
@@ -766,6 +1045,8 @@ mod tests {
object: StdArc::from("o"),
version_id: None,
kind: MrfKind::MetadataCorruption,
scope: None,
lease: None,
enqueued_at_ms: 0,
attempts: 0,
});
@@ -777,6 +1058,8 @@ mod tests {
object: StdArc::from("o"),
version_id: None,
kind: MrfKind::PartialWrite,
scope: None,
lease: None,
enqueued_at_ms: 0,
attempts: 0,
});
+72 -2
View File
@@ -35,6 +35,7 @@ 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()
@@ -79,14 +80,18 @@ fn journal_record(kind: u8, bucket: &str, object: &str, version: Option<[u8; 16]
body
}
fn write_journal_to_disks(disk_paths: &[std::path::PathBuf], data: &[u8]) {
fn write_journal_path_to_disks(disk_paths: &[std::path::PathBuf], relative_path: &str, data: &[u8]) {
for path in disk_paths {
let journal = path.join(META_BUCKET).join(JOURNAL_REL);
let journal = path.join(META_BUCKET).join(relative_path);
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,
@@ -187,8 +192,73 @@ 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()
}));
}
+25 -86
View File
@@ -16,14 +16,14 @@
//!
//! `scripts/test/vault_ha_kms_live.sh` owns the official Vault containers and
//! kills the active node while this test continuously decrypts through a
//! surviving standby. KV2 and Transit must recover after the bounded circuit
//! interval, use a bounded number of attempts, and leave the circuit and
//! in-flight gauges at zero after a new leader is elected.
//! surviving standby. KV2 and Transit requests must remain successful, use a
//! bounded number of attempts, and leave the circuit and in-flight gauges at
//! zero after a new leader is elected.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use metrics_util::MetricKind;
@@ -43,11 +43,6 @@ const OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
const IN_FLIGHT: &str = "rustfs_kms_backend_in_flight";
const CIRCUIT_OPEN: &str = "rustfs_kms_backend_circuit_open";
const MAX_ATTEMPTS: u32 = 10;
const ATTEMPT_TIMEOUT: Duration = Duration::from_secs(2);
const HEALTHY_PROGRESS_TIMEOUT: Duration = Duration::from_secs(20);
// The circuit remains open for 30s after five failed attempts.
const POST_FAILOVER_PROGRESS_TIMEOUT: Duration = Duration::from_secs(35);
const FAILOVER_ERROR_POLL_INTERVAL: Duration = Duration::from_millis(100);
type MetricEntry = (
metrics_util::CompositeKey,
@@ -69,7 +64,7 @@ fn config(backend: KmsBackend, backend_config: BackendConfig) -> KmsConfig {
backend,
backend_config,
allow_insecure_dev_defaults: true,
timeout: ATTEMPT_TIMEOUT,
timeout: Duration::from_secs(2),
retry_attempts: MAX_ATTEMPTS,
enable_cache: false,
..KmsConfig::default()
@@ -169,31 +164,14 @@ fn retryable_failures(snapshot: &[MetricEntry], operation: &str) -> u64 {
.sum()
}
async fn wait_for_count(
counter: &AtomicU64,
failure: &Mutex<Option<String>>,
minimum: u64,
description: &str,
timeout: Duration,
) {
tokio::time::timeout(timeout, async {
async fn wait_for_count(counter: &AtomicU64, minimum: u64, description: &str) {
tokio::time::timeout(Duration::from_secs(20), async {
while counter.load(Ordering::SeqCst) < minimum {
if let Some(error) = failure.lock().expect("decrypt failure lock poisoned").as_ref() {
panic!(
"{description} worker failed after {} successful decrypts: {error}",
counter.load(Ordering::SeqCst)
);
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.unwrap_or_else(|_| {
panic!(
"timed out after {timeout:?} waiting for {description}: completed {}, expected {minimum}",
counter.load(Ordering::SeqCst)
)
});
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
}
async fn wait_for_file(path: &Path, description: &str) {
@@ -211,8 +189,7 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
request: DecryptRequest,
expected: Vec<u8>,
completed: Arc<AtomicU64>,
allow_failover_errors: Arc<AtomicBool>,
failure: Arc<Mutex<Option<String>>>,
failed: Arc<AtomicBool>,
stop: CancellationToken,
) {
while !stop.is_cancelled() {
@@ -220,18 +197,8 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
Ok(response) if response.plaintext == expected => {
completed.fetch_add(1, Ordering::SeqCst);
}
Ok(_) => {
*failure.lock().expect("decrypt failure lock poisoned") =
Some("decrypt returned unexpected plaintext".to_string());
return;
}
Err(rustfs_kms::KmsError::BackendError { .. } | rustfs_kms::KmsError::OperationTimedOut { .. })
if allow_failover_errors.load(Ordering::SeqCst) =>
{
tokio::time::sleep(FAILOVER_ERROR_POLL_INTERVAL).await;
}
Err(error) => {
*failure.lock().expect("decrypt failure lock poisoned") = Some(error.to_string());
Ok(_) | Err(_) => {
failed.store(true, Ordering::SeqCst);
return;
}
}
@@ -329,9 +296,7 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
);
let stop = CancellationToken::new();
let allow_failover_errors = Arc::new(AtomicBool::new(false));
let kv2_failure = Arc::new(Mutex::new(None));
let transit_failure = Arc::new(Mutex::new(None));
let failed = Arc::new(AtomicBool::new(false));
let kv2_completed = Arc::new(AtomicU64::new(0));
let transit_completed = Arc::new(AtomicU64::new(0));
let kv2_worker = tokio::spawn(decrypt_loop(
@@ -339,8 +304,7 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
kv2_request,
kv2_data_key.plaintext_key,
Arc::clone(&kv2_completed),
Arc::clone(&allow_failover_errors),
Arc::clone(&kv2_failure),
Arc::clone(&failed),
stop.clone(),
));
let transit_worker = tokio::spawn(decrypt_loop(
@@ -348,21 +312,12 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
transit_request,
transit_data_key.plaintext_key,
Arc::clone(&transit_completed),
Arc::clone(&allow_failover_errors),
Arc::clone(&transit_failure),
Arc::clone(&failed),
stop.clone(),
));
wait_for_count(&kv2_completed, &kv2_failure, 2, "two healthy KV2 decrypts", HEALTHY_PROGRESS_TIMEOUT).await;
wait_for_count(
&transit_completed,
&transit_failure,
2,
"two healthy Transit decrypts",
HEALTHY_PROGRESS_TIMEOUT,
)
.await;
allow_failover_errors.store(true, Ordering::SeqCst);
wait_for_count(&kv2_completed, 2, "two healthy KV2 decrypts").await;
wait_for_count(&transit_completed, 2, "two healthy Transit decrypts").await;
std::fs::write(&marker, b"ready").expect("publish failover readiness marker");
wait_for_file(&elected, "the replacement Vault leader").await;
@@ -371,39 +326,18 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
let kv2_after_election = kv2_completed.load(Ordering::SeqCst) + 2;
let transit_after_election = transit_completed.load(Ordering::SeqCst) + 2;
wait_for_count(
&kv2_completed,
&kv2_failure,
kv2_after_election,
"post-failover KV2 decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(
&transit_completed,
&transit_failure,
transit_after_election,
"post-failover Transit decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(&kv2_completed, kv2_after_election, "post-failover KV2 decrypts").await;
wait_for_count(&transit_completed, transit_after_election, "post-failover Transit decrypts").await;
stop.cancel();
kv2_worker.await.expect("KV2 decrypt worker must join");
transit_worker.await.expect("Transit decrypt worker must join");
assert!(
kv2_failure.lock().expect("KV2 failure lock poisoned").is_none(),
"no KV2 decrypt may fail or return different plaintext"
);
assert!(
transit_failure.lock().expect("Transit failure lock poisoned").is_none(),
"no Transit decrypt may fail or return different plaintext"
);
assert!(!failed.load(Ordering::SeqCst), "no decrypt may fail or return different plaintext");
}
#[test]
#[ignore = "requires a real three-node Vault Raft cluster; run scripts/test/vault_ha_kms_live.sh"]
fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
@@ -415,6 +349,11 @@ fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
});
let snapshot = snapshotter.snapshot().into_vec();
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "circuit_open")]),
0,
"a bounded leader election must not open the circuit"
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]),
0,
+4 -33
View File
@@ -54,7 +54,6 @@ const ERR_LIFECYCLE_INVALID_EXPIRED_OBJECT_ALL_VERSIONS: &str =
"Days must be a positive integer and Date must not be specified inside Expiration with ExpiredObjectAllVersions";
const ERR_LIFECYCLE_INVALID_DEL_MARKER_EXPIRATION_DAYS: &str = "Days must be a positive integer with DelMarkerExpiration";
const ERR_LIFECYCLE_INVALID_RULE_ID_TOO_LONG: &str = "Rule ID must be at most 255 characters";
const ERR_LIFECYCLE_INVALID_RULE_ID_EMPTY: &str = "Rule ID must not be empty";
const ERR_LIFECYCLE_INVALID_RULE_STATUS: &str = "Rule status must be either Enabled or Disabled";
const ERR_LIFECYCLE_DEL_MARKER_WITH_TAGS: &str = "Rule with DelMarkerExpiration cannot have tags based filtering";
const ERR_LIFECYCLE_EXPIRED_OBJECT_DELETE_MARKER_WITH_TAGS: &str =
@@ -403,13 +402,10 @@ impl Lifecycle for BucketLifecycleConfiguration {
NoncurrentVersionTransitionOps::validate(transition)?;
}
}
if let Some(id) = &r.id {
if id.is_empty() {
return Err(std::io::Error::other(ERR_LIFECYCLE_INVALID_RULE_ID_EMPTY));
}
if id.len() > 255 {
return Err(std::io::Error::other(ERR_LIFECYCLE_INVALID_RULE_ID_TOO_LONG));
}
if let Some(id) = &r.id
&& id.len() > 255
{
return Err(std::io::Error::other(ERR_LIFECYCLE_INVALID_RULE_ID_TOO_LONG));
}
r.validate()?;
if let Some(object_lock_enabled) = lr.object_lock_enabled.as_ref()
@@ -3734,31 +3730,6 @@ mod tests {
.expect("empty prefix with filter should be valid");
}
#[tokio::test]
async fn validate_rejects_empty_rule_id() {
let lc = BucketLifecycleConfiguration {
expiry_updated_at: None,
rules: vec![LifecycleRule {
status: ExpirationStatus::from_static(ExpirationStatus::ENABLED),
expiration: Some(LifecycleExpiration {
days: Some(30),
..Default::default()
}),
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
filter: None,
id: Some(String::new()),
noncurrent_version_expiration: None,
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}],
};
let error = lc.validate(&ObjectLockConfiguration::default()).await.unwrap_err();
assert_eq!(error.to_string(), ERR_LIFECYCLE_INVALID_RULE_ID_EMPTY);
}
// --- TASK-004 tests: ExpiredObjectAllVersions ---
#[tokio::test]
+5 -50
View File
@@ -29,8 +29,7 @@ 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, SizeReconciliationEntry, SizeReconciliationScope, SizeSummary,
TierStats, hash_path, prefix_usage_in_cache,
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};
@@ -193,10 +192,6 @@ const MAX_DATA_USAGE_CACHE_DEPTH: usize = 1024;
pub trait ScannerSizeSummaryExt {
/// Fold one object's contribution into the summary, including its tier.
fn actions_accounting(&mut self, oi: &ObjectInfo, size: i64, actual_size: i64);
/// Fold counters and physical tier usage for an object whose metadata is
/// valid but whose logical size is currently unavailable. Logical totals
/// stay unchanged.
fn actions_accounting_unknown(&mut self, oi: &ObjectInfo);
}
impl ScannerSizeSummaryExt for SizeSummary {
@@ -230,34 +225,6 @@ impl ScannerSizeSummaryExt for SizeSummary {
});
}
}
fn actions_accounting_unknown(&mut self, oi: &ObjectInfo) {
if oi.delete_marker {
self.delete_markers = self.delete_markers.saturating_add(1);
return;
}
if oi.version_id.is_some_and(|v| !v.is_nil()) {
self.versions = self.versions.saturating_add(1);
}
if oi.transitioned_object.free_version {
return;
}
let tier = if oi.transitioned_object.status == TRANSITION_COMPLETE {
oi.transitioned_object.tier.clone()
} else {
oi.storage_class.clone().unwrap_or_else(|| storageclass::STANDARD.to_string())
};
if let Some(tier_stats) = self.tier_stats.get_mut(&tier) {
*tier_stats = tier_stats.add(&TierStats {
total_size: u64::try_from(oi.size).unwrap_or(0),
num_versions: 1,
num_objects: u64::from(oi.is_latest),
});
}
}
}
// ===== Cache-related data structures =====
@@ -377,10 +344,6 @@ pub struct DataUsageCacheInfo {
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
#[serde(default)]
pub cache_key_format: u16,
/// Bounded durable debts for versions whose logical size was not trusted.
/// The map key is an identity key, never a user-controlled metric label.
#[serde(default)]
pub size_reconciliation: HashMap<String, SizeReconciliationEntry>,
}
impl Serialize for DataUsageCacheInfo {
@@ -390,8 +353,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 field_count = 16 + usize::from(!self.size_reconciliation.is_empty());
let mut state = serializer.serialize_map(Some(field_count))?;
let mut state = serializer.serialize_map(Some(16))?;
state.serialize_entry("name", &self.name)?;
state.serialize_entry("next_cycle", &self.next_cycle)?;
state.serialize_entry("leader_epoch", &self.leader_epoch)?;
@@ -408,9 +370,6 @@ 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)?;
if !self.size_reconciliation.is_empty() {
state.serialize_entry("size_reconciliation", &self.size_reconciliation)?;
}
state.end()
}
}
@@ -469,18 +428,14 @@ impl DataUsageCache {
self.checked_flatten(name).is_some()
});
if !reusable {
let (pending_heals, size_reconciliation) = if self.info.name == name {
(
std::mem::take(&mut self.info.pending_heals),
std::mem::take(&mut self.info.size_reconciliation),
)
let pending_heals = if self.info.name == name {
std::mem::take(&mut self.info.pending_heals)
} else {
(Vec::new(), HashMap::new())
Vec::new()
};
*self = Self::default();
self.info.name = name.to_string();
self.info.pending_heals = pending_heals;
self.info.size_reconciliation = size_reconciliation;
}
self.info.next_cycle = next_cycle;
@@ -673,34 +673,6 @@ fn size_summary_actions_accounting_accumulates_tier_stats() {
);
}
#[test]
fn size_summary_unknown_accounting_keeps_physical_tier_and_version_only() {
let mut summary = SizeSummary::default();
summary
.tier_stats
.insert(storageclass::STANDARD.to_string(), TierStats::default());
let object = ObjectInfo {
size: 12,
storage_class: Some(storageclass::STANDARD.to_string()),
version_id: Some(uuid::Uuid::new_v4()),
is_latest: true,
..Default::default()
};
summary.actions_accounting_unknown(&object);
assert_eq!(summary.total_size, 0, "unknown logical size must not become zero or physical bytes");
assert_eq!(summary.versions, 1);
assert_eq!(
summary.tier_stats.get(storageclass::STANDARD),
Some(&TierStats {
total_size: 12,
num_versions: 1,
num_objects: 1,
})
);
}
#[test]
fn test_data_usage_entry_merge_sums_failed_objects() {
let mut left = DataUsageEntry {
@@ -1107,16 +1079,6 @@ fn data_usage_cache_prepare_for_scan_preserves_pending_heal_only_progress() {
scan_plan_digest: Some(TEST_PLAN_DIGEST),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
pending_heals: vec![pending_heal.clone()],
size_reconciliation: HashMap::from([(
"size-key".to_string(),
SizeReconciliationEntry {
key: "size-key".to_string(),
bucket: "bucket".to_string(),
object: "prefix/object".to_string(),
reason: "invalid_declared_size".to_string(),
..Default::default()
},
)]),
..Default::default()
},
..Default::default()
@@ -1126,7 +1088,6 @@ fn data_usage_cache_prepare_for_scan_preserves_pending_heal_only_progress() {
assert_eq!(outcome, DataUsageCachePrepareOutcome::Reused);
assert_eq!(cache.info.pending_heals, vec![pending_heal]);
assert!(cache.info.size_reconciliation.contains_key("size-key"));
assert!(cache.cache.is_empty());
assert!(!cache.info.snapshot_complete);
}
+8 -131
View File
@@ -20,9 +20,8 @@ use std::time::{Duration, Instant, SystemTime};
use crate::ReplTargetSizeSummary;
use crate::data_usage_define::{
DATA_USAGE_SCAN_CHECKPOINT_VERSION, DataUsageCache, DataUsageCacheInfo, DataUsageEntry, DataUsageHash, DataUsageHashMap,
DataUsageScanCheckpoint, DataUsageScanCheckpointReason, PendingScannerHeal, PendingScannerHealKind, ScannerSizeSummaryExt,
SizeReconciliationEntry, SizeSummary, hash_path,
DATA_USAGE_SCAN_CHECKPOINT_VERSION, DataUsageCache, DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageScanCheckpoint,
DataUsageScanCheckpointReason, PendingScannerHeal, PendingScannerHealKind, ScannerSizeSummaryExt, SizeSummary, hash_path,
};
use crate::error::ScannerError;
use crate::runtime_config::{
@@ -98,9 +97,6 @@ const METRIC_SCANNER_EXCESS_FOLDERS_TOTAL: &str = "rustfs_scanner_excess_folders
const METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL: &str = "rustfs_scanner_pending_heal_prune_total";
const METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL: &str = "rustfs_scanner_pending_heal_malformed_total";
const MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET: usize = 128;
const MAX_SIZE_RECONCILIATION_ENTRIES_PER_BUCKET: usize = 10_000;
const MAX_SIZE_RECONCILIATION_BYTES_PER_BUCKET: usize = 8 * 1024 * 1024;
const MAX_SIZE_RECONCILIATION_AGE_SECS: u64 = 7 * 24 * 60 * 60;
// --- scanner excess alerts as S3 notification events (rustfs/backlog#1868) --
//
@@ -368,7 +364,7 @@ impl PendingScannerAccounting<'_> {
fn apply(self, size_summary: &mut SizeSummary, cumulative_size: &mut i64, queued: bool) {
let size = if queued { self.expired_size } else { self.retained_size };
size_summary.actions_accounting(self.object, size, self.retained_size);
*cumulative_size = cumulative_size.saturating_add(size);
*cumulative_size += size;
}
}
@@ -675,65 +671,10 @@ pub struct FolderScanner {
skip_heal: Arc<std::sync::atomic::AtomicBool>,
local_disk: Arc<Disk>,
pending_heals_changed: bool,
pending_size_reconciliation_keys: HashSet<String>,
pending_size_reconciliation_scopes: HashSet<String>,
pending_size_reconciliation_truncated: bool,
#[cfg(test)]
list_path_raw_options_observer: Option<mpsc::UnboundedSender<ListPathRawTimeoutSnapshot>>,
}
fn size_reconciliation_entry_bytes(entry: &SizeReconciliationEntry) -> usize {
entry.key.len()
+ entry.bucket.len()
+ entry.object.len()
+ entry.version_id.as_deref().map_or(0, str::len)
+ entry.generation.as_deref().map_or(0, str::len)
+ entry.reason.len()
+ std::mem::size_of::<u64>()
+ std::mem::size_of::<u32>()
}
fn size_reconciliation_scope_key(bucket: &str, object: &str) -> String {
format!("{}:{}|{}:{}", bucket.len(), bucket, object.len(), object)
}
fn prune_size_reconciliation(info: &mut DataUsageCacheInfo, now: u64) {
info.size_reconciliation.retain(|key, entry| {
if entry.first_seen == 0 || entry.first_seen > now {
entry.first_seen = now;
}
key == &entry.key
&& entry.key.len() <= 4096
&& entry.bucket.len() <= 512
&& entry.object.len() <= 512
&& entry.version_id.as_deref().is_none_or(|value| value.len() <= 64)
&& entry.generation.as_deref().is_none_or(|value| value.len() <= 64)
&& entry.reason.len() <= 64
&& now.saturating_sub(entry.first_seen) <= MAX_SIZE_RECONCILIATION_AGE_SECS
});
while info.size_reconciliation.len() > MAX_SIZE_RECONCILIATION_ENTRIES_PER_BUCKET
|| info
.size_reconciliation
.values()
.map(size_reconciliation_entry_bytes)
.sum::<usize>()
> MAX_SIZE_RECONCILIATION_BYTES_PER_BUCKET
{
let oldest = info
.size_reconciliation
.iter()
.min_by(|(left_key, left), (right_key, right)| {
left.first_seen.cmp(&right.first_seen).then_with(|| left_key.cmp(right_key))
})
.map(|(key, _)| key.clone());
let Some(oldest) = oldest else {
break;
};
info.size_reconciliation.remove(&oldest);
}
}
impl FolderScanner {
fn now_secs() -> u64 {
SystemTime::now()
@@ -807,60 +748,6 @@ impl FolderScanner {
}
}
/// Apply the per-object size-resolution ledger updates in one place. The
/// scanner cache is the durable boundary; both working copies are updated
/// so an incremental publication cannot lose a debt or its resolution.
fn apply_size_reconciliation(&mut self, summary: &SizeSummary) {
let now = Self::now_secs();
self.pending_size_reconciliation_keys
.extend(summary.size_reconciliation.iter().map(|entry| entry.key.clone()));
self.pending_size_reconciliation_scopes.extend(
summary
.reconciliation_scopes
.iter()
.map(|scope| size_reconciliation_scope_key(&scope.bucket, &scope.object)),
);
self.pending_size_reconciliation_truncated |= summary.size_reconciliation_truncated;
for info in [&mut self.new_cache.info, &mut self.update_cache.info] {
for incoming in &summary.size_reconciliation {
if let Some(existing) = info.size_reconciliation.get_mut(&incoming.key) {
existing.reason = incoming.reason.clone();
existing.physical_size = incoming.physical_size;
existing.generation = incoming.generation.clone();
existing.version_id = incoming.version_id.clone();
existing.attempts = existing.attempts.saturating_add(1);
continue;
}
if size_reconciliation_entry_bytes(incoming) > MAX_SIZE_RECONCILIATION_BYTES_PER_BUCKET {
continue;
}
let mut entry = incoming.clone();
entry.first_seen = now;
entry.attempts = 1;
info.size_reconciliation.insert(entry.key.clone(), entry);
}
}
}
fn finish_size_reconciliation_batch(&mut self) {
let now = Self::now_secs();
let current_keys = std::mem::take(&mut self.pending_size_reconciliation_keys);
let scopes = std::mem::take(&mut self.pending_size_reconciliation_scopes);
let truncated = std::mem::replace(&mut self.pending_size_reconciliation_truncated, false);
for info in [&mut self.new_cache.info, &mut self.update_cache.info] {
if !truncated {
info.size_reconciliation.retain(|key, entry| {
!scopes.contains(&size_reconciliation_scope_key(&entry.bucket, &entry.object)) || current_keys.contains(key)
});
}
prune_size_reconciliation(info, now);
}
}
fn record_scan_resume_hint(&mut self, folder: &str) {
self.new_cache.info.scan_resume_after = Some(folder.to_string());
self.update_cache.info.scan_resume_after = Some(folder.to_string());
@@ -1488,13 +1375,14 @@ impl FolderScanner {
// Single-flight (backlog#1894 axis A) — the
// recording mode and its guarantees are pinned by
// corrupt_metadata_recording below.
let mrf_accepted = rustfs_common::mrf_channel::try_send_mrf_intent(
let mrf_result = rustfs_common::mrf_channel::try_send_mrf_intent_typed(
rustfs_common::mrf_channel::MrfKind::MetadataCorruption,
&item.bucket,
&object,
None,
None,
);
match corrupt_metadata_recording(mrf_accepted) {
match corrupt_metadata_recording(mrf_result) {
CorruptMetadataRecording::LedgerOnly => {
// Recorded as Full (retry-later): admission
// for this target happens in the MRF
@@ -1539,7 +1427,6 @@ impl FolderScanner {
abandoned_children.remove(&path_join_buf(&[&item.bucket, &item.object_path()]));
apply_scanner_size_summary(into, &sz);
self.apply_size_reconciliation(&sz);
into.objects += 1;
object_count += 1;
self.budget.record_object_scanned();
@@ -2219,7 +2106,6 @@ impl FolderScanner {
}
}
self.finish_size_reconciliation_batch();
done_folder();
let scanned_objects = u64::try_from(into.objects).unwrap_or(u64::MAX);
emit_scanner_folder_trace(&self.root, &folder.name, scanned_objects, trace_started_at, "completed");
@@ -2305,17 +2191,10 @@ pub async fn scan_data_folder(
skip_heal,
local_disk,
pending_heals_changed: false,
pending_size_reconciliation_keys: HashSet::new(),
pending_size_reconciliation_scopes: HashSet::new(),
pending_size_reconciliation_truncated: false,
#[cfg(test)]
list_path_raw_options_observer: None,
};
let now = FolderScanner::now_secs();
prune_size_reconciliation(&mut scanner.new_cache.info, now);
prune_size_reconciliation(&mut scanner.update_cache.info, now);
// Check if context is cancelled
if ctx.is_cancelled() {
return Err(ScannerError::Other("Operation cancelled".to_string()));
@@ -2339,9 +2218,7 @@ pub async fn scan_data_folder(
new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN);
new_cache.info.last_update = Some(SystemTime::now());
new_cache.info.next_cycle = cache.info.next_cycle;
let unresolved_objects = root.failed_objects > 0
|| !new_cache.info.failed_objects.is_empty()
|| !new_cache.info.size_reconciliation.is_empty();
let unresolved_objects = root.failed_objects > 0 || !new_cache.info.failed_objects.is_empty();
new_cache.info.snapshot_complete = !unresolved_objects;
let had_scan_checkpoint = cache.info.scan_checkpoint.is_some() || new_cache.info.scan_checkpoint.is_some();
new_cache.info.scan_resume_after = None;
@@ -2369,7 +2246,7 @@ pub async fn scan_data_folder(
if root_has_progress {
new_cache.replace_hashed(&root_hash, &None, &root);
}
if partial_cache_is_useful(&root, pending_heals_changed) || !new_cache.info.size_reconciliation.is_empty() {
if partial_cache_is_useful(&root, pending_heals_changed) {
if new_cache.root().is_some() {
new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN);
}
+71 -814
View File
@@ -13,7 +13,6 @@
// limitations under the License.
/// Per-object scan actions: ScannerItem, the get-size failure policy, and the heal/ILM admission helpers.
use super::*;
use sha2::{Digest as _, Sha256};
/// Cached folder information for scanning
#[derive(Clone, Debug)]
@@ -33,263 +32,6 @@ pub(super) enum GetSizeFailureAction {
HealMetadata { object: String },
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum SizeResolutionReason {
CompressedSizeUnknown,
InvalidPhysicalSize,
UnsupportedCompression,
InvalidObjectSize,
InvalidPartSize,
InvalidDeclaredSize,
SizeOverflowOrMismatch,
}
impl SizeResolutionReason {
fn as_str(self) -> &'static str {
match self {
Self::CompressedSizeUnknown => "compressed_size_unknown",
Self::InvalidPhysicalSize => "invalid_physical_size",
Self::UnsupportedCompression => "unsupported_compression",
Self::InvalidObjectSize => "invalid_object_size",
Self::InvalidPartSize => "invalid_part_size",
Self::InvalidDeclaredSize => "invalid_declared_size",
Self::SizeOverflowOrMismatch => "size_overflow_or_mismatch",
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(super) enum SizeResolution {
Known { logical: i64, physical: i64 },
Unknown { physical: i64, reason: SizeResolutionReason },
Corrupt { physical: i64, reason: SizeResolutionReason },
}
impl SizeResolution {
fn known_size(&self) -> Option<i64> {
match self {
Self::Known { logical, .. } => Some(*logical),
Self::Unknown { .. } | Self::Corrupt { .. } => None,
}
}
}
fn size_reconciliation_key(oi: &ObjectInfo, reason: SizeResolutionReason) -> String {
let version = oi
.version_id
.filter(|version| !version.is_nil())
.map(|version| version.to_string())
.unwrap_or_default();
let generation = oi
.data_dir
.filter(|generation| !generation.is_nil())
.map(|generation| generation.to_string())
.unwrap_or_default();
// Length-prefix each component so an object key containing the separator
// cannot alias another identity. S3 keys are bounded in normal operation;
// oversized persisted values use a digest so a corrupt metadata record
// cannot grow the ledger without bound.
fn component(value: &str) -> String {
const MAX_COMPONENT_LEN: usize = 512;
if value.len() <= MAX_COMPONENT_LEN {
return format!("{}:{}", value.len(), value);
}
let digest = Sha256::digest(value.as_bytes());
let digest = hex_simd::encode_to_string(digest, hex_simd::AsciiCase::Lower);
format!("hash:{}:{}", value.len(), digest)
}
format!(
"{}|{}|{}|{}|{}",
component(&oi.bucket),
component(&oi.name),
component(&version),
component(&generation),
component(reason.as_str())
)
}
pub(super) fn bounded_reconciliation_field(value: &str) -> String {
const MAX_FIELD_LEN: usize = 512;
if value.len() <= MAX_FIELD_LEN {
return value.to_string();
}
let digest = hex_simd::encode_to_string(Sha256::digest(value.as_bytes()), hex_simd::AsciiCase::Lower);
let prefix_len = MAX_FIELD_LEN - 65;
let prefix = value
.char_indices()
.take_while(|(offset, ch)| offset.saturating_add(ch.len_utf8()) <= prefix_len)
.map(|(_, ch)| ch)
.collect::<String>();
format!("{}~{}", prefix, digest)
}
fn record_size_resolution(summary: &mut SizeSummary, oi: &ObjectInfo, resolution: &SizeResolution) {
match resolution {
SizeResolution::Known { .. } => {}
SizeResolution::Unknown { physical, reason } | SizeResolution::Corrupt { physical, reason } => {
summary.record_size_reconciliation(SizeReconciliationEntry {
key: size_reconciliation_key(oi, *reason),
bucket: bounded_reconciliation_field(&oi.bucket),
object: bounded_reconciliation_field(&oi.name),
version_id: oi
.version_id
.filter(|version| !version.is_nil())
.map(|version| version.to_string()),
generation: oi
.data_dir
.filter(|generation| !generation.is_nil())
.map(|generation| generation.to_string()),
reason: reason.as_str().to_string(),
physical_size: u64::try_from(*physical).ok(),
first_seen: 0,
attempts: 0,
});
}
}
}
/// Resolve the size metadata once at the scanner trust boundary. A compressed
/// -1 sentinel is valid legacy metadata, but it cannot participate in normal
/// logical-size accounting or size-filtered lifecycle rules.
pub(super) fn resolve_size(oi: &ObjectInfo) -> SizeResolution {
let physical = oi.size;
if physical < 0 {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::InvalidPhysicalSize,
};
}
let compressed = match oi.compression_read_plan() {
Ok((_, _, compressed)) => compressed,
Err(_) => {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::UnsupportedCompression,
};
}
};
if oi.actual_size < -1 || (oi.actual_size == -1 && !compressed) {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::InvalidObjectSize,
};
}
// Match ObjectInfo::get_actual_size: a positive in-memory value is the
// authoritative decoded size. Stale declared/part metadata must not turn
// an otherwise valid object into a false corruption report.
if oi.actual_size > 0 {
return SizeResolution::Known {
logical: oi.actual_size,
physical,
};
}
if oi
.parts
.iter()
.any(|part| part.actual_size < -1 || (part.actual_size < 0 && !compressed))
{
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::InvalidPartSize,
};
}
let declared = rustfs_utils::http::get_str(&oi.user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE);
let declared = match declared {
Some(value) if value.is_empty() => {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::InvalidDeclaredSize,
};
}
Some(value) => match value.parse::<i64>() {
Ok(value) if value >= 0 => Some(value),
_ => {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::InvalidDeclaredSize,
};
}
},
None => None,
};
let logical = match oi.get_actual_size() {
Ok(size) if size == -1 && compressed && declared.is_none() => {
return SizeResolution::Unknown {
physical,
reason: SizeResolutionReason::CompressedSizeUnknown,
};
}
Ok(size) if size >= 0 => size,
Ok(_) | Err(_) => {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::SizeOverflowOrMismatch,
};
}
};
if compressed && logical == 0 && physical != 0 && oi.parts.is_empty() && declared.is_none() {
return SizeResolution::Corrupt {
physical,
reason: SizeResolutionReason::SizeOverflowOrMismatch,
};
}
SizeResolution::Known { logical, physical }
}
fn resolve_sizes(object_infos: &[ObjectInfo]) -> Vec<SizeResolution> {
object_infos.iter().map(resolve_size).collect()
}
fn lifecycle_rule_has_size_filter(lifecycle: &BucketLifecycleConfiguration, rule_id: &str) -> bool {
let filter_has_size = |filter: &s3s::dto::LifecycleRuleFilter| {
filter.object_size_greater_than.is_some()
|| filter.object_size_less_than.is_some()
|| filter
.and
.as_ref()
.is_some_and(|and| and.object_size_greater_than.is_some() || and.object_size_less_than.is_some())
};
lifecycle
.rules
.iter()
.find(|rule| {
if rule_id.is_empty() {
rule.id.as_deref().is_none_or(str::is_empty)
} else {
rule.id.as_deref() == Some(rule_id)
}
})
.and_then(|rule| rule.filter.as_ref())
.is_some_and(filter_has_size)
}
fn lifecycle_event_allowed(resolution: &SizeResolution, event: &Event, lifecycle: &BucketLifecycleConfiguration) -> bool {
match resolution {
// Missing or invalid logical size only defers actions whose selected
// rule actually depends on that size. Time/version-only actions retain
// their existing semantics, including intrinsic events without a rule ID.
SizeResolution::Unknown { .. } | SizeResolution::Corrupt { .. } => {
!lifecycle_rule_has_size_filter(lifecycle, &event.rule_id)
}
SizeResolution::Known { .. } => true,
}
}
/// A successful newer-noncurrent batch consumes both known and unresolved
/// versions from the retained-version alert count. The two accounting paths
/// are separate because only known sizes can contribute byte totals.
fn remaining_versions_after_queued_noncurrent(remaining_versions: usize, known_count: usize, unknown_count: usize) -> usize {
remaining_versions.saturating_sub(known_count.saturating_add(unknown_count))
}
/// How the corrupt-metadata branch records the repair after attempting an
/// MRF intent (backlog#1894 axis A).
#[derive(Debug, PartialEq, Eq)]
@@ -308,11 +50,12 @@ pub(super) enum CorruptMetadataRecording {
ImmediateAndLedger,
}
pub(super) fn corrupt_metadata_recording(mrf_accepted: bool) -> CorruptMetadataRecording {
if mrf_accepted {
CorruptMetadataRecording::LedgerOnly
} else {
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,
}
}
@@ -577,48 +320,34 @@ impl ScannerItem {
"Scanner lifecycle evaluation started"
);
let resolved_sizes = resolve_sizes(&object_infos);
if let Some(first) = object_infos.first() {
size_summary.record_reconciliation_scope(
&bounded_reconciliation_field(&first.bucket),
&bounded_reconciliation_field(&first.name),
);
}
for (oi, resolution) in object_infos.iter().zip(resolved_sizes.iter()) {
record_size_resolution(size_summary, oi, resolution);
}
let has_corrupt_size = resolved_sizes
.iter()
.any(|resolution| matches!(resolution, SizeResolution::Corrupt { .. }));
// `versioning_config` is resolved once per object by the caller
// (`get_size`) and handed in; only `prefix_enabled` is consulted here.
let Some(lifecycle) = self.lifecycle.clone() else {
let mut cumulative_size: i64 = 0;
for (oi, resolved_size) in object_infos.iter().zip(resolved_sizes.iter()) {
let accounting_size = match resolved_size {
SizeResolution::Known { logical, .. } => *logical,
// A valid compressed legacy sentinel has no logical size,
// but heal and replication still need to run. The
// physical size is only an input to those operations; it
// is not folded into the logical total below.
SizeResolution::Unknown { physical, .. } => {
self.heal_actions(oi, *physical, size_summary).await;
size_summary.actions_accounting_unknown(oi);
continue;
}
SizeResolution::Corrupt { .. } => {
size_summary.actions_accounting_unknown(oi);
let Some(lifecycle) = self.lifecycle.as_ref() else {
let mut cumulative_size = 0;
for oi in object_infos.iter() {
let actual_size = match oi.get_actual_size() {
Ok(size) => size,
Err(_) => {
warn!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_LIFECYCLE_ACTION,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %self.bucket,
object = %oi.name,
state = "size_lookup_failed",
"Scanner lifecycle action used fallback size"
);
continue;
}
};
let size = self.heal_actions(oi, accounting_size, size_summary).await;
let size = self.heal_actions(oi, actual_size, size_summary).await;
size_summary.actions_accounting(oi, size, accounting_size);
size_summary.actions_accounting(oi, size, actual_size);
cumulative_size = cumulative_size.saturating_add(size);
cumulative_size += size;
}
self.alert_excessive_versions(object_infos.len(), cumulative_size);
@@ -672,108 +401,25 @@ impl ScannerItem {
let mut to_delete_objs: Vec<ObjectToDelete> = Vec::new();
let mut noncurrent_events: Vec<Event> = Vec::new();
let mut noncurrent_accounting: Vec<PendingScannerAccounting<'_>> = Vec::new();
let mut noncurrent_unknown: Vec<&ObjectInfo> = Vec::new();
let mut cumulative_size = 0;
let mut remaining_versions = object_infos.len();
'eventLoop: {
for (i, event) in events.iter().enumerate() {
let oi = &object_infos[i];
let known_size = resolved_sizes[i].known_size();
if has_corrupt_size
&& matches!(
event.action,
IlmAction::DeleteAllVersionsAction | IlmAction::DelMarkerDeleteAllVersionsAction
)
{
// An all-version delete would also remove a corrupt
// sibling that could not be reconciled safely.
continue;
}
if !lifecycle_event_allowed(&resolved_sizes[i], event, &lifecycle) {
// An unknown logical size must not make an otherwise
// non-destructive scan disappear from heal/physical-tier
// accounting. Size-filtered or deferred events remain
// pending, so retain the version-only physical counters.
if let SizeResolution::Unknown { physical, .. } = &resolved_sizes[i] {
self.heal_actions(oi, *physical, size_summary).await;
size_summary.actions_accounting_unknown(oi);
}
continue;
}
let actual_size = match known_size {
Some(size) => size,
None => {
match event.action {
IlmAction::DeleteAction
| IlmAction::DeleteRestoredAction
| IlmAction::DeleteRestoredVersionAction
| IlmAction::DeleteAllVersionsAction
| IlmAction::DelMarkerDeleteAllVersionsAction => {
let done_ilm = Metrics::time_ilm(event.action);
let trace_started_at = trace_start_instant();
let queued = apply_expiry_rule(event, &LcEventSrc::Scanner, oi).await;
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
done_ilm(1)();
if event.action == IlmAction::DeleteAllVersionsAction
|| event.action == IlmAction::DelMarkerDeleteAllVersionsAction
{
remaining_versions = 0;
}
} else if matches!(
event.action,
IlmAction::DeleteAction
| IlmAction::DeleteRestoredAction
| IlmAction::DeleteRestoredVersionAction
) {
size_summary.actions_accounting_unknown(oi);
} else {
size_summary.actions_accounting_unknown(oi);
for (j, retained) in object_infos.iter().enumerate().skip(i + 1) {
match &resolved_sizes[j] {
SizeResolution::Known { logical, .. } => PendingScannerAccounting {
object: retained,
retained_size: *logical,
expired_size: 0,
}
.apply(size_summary, &mut cumulative_size, false),
SizeResolution::Unknown { .. } => {
size_summary.actions_accounting_unknown(retained);
}
SizeResolution::Corrupt { .. } => {}
}
}
}
}
IlmAction::DeleteVersionAction => {
if let Some(opt) = object_opts.get(i) {
to_delete_objs.push(ObjectToDelete {
object_name: opt.name.clone(),
version_id: opt.version_id,
..Default::default()
});
noncurrent_events.push(event.clone());
noncurrent_unknown.push(oi);
}
}
IlmAction::TransitionAction | IlmAction::TransitionVersionAction => {
let trace_started_at = trace_start_instant();
let queued = apply_transition_rule(event, &LcEventSrc::Scanner, oi).await;
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
let done_ilm = Metrics::time_ilm(event.action);
done_ilm(1)();
}
size_summary.actions_accounting_unknown(oi);
}
IlmAction::NoneAction | IlmAction::ActionCount => {
if let SizeResolution::Unknown { physical, .. } = &resolved_sizes[i] {
self.heal_actions(oi, *physical, size_summary).await;
}
size_summary.actions_accounting_unknown(oi);
}
}
continue;
let actual_size = match oi.get_actual_size() {
Ok(size) => size,
Err(_) => {
warn!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_LIFECYCLE_ACTION,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %self.bucket,
object = %oi.name,
state = "size_lookup_failed",
"Scanner lifecycle action used fallback size"
);
0
}
};
@@ -801,24 +447,36 @@ impl ScannerItem {
done_ilm(1)();
remaining_versions = 0;
} else {
if let Some(actual_size) = known_size {
PendingScannerAccounting {
object: oi,
retained_size: actual_size,
expired_size: 0,
}
.apply(size_summary, &mut cumulative_size, false);
for retained in object_infos.iter().skip(i + 1) {
let retained_size = match retained.get_actual_size() {
Ok(size) => size,
Err(_) => {
warn!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_LIFECYCLE_ACTION,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_LIFECYCLE,
bucket = %self.bucket,
object = %retained.name,
state = "size_lookup_failed",
"Scanner lifecycle action used fallback size"
);
0
}
};
PendingScannerAccounting {
object: oi,
retained_size: actual_size,
object: retained,
retained_size,
expired_size: 0,
}
.apply(size_summary, &mut cumulative_size, false);
}
for (j, retained) in object_infos.iter().enumerate().skip(i + 1) {
if let Some(retained_size) = resolved_sizes[j].known_size() {
PendingScannerAccounting {
object: retained,
retained_size,
expired_size: 0,
}
.apply(size_summary, &mut cumulative_size, false);
}
}
}
break 'eventLoop;
}
@@ -854,13 +512,11 @@ impl ScannerItem {
version_id: opt.version_id,
..Default::default()
});
if let Some(actual_size) = known_size {
noncurrent_accounting.push(PendingScannerAccounting {
object: oi,
retained_size: actual_size,
expired_size: 0,
});
}
noncurrent_accounting.push(PendingScannerAccounting {
object: oi,
retained_size: actual_size,
expired_size: 0,
});
account_now = false;
}
noncurrent_events.push(event.clone());
@@ -893,7 +549,7 @@ impl ScannerItem {
if account_now {
size_summary.actions_accounting(oi, size, actual_size);
cumulative_size = cumulative_size.saturating_add(size);
cumulative_size += size;
}
}
}
@@ -921,20 +577,11 @@ impl ScannerItem {
}
if record_scanner_ilm_action_if_queued(global_metrics(), action, count, queued) {
done_ilm(count)();
remaining_versions = remaining_versions_after_queued_noncurrent(
remaining_versions,
noncurrent_accounting.len(),
noncurrent_unknown.len(),
);
remaining_versions = remaining_versions.saturating_sub(noncurrent_accounting.len());
}
for pending in noncurrent_accounting {
pending.apply(size_summary, &mut cumulative_size, queued);
}
if !queued {
for object in noncurrent_unknown {
size_summary.actions_accounting_unknown(object);
}
}
}
self.alert_excessive_versions(remaining_versions, cumulative_size);
}
@@ -1283,394 +930,4 @@ mod tests {
assert_eq!(item.object_name, "object");
assert_eq!(item.object_path(), "object");
}
#[test]
fn size_resolution_rejects_negative_overflow_and_unknown_compression() {
let compressed = |actual_size: i64, declared: Option<&str>| {
let mut user_defined = HashMap::new();
rustfs_utils::http::insert_str(&mut user_defined, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
if let Some(declared) = declared {
rustfs_utils::http::insert_str(&mut user_defined, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, declared.to_string());
}
ObjectInfo {
size: 12,
actual_size,
user_defined: Arc::new(user_defined),
..Default::default()
}
};
let normal = ObjectInfo {
size: 12,
actual_size: 10,
..Default::default()
};
assert_eq!(
resolve_size(&normal),
SizeResolution::Known {
logical: 10,
physical: 12
}
);
let stale_declared_metadata = ObjectInfo {
size: 12,
actual_size: 10,
user_defined: Arc::new(HashMap::from([("x-rustfs-internal-actual-size".to_string(), "not-a-size".to_string())])),
parts: Arc::new(vec![rustfs_filemeta::ObjectPartInfo {
actual_size: -2,
..Default::default()
}]),
..Default::default()
};
assert_eq!(
resolve_size(&stale_declared_metadata),
SizeResolution::Known {
logical: 10,
physical: 12
}
);
assert_eq!(
resolve_size(&compressed(0, Some("9"))),
SizeResolution::Known {
logical: 9,
physical: 12
}
);
assert_eq!(
resolve_size(&compressed(-1, None)),
SizeResolution::Unknown {
physical: 12,
reason: SizeResolutionReason::CompressedSizeUnknown,
}
);
assert!(matches!(
resolve_size(&compressed(0, Some("not-a-size"))),
SizeResolution::Corrupt {
reason: SizeResolutionReason::InvalidDeclaredSize,
..
}
));
assert!(matches!(
resolve_size(&ObjectInfo {
size: 12,
actual_size: -2,
..Default::default()
}),
SizeResolution::Corrupt { .. }
));
assert!(matches!(resolve_size(&compressed(0, Some("-1"))), SizeResolution::Corrupt { .. }));
assert!(matches!(resolve_size(&compressed(0, Some(""))), SizeResolution::Corrupt { .. }));
let unsupported = {
let mut object = compressed(0, None);
let mut metadata = (*object.user_defined).clone();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "unsupported".to_string());
object.user_defined = Arc::new(metadata);
object
};
assert!(matches!(resolve_size(&unsupported), SizeResolution::Corrupt { .. }));
let invalid_part = {
let mut object = compressed(0, None);
object.parts = Arc::new(vec![rustfs_filemeta::ObjectPartInfo {
size: 12,
actual_size: -2,
..Default::default()
}]);
object
};
assert!(matches!(resolve_size(&invalid_part), SizeResolution::Corrupt { .. }));
let overflow = {
let mut object = compressed(0, None);
object.parts = Arc::new(vec![
rustfs_filemeta::ObjectPartInfo {
size: 1,
actual_size: i64::MAX,
..Default::default()
},
rustfs_filemeta::ObjectPartInfo {
size: 1,
actual_size: 1,
..Default::default()
},
]);
object
};
assert!(matches!(resolve_size(&overflow), SizeResolution::Corrupt { .. }));
let mismatch = compressed(0, None);
assert!(matches!(resolve_size(&mismatch), SizeResolution::Corrupt { .. }));
assert_eq!(
resolve_size(&ObjectInfo {
size: 0,
actual_size: 0,
..Default::default()
}),
SizeResolution::Known { logical: 0, physical: 0 }
);
}
#[test]
fn size_resolution_records_and_replays_one_identity() {
let version_id = uuid::Uuid::new_v4();
let generation = uuid::Uuid::new_v4();
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, "not-a-number".to_string());
let corrupt = ObjectInfo {
bucket: "bucket".to_string(),
name: "object".to_string(),
size: 12,
version_id: Some(version_id),
data_dir: Some(generation),
user_defined: Arc::new(metadata),
..Default::default()
};
let mut summary = SizeSummary::default();
let resolution = resolve_size(&corrupt);
record_size_resolution(&mut summary, &corrupt, &resolution);
record_size_resolution(&mut summary, &corrupt, &resolution);
assert_eq!(summary.size_reconciliation.len(), 1);
assert_eq!(summary.size_reconciliation[0].reason, "invalid_declared_size");
assert_eq!(summary.size_reconciliation[0].physical_size, Some(12));
let known = ObjectInfo {
actual_size: 12,
user_defined: Arc::new(HashMap::new()),
..corrupt.clone()
};
record_size_resolution(&mut summary, &known, &resolve_size(&known));
summary.record_reconciliation_scope(&known.bucket, &known.name);
assert_eq!(summary.reconciliation_scopes.len(), 1);
assert_eq!(summary.reconciliation_scopes[0].bucket, "bucket");
}
#[test]
fn malformed_size_has_same_ilm_accounting() {
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_ACTUAL_SIZE, "invalid".to_string());
let object = ObjectInfo {
bucket: "bucket".to_string(),
name: "object".to_string(),
size: 12,
user_defined: Arc::new(metadata),
..Default::default()
};
let resolution = resolve_size(&object);
let mut without_ilm = SizeSummary::default();
let mut with_ilm = SizeSummary::default();
record_size_resolution(&mut without_ilm, &object, &resolution);
record_size_resolution(&mut with_ilm, &object, &resolution);
assert_eq!(without_ilm.size_reconciliation, with_ilm.size_reconciliation);
assert_eq!(without_ilm.total_size, 0);
assert_eq!(with_ilm.total_size, 0);
assert!(without_ilm.tier_stats.is_empty());
assert!(with_ilm.tier_stats.is_empty());
}
#[test]
fn size_resolution_parses_once_per_version() {
let objects = vec![
ObjectInfo {
bucket: "bucket".to_string(),
name: "one".to_string(),
size: 1,
actual_size: 1,
..Default::default()
},
ObjectInfo {
bucket: "bucket".to_string(),
name: "two".to_string(),
size: 2,
actual_size: -2,
..Default::default()
},
];
let resolutions = resolve_sizes(&objects);
assert_eq!(resolutions.len(), objects.len());
assert!(matches!(resolutions[0], SizeResolution::Known { logical: 1, .. }));
assert!(matches!(resolutions[1], SizeResolution::Corrupt { .. }));
}
#[test]
fn queued_unknown_noncurrent_versions_are_removed_from_alert_count() {
assert_eq!(remaining_versions_after_queued_noncurrent(3, 1, 2), 0);
assert_eq!(remaining_versions_after_queued_noncurrent(7, 2, 1), 4);
assert_eq!(remaining_versions_after_queued_noncurrent(usize::MAX, usize::MAX, usize::MAX), 0);
}
#[test]
fn malformed_size_blocks_size_dependent_transition_but_allows_time_only_expiry() {
let size_filtered = BucketLifecycleConfiguration {
rules: vec![s3s::dto::LifecycleRule {
status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED),
expiration: None,
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
id: Some("size".to_string()),
filter: Some(s3s::dto::LifecycleRuleFilter {
object_size_greater_than: Some(1),
..Default::default()
}),
noncurrent_version_expiration: None,
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}],
..Default::default()
};
let unknown = SizeResolution::Unknown {
physical: 12,
reason: SizeResolutionReason::CompressedSizeUnknown,
};
let size_event = Event {
action: IlmAction::DeleteAction,
rule_id: "size".to_string(),
..Default::default()
};
assert!(!lifecycle_event_allowed(&unknown, &size_event, &size_filtered));
assert!(!lifecycle_event_allowed(
&unknown,
&Event {
action: IlmAction::TransitionAction,
rule_id: "size".to_string(),
..Default::default()
},
&size_filtered
));
let mixed_filters = BucketLifecycleConfiguration {
rules: vec![
size_filtered.rules[0].clone(),
s3s::dto::LifecycleRule {
status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED),
expiration: None,
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
id: Some("time".to_string()),
filter: None,
noncurrent_version_expiration: None,
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
},
],
..Default::default()
};
assert!(lifecycle_event_allowed(
&unknown,
&Event {
action: IlmAction::DeleteAction,
rule_id: "time".to_string(),
..Default::default()
},
&mixed_filters
));
assert!(lifecycle_event_allowed(
&unknown,
&Event {
action: IlmAction::TransitionAction,
..Default::default()
},
&BucketLifecycleConfiguration::default()
));
assert!(lifecycle_event_allowed(
&SizeResolution::Corrupt {
physical: 12,
reason: SizeResolutionReason::InvalidDeclaredSize,
},
&Event {
action: IlmAction::DeleteAction,
..Default::default()
},
&BucketLifecycleConfiguration::default()
));
assert!(!lifecycle_event_allowed(
&SizeResolution::Corrupt {
physical: 12,
reason: SizeResolutionReason::InvalidDeclaredSize,
},
&Event {
action: IlmAction::DeleteAction,
rule_id: "size".to_string(),
..Default::default()
},
&size_filtered
));
assert!(lifecycle_rule_has_size_filter(
&BucketLifecycleConfiguration {
rules: vec![s3s::dto::LifecycleRule {
status: s3s::dto::ExpirationStatus::from_static(s3s::dto::ExpirationStatus::ENABLED),
expiration: None,
abort_incomplete_multipart_upload: None,
del_marker_expiration: None,
id: None,
filter: Some(s3s::dto::LifecycleRuleFilter {
object_size_greater_than: Some(1),
..Default::default()
}),
noncurrent_version_expiration: None,
noncurrent_version_transitions: None,
prefix: None,
transitions: None,
}],
..Default::default()
},
""
));
assert!(lifecycle_event_allowed(
&SizeResolution::Known {
logical: 10,
physical: 12,
},
&Event {
action: IlmAction::DeleteAllVersionsAction,
..Default::default()
},
&BucketLifecycleConfiguration::default()
));
assert!(lifecycle_event_allowed(
&unknown,
&Event {
action: IlmAction::DeleteAction,
rule_id: "time-only".to_string(),
..Default::default()
},
&BucketLifecycleConfiguration::default()
));
}
#[tokio::test]
async fn long_object_size_reconciliation_scope_uses_bounded_identity() {
let object_name = "o".repeat(600);
let mut item = scanner_item_with_prefix("");
item.object_name = object_name.clone();
let mut metadata = HashMap::new();
rustfs_utils::http::insert_str(&mut metadata, rustfs_utils::http::SUFFIX_COMPRESSION, "zstd".to_string());
let object = ObjectInfo {
bucket: item.bucket.clone(),
name: object_name.clone(),
size: 12,
actual_size: -1,
version_id: Some(uuid::Uuid::new_v4()),
user_defined: Arc::new(metadata),
..Default::default()
};
let mut summary = SizeSummary::default();
item.apply_actions(vec![object], None, VersioningConfiguration::default(), &mut summary)
.await;
let bounded_bucket = bounded_reconciliation_field(&item.bucket);
let bounded_object = bounded_reconciliation_field(&object_name);
assert_eq!(summary.reconciliation_scopes[0].bucket, bounded_bucket);
assert_eq!(summary.reconciliation_scopes[0].object, bounded_object);
assert_eq!(summary.size_reconciliation[0].object, bounded_object);
assert_eq!(summary.versions, 1);
assert_eq!(summary.total_size, 0);
}
}
+14 -65
View File
@@ -48,8 +48,20 @@ 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(true), CorruptMetadataRecording::LedgerOnly);
assert_eq!(corrupt_metadata_recording(false), CorruptMetadataRecording::ImmediateAndLedger);
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
);
}
fn cooldown_map_len() -> usize {
@@ -326,9 +338,6 @@ async fn build_test_scanner() -> (FolderScanner, std::path::PathBuf) {
skip_heal: Arc::new(AtomicBool::new(false)),
local_disk: disk,
pending_heals_changed: false,
pending_size_reconciliation_keys: HashSet::new(),
pending_size_reconciliation_scopes: HashSet::new(),
pending_size_reconciliation_truncated: false,
list_path_raw_options_observer: None,
};
@@ -391,66 +400,6 @@ async fn test_record_failed_ttl_zero_noop() {
assert!(!scanner.should_skip_failed("path2"));
}
#[tokio::test]
async fn malformed_size_reconciliation_replays_after_restart() {
let (mut scanner, temp_dir) = build_test_scanner().await;
let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir);
let entry = SizeReconciliationEntry {
key: "1:b|6:object|0:|0:".to_string(),
bucket: "b".to_string(),
object: "object".to_string(),
reason: "invalid_declared_size".to_string(),
physical_size: Some(12),
..Default::default()
};
let mut summary = SizeSummary::default();
summary.record_size_reconciliation(entry.clone());
summary.record_reconciliation_scope("b", "object");
scanner.apply_size_reconciliation(&summary);
scanner.apply_size_reconciliation(&summary);
assert_eq!(scanner.new_cache.info.size_reconciliation.len(), 1);
assert_eq!(scanner.update_cache.info.size_reconciliation.len(), 1);
assert_eq!(scanner.new_cache.info.size_reconciliation[&entry.key].attempts, 2);
let encoded = rmp_serde::to_vec_named(&scanner.new_cache.info).expect("size ledger should encode");
let decoded: crate::data_usage_define::DataUsageCacheInfo =
rmp_serde::from_slice(&encoded).expect("size ledger should decode");
assert_eq!(decoded.size_reconciliation.len(), 1);
assert_eq!(decoded.size_reconciliation[&entry.key].reason, "invalid_declared_size");
let mut resolved = SizeSummary::default();
resolved.record_reconciliation_scope("b", "object");
scanner.apply_size_reconciliation(&resolved);
assert!(scanner.new_cache.info.size_reconciliation.is_empty());
assert!(scanner.update_cache.info.size_reconciliation.is_empty());
}
#[tokio::test]
async fn malformed_size_reconciliation_clears_bounded_long_object_scope() {
let (mut scanner, temp_dir) = build_test_scanner().await;
let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir);
let long_object = "o".repeat(600);
let bounded_object = item_actions::bounded_reconciliation_field(&long_object);
let entry = SizeReconciliationEntry {
key: "long-object-key".to_string(),
bucket: "b".to_string(),
object: bounded_object,
reason: "invalid_declared_size".to_string(),
..Default::default()
};
let mut summary = SizeSummary::default();
summary.record_size_reconciliation(entry);
scanner.apply_size_reconciliation(&summary);
assert_eq!(scanner.new_cache.info.size_reconciliation.len(), 1);
let mut resolved = SizeSummary::default();
resolved.record_reconciliation_scope("b", &long_object);
scanner.apply_size_reconciliation(&resolved);
assert!(scanner.new_cache.info.size_reconciliation.is_empty());
}
#[test]
fn test_classify_get_size_failure_marks_metadata_heal_object_path() {
let temp_dir = std::env::temp_dir();
+3
View File
@@ -43,6 +43,9 @@ allow-git = [
# RustFS fork carrying presigned expiry and constant-time authentication fixes.
# owner: rustfs-maintainers review: 2026-10
"https://github.com/rustfs/s3s.git",
# MiMalloc fork pinned for hotpath allocation counting support.
# owner: houseme review: 2026-10
"https://github.com/xonatius/mimalloc_rust.git",
]
[bans]
-149
View File
@@ -1,149 +0,0 @@
# CI gate matrix
This file is the source of truth for which validation runs on each event, its
configured wall-clock budget, and whether it can block a merge. Test taxonomy,
naming, and nextest serialization rules remain in [README.md](README.md); e2e
membership and counts remain in
[e2e-suite-inventory.md](e2e-suite-inventory.md).
The distinction between **required** and **report-only** is load-bearing:
a failing job blocks a merge only when its exact check name is present in the
live `main` ruleset. A workflow name, a `merge_group` trigger, or a red PR check
does not make a job required by itself.
## Required merge checks
The live `main` ruleset (`6436880`) currently requires exactly these contexts:
| Required context | Producer | Validation |
|---|---|---|
| `CLA Check` | `.github/workflows/cla.yml` | Contributor agreement |
| `Quick Checks` | `.github/workflows/ci.yml` | Formatting and repository guard scripts |
| `Test and Lint` | `.github/workflows/ci.yml` | Clippy, workspace nextest excluding `e2e_test`, doctests, and migration proofs |
For pull requests limited to the paths excluded by the main CI workflow,
`.github/workflows/ci-docs-only.yml` reports `Quick Checks` and
`Test and Lint` under the same names. It runs the real quick checks and the
planning-document guard; it does not claim that Rust compilation or runtime
tests ran. Despite the workflow name, these paths also include selected deploy,
workflow, and lock files.
Verify the live rule rather than trusting this snapshot before changing merge
policy:
```bash
gh api repos/rustfs/rustfs/rulesets/6436880 \
--jq '.rules[] | select(.type == "required_status_checks") | .parameters'
```
The ruleset currently has `strict_required_status_checks_policy=false`.
`Continuous Integration` accepts `merge_group` events and runs `e2e-full` for
them, but `End-to-End Tests (full merge gate)` is not currently a required
context. Therefore the repository is prepared to test a merge-queue SHA, but
the workflow alone does not prove that every merge passed that lane.
## Pull request and merge matrix
Budgets below are job `timeout-minutes`, not typical runtimes. “Report-only”
means the result is visible and actionable but is not in the live required
context list.
| Event | Validation | Budget | Merge status | Reproduction |
|---|---|---:|---|---|
| PR, non-doc change | `Quick Checks` | 10 min | Required | `make pre-commit` (broader local umbrella) |
| PR, non-doc change | `Test and Lint` | 90 min | Required | `cargo nextest run --profile ci --all --exclude e2e_test` |
| PR, non-doc change | `Typos` | 10 min | Report-only | `typos` |
| PR, non-doc change | `ILM Integration (serial)` | 90 min | Report-only | Use the exact command in `.github/workflows/ci.yml` |
| PR, non-doc change | rio-v2 / swift / sftp test-and-lint variants | 90 min each | Report-only | `cargo nextest run` with the workflow's feature set |
| PR, non-doc change | `Build RustFS Debug Binary` | 30 min | Report-only; prerequisite for black-box lanes | `cargo build -p rustfs --bins` |
| PR, non-doc change | `io_uring Integration (real)` | 30 min | Report-only | `cargo test -p rustfs-ecstore --lib uring_ -- --test-threads=1 --nocapture` |
| PR, non-doc change | `End-to-End Tests` (`e2e-smoke` plus `s3s-e2e`) | 30 min | Report-only | `cargo nextest run --profile e2e-smoke -p e2e_test`; then `./scripts/e2e-run.sh ./target/debug/rustfs <data-dir>` |
| PR, non-doc change | `S3 Implemented Tests` | 60 min | Report-only | Build `rustfs`, then run `scripts/s3-tests/run.sh` with `DEPLOY_MODE=binary`, `TEST_MODE=single`, and `MAXFAIL=0` |
| PR, non-doc change | `S3 Lifecycle Behavior Tests` | 30 min | Report-only | Use the accelerated scanner environment in `.github/workflows/ci.yml` with `scripts/s3-tests/run.sh` |
| PR touching dependency or workflow inputs | Cargo Deny / Workflow Pin Report / Dependency Review | 20 / 5 / 30 min | Report-only | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
| PR touching architecture rules or architecture docs | `Architecture Migration Rules` | 10 min | Report-only | `scripts/check_architecture_migration_rules.sh` |
| PR touching Nix or workspace manifests | `Nix Build & Check` | 60 min | Report-only | `nix flake check` |
| PR limited to main-CI-excluded paths | companion `Quick Checks` and `Test and Lint` | 10 min each | Required | `git diff --check`; `make doc-paths-check` when documentation paths changed |
| `merge_group` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Standard required contexts only; `e2e-full` report-only | `cargo nextest run --profile e2e-full -p e2e_test` |
| Push to `main` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Post-merge detection | Same as `merge_group` |
| PR touching fuzz inputs or harness paths | Build plus five 60-second fuzz smoke targets | 60 min build; 30 min per target | Report-only | `MAX_TOTAL_TIME=60 ./scripts/fuzz/run.sh` |
| PR touching selected ecstore disk/format paths | `Rename Safety` on Windows | 60 min | Report-only | Run the four `cargo test -p rustfs-ecstore --lib <filter>` commands in `windows-filesystem.yml` on Windows |
The authoritative e2e filters live in `.config/nextest.toml`; extend a profile
instead of adding a second ad-hoc selector. Before a profile runs,
`scripts/check_test_wiring.py` compares its exact membership to the committed
digest so a silent test drop fails closed.
## Scheduled and manual validation
Scheduled lanes are independent fault domains. They do not block a pull
request, but their workflow-local gate can fail the run and scheduled failures
are routed to the shared failure-issue action. The scheduled-validation
watchdog and freshness workflow separately detect incomplete runs and missing
schedules.
| Cadence (UTC unless noted) | Workflow / validation | Budget | Verdict and artifacts | Reproduction |
|---|---|---:|---|---|
| Daily 02:17 | Fuzz: five nightly corpus targets | 60 min build; 60 min per target | Gate; corpus/crash artifacts, scheduled failure alert | `MAX_TOTAL_TIME=<seconds> ./scripts/fuzz/run.sh` |
| Daily 03:17 | MinIO interop (EC + SSE read parity) | 40 min | Gate; scheduled failure alert | Dispatch `minio-interop.yml` or follow its pinned Docker fixture steps |
| Daily 04:29 | Replication / cluster-fault / protocol e2e | 45 / 90 / 90 min | Three independent gates; JUnit, membership, and server logs | `cargo nextest run --profile e2e-repl-nightly -p e2e_test`; `--profile e2e-nightly`; `-j 1 --profile e2e-protocols` |
| Daily 06:31 | Warp performance A/B | 180 min | Regression budget gate; A/B summaries and server logs | `bash scripts/run_hotpath_warp_abba.sh --help` |
| Daily 00:07 Asia/Shanghai (16:07 UTC previous day) | Nightly GNU build and Vault lanes | 150 / 90 / 60 min | Build, live Vault, and HA failover gates | Use the commands and pinned Vault images in `nightly-gnu.yml` |
| Daily 03:23 | Security Audit | 20 / 5 min, plus 30 min on PR dependency review | Cargo Deny and workflow-pin gates; scheduled failure alert | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
| Daily 23:47 | Scheduled Validation Freshness | 10 min | Fails when a critical schedule was never created or is stale | Dispatch `scheduled-validation-freshness.yml` |
| Sunday 00:11 | Full `Continuous Integration` matrix | Per-job budgets above | Weekly variant coverage, including dormant rio-v2 binary/e2e lanes | Dispatch `ci.yml` |
| Sunday 01:13 | Seven-platform build matrix | 150 min per platform | Build/package integrity; scheduled failure alert | Dispatch `build.yml` with an exact platform set |
| Sunday 02:19 | Ceph s3-tests full sweep: single and real four-node, four shards each | 180 min per shard | Compatibility gate; report, JUnit, exact node IDs, and server logs | `scripts/s3-tests/run.sh` against an existing single or distributed target |
| Sunday 06:41 | Mint | 120 min | **Report-only by design**; per-suite PASS/FAIL/NA and raw `log.json` | Reproduce the pinned Docker sequence in `mint.yml` or dispatch it |
| Sunday 07:43 | Workspace line coverage | 120 min | Report-only trend; lcov and JSON retained 90 days | `make coverage` |
| Monthly, day 1 06:37 | Runner Hygiene | 15 min | Validates runner ephemerality; scheduled failure alert | Dispatch `runner-hygiene.yml` |
Manual `workflow_dispatch` exists for the scheduled workflows above. Manual
runs are debugging evidence and intentionally do not open scheduled-failure
issues. A manual performance run may explicitly allow a known regression; that
override must not be treated as an ordinary passing baseline.
## Release validation
Release validation is post-merge and tag-driven; it does not substitute for a
pull-request gate.
| Event | Validation | Budget | Result |
|---|---|---:|---|
| Push to `main` or weekly schedule | `Build and Release` platform matrix | 150 min per platform | Build artifacts for all selected targets; no release publication on a main push |
| Valid release or preview tag | `Build and Release` plus asset checks | 150 min per platform | Draft release, checksummed assets, and publish step |
| Successful non-preview release-tag build | Docker image build and image scan | 60 min build; 30 min scan | Multi-architecture images plus vulnerability report |
| Successful release-tag build | DEB/RPM packaging | 30 min per architecture | Packages and checksum files uploaded to the release |
| Successful non-preview release-tag build | Helm template test and package | 30 min build; 30 min publish | Versioned chart and repository index |
Use an exact preview tag for end-to-end release rehearsal. Manual dispatches
are backfill/debug paths and do not prove the automatic `workflow_run` chain.
## Evidence requirements
A green check is useful only when it proves the intended behavior ran:
- Record the exact commit SHA and run URL.
- Separate product failure from runner prerequisites, service readiness, and
cancellation. Repair the precondition, then rerun the exact workload.
- Preserve membership manifests, JUnit, raw compatibility logs, seeds, and
server logs where the workflow provides them.
- For a bug fix or a new fault checker, provide sensitivity evidence: the old
behavior or an intentional mutation must fail the new oracle, and the fixed
behavior must pass it.
- Never promote a report-only lane to required from one green run. Require at
least 14 days and 30 representative pull requests with at least 99% complete
execution, then update the ruleset and this table together.
## Change checklist
Update this file in the same pull request when any of these change:
- workflow triggers, job names, timeouts, or nextest profile ownership;
- required status contexts or strict/merge-queue policy;
- scheduled cadence, alert routing, artifact contract, or local reproduction;
- report-only versus gating semantics.
Do not copy per-module test counts here. Update
[e2e-suite-inventory.md](e2e-suite-inventory.md) and its enforced membership
digest instead.
+2 -2
View File
@@ -336,13 +336,13 @@ opentelemetry = { workspace = true }
tracing-opentelemetry = { workspace = true }
# Data structures
hashbrown = { workspace = true, features = ["serde", "rayon"] }
rustfs-mimalloc = { workspace = true }
mimalloc = { workspace = true }
[target.'cfg(target_os = "linux")'.dependencies]
libsystemd.workspace = true
[target.'cfg(not(target_os = "windows"))'.dependencies]
rustfs-mimalloc-sys.workspace = true
libmimalloc-sys.workspace = true
[dev-dependencies]
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
+7 -1
View File
@@ -369,8 +369,14 @@ pub fn allocator_reclaim_controller_snapshot(ctx: &CancellationToken) -> Allocat
}
#[cfg(not(target_os = "windows"))]
#[allow(unsafe_code)]
fn collect_allocator_memory(force: bool) -> Result<(), String> {
rustfs_mimalloc::MiMalloc::collect(force);
// SAFETY: `mi_collect` is provided by the active global allocator backend
// on this target family. It is explicitly intended to reclaim retained
// pages/segments and does not require additional invariants from the caller.
unsafe {
libmimalloc_sys::mi_collect(force);
}
Ok(())
}
+8 -10
View File
@@ -26,22 +26,22 @@ struct MiMallocAllocator;
unsafe impl GlobalAlloc for MiMallocAllocator {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { rustfs_mimalloc::MiMalloc.alloc(layout) }
unsafe { mimalloc::MiMalloc.alloc(layout) }
}
unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { rustfs_mimalloc::MiMalloc.alloc_zeroed(layout) }
unsafe { mimalloc::MiMalloc.alloc_zeroed(layout) }
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { rustfs_mimalloc::MiMalloc.dealloc(ptr, layout) }
unsafe { mimalloc::MiMalloc.dealloc(ptr, layout) }
}
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { rustfs_mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
unsafe { mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
}
}
@@ -51,7 +51,7 @@ static GLOBAL: hotpath::CountingAllocator<MiMallocAllocator> = hotpath::Counting
#[cfg(not(all(feature = "hotpath", feature = "hotpath-alloc")))]
#[global_allocator]
static GLOBAL: rustfs_mimalloc::MiMalloc = rustfs_mimalloc::MiMalloc;
static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc;
fn main() {
let _hotpath_guard = hotpath::HotpathGuardBuilder::new("main").build();
@@ -71,9 +71,8 @@ mod tests {
allocation.extend_from_slice(&[7_u8; 64]);
assert_eq!(allocation.len(), 64);
let heap = rustfs_mimalloc::heap::Heap::main();
// SAFETY: the live Vec pointer is valid to inspect for heap ownership.
assert!(unsafe { heap.contains(allocation.as_ptr()) });
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(allocation.as_ptr().cast()) });
}
#[test]
@@ -86,13 +85,12 @@ mod tests {
let layout = Layout::from_size_align(32, 8).expect("valid test allocation layout");
let grown_layout = Layout::from_size_align(64, 8).expect("valid grown test allocation layout");
let allocator = super::MiMallocAllocator;
let heap = rustfs_mimalloc::heap::Heap::main();
// SAFETY: The pointer is checked for null before use and later released
// through the same allocator with the corresponding layout.
let ptr = unsafe { allocator.alloc_zeroed(layout) };
assert!(!ptr.is_null());
assert!(unsafe { heap.contains(ptr) });
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(ptr.cast()) });
assert!(unsafe { std::slice::from_raw_parts(ptr, 32).iter().all(|byte| *byte == 0) });
// SAFETY: `ptr` was allocated by `allocator` with `layout`; on failure
@@ -104,7 +102,7 @@ mod tests {
panic!("mimalloc realloc failed in allocator smoke test");
}
assert!(unsafe { heap.contains(grown_ptr) });
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(grown_ptr.cast()) });
// SAFETY: `grown_ptr` was reallocated by `allocator` and is released
// with the matching grown layout.
unsafe { allocator.dealloc(grown_ptr, grown_layout) };
+36 -19
View File
@@ -17,7 +17,10 @@ use rustfs_io_metrics::{
record_cpu_usage, record_memory_usage, record_process_memory_split,
};
use serde::Serialize;
#[cfg(any(test, not(target_os = "windows")))]
use serde_json::Value;
#[cfg(not(target_os = "windows"))]
use std::ffi::CStr;
use std::path::Path;
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
@@ -228,18 +231,7 @@ fn read_cgroup_memory_snapshot() -> Option<CgroupMemorySnapshot> {
read_cgroup_v2().or_else(read_cgroup_v1)
}
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
let json = rustfs_mimalloc::MiMalloc::stats_json();
if json.is_empty() {
return None;
}
let observation = parse_mimalloc_stats_json(&json)?;
Some(AllocatorMemorySnapshot {
backend: crate::allocator_reclaim::allocator_backend(),
observation,
})
}
#[cfg(any(test, not(target_os = "windows")))]
fn numeric_json_value(value: &Value) -> Option<u64> {
match value {
Value::Number(number) => number
@@ -250,6 +242,7 @@ fn numeric_json_value(value: &Value) -> Option<u64> {
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
match value {
Value::Object(fields) => fields
@@ -261,6 +254,7 @@ fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64> {
match value {
Value::Object(fields) => {
@@ -277,10 +271,12 @@ fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64>
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_current(value: &Value, metric: &str) -> Option<u64> {
mimalloc_stat_field(value, metric, "current")
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64> {
metrics
.iter()
@@ -289,6 +285,7 @@ fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64
.filter(|value| *value > 0)
}
#[cfg(any(test, not(target_os = "windows")))]
fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservation> {
let value = serde_json::from_str::<Value>(stats_json).ok()?;
let malloc_metrics = ["malloc_normal", "malloc_huge"];
@@ -315,6 +312,33 @@ fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservat
}
}
#[cfg(not(target_os = "windows"))]
#[allow(unsafe_code)]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
// SAFETY: `mi_stats_get_json` returns a null-terminated JSON buffer owned by
// mimalloc when called with a null input buffer. The mimalloc API requires
// freeing that buffer with `mi_free`; parsing finishes before the buffer is freed.
let observation = unsafe {
let stats_ptr = libmimalloc_sys::mi_stats_get_json(0, std::ptr::null_mut());
if stats_ptr.is_null() {
return None;
}
let observation = CStr::from_ptr(stats_ptr).to_str().ok().and_then(parse_mimalloc_stats_json);
libmimalloc_sys::mi_free(stats_ptr.cast());
observation?
};
Some(AllocatorMemorySnapshot {
backend: crate::allocator_reclaim::allocator_backend(),
observation,
})
}
#[cfg(target_os = "windows")]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
None
}
fn configured_memory_observability_interval_secs() -> u64 {
rustfs_utils::get_env_u64(ENV_MEMORY_OBSERVABILITY_INTERVAL_SECS, DEFAULT_MEMORY_OBSERVABILITY_INTERVAL_SECS).max(1)
}
@@ -542,13 +566,6 @@ mod tests {
assert_eq!(parse_mimalloc_stats_json(r#"{ "allocator": "unknown" }"#), None);
}
#[test]
fn read_allocator_memory_snapshot_uses_mimalloc_stats_json() {
let snapshot = super::read_allocator_memory_snapshot();
#[cfg(not(target_os = "windows"))]
assert!(snapshot.is_some(), "allocator snapshot should be available on non-Windows");
}
#[test]
fn memory_observability_snapshot_reports_disabled_when_metrics_are_disabled() {
let snapshot = build_memory_observability_status_snapshot(false, 15, false);
+1 -1
View File
@@ -241,7 +241,7 @@ env \
RUSTFS_TEST_VAULT_FAILOVER_MARKER="$MARKER" \
RUSTFS_TEST_VAULT_OLD_LEADER="$OLD_LEADER" \
cargo test -p rustfs-kms --test vault_ha_failover_live \
vault_raft_leader_failure_recovers_kv2_and_transit_decrypts -- \
vault_raft_leader_failure_preserves_kv2_and_transit_decrypts -- \
--ignored --nocapture --test-threads=1 &
TEST_PID=$!