Compare commits

..

7 Commits

Author SHA1 Message Date
houseme ab91a05ef5 Merge branch 'main' into hs01-mrf-wiring 2026-08-18 10:18:18 +08:00
houseme e5ec4b95cf fix: include mrf heal source counts
Co-Authored-By: heihutu <heihutu@gmail.com>
2026-08-18 10:16:23 +08:00
houseme 9a2d06b370 test(heal): lock heal vs delete/overwrite race invariants (HS-12) (#6183)
* test(heal): add concurrency invariants for heal vs delete/overwrite races (HS-12)

Audit conclusion for backlog#1874: RustFS does not need a persistent
object-level healing marker (MinIO x-minio-healing) because every path
that can touch the same (bucket, object) commit surface serializes on
the same namespace write lock, and the heal lock guard spans the whole
rename commit including the HEAL_RENAME_INCOMPLETE partial path.

Lock the conclusion in with two race regression tests:

- heal_racing_version_delete_never_resurrects_the_deleted_version:
  shard damage is injected on the doomed version so a Deep heal has real
  reconstruction work while a versioned DELETE runs concurrently; the
  deleted version must stay deleted and the survivor intact.
- heal_racing_unversioned_overwrites_preserves_the_last_commit:
  unversioned overwrites (activating the post-commit tail that deletes
  the replaced data dir without the ns lock) race a Deep heal in a
  loop; the final current version must be exactly the last commit.

Also adds docs/operations/heal-concurrency-safety-notes-zh.md with the
full intersection matrix (17 intersections), lock-coverage argument,
and the residual-window classification (commit tail races are
fail-into-retry safe; bare prefix delete has zero production callers;
admin no_lock is an explicit operator opt-in).

Co-Authored-By: heihutu <heihutu@gmail.com>

* test: remove redundant heal etag clone

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-18 02:01:04 +00:00
Zhengchao An 00de43528c fix(ecstore): describe peer bucket RPC failures with no details (#6190)
heal_bucket, list_bucket, get_bucket_info and delete_bucket returned Error::other("") when a peer answered success=false without an error payload, so operators saw a bare "io error " after quorum reduction. Route all five bucket RPCs through peer_failure_without_details, which names the operation and bucket while staying identical across the peers of one operation so reduce_errs keeps grouping them into a single dominant error.
2026-08-18 01:55:29 +00:00
houseme 577de92c02 feat(ecstore,scanner): deliver MRF intents from error paths (HS-01)
Wire the three production delivery points, each a single non-blocking
try_send next to the existing in-memory heal paths, which stay as the
fast path:

- read.rs decode-error branch: DecodeFailure intent beside the existing
  read-repair submit, so an Urgent ECDecode request survives restarts
  even when the Low-priority read-repair request was dropped or lost.
- add_partial: PartialWrite intent, giving partial-write recovery a
  durable Normal-priority object heal across restarts.
- scanner_folder metadata-corruption classification: MetadataCorruption
  intent beside the existing High-priority scanner heal request.

All three are on error paths only: zero cost on healthy IO.

Part of backlog#1865 (option a).

Co-Authored-By: heihutu <heihutu@gmail.com>
2026-08-18 09:03:18 +08:00
houseme 06eeb886f8 feat(heal): add MRF queue, durable journal, and intent consumer (HS-01)
Consumer half of the mission repair feed: a bounded pending queue
(100k intents / 8 MiB dual ceiling, drop-newest on overflow), a durable
journal at buckets/.heal/mrf/journal.bin holding the unaccepted pending
snapshot, and a consumer task that batches intents off the global
channel, translates them into prioritized heal requests (decode
failure -> Urgent ECDecode, metadata corruption -> High Metadata,
partial write -> Normal object heal), and retries full admissions with
a 5s backoff and a 3-attempt ceiling.

Durability: every journal record carries its own CRC32 and a
format/version header, so a torn tail truncates cleanly at replay; the
journal is deleted after a successful replay and when the pending set
drains (mirroring MinIO's post-replay list.bin unlink). Losing the last
500 ms flush window is acceptable: replayed duplicates merge via the
manager dedup key and read-repair remains the safety net.

Metrics: rustfs_heal_mrf_queue_depth/_queue_bytes, _dropped_total
{reason}, _replayed_total, _journal_bytes, _journal_fsync_total.
The consumer is wired at heal runtime bootstrap right after manager
start, honoring RUSTFS_HEAL_MRF_ENABLE (default on, rollback = off).

Tests: unit tests for the dual ceiling, record roundtrip, torn-tail
truncation, and the priority mapping; integration tests against a real
4-disk ECStore proving channel intents reach the manager queue as
Urgent/mrf-attributed requests and journal replay arms intents, drops
torn tails, and removes the file.

Part of backlog#1865 (option a).

Co-Authored-By: heihutu <heihutu@gmail.com>
2026-08-18 09:03:11 +08:00
houseme f477d27861 feat(common): add MRF intent channel and Mrf request source (HS-01)
Introduce the producer-facing half of the mission repair feed: a global
bounded (8192) channel carrying lightweight MrfIntent values from IO
error paths, plus the RUSTFS_HEAL_MRF_ENABLE delivery kill-switch and
config constants for queue/journal sizing. Delivery is strictly
non-blocking (try_send, drop-on-full) so it can sit on decode-failure
and partial-write paths without adding latency. HealRequestSource grows
a 'mrf' variant so admission accounting can attribute replayed intents.

Part of backlog#1865 (option a: wire HealEvent-style intents with a
durable retry ledger).

Co-Authored-By: heihutu <heihutu@gmail.com>
2026-08-18 09:03:02 +08:00
21 changed files with 1554 additions and 35 deletions
+4
View File
@@ -287,6 +287,9 @@ pub enum HealRequestSource {
Scanner,
AutoHeal,
ReadRepair,
/// Mission Repair Feed: intents delivered by error paths and replayed
/// from the durable MRF journal.
Mrf,
}
impl HealRequestSource {
@@ -297,6 +300,7 @@ impl HealRequestSource {
Self::Scanner => "scanner",
Self::AutoHeal => "auto_heal",
Self::ReadRepair => "read_repair",
Self::Mrf => "mrf",
}
}
}
+1
View File
@@ -17,6 +17,7 @@ pub mod globals;
pub mod heal_channel;
pub mod last_minute;
pub mod metrics;
pub mod mrf_channel;
mod readiness;
pub mod table_catalog;
pub mod trace_bus;
+203
View File
@@ -0,0 +1,203 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Mission Repair Feed (MRF) intent channel.
//!
//! Producers on error paths (read decode failure, scanner metadata
//! corruption, partial-write recovery) hand a lightweight [`MrfIntent`] to the
//! heal crate through a global bounded channel. Delivery is strictly
//! non-blocking: `try_send_mrf_intent` never awaits and drops the intent
//! (counting it) when the channel is full or uninitialized — losing one heal
//! hint is always preferred over stalling an IO path. Durable replay of
//! unconsumed intents is the consumer's job (see `rustfs-heal`
//! `heal::mrf_queue`), mirroring MinIO's `.heal/mrf/list.bin`.
use std::sync::{
Arc, OnceLock,
atomic::{AtomicBool, Ordering},
};
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;
/// 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)]
pub enum MrfKind {
/// Erasure decode failed while serving a read (read path).
DecodeFailure,
/// Scanner classified object metadata as corrupt.
MetadataCorruption,
/// A write left the object with fewer committed shards than the set size.
PartialWrite,
}
impl MrfKind {
pub const fn as_str(self) -> &'static str {
match self {
MrfKind::DecodeFailure => "decode-failure",
MrfKind::MetadataCorruption => "metadata-corruption",
MrfKind::PartialWrite => "partial-write",
}
}
}
/// One repair intent. Kept deliberately small so the in-memory queue and the
/// journal stay bounded; `bucket`/`object` are `Arc<str>` so re-arming an
/// intent never re-allocates the strings.
#[derive(Clone, Debug)]
pub struct MrfIntent {
pub bucket: Arc<str>,
pub object: Arc<str>,
/// Version the intent targets, as raw UUID bytes.
pub version_id: Option<[u8; 16]>,
pub kind: MrfKind,
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,
}
/// Consumer-side retry ceiling before an intent is given up on.
pub const MRF_MAX_ATTEMPTS: u8 = 3;
impl MrfIntent {
/// Rough in-memory footprint used by the queue's byte budget.
pub fn estimated_bytes(&self) -> usize {
// Struct + strings + version bytes; buckets and objects are usually
// far below this bound, so rounding up keeps the budget conservative.
64 + self.bucket.len() + self.object.len()
}
}
static GLOBAL_MRF_SENDER: OnceLock<mpsc::Sender<MrfIntent>> = OnceLock::new();
/// 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.
static MRF_DELIVERY_ENABLED: AtomicBool = AtomicBool::new(true);
/// Override delivery (used at heal-runtime startup from configuration).
pub fn set_mrf_delivery_enabled(enabled: bool) {
MRF_DELIVERY_ENABLED.store(enabled, Ordering::Relaxed);
}
/// Whether producers currently deliver intents.
pub fn mrf_delivery_enabled() -> bool {
MRF_DELIVERY_ENABLED.load(Ordering::Relaxed)
}
/// Create the global MRF channel and return the consumer half. Fails if the
/// channel is already initialized (the heal runtime is a singleton).
pub fn init_mrf_channel() -> Result<mpsc::Receiver<MrfIntent>, &'static str> {
let (sender, receiver) = mpsc::channel(MRF_CHANNEL_CAPACITY);
GLOBAL_MRF_SENDER
.set(sender)
.map_err(|_| "MRF channel sender already initialized")?;
Ok(receiver)
}
/// Best-effort, non-blocking intent delivery from an error path.
///
/// Returns `true` when the intent was accepted into the channel. `false`
/// means the intent was dropped (feature disabled, channel not yet
/// initialized, or channel full) — callers must not retry or await; the
/// existing read-repair / scanner heal paths remain the safety net.
///
/// 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 {
if !mrf_delivery_enabled() {
return false;
}
let Some(sender) = GLOBAL_MRF_SENDER.get() else {
return false;
};
let intent = MrfIntent {
bucket: Arc::from(bucket),
object: Arc::from(object),
version_id: version_id.map(|vid| *vid.as_bytes()),
kind,
enqueued_at_ms: unix_now_ms(),
attempts: 0,
};
sender.try_send(intent).is_ok()
}
fn unix_now_ms() -> u64 {
// Kept trivial: the timestamp is diagnostic metadata only; wall-clock
// 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)
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn intents_estimate_is_conservative() {
let intent = MrfIntent {
bucket: Arc::from("bucket"),
object: Arc::from("object"),
version_id: Some([0u8; 16]),
kind: MrfKind::DecodeFailure,
enqueued_at_ms: 0,
attempts: 0,
};
assert!(intent.estimated_bytes() >= intent.bucket.len() + intent.object.len());
}
#[tokio::test]
async fn try_send_delivers_and_respects_capacity() {
let mut receiver = init_mrf_channel().expect("first initialization should succeed");
assert!(init_mrf_channel().is_err(), "double initialization must fail");
assert!(try_send_mrf_intent(MrfKind::DecodeFailure, "b", "o", Some(Uuid::nil())));
let intent = receiver.recv().await.expect("intent should arrive");
assert_eq!(intent.kind, MrfKind::DecodeFailure);
assert_eq!(intent.bucket.as_ref(), "b");
// Disable delivery: producers become no-ops.
set_mrf_delivery_enabled(false);
assert!(!try_send_mrf_intent(MrfKind::PartialWrite, "b", "o", None));
set_mrf_delivery_enabled(true);
// 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) {
accepted += 1;
}
}
assert_eq!(accepted, MRF_CHANNEL_CAPACITY);
}
#[test]
fn try_send_without_channel_is_false() {
// This test may run after the tokio test above in the same process;
// the singleton semantics make a clean "uninitialized" case hard, so
// assert the flag-off behavior only.
set_mrf_delivery_enabled(false);
assert!(!try_send_mrf_intent(MrfKind::MetadataCorruption, "b", "o", None));
set_mrf_delivery_enabled(true);
}
}
+28
View File
@@ -177,3 +177,31 @@ pub const DEFAULT_HEAL_MAINLINE_WRITE_UTILIZATION_HIGH_PERCENT: usize = 80;
/// Default foreground pressure recheck delay for heal scheduler, in milliseconds.
pub const DEFAULT_HEAL_MAINLINE_MAX_SLEEP_MS: u64 = 250;
/// Environment variable that toggles the MRF (mission repair feed) intent
/// pipeline: error paths deliver repair intents to the heal runtime, and
/// unconsumed intents are replayed from the durable journal after a restart.
pub const ENV_HEAL_MRF_ENABLE: &str = "RUSTFS_HEAL_MRF_ENABLE";
/// Environment variable for the MRF in-memory queue capacity (intent count).
pub const ENV_HEAL_MRF_QUEUE_SIZE: &str = "RUSTFS_HEAL_MRF_QUEUE_SIZE";
/// Environment variable for the MRF journal byte budget. The journal is
/// compacted once its on-disk size crosses this bound.
pub const ENV_HEAL_MRF_JOURNAL_MAX_BYTES: &str = "RUSTFS_HEAL_MRF_JOURNAL_MAX_BYTES";
/// Environment variable for the MRF journal replay batch size (intents per
/// replay push round).
pub const ENV_HEAL_MRF_REPLAY_BATCH: &str = "RUSTFS_HEAL_MRF_REPLAY_BATCH";
/// Default behavior keeps the MRF intent pipeline enabled.
pub const DEFAULT_HEAL_MRF_ENABLE: bool = true;
/// Default MRF queue capacity (matches MinIO's 100k MRF list ceiling).
pub const DEFAULT_HEAL_MRF_QUEUE_SIZE: usize = 100_000;
/// Default MRF journal byte budget (8 MiB), mirroring the channel payload cap.
pub const DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES: usize = 8 * 1024 * 1024;
/// Default MRF replay batch size.
pub const DEFAULT_HEAL_MRF_REPLAY_BATCH: usize = 256;
+1 -2
View File
@@ -20,7 +20,6 @@
#![allow(clippy::all)]
use http::{HeaderMap, HeaderName, HeaderValue};
use rustfs_utils::http::headers::AMZ_CHECKSUM_MODE;
use std::collections::HashMap;
use time::OffsetDateTime;
use tracing::warn;
@@ -77,7 +76,7 @@ impl GetObjectOptions {
}
}
if self.checksum {
headers.insert(HeaderName::from_static(AMZ_CHECKSUM_MODE), HeaderValue::from_static("ENABLED"));
headers.insert(HeaderName::from_static("x-amz-checksum-mode"), HeaderValue::from_static("ENABLED"));
}
headers
}
+6 -10
View File
@@ -54,10 +54,6 @@ use rustfs_config::MAX_S3_CLIENT_RESPONSE_SIZE;
use rustfs_rio::HashReader;
use rustfs_utils::HashAlgorithm;
use rustfs_utils::{
http::headers::{
AMZ_CHECKSUM_CRC32, AMZ_CHECKSUM_CRC32C, AMZ_CHECKSUM_CRC64NVME, AMZ_CHECKSUM_MODE, AMZ_CHECKSUM_SHA1,
AMZ_CHECKSUM_SHA256,
},
net::get_endpoint_url,
retry::{DEFAULT_RETRY_CAP, DEFAULT_RETRY_UNIT, MAX_JITTER, MAX_RETRY, RetryTimer},
};
@@ -1387,12 +1383,12 @@ pub(crate) fn to_object_info_for_provider(
};
// Extract checksums
let checksum_crc32 = get_header(AMZ_CHECKSUM_CRC32);
let checksum_crc32c = get_header(AMZ_CHECKSUM_CRC32C);
let checksum_sha1 = get_header(AMZ_CHECKSUM_SHA1);
let checksum_sha256 = get_header(AMZ_CHECKSUM_SHA256);
let checksum_crc64nvme = get_header(AMZ_CHECKSUM_CRC64NVME);
let checksum_mode = get_header(AMZ_CHECKSUM_MODE);
let checksum_crc32 = get_header("x-amz-checksum-crc32");
let checksum_crc32c = get_header("x-amz-checksum-crc32c");
let checksum_sha1 = get_header("x-amz-checksum-sha1");
let checksum_sha256 = get_header("x-amz-checksum-sha256");
let checksum_crc64nvme = get_header("x-amz-checksum-crc64nvme");
let checksum_mode = get_header("x-amz-checksum-mode");
// Build and return the ObjectInfo struct
Ok(ObjectInfo {
@@ -214,6 +214,19 @@ fn pool_write_quorum(participant_count: usize) -> usize {
(participant_count / 2) + 1
}
/// Error for a peer that reported `success = false` without an error payload.
///
/// The message must stay identical across the peers of one operation: `reduce_errs`
/// buckets `Error::Io` by kind plus rendered message, so any per-peer detail (address,
/// timing) would split one shared failure into single-count buckets and downgrade a real
/// dominant error into `ErasureWriteQuorum`.
fn peer_failure_without_details(op: &str, bucket: Option<&str>) -> Error {
match bucket {
Some(bucket) => Error::other(format!("{op}({bucket}): peer returned failure without error details")),
None => Error::other(format!("{op}: peer returned failure without error details")),
}
}
fn reduce_pool_write_quorum_errs(per_pool_errs: &[Option<Error>]) -> Option<Error> {
if per_pool_errs.is_empty() {
return Some(Error::ErasureWriteQuorum);
@@ -1078,7 +1091,7 @@ impl PeerS3Client for RemotePeerS3Client {
return if let Some(err) = response.error {
Err(err.into())
} else {
Err(Error::other(""))
Err(peer_failure_without_details("heal_bucket", Some(bucket)))
};
}
@@ -1105,7 +1118,7 @@ impl PeerS3Client for RemotePeerS3Client {
return if let Some(err) = response.error {
Err(err.into())
} else {
Err(Error::other(""))
Err(peer_failure_without_details("list_bucket", None))
};
}
let bucket_infos = response
@@ -1136,9 +1149,7 @@ impl PeerS3Client for RemotePeerS3Client {
return if let Some(err) = response.error {
Err(err.into())
} else {
Err(Error::other(format!(
"make_bucket({bucket}): peer returned failure without error details"
)))
Err(peer_failure_without_details("make_bucket", Some(bucket)))
};
}
@@ -1162,7 +1173,7 @@ impl PeerS3Client for RemotePeerS3Client {
return if let Some(err) = response.error {
Err(err.into())
} else {
Err(Error::other(""))
Err(peer_failure_without_details("get_bucket_info", Some(bucket)))
};
}
let bucket_info = serde_json::from_str::<BucketInfo>(&response.bucket_info)?;
@@ -1190,7 +1201,7 @@ impl PeerS3Client for RemotePeerS3Client {
return if let Some(err) = response.error {
Err(err.into())
} else {
Err(Error::other(""))
Err(peer_failure_without_details("delete_bucket", Some(bucket)))
};
}
@@ -2314,4 +2325,37 @@ mod tests {
.collect::<Vec<_>>();
assert_eq!(calls, vec![1, 1, 0, 0, 0, 0, 0, 0]);
}
#[test]
fn peer_failure_without_details_names_operation_and_bucket() {
for op in ["heal_bucket", "make_bucket", "get_bucket_info", "delete_bucket"] {
let message = peer_failure_without_details(op, Some("ops-bucket")).to_string();
assert!(message.contains(op), "{op} message must name the operation: {message}");
assert!(message.contains("ops-bucket"), "{op} message must name the bucket: {message}");
}
let message = peer_failure_without_details("list_bucket", None).to_string();
assert!(message.contains("list_bucket"), "cluster-wide message must name the operation");
assert!(!message.trim().is_empty());
}
#[test]
fn peer_failure_without_details_keeps_one_reduce_errs_bucket_per_operation() {
// reduce_errs groups Io errors by kind plus rendered message: peers failing the
// same operation on the same bucket must still reach quorum as one dominant error.
let per_pool_errs = vec![
Some(peer_failure_without_details("delete_bucket", Some("shared"))),
Some(peer_failure_without_details("delete_bucket", Some("shared"))),
Some(peer_failure_without_details("delete_bucket", Some("shared"))),
];
assert_eq!(
reduce_pool_write_quorum_errs(&per_pool_errs),
Some(peer_failure_without_details("delete_bucket", Some("shared")))
);
assert_ne!(
peer_failure_without_details("delete_bucket", Some("shared")),
peer_failure_without_details("get_bucket_info", Some("shared"))
);
}
}
+219
View File
@@ -3297,4 +3297,223 @@ mod heal_result_report_tests {
assert!(result.detail.contains("part 1"));
assert!(result.detail.contains("bitrot_failure=true"));
}
// HS-12 (backlog#1874): a versioned DELETE racing an object heal must never
// resurrect the deleted version. The heal has real reconstruction work (a
// shard of the doomed version is removed), so both sides touch the same
// (bucket, object, data_dir); whichever order the ns write lock serializes
// them in, the committed delete must win.
#[tokio::test]
#[serial_test::serial]
async fn heal_racing_version_delete_never_resurrects_the_deleted_version() {
let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await;
let bucket = "heal-race-delete-no-resurrect";
let object = "object.bin";
set.make_bucket(
bucket,
&MakeBucketOptions {
versioning_enabled: true,
..Default::default()
},
)
.await
.expect("versioned bucket should be created");
let mut first_reader = PutObjReader::from_vec(vec![0x11; 1024 * 1024]);
let first_info = set
.put_object(
bucket,
object,
&mut first_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("first version should be written");
let first_version = first_info
.version_id
.expect("versioned put should return the first version id")
.to_string();
let mut second_reader = PutObjReader::from_vec(vec![0x22; 1024 * 1024]);
let second_info = set
.put_object(
bucket,
object,
&mut second_reader,
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("second version should be written");
let second_version = second_info
.version_id
.expect("versioned put should return the second version id")
.to_string();
// Damage one shard of the doomed version so the racing heal performs an
// actual reconstruction over its data dir instead of an early exit.
let doomed_source = disks[0]
.read_version("", bucket, object, &first_version, &ReadOptions::default())
.await
.expect("doomed version metadata should be readable");
let doomed_data_dir = doomed_source
.data_dir
.expect("non-inline version should have a data directory");
tokio::fs::remove_file(
temp_dirs[1]
.path()
.join(bucket)
.join(object)
.join(doomed_data_dir.to_string())
.join("part.1"),
)
.await
.expect("shard damage should be injected before the race");
let delete_set = set.clone();
let (delete_res, heal_res) = tokio::join!(
async {
delete_set
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(first_version.clone()),
object_lock_config_snapshot: Some(Arc::new(crate::set_disk::ObjectLockConfigSnapshot::new(
crate::bucket::metadata_sys::ObjectLockConfigState::ConfirmedAbsent,
))),
..Default::default()
},
)
.await
},
async {
set.heal_object(
bucket,
object,
"",
&HealOpts {
scan_mode: HealScanMode::Deep,
..Default::default()
},
)
.await
},
);
delete_res.expect("version delete must succeed under lock serialization");
// The heal may legitimately report a transient failure when the version
// it was rebuilding disappears mid-flight; only the end state matters.
drop(heal_res);
let resurrected = set
.get_object_info(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: Some(first_version.clone()),
..Default::default()
},
)
.await;
assert!(
matches!(&resurrected, Err(Error::FileVersionNotFound) | Err(Error::ObjectNotFound(..))),
"a racing heal must not resurrect the deleted version: {resurrected:?}"
);
let survivor = set
.get_object_info(
bucket,
object,
&ObjectOptions {
versioned: true,
version_id: Some(second_version.clone()),
..Default::default()
},
)
.await
.expect("surviving version must remain readable after the race");
assert_eq!(survivor.size, 1024 * 1024, "survivor size must be intact");
}
// HS-12 (backlog#1874): unversioned overwrite commits race a Deep heal on
// the same object. The overwrite's post-commit tail deletes the replaced
// data dir without the ns lock (object.rs commit tail), which is exactly
// the intersection the audit flagged: the heal must tolerate the tail race
// (retryable outcome) and every committed overwrite must survive — the
// final current version is exactly the last payload written.
#[tokio::test]
#[serial_test::serial]
async fn heal_racing_unversioned_overwrites_preserves_the_last_commit() {
let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await;
let bucket = "heal-race-put-overwrite";
let object = "object.bin";
set.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("bucket should be created");
const ROUNDS: usize = 8;
const PAYLOAD_SIZE: usize = 256 * 1024;
let mut last_etag = String::new();
for round in 0..ROUNDS {
// Give the heal something to rebuild on alternating rounds: remove a
// shard of the current data dir right before the race.
if round % 2 == 1 {
let current = disks[2]
.read_version("", bucket, object, "", &ReadOptions::default())
.await
.expect("current metadata should be readable");
if let Some(data_dir) = current.data_dir {
let shard = temp_dirs[3]
.path()
.join(bucket)
.join(object)
.join(data_dir.to_string())
.join("part.1");
if shard.exists() {
tokio::fs::remove_file(&shard)
.await
.expect("shard damage should be injectable mid-race");
}
}
}
let payload = vec![round as u8; PAYLOAD_SIZE];
let mut put_reader = PutObjReader::from_vec(payload);
let put_opts = ObjectOptions::default();
let heal_opts = HealOpts {
scan_mode: HealScanMode::Deep,
..Default::default()
};
let (put_res, heal_res) = tokio::join!(
set.put_object(bucket, object, &mut put_reader, &put_opts),
set.heal_object(bucket, object, "", &heal_opts),
);
let put_info = put_res.expect("overwrite must succeed under lock serialization");
last_etag = put_info.etag.clone().unwrap_or_default();
// Heal outcome is unconstrained (may hit the tail race and report a
// retryable error); the invariant is checked on the end state.
drop(heal_res);
}
let final_info = set
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("object must remain readable after the race loop");
assert_eq!(
final_info.size, PAYLOAD_SIZE as i64,
"final current version must be the last committed overwrite"
);
assert_eq!(
final_info.etag.unwrap_or_default(),
last_etag,
"the racing heal loop must never leave a stale or resurrected current version"
);
}
}
@@ -5845,6 +5845,14 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks {
#[tracing::instrument(skip(self))]
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 mut request = rustfs_common::heal_channel::create_heal_request_with_options(
bucket.to_string(),
Some(object.to_string()),
+9
View File
@@ -1077,6 +1077,15 @@ impl SetDisks {
"Recoverable decode error triggered read repair"
);
let version_id = fi.version_id.as_ref().map(ToString::to_string);
// MRF journal intent: keeps a durable Urgent ECDecode
// request alive across restarts even when the in-memory
// read-repair request is dropped or lost (HS-01).
rustfs_common::mrf_channel::try_send_mrf_intent(
rustfs_common::mrf_channel::MrfKind::DecodeFailure,
bucket,
object,
fi.version_id,
);
submit_read_repair_heal(
bucket,
object,
+2
View File
@@ -89,6 +89,8 @@ async-trait = { workspace = true }
futures = { workspace = true }
metrics = { workspace = true }
base64 = { workspace = true }
bytes = { workspace = true }
crc-fast = { workspace = true }
[dev-dependencies]
serde_json = { workspace = true, features = ["raw_value"] }
+2 -1
View File
@@ -612,7 +612,8 @@ impl HealChannelProcessor {
HealRequestSource::Admin
| HealRequestSource::AutoHeal
| HealRequestSource::Internal
| HealRequestSource::ReadRepair => true,
| HealRequestSource::ReadRepair
| HealRequestSource::Mrf => true,
});
// Build HealOptions with all available fields
+2
View File
@@ -270,6 +270,7 @@ pub struct HealSourceCounts {
pub auto_heal: u64,
pub internal: u64,
pub read_repair: u64,
pub mrf: u64,
}
impl HealSourceCounts {
@@ -280,6 +281,7 @@ impl HealSourceCounts {
HealRequestSource::AutoHeal => self.auto_heal += 1,
HealRequestSource::Internal => self.internal += 1,
HealRequestSource::ReadRepair => self.read_repair += 1,
HealRequestSource::Mrf => self.mrf += 1,
}
}
}
+1
View File
@@ -16,6 +16,7 @@ pub mod channel;
pub mod erasure_healer;
pub mod event;
pub mod manager;
pub mod mrf_queue;
pub mod progress;
pub(crate) mod replacement_readiness;
pub mod resume;
+682
View File
@@ -0,0 +1,682 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Mission Repair Feed (MRF) queue, journal, and consumer.
//!
//! Intents arriving on the global channel (see `rustfs_common::mrf_channel`)
//! are buffered in a bounded in-memory queue, translated into prioritized
//! heal requests, and — while they are not yet accepted by the heal manager —
//! mirrored into a durable journal so a crash or restart can replay them.
//! This is the RustFS counterpart of MinIO's `.heal/mrf/list.bin` replay,
//! layered on top of (not replacing) read-repair and scanner heal.
//!
//! Durability model: the journal is a snapshot of the *unaccepted* pending
//! set, rewritten on a group-commit cadence (every flush interval or flush
//! threshold new intents). A rewrite is atomic at the record level only — a
//! torn tail simply truncates during replay because every record carries its
//! own CRC32. Losing the last flush window (≤500 ms) is acceptable: replayed
//! duplicates are merged by the manager's dedup key, and read-repair remains
//! the safety net.
use super::{DiskStore, HealDiskExt as _, local_disk_map_read};
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::sync::Arc;
use std::time::Duration;
use tokio::sync::mpsc;
use uuid::Uuid;
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";
/// Record format tag.
const MRF_JOURNAL_FORMAT: u8 = 1;
/// Record layout version.
const MRF_JOURNAL_VERSION: u8 = 1;
/// 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;
#[derive(Debug, Clone)]
pub(crate) struct MrfConsumerConfig {
/// In-memory queue capacity in intents.
pub queue_capacity: usize,
/// Journal byte budget; a pending snapshot above this bound is rejected
/// oldest-first so the journal can never grow unbounded.
pub journal_max_bytes: usize,
/// How many journal intents to re-arm per replay round.
pub replay_batch: usize,
/// Group-commit cadence for the journal snapshot.
pub flush_interval: Duration,
/// New intents between flushes that force an early snapshot.
pub flush_threshold: usize,
/// Backoff after the heal manager reports a full admission.
pub admission_backoff: Duration,
}
impl Default for MrfConsumerConfig {
fn default() -> Self {
Self {
queue_capacity: rustfs_utils::get_env_usize(
rustfs_config::ENV_HEAL_MRF_QUEUE_SIZE,
rustfs_config::DEFAULT_HEAL_MRF_QUEUE_SIZE,
),
journal_max_bytes: rustfs_utils::get_env_usize(
rustfs_config::ENV_HEAL_MRF_JOURNAL_MAX_BYTES,
rustfs_config::DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES,
),
replay_batch: rustfs_utils::get_env_usize(
rustfs_config::ENV_HEAL_MRF_REPLAY_BATCH,
rustfs_config::DEFAULT_HEAL_MRF_REPLAY_BATCH,
),
flush_interval: Duration::from_millis(500),
flush_threshold: 1000,
admission_backoff: Duration::from_secs(5),
}
}
}
/// Bounded pending set with count and byte ceilings. Overflow drops the
/// incoming intent (never a resident one) and counts the loss.
pub(crate) struct MrfQueue {
pending: VecDeque<MrfIntent>,
bytes: usize,
capacity: usize,
byte_budget: usize,
}
impl MrfQueue {
pub(crate) fn new(capacity: usize, byte_budget: usize) -> Self {
Self {
pending: VecDeque::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 {
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;
}
self.bytes += cost;
self.pending.push_back(intent);
true
}
pub(crate) fn pop_front(&mut self) -> Option<MrfIntent> {
let intent = self.pending.pop_front()?;
self.bytes = self.bytes.saturating_sub(intent.estimated_bytes());
Some(intent)
}
pub(crate) fn push_back(&mut self, intent: MrfIntent) {
self.bytes += intent.estimated_bytes();
self.pending.push_back(intent);
}
pub(crate) fn depth(&self) -> usize {
self.pending.len()
}
pub(crate) fn bytes(&self) -> usize {
self.bytes
}
pub(crate) fn intents(&self) -> impl Iterator<Item = &MrfIntent> {
self.pending.iter()
}
}
// ---------------------------------------------------------------------------
// Journal record codec
// ---------------------------------------------------------------------------
/// Append one encoded record to `out`.
pub(crate) fn encode_intent(intent: &MrfIntent, out: &mut Vec<u8>) {
let start = out.len();
out.push(MRF_JOURNAL_FORMAT);
out.push(MRF_JOURNAL_VERSION);
out.push(match intent.kind {
rustfs_common::mrf_channel::MrfKind::DecodeFailure => 1,
rustfs_common::mrf_channel::MrfKind::MetadataCorruption => 2,
rustfs_common::mrf_channel::MrfKind::PartialWrite => 3,
});
out.push(intent.attempts);
out.extend_from_slice(&intent.enqueued_at_ms.to_le_bytes());
match intent.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());
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());
}
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 {
return None;
}
let kind = match data[2] {
1 => rustfs_common::mrf_channel::MrfKind::DecodeFailure,
2 => rustfs_common::mrf_channel::MrfKind::MetadataCorruption,
3 => rustfs_common::mrf_channel::MrfKind::PartialWrite,
_ => return None,
};
let attempts = data[3];
let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().expect("slice length checked"));
let has_version = data[12] != 0;
let mut cursor = MRF_RECORD_FIXED_HEAD;
let version_id = if has_version {
if data.len() < cursor + 16 {
return None;
}
let bytes: [u8; 16] = data[cursor..cursor + 16].try_into().expect("slice length checked");
cursor += 16;
Some(bytes)
} 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;
cursor += 8;
let body_end = cursor.checked_add(bucket_len)?.checked_add(object_len)?;
let record_end = body_end.checked_add(4)?;
if data.len() < record_end {
return None;
}
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")) {
return None;
}
let bucket = std::sync::Arc::from(std::str::from_utf8(&data[cursor..cursor + bucket_len]).ok()?);
let object = std::sync::Arc::from(std::str::from_utf8(&data[cursor + bucket_len..body_end]).ok()?);
Some((
MrfIntent {
bucket,
object,
version_id,
kind,
enqueued_at_ms,
attempts,
},
record_end,
))
}
/// Decode a whole journal, stopping at the first torn or corrupt record.
/// Returns the decoded intents and the number of trailing bytes discarded.
pub(crate) fn decode_journal(data: &[u8]) -> (Vec<MrfIntent>, usize) {
let mut intents = Vec::new();
let mut cursor = 0usize;
while cursor < data.len() {
match decode_one(&data[cursor..]) {
Some((intent, consumed)) => {
intents.push(intent);
cursor += consumed;
}
None => break,
}
}
let truncated = data.len() - cursor;
(intents, truncated)
}
// ---------------------------------------------------------------------------
// Journal disk IO (all local disks, first successful read wins)
// ---------------------------------------------------------------------------
async fn journal_disks() -> Vec<DiskStore> {
let map = local_disk_map_read().await;
map.values().flatten().cloned().collect()
}
async fn read_journal() -> Option<Vec<u8>> {
for disk in journal_disks().await {
match disk.read_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH).await {
Ok(bytes) => return Some(bytes.to_vec()),
Err(_) => continue,
}
}
None
}
async fn write_journal(data: &[u8]) {
let payload = bytes::Bytes::copy_from_slice(data);
for disk in journal_disks().await {
if let Err(err) = disk
.write_all(super::RUSTFS_META_BUCKET, MRF_JOURNAL_PATH, payload.clone())
.await
{
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);
}
async fn delete_journal() {
for disk in journal_disks().await {
let _ = disk
.delete(
super::RUSTFS_META_BUCKET,
MRF_JOURNAL_PATH,
crate::heal::storage_api::owner::EcstoreDeleteOptions::default(),
)
.await;
}
}
fn warn_mrf_journal_write(err: &super::DiskError) {
tracing::warn!(
target: "rustfs::heal::mrf",
error = %err,
"MRF journal write failed; unconsumed intents may be lost on restart"
);
}
// ---------------------------------------------------------------------------
// Consumer
// ---------------------------------------------------------------------------
/// Translate an intent into the prioritized heal request the issue specifies:
/// decode failures go Urgent ECDecode, metadata corruption goes High
/// Metadata, partial writes go Normal object heal.
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 (heal_type, priority) = match intent.kind {
rustfs_common::mrf_channel::MrfKind::DecodeFailure => (
HealType::ECDecode {
bucket,
object,
version_id,
},
HealPriority::Urgent,
),
rustfs_common::mrf_channel::MrfKind::MetadataCorruption => (HealType::Metadata { bucket, object }, HealPriority::High),
rustfs_common::mrf_channel::MrfKind::PartialWrite => (
HealType::Object {
bucket,
object,
version_id,
},
HealPriority::Normal,
),
};
let mut request = HealRequest::new(heal_type, HealOptions::default(), priority);
request.source = rustfs_common::heal_channel::HealRequestSource::Mrf;
request
}
struct MrfRuntime {
queue: MrfQueue,
config: MrfConsumerConfig,
new_since_flush: usize,
/// True while a journal snapshot exists on disk that no longer reflects
/// an all-consumed pending set; the next idle tick removes it (MinIO
/// deletes its `list.bin` after replay for the same reason).
journal_on_disk: bool,
/// Earliest instant a full-admission retry may proceed.
backoff_until: Option<tokio::time::Instant>,
}
impl MrfRuntime {
fn record_accept(&mut self) {
// Accepted intents leave the pending set; the next flush persists the
// smaller snapshot, which is the journal's compaction.
}
fn snapshot(&self) -> Vec<u8> {
let mut buf = Vec::new();
for intent in self.queue.intents() {
encode_intent(intent, &mut buf);
}
buf
}
async fn flush(&mut self) {
write_journal(&self.snapshot()).await;
self.new_since_flush = 0;
self.journal_on_disk = true;
}
/// Drain pending intents into the heal manager until it is full, the
/// queue empties, or attempts are exhausted.
async fn dispatch(&mut self, manager: &HealManager) {
if let Some(until) = self.backoff_until {
if tokio::time::Instant::now() < until {
return;
}
self.backoff_until = None;
}
while let Some(mut intent) = self.queue.pop_front() {
let request = build_heal_request(&intent);
match manager.submit_heal_request(request).await {
Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => self.record_accept(),
Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => {
intent.attempts = intent.attempts.saturating_add(1);
if intent.attempts >= MRF_MAX_ATTEMPTS {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "attempts_exhausted").increment(1);
continue;
}
self.queue.push_back(intent);
self.backoff_until = Some(tokio::time::Instant::now() + self.config.admission_backoff);
break;
}
Ok(HealAdmissionResult::Dropped(_)) => {
counter!("rustfs_heal_mrf_dropped_total", "reason" => "admission_policy").increment(1);
}
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);
continue;
}
self.queue.push_back(intent);
self.backoff_until = Some(tokio::time::Instant::now() + self.config.admission_backoff);
break;
}
}
}
gauge!("rustfs_heal_mrf_queue_depth").set(self.queue.depth() as f64);
gauge!("rustfs_heal_mrf_queue_bytes").set(self.queue.bytes() as f64);
}
}
/// Initialize the global MRF channel (honoring `RUSTFS_HEAL_MRF_ENABLE`) and
/// spawn the consumer task. Called once from the heal runtime bootstrap right
/// after the manager started; a disabled feature or a double call is a no-op.
/// Public for integration tests that drive the real consumer loop.
pub fn spawn_mrf_consumer(manager: Arc<HealManager>) {
let enabled = rustfs_utils::get_env_bool(rustfs_config::ENV_HEAL_MRF_ENABLE, rustfs_config::DEFAULT_HEAL_MRF_ENABLE);
rustfs_common::mrf_channel::set_mrf_delivery_enabled(enabled);
if !enabled {
tracing::info!(
target: "rustfs::heal::mrf",
"MRF intent pipeline disabled by configuration; producers will not deliver"
);
return;
}
let receiver = match rustfs_common::mrf_channel::init_mrf_channel() {
Ok(receiver) => receiver,
Err(err) => {
tracing::warn!(
target: "rustfs::heal::mrf",
error = err,
"MRF channel initialization failed; intents will be dropped at producers"
);
return;
}
};
tokio::spawn(async move {
run_mrf_consumer(manager, receiver).await;
});
tracing::info!(target: "rustfs::heal::mrf", "MRF intent consumer started");
}
/// Replay the durable journal into a fresh pending queue and submit whatever
/// it armed. Returns the number of intact intents replayed. Duplicates are
/// merged by the manager's dedup key; the journal file is removed once read
/// (torn tails truncate via the per-record CRC). Public for integration tests;
/// the live consumer invokes this through [`replay_into`] at startup.
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
}
/// Shared replay core: read + decode + re-arm + delete, then drain what fits.
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;
};
let (intents, truncated) = decode_journal(&data);
if truncated > 0 {
tracing::warn!(
target: "rustfs::heal::mrf",
truncated_bytes = truncated,
"MRF journal had a torn tail; truncated records were discarded"
);
}
counter!("rustfs_heal_mrf_replayed_total").increment(intents.len() as u64);
let replayed = intents.len();
for intent in intents {
queue.try_push(intent);
}
delete_journal().await;
// Drain the replayed intents immediately; whatever the manager refuses
// stays armed in `queue` for the consumer's retry loop.
if backoff_until.is_none() {
while let Some(mut intent) = queue.pop_front() {
let request = build_heal_request(&intent);
match manager.submit_heal_request(request).await {
Ok(HealAdmissionResult::Accepted) | Ok(HealAdmissionResult::Merged) => {}
Ok(HealAdmissionResult::Full) | Ok(HealAdmissionResult::Dropped(HealAdmissionDropReason::QueueFull)) => {
intent.attempts = intent.attempts.saturating_add(1);
if intent.attempts < MRF_MAX_ATTEMPTS {
queue.push_back(intent);
*backoff_until = Some(tokio::time::Instant::now());
}
break;
}
Ok(HealAdmissionResult::Dropped(_)) | Err(_) => {}
}
}
}
replayed
}
/// Replay the journal, then keep draining the channel into the heal manager
/// while persisting the pending snapshot.
async fn run_mrf_consumer(manager: Arc<HealManager>, mut receiver: mpsc::Receiver<MrfIntent>) {
let config = MrfConsumerConfig::default();
let mut runtime = MrfRuntime {
queue: MrfQueue::new(config.queue_capacity, config.journal_max_bytes),
config: config.clone(),
new_since_flush: 0,
journal_on_disk: false,
backoff_until: None,
};
// 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 mut flush_tick = tokio::time::interval(runtime.config.flush_interval);
flush_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut batch: Vec<MrfIntent> = Vec::with_capacity(runtime.config.replay_batch);
loop {
tokio::select! {
received = receiver.recv_many(&mut batch, runtime.config.replay_batch) => {
if received == 0 {
// Channel closed: flush once more and stop.
runtime.flush().await;
tracing::info!(
target: "rustfs::heal::mrf",
"MRF channel closed; consumer stopped after final flush"
);
return;
}
for intent in batch.drain(..) {
runtime.queue.try_push(intent);
runtime.new_since_flush += 1;
}
runtime.dispatch(manager.as_ref()).await;
if runtime.new_since_flush >= runtime.config.flush_threshold {
runtime.flush().await;
}
}
_ = flush_tick.tick() => {
if runtime.new_since_flush > 0 || runtime.queue.depth() > 0 {
runtime.flush().await;
runtime.dispatch(manager.as_ref()).await;
} else if runtime.journal_on_disk {
// 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);
}
gauge!("rustfs_heal_mrf_queue_depth").set(runtime.queue.depth() as f64);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use rustfs_common::mrf_channel::{MrfIntent, MrfKind};
use std::sync::Arc as StdArc;
fn intent(bucket: &str, object: &str, attempts: u8) -> MrfIntent {
MrfIntent {
bucket: StdArc::from(bucket),
object: StdArc::from(object),
version_id: Some([7u8; 16]),
kind: MrfKind::DecodeFailure,
enqueued_at_ms: 1_700_000_000_000,
attempts,
}
}
#[test]
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");
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)),
"byte budget must drop before the second intent fits"
);
}
#[test]
fn journal_roundtrip_preserves_intents() {
let intents = vec![
intent("bucket-a", "object/a", 0),
intent("bucket-b", "object/b", 2),
MrfIntent {
bucket: StdArc::from("bucket-c"),
object: StdArc::from("object/c"),
version_id: None,
kind: MrfKind::MetadataCorruption,
enqueued_at_ms: 5,
attempts: 1,
},
];
let mut buf = Vec::new();
for intent in &intents {
encode_intent(intent, &mut buf);
}
let (decoded, truncated) = decode_journal(&buf);
assert_eq!(truncated, 0);
assert_eq!(decoded.len(), intents.len());
for (left, right) in decoded.iter().zip(intents.iter()) {
assert_eq!(left.bucket, right.bucket);
assert_eq!(left.object, right.object);
assert_eq!(left.version_id, right.version_id);
assert_eq!(left.kind, right.kind);
assert_eq!(left.attempts, right.attempts);
}
}
#[test]
fn journal_torn_tail_is_truncated() {
let mut buf = Vec::new();
encode_intent(&intent("b", "o", 0), &mut buf);
let mut torn = buf.clone();
torn.extend_from_slice(&buf[..buf.len() / 2]);
let (decoded, truncated) = decode_journal(&torn);
assert_eq!(decoded.len(), 1, "the intact record must survive");
assert!(truncated > 0, "the partial tail must be discarded");
// A corrupted body (CRC mismatch) also truncates from that record on.
let mut corrupt = buf.clone();
let mid = MRF_RECORD_FIXED_HEAD + 4;
corrupt[mid] ^= 0xff;
let (decoded, truncated) = decode_journal(&corrupt);
assert!(decoded.is_empty());
assert_eq!(truncated, corrupt.len());
}
#[test]
fn heal_request_mapping_follows_priority_matrix() {
let decode = build_heal_request(&intent("b", "o", 0));
assert!(matches!(decode.heal_type, HealType::ECDecode { .. }));
assert_eq!(decode.priority, HealPriority::Urgent);
let metadata = build_heal_request(&MrfIntent {
bucket: StdArc::from("b"),
object: StdArc::from("o"),
version_id: None,
kind: MrfKind::MetadataCorruption,
enqueued_at_ms: 0,
attempts: 0,
});
assert!(matches!(metadata.heal_type, HealType::Metadata { .. }));
assert_eq!(metadata.priority, HealPriority::High);
let partial = build_heal_request(&MrfIntent {
bucket: StdArc::from("b"),
object: StdArc::from("o"),
version_id: None,
kind: MrfKind::PartialWrite,
enqueued_at_ms: 0,
attempts: 0,
});
assert!(matches!(partial.heal_type, HealType::Object { .. }));
assert_eq!(partial.priority, HealPriority::Normal);
}
}
+4
View File
@@ -158,6 +158,10 @@ pub async fn init_heal_manager_with_workload_provider(
return Err(err);
}
// Start the MRF intent consumer (error-path repair intents + durable
// journal replay) now that the manager can accept submissions.
heal::mrf_queue::spawn_mrf_consumer(heal_manager.clone());
#[cfg(test)]
test_hook_after_manager_start().await;
+189
View File
@@ -0,0 +1,189 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! HS-01 (rustfs/backlog#1865): MRF intent pipeline integration tests.
//!
//! Drives the real consumer loop (`spawn_mrf_consumer`) against a real
//! 4-disk `ECStore` heal storage and a `HealManager` that has not started its
//! scheduler, so submitted intents stay observable in the admission queue.
//! Under `cargo nextest` each test runs in its own process, which keeps the
//! process-global MRF channel singleton safe.
use rustfs_common::mrf_channel::{self, MrfKind};
use rustfs_heal::heal::{
manager::{HealConfig, HealManager},
mrf_queue,
storage::{ECStoreHealStorage, HealStorageAPI},
};
use serial_test::serial;
use std::{path::Path, sync::Arc, time::Duration};
mod storage_api;
use storage_api::endpoint_index::{Endpoint, EndpointServerPools, Endpoints, PoolEndpoints, init_local_disks};
const META_BUCKET: &str = ".rustfs.sys";
const JOURNAL_REL: &str = "buckets/.heal/mrf/journal.bin";
async fn heal_env() -> (Vec<std::path::PathBuf>, Arc<dyn HealStorageAPI>) {
let env = rustfs_test_utils::TestECStoreEnv::builder()
.prefix("rustfs_heal_mrf_test")
.build()
.await;
let heal_storage: Arc<dyn HealStorageAPI> = Arc::new(ECStoreHealStorage::new(env.ecstore.clone()));
(env.disk_paths, heal_storage)
}
fn make_manager(storage: Arc<dyn HealStorageAPI>) -> Arc<HealManager> {
Arc::new(HealManager::new(
storage,
Some(HealConfig {
// Keep the scheduler from draining the queue before assertions.
heal_interval: Duration::from_secs(3600),
enable_auto_heal: false,
..Default::default()
}),
))
}
/// Encode one journal record independently of the implementation, so a format
/// drift between writer and this fixture fails loudly here.
fn journal_record(kind: u8, bucket: &str, object: &str, version: Option<[u8; 16]>, attempts: u8) -> Vec<u8> {
let mut body = vec![1u8, 1, kind, attempts];
body.extend_from_slice(&1_700_000_000_000u64.to_le_bytes());
match version {
Some(bytes) => {
body.push(1);
body.extend_from_slice(&bytes);
}
None => body.push(0),
}
body.extend_from_slice(&(bucket.len() as u32).to_le_bytes());
body.extend_from_slice(&(object.len() as u32).to_le_bytes());
body.extend_from_slice(bucket.as_bytes());
body.extend_from_slice(object.as_bytes());
let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc);
hasher.update(&body);
body.extend_from_slice(&(hasher.finalize() as u32).to_le_bytes());
body
}
fn write_journal_to_disks(disk_paths: &[std::path::PathBuf], data: &[u8]) {
for path in disk_paths {
let journal = path.join(META_BUCKET).join(JOURNAL_REL);
std::fs::create_dir_all(journal.parent().expect("journal parent")).expect("create journal dir");
std::fs::write(&journal, data).expect("write journal fixture");
}
}
async fn wait_until<F, Fut>(deadline: Duration, mut probe: F) -> bool
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = bool>,
{
let start = std::time::Instant::now();
while start.elapsed() < deadline {
if probe().await {
return true;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
false
}
/// A decode-failure intent delivered on the global channel must surface in the
/// heal manager as an Urgent request attributed to the MRF source.
#[tokio::test]
#[serial]
async fn decode_failure_intent_maps_to_urgent_mrf_heal_request() {
let (_disk_paths, storage) = heal_env().await;
let manager = make_manager(storage);
mrf_queue::spawn_mrf_consumer(manager.clone());
assert!(
mrf_channel::try_send_mrf_intent(MrfKind::DecodeFailure, "mrf-bucket", "mrf-object", None),
"intent should be accepted while the consumer holds the channel"
);
let appeared = wait_until(Duration::from_secs(10), || async {
let snapshot = manager.operations_snapshot().await;
snapshot.queued_by_source.mrf >= 1 && snapshot.queued_by_priority.urgent >= 1
})
.await;
assert!(
appeared,
"MRF intent must reach the manager queue as an Urgent request (snapshot: {:?})",
manager.operations_snapshot().await
);
}
/// A journal left behind by a previous process must be replayed into the
/// manager queue and then removed, and a torn tail must not block replay of
/// the intact records.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial]
async fn journal_replay_arms_intents_and_deletes_the_file() {
let (disk_paths, storage) = heal_env().await;
// The journal reader resolves disks through the process-local disk map;
// register the environment's disks the same way server startup does.
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-test".to_string(),
platform: String::new(),
};
init_local_disks(EndpointServerPools::from(vec![pool]))
.await
.expect("local disks should register");
let mut journal = journal_record(1, "replay-bucket", "replay-object", Some([9u8; 16]), 0);
journal.extend(journal_record(3, "replay-bucket", "partial-object", None, 1));
// Torn tail: a third record truncated mid-way must not block the two
// intact records above.
journal.extend_from_slice(&journal_record(2, "replay-bucket", "metadata-object", None, 0)[..8]);
write_journal_to_disks(&disk_paths, &journal);
let manager = make_manager(storage);
// Replay directly (not via the process-global channel consumer, which the
// sibling test already claimed in this process under plain `cargo test`).
let replayed = mrf_queue::replay_journal_once(&manager).await;
assert_eq!(replayed, 2, "the two intact records must be replayed");
let snapshot = manager.operations_snapshot().await;
assert_eq!(snapshot.queued_by_source.mrf, 2, "replayed intents must be attributed to the MRF source");
assert!(
disk_paths
.iter()
.all(|path| !Path::new(path).join(META_BUCKET).join(JOURNAL_REL).exists()),
"the journal file must be removed after a successful 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");
}
+9
View File
@@ -2478,6 +2478,15 @@ impl FolderScanner {
}
if let GetSizeFailureAction::HealMetadata { object } = failure_action {
// MRF journal intent: durable High-priority Metadata
// heal across restarts (HS-01); the scanner heal
// request below stays as the immediate path.
rustfs_common::mrf_channel::try_send_mrf_intent(
rustfs_common::mrf_channel::MrfKind::MetadataCorruption,
&item.bucket,
&object,
None,
);
self.send_required_scanner_heal_request(
PendingScannerHealKind::Object,
item.bucket.clone(),
@@ -0,0 +1,115 @@
# Heal 并发安全说明(对象级 healing 标记对标审计结论)
对应 backlog rustfs/backlog#1874(父 #1862,HS-12)。本文回答一个问题:MinIO 在 heal
期间对对象打 `x-minio-healing:true` 元数据标记以防"heal 提交与并发删除/版本清理互毁"
cmd/xl-storage.go RenameData 的 healing 分支),RustFS 是否需要同款防御。
**结论:不需要。** RustFS 不存在 MinIO 用 healing 标记防御的那类竞争:所有会触达同一
`(bucket, object)` 提交面的路径都在同一把对象级 namespace 写锁上互斥,且 heal 的锁
guard 覆盖 rename 提交全程;MinIO 需要标记的根因(RenameData 提交内部与版本清理逻辑
交错)在 RustFS 的提交模型中不存在。RustFS 已有一个瞬态 healing 旗标用于另一目的
(见下文 §2),并有并发不变量回归测试锁定本结论(§5)。
## 1. 两个防御模型的对照
MinIOheal 时对对象写 `x-minio-healing:true`(持久元数据标记),后续任何 RenameData
提交看到该标记就跳过版本清理/legacy purge 逻辑——防御发生在锁外,靠元数据让路。
RustFS:三层防御,全部不依赖持久对象标记:
1. **锁内互斥**:heal 与一切前台/后台写路径的提交点在同一把 `(bucket, object)` ns 写锁
上串行(分布式部署为 quorum 锁 RPC,单机为进程内锁管理器;锁粒度是对象级,version
恒为 None)。
2. **提交模型隔离**rename_data 提交内没有会与 heal 交错的版本清理逻辑;被替换旧版本
的 data_dir 物理删除被移出提交临界区(commit tail),且只删已被新提交替换的 unshared
目录。
3. **瞬态 healing 旗标**`FileInfo::set_healing`crates/filemeta/src/fileinfo.rs)在
heal 提交的内存 FileInfo 上打 `"healing"` 内部键,rename_data 据此允许先清空 stale
目标 data_dir 再 rename——解决 heal 复用 data_dir 做 in-place 修复时 rename(2) 无法
替换非空目录的文件系统语义冲突(EEXIST/ENOTEMPTY)。该键是瞬态的,不落盘
`is_skip_meta_key`),与 MinIO 的持久标记目的不同。非 heal 提交撞上非空目标
data_dir 会显式失败,有测试锁定两个方向的行为。
## 2. 交点矩阵
中心路径:`heal_object_with_explicit_version_regen`crates/ecstore/src/set_disk/ops/heal.rs
下称 heal.rs)在入口取 `(bucket, object)` ns 写锁,guard 绑定到函数作用域末尾,覆盖
quorum 元数据读取 → EC 重建 → 逐盘 rename 提交 → tmp 清理 → HEAL_RENAME_INCOMPLETE
部分提交返回 → 孤儿 data_dir 回收的全过程。并发侧逐交点判定:
| # | 并发路径 | 并发侧锁 | 判定 | 关键证据 |
|---|---|---|---|---|
| 1 | PUT 对象提交 | `put_object_commit` 对象写锁,rename_data 在锁内 | 同锁串行 | ops/object.rs 提交锁段 + rename 调用点 |
| 2 | PUT 旧 data_dir tail 清理 | drop 对象锁后的 `commit_rename_data_dir`,无锁 | 无锁并发,语义安全(见 §3.1 | object.rs drop 后 tail 段;io_primitives.rs |
| 3 | DELETE 单对象/版本 | `delete_object` 对象写锁,delete_version 在锁内 | 同锁串行 | object.rs delete_object 锁段 |
| 4 | DELETE 批量 | 批量逐对象写锁(dist 走批量锁 RPC | 同锁串行 | object.rs delete_objects 锁段 |
| 5 | CompleteMultipart | 对象写锁 + upload 路径锁双锁,rename 在锁内 | 同锁串行 | ops/multipart.rs 提交锁段 |
| 6 | CompleteMultipart tail 清理 | drop 对象锁后的旧 data_dir 删除 | 无锁并发,语义安全(见 §3.1 | multipart.rs drop 后 tail 段 |
| 7 | AbortMultipart | 仅 multipart bucket 的 upload 路径锁 | 锁 key 不相交,但资源不相交(abort 不触对象 data_dir/xl.meta)→ 无实际交点 | multipart.rs abort 锁段 |
| 8 | ILM expiry(含 DeleteAllVersions | DeleteAllVersions 走 `delete_prefix_object=true` → 仍取对象锁;FreeVersionTask 显式取锁;noncurrent 批量走批量锁 | 同锁串行 | bucket_lifecycle_ops.rs 消费端链路 |
| 9 | 纯 prefix 删除(绕锁能力面) | `delete_prefix`-only 不取子对象锁 | 无锁并发,但生产调用方为零(见 §3.2 | object.rs delete_object 锁条件 |
| 10 | 孤儿 data_dir 回收 reclaim_orphan_data_dirs | 函数本体无锁;唯一生产调用方在 heal 锁内 | heal 流程内=锁内串行 | heal.rs 收尾调用;io_primitives.rs |
| 11 | 旧清理 receipt 对账 reconcile_old_data_cleanup_receipts | 函数本体无锁;调用点在 heal 锁内 + epoch fence 防误删 | 锁内串行 | object.rs 对账函数 |
| 12 | replication | 数据面为远端 HTTP 写(不落本地盘);本地元数据回写走对象锁 | 同锁串行 / 无交点 | replication_resyncer.rs 链路 |
| 13 | data_movement / rebalance / decommission 源清理 | 显式取对象锁 + 版本未变复核 + guard 复用(no_lock 只是复用已持锁) | 同锁串行 | data_movement/mod.rs 源清理 |
| 14 | copy_object | 目标对象锁 / 走 put 链锁 | 同锁串行 | object.rs copy_object 锁段 |
| 15 | 另一 heal 任务(跨 HealType/force_start | dedup key 跨类型不相交 + force_start 跳过去重 → 任务级可并发 | 最终在 ns 写锁上串行 | heal/manager.rs dedup key 构成 |
| 16 | admin `no_lock=true` heal | 客户端可控绕锁 | 无锁并发,明示运维选项(见 §3.3 | admin/handlers/heal.rs 透传 |
| 17 | stale multipart 清理 | multipart bucket 的 upload 路径锁 | 资源不相交 → 无交点 | bucket_lifecycle_ops.rs 清理链路 |
## 3. 残留窗口定性
### 3.1 PUT/CompleteMultipart commit tail(交点 2/6
写路径提交成功、释放对象锁之后,才 best-effort 删除被替换的旧 data_dir(注释明示有意
不阻塞下一操作)。该删除与并发 heal 对同一旧 data_dir 的读取/重建存在竞态窗口,但语义
安全:
- 删除目标是已被新提交替换的 unshared data_dirheal 的 canonical 元数据来自 quorum
仲裁(ETag/mod_time),此时 quorum 已指向新版本,heal 不会把已替换版本当作 canonical
复活;
- 竞态最坏后果 = heal 当轮对旧版本的一次 transient 失败/空转,重试轮自然收敛;清理
residue 会上报并重新入队 heal`report_old_data_dir_cleanup`);
- 换盘重建等长 heal 走 per-version 显式版本请求,quorum 元数据在锁内读取,不受 tail
影响。
### 3.2 纯 prefix 删除(交点 9
`delete_prefix && !delete_prefix_object` 的路径不取子对象锁(对象名空间锁无法保护前缀
递归删除),与并发 heal 存在理论复活窗口(heal 在 prefix 删除进行中依据旧 quorum 元
数据重建某版本)。全仓库核对结论:该路径的**生产调用方为零**——所有生产 `delete_prefix:
true` 调用点均同时设置 `delete_prefix_object: true`(从而取对象锁)或在测试模块内。这
是 API 能力面的暴露而非行为风险。若未来有调用方需要纯 prefix 删除,须在调用点证明与
heal/scanner 的隔离(例如 bucket 级停扫围栏)。
### 3.3 admin `no_lock=true`(交点 16
admin heal 请求可透传客户端 `nolock` 参数绕过 ns 锁(与 MinIO madmin 的同名选项对齐)。
这是运维明示选项:使用即自负与并发写的竞争责任。文档化即可,不建议收紧。
## 4. heal 侧自身的不变量保障
- dedup key 跨 HealType 不相交(object/metadata/mrf/ecdecode/prefix 各自键面)+ admin
`force_start` 可跳过去重 → 同对象可能同时存在多个 heal 任务,但它们的执行体全部在
`heal_object` 入口的 ns 写锁上串行(生产入口均 `no_lock=false`);
- read-repair 的本地 TTL 预留只去重自身来源,不拦截其他来源的 heal——同样由 ns 锁兜底;
- healing 旗标不落盘,故不存在"标记残留导致后续提交错误让路"的反向风险。
## 5. 回归测试
以下两个并发不变量测试随本审计加入 `crates/ecstore/src/set_disk/ops/heal.rs` 测试模块:
- `heal_racing_version_delete_never_resurrects_the_deleted_version`:注入 doomed 版本
shard 损坏后,版本化 DELETE 与 Deep heal 真并发(同一把锁争用),断言已删除版本不被
复活、存活版本完好;
- `heal_racing_unversioned_overwrites_preserves_the_last_commit`:非版本化覆盖提交(激活
commit tail 旧 data_dir 删除)与 Deep heal 循环竞态,断言最终 current 恰为最后一次
提交(etag 级一致)。
## 6. 结论
MinIO 的 `x-minio-healing` 是锁外元数据防御,前提是其 RenameData 提交内部存在与 heal
交错的版本清理逻辑;RustFS 的提交模型把这类交错从根上消除(提交面锁内互斥 + 清理外
移到 tail + tail 只删 unshared 旧目录),因此引入持久对象级 healing 标记没有对应的竞争
可防,反而会引入 FileInfo 落盘格式变更与标记残留清理两类新成本。维持现状,本对标疑点
关闭。
+2
View File
@@ -317,6 +317,7 @@ fn add_source_counts(total: &mut rustfs_heal::HealSourceCounts, next: rustfs_hea
total.auto_heal = total.auto_heal.saturating_add(next.auto_heal);
total.internal = total.internal.saturating_add(next.internal);
total.read_repair = total.read_repair.saturating_add(next.read_repair);
total.mrf = total.mrf.saturating_add(next.mrf);
}
fn add_operations(total: &mut rustfs_heal::HealOperationsSnapshot, next: rustfs_heal::HealOperationsSnapshot) {
@@ -2353,6 +2354,7 @@ mod tests {
auto_heal: value,
internal: value,
read_repair: value,
mrf: value,
};
let operations = |value| rustfs_heal::HealOperationsSnapshot {
queue_length: value,
+16 -15
View File
@@ -67,9 +67,6 @@ use rustfs_policy::policy::action::{Action, S3Action};
use rustfs_s3_types::EventName;
use rustfs_signer::pre_sign_v4;
use rustfs_utils::egress::{OutboundDnsResolver, OutboundPolicy};
use rustfs_utils::http::headers::{
AMZ_CHECKSUM_CRC32, AMZ_CHECKSUM_CRC32C, AMZ_CHECKSUM_CRC64NVME, AMZ_CHECKSUM_SHA1, AMZ_CHECKSUM_SHA256, AMZ_CHECKSUM_TYPE,
};
use rustfs_utils::http::{
SUFFIX_SOURCE_DELETEMARKER, SUFFIX_SOURCE_MTIME, SUFFIX_SOURCE_REPLICATION_CHECK, SUFFIX_SOURCE_REPLICATION_REQUEST,
SUFFIX_SOURCE_VERSION_ID, get_source_scheme, insert_header,
@@ -1034,24 +1031,28 @@ fn build_get_object_response_headers(output: &GetObjectOutput, base_headers: &He
)?;
}
if let Some(checksum_crc32) = &output.checksum_crc32 {
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_CRC32), checksum_crc32.clone())?;
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-crc32"), checksum_crc32.clone())?;
}
if let Some(checksum_crc32c) = &output.checksum_crc32c {
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_CRC32C), checksum_crc32c.clone())?;
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-crc32c"), checksum_crc32c.clone())?;
}
if let Some(checksum_crc64nvme) = &output.checksum_crc64nvme {
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_CRC64NVME), checksum_crc64nvme.clone())?;
insert_string_header(
&mut headers,
HeaderName::from_static("x-amz-checksum-crc64nvme"),
checksum_crc64nvme.clone(),
)?;
}
if let Some(checksum_sha1) = &output.checksum_sha1 {
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_SHA1), checksum_sha1.clone())?;
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-sha1"), checksum_sha1.clone())?;
}
if let Some(checksum_sha256) = &output.checksum_sha256 {
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_SHA256), checksum_sha256.clone())?;
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-sha256"), checksum_sha256.clone())?;
}
if let Some(checksum_type) = &output.checksum_type {
insert_string_header(
&mut headers,
HeaderName::from_static(AMZ_CHECKSUM_TYPE),
HeaderName::from_static("x-amz-checksum-type"),
checksum_type.as_str().to_string(),
)?;
}
@@ -1113,12 +1114,12 @@ fn clear_object_lambda_variant_headers(headers: &mut HeaderMap) {
http::header::ETAG,
http::header::LAST_MODIFIED,
http::header::EXPIRES,
HeaderName::from_static(AMZ_CHECKSUM_CRC32),
HeaderName::from_static(AMZ_CHECKSUM_CRC32C),
HeaderName::from_static(AMZ_CHECKSUM_CRC64NVME),
HeaderName::from_static(AMZ_CHECKSUM_SHA1),
HeaderName::from_static(AMZ_CHECKSUM_SHA256),
HeaderName::from_static(AMZ_CHECKSUM_TYPE),
HeaderName::from_static("x-amz-checksum-crc32"),
HeaderName::from_static("x-amz-checksum-crc32c"),
HeaderName::from_static("x-amz-checksum-crc64nvme"),
HeaderName::from_static("x-amz-checksum-sha1"),
HeaderName::from_static("x-amz-checksum-sha256"),
HeaderName::from_static("x-amz-checksum-type"),
HeaderName::from_static("x-amz-tagging-count"),
HeaderName::from_static("x-amz-request-route"),
HeaderName::from_static("x-amz-request-token"),