Compare commits

..

1 Commits

Author SHA1 Message Date
马登山 f4f9b315a6 fix(heal): coalesce duplicate MRF intents 2026-08-23 09:05:45 +08:00
25 changed files with 1012 additions and 165 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=52a05fdfae8bcf6f5828cc2b1e91b2a268139d3f7e1fc47d7b875b55fffb3995
sha256-linux=a22d8af72e250595ac4445e8c880f3f9706202e09ed196e5b7baac632dead8d8
sha256-darwin=9f767b37ed8b1c82da62ea441462d75487785c8086e56f08fb6f6cd89c6e2e52
sha256-linux=fbdaf42b220958d4b1e8880e0f8b5a7992d38e21051bb60596dd4538424757d6
+3 -4
View File
@@ -107,10 +107,9 @@ filter = 'package(e2e_test) & test(/^inline_fast_path_cluster_test::/)'
test-group = 'e2e-inline-boundaries'
# Vault KMS tests share the fixed dev-server port 8200. serial_test's #[serial]
# does not cross nextest process boundaries, so keep every Vault-backed test in
# one group.
# does not cross nextest process boundaries, so keep these tests in one group.
[[profile.default.overrides]]
filter = 'package(e2e_test) & (test(/^kms::kms_vault_test::/) | test(/^kms::configured_roundtrip_test::test_configured_vault_kms_admin_and_versioned_cleanup$/))'
filter = 'package(e2e_test) & test(/^kms::kms_vault_test::/)'
test-group = 'e2e-vault'
# ---------------------------------------------------------------------------
@@ -444,5 +443,5 @@ filter = 'package(e2e_test) & test(/^inline_fast_path_cluster_test::/)'
test-group = 'e2e-inline-boundaries'
[[profile.e2e-full.overrides]]
filter = 'package(e2e_test) & (test(/^kms::kms_vault_test::/) | test(/^kms::configured_roundtrip_test::test_configured_vault_kms_admin_and_versioned_cleanup$/))'
filter = 'package(e2e_test) & test(/^kms::kms_vault_test::/)'
test-group = 'e2e-vault'
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;
}
}
@@ -432,6 +432,7 @@ async fn test_configured_local_kms_admin_and_versioned_cleanup() -> TestResult {
}
#[tokio::test]
#[ignore = "requires a Vault binary"]
async fn test_configured_vault_kms_admin_and_versioned_cleanup() -> TestResult {
let mut env = VaultTestEnvironment::new().await?;
env.start_vault().await?;
@@ -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()
}));
}
+3 -2
View File
@@ -1375,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
@@ -50,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,
}
}
+14 -2
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 {
+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]
+2 -2
View File
@@ -58,7 +58,7 @@
| heal_erasure_disk_rebuild_test | 4 | 🌙 |
| inline_fast_path_cluster_test | 16 | |
| internode_rpc_signature_e2e_test | 5 | |
| kms | 47 | |
| kms | 46 | |
| leading_slash_key_test | 2 | ✅ |
| lifecycle_regression_test | 4 | |
| list_buckets_auth_test | 1 | ✅ |
@@ -99,4 +99,4 @@
| tls_hot_reload_test | 1 | ✅ |
| version_id_regression_test | 10 | ✅ |
**Total listed: 576 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 454 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.
**Total listed: 575 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 453 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.
+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);