mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-18 02:33:15 +00:00
Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ab91a05ef5 | |||
| e5ec4b95cf | |||
| 9a2d06b370 | |||
| 00de43528c | |||
| 577de92c02 | |||
| 06eeb886f8 | |||
| f477d27861 |
@@ -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",
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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"))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()),
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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"] }
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
@@ -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. 两个防御模型的对照
|
||||
|
||||
MinIO:heal 时对对象写 `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_dir;heal 的 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 落盘格式变更与标记残留清理两类新成本。维持现状,本对标疑点
|
||||
关闭。
|
||||
@@ -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
@@ -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"),
|
||||
|
||||
Reference in New Issue
Block a user