Compare commits

...

4 Commits

Author SHA1 Message Date
overtrue 31c740414c refactor: import x-amz-checksum header names from the shared constants 2026-08-18 09:39:20 +08:00
houseme 35a30cd614 feat(scanner): emit excess alerts as S3 notification events (HS-04) (#6176)
* feat(scanner): emit excess alerts as S3 notification events

The excess-versions / excess-version-size / excess-folders alerts were
metrics-and-logs only; consoles and external auditors had no way to hear
them (rustfs/backlog#1868, HS-04). MinIO emits s3:ObjectManyVersions /
s3:ObjectLargeVersions / s3:PrefixManyFolders for the same conditions —
RustFS carries those as EventName::Scanner* with s3:Scanner:* wire names
that already existed unpublished.

The three alert sites now also dispatch through the standard event
pipeline (send_event via the storage_api owner facade), carrying the
actual values and thresholds in req_params and UserAgent "Scanner".
Without a cooldown a single over-threshold object would re-emit on every
~60s scan cycle, so emissions are edge-held per (kind, bucket, object)
for 24h (RUSTFS_SCANNER_ALERT_COOLDOWN_SECS, 0 = every cycle), backed by
a process-global map with a 4096-key hard cap that clears rather than
grows. Metrics and structured logs stay level-triggered every cycle;
only the notification events are held back. A restart resets the
cooldown deliberately: one re-emission per still-hot key buys back
visibility after the restarts that accompany incident response.

Tests pin the edge-hold semantics (first fires, immediate re-check held,
independent keys, cooldown expiry re-fires, zero cooldown always emits,
hard bound) in one sequential test for the process-global map, and pin
the emitted wire names against EventName's canonical string forms so a
subscribed bucket notification can never silently stop matching.
docs/operations/scanner-excess-alerts.md documents the three events,
the metric-vs-event cadence difference, and the HS-15 threshold deltas
(alert_excess_folders 65538 vs MinIO 50000 is deliberate: Proxmox
Backup Server chunk layout compatibility).

Closes rustfs/backlog#1868.

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

* docs(operations): split scanner excess alerts into English and Chinese pages

The page shipped Chinese-only; keep it as scanner-excess-alerts_zh.md and
add a faithful English translation at the original path, cross-linked at
the top of both.

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

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-18 08:46:32 +08:00
houseme 360bceafce feat(heal): add progress and trace observability (#6179)
* feat(heal): track erasure set progress baseline

Record erasure-set heal byte progress from per-object results and seed progress totals from complete usage-cache snapshots when available.

Keep usage-cache failures observational so heal execution continues without a baseline.

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

* feat(heal): skip filtered erasure set versions

Skip erasure-set versions written after the durable heal start time, and queue lifecycle-expired versions for expiry before skipping them.

Track new-version and ILM-expired skips separately so progress can explain completed baseline work without treating these skips as retry-blocking failures.

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

* feat(heal): wire abandoned data-dir cleanup check

Connect check_abandoned_parts through ECStore, pool, and set layers so heal can invoke the existing orphan data-dir reclaim path instead of returning NotImplemented.

Add dry-run support to the reclaim scan and cover dry-run plus scoped set behavior with regression tests.

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

* feat(obs): add heal scanner trace bus

Introduce an in-process broadcast trace bus with typed heal and scanner events, lazy event construction, and bounded lagged-subscriber behavior.

Cover zero-subscriber publishing, subscription delivery, drop accounting, and lagged receivers with focused common-crate tests.

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

* feat(obs): stream heal trace events from admin API

Wire the admin trace endpoint to the common trace bus for heal/scanner events, including kind, regex, and threshold filtering.

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

* feat(obs): emit heal trace events

Publish heal task lifecycle and abandoned-parts cleanup events through the common trace bus so the admin trace stream has live heal diagnostics.

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

* feat(obs): emit scanner trace events

Publish scanner folder, lifecycle action, and heal-candidate events through the common trace bus for live admin scanner diagnostics.

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

* fix(heal): route data usage loader through storage api

Keep ECStore data-usage facade access behind the heal storage_api boundary so architecture migration guards can validate the heal progress path.

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

* perf(heal): avoid lifecycle snapshots on ordinary heal pages

Only request lifecycle object snapshots when the heal pass has lifecycle expiry context. This keeps ordinary listing and disk-walk pages from cloning FileInfo/ObjectInfo payloads while preserving the skip path that queues expired versions.

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

* test(heal): update bug-fix mocks for lifecycle snapshots

Carry the lifecycle snapshot opt-in argument through the remaining heal bug-fix test mocks so all-targets clippy covers the updated storage trait.

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

* test(rustfs): sync heal storage mock signature

Update the rustfs storage RPC test mock for the lifecycle snapshot opt-in argument and cover it with rustfs all-targets clippy.

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

* test(e2e): allocate smoke ports across nextest processes

Serialize E2E port selection with a small /tmp allocator so nextest workers do not reuse the same just-released ephemeral port before RustFS binds it.

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

---------

Co-authored-by: heihutu <heihutu@gmail.com>
2026-08-18 08:29:29 +08:00
Zhengchao An 7cb91a0190 chore(ecstore): adjudicate 32 bare dead_code allows (#6173)
Replace every bare `#[allow(dead_code)]` in ecstore with either a deletion or a per-item allow carrying a `reason`. Blanket allows at module, struct, and impl level silence the lint for future members too, so each is narrowed to the members that are actually dead.

Delete the dead cluster in `config/heal.rs` (`Config`, its three methods, `RUSTFS_BITROT_CYCLE_IN_MONTHS`, `parse_bitrot_config`) rather than annotate it: it has no callers and is unreachable outside the crate, and `parse_bitrot_config` would panic on its disabled path via `Duration::from_secs_f64(-1.0)`. `DEFAULT_KVS` stays, since the config registry uses it.

Correct two `reason` strings on `Checksum::new` and `PutObjReader::md5_current_hex_string`, which are methods but carried a field-only rationale.

Refs backlog#1823

Co-authored-by: houseme <housemecn@gmail.com>
2026-08-17 23:51:56 +00:00
55 changed files with 2777 additions and 234 deletions
Generated
+2
View File
@@ -9280,6 +9280,7 @@ dependencies = [
"s3s",
"serde",
"serde_json",
"smallvec",
"tokio",
"tonic",
"tracing",
@@ -10251,6 +10252,7 @@ dependencies = [
"rustfs-ecstore",
"rustfs-filemeta",
"rustfs-lock",
"rustfs-s3-types",
"rustfs-storage-api",
"rustfs-utils",
"s3s",
+1
View File
@@ -42,6 +42,7 @@ chrono = { workspace = true, features = ["serde"] }
jiff = { workspace = true, features = ["serde"] }
metrics = { workspace = true }
serde = { workspace = true, features = ["derive"] }
smallvec = { workspace = true }
rmp-serde = { workspace = true }
s3s = { workspace = true, features = ["minio"] }
tracing = { workspace = true }
+1
View File
@@ -19,6 +19,7 @@ pub mod last_minute;
pub mod metrics;
mod readiness;
pub mod table_catalog;
pub mod trace_bus;
pub use globals::*;
pub use readiness::{GlobalReadiness, SystemStage};
+333
View File
@@ -0,0 +1,333 @@
// 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.
use smallvec::SmallVec;
use std::{
sync::{
Arc, OnceLock,
atomic::{AtomicUsize, Ordering},
},
time::{Duration, SystemTime},
};
use tokio::sync::broadcast;
const DEFAULT_TRACE_BUS_CAPACITY: usize = 1024;
const TRACE_ATTR_INLINE_CAPACITY: usize = 8;
static GLOBAL_TRACE_BUS: OnceLock<TraceBus> = OnceLock::new();
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TraceKind {
Heal,
Scanner,
}
impl TraceKind {
pub const fn as_str(self) -> &'static str {
match self {
Self::Heal => "heal",
Self::Scanner => "scanner",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TraceFunc {
HealTask,
HealBucket,
HealObject,
HealCheckAbandonedParts,
HealErasureSetPage,
ScannerFolder,
ScannerIlmAction,
ScannerHealCandidate,
Dropped,
}
impl TraceFunc {
pub const fn as_str(self) -> &'static str {
match self {
Self::HealTask => "heal.Task",
Self::HealBucket => "heal.Bucket",
Self::HealObject => "heal.Object",
Self::HealCheckAbandonedParts => "heal.CheckAbandonedParts",
Self::HealErasureSetPage => "heal.ErasureSetPage",
Self::ScannerFolder => "scanner.Folder",
Self::ScannerIlmAction => "scanner.IlmAction",
Self::ScannerHealCandidate => "scanner.HealCandidate",
Self::Dropped => "trace.Dropped",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum TraceVal {
Bool(bool),
U64(u64),
I64(i64),
Str(Arc<str>),
}
impl From<bool> for TraceVal {
fn from(value: bool) -> Self {
Self::Bool(value)
}
}
impl From<u64> for TraceVal {
fn from(value: u64) -> Self {
Self::U64(value)
}
}
impl From<i64> for TraceVal {
fn from(value: i64) -> Self {
Self::I64(value)
}
}
impl From<&str> for TraceVal {
fn from(value: &str) -> Self {
Self::Str(Arc::from(value))
}
}
impl From<String> for TraceVal {
fn from(value: String) -> Self {
Self::Str(Arc::from(value))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TraceAttr {
pub key: &'static str,
pub value: TraceVal,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TraceEvent {
pub kind: TraceKind,
pub func: TraceFunc,
pub time: SystemTime,
pub bucket: Option<Arc<str>>,
pub object: Option<Arc<str>>,
pub duration: Duration,
pub bytes: u64,
pub attrs: SmallVec<[TraceAttr; TRACE_ATTR_INLINE_CAPACITY]>,
}
impl TraceEvent {
pub fn new(kind: TraceKind, func: TraceFunc) -> Self {
Self {
kind,
func,
time: SystemTime::now(),
bucket: None,
object: None,
duration: Duration::ZERO,
bytes: 0,
attrs: SmallVec::new(),
}
}
pub fn with_bucket(mut self, bucket: impl Into<Arc<str>>) -> Self {
self.bucket = Some(bucket.into());
self
}
pub fn with_object(mut self, object: impl Into<Arc<str>>) -> Self {
self.object = Some(object.into());
self
}
pub fn with_duration(mut self, duration: Duration) -> Self {
self.duration = duration;
self
}
pub fn with_bytes(mut self, bytes: u64) -> Self {
self.bytes = bytes;
self
}
pub fn with_attr(mut self, key: &'static str, value: impl Into<TraceVal>) -> Self {
self.attrs.push(TraceAttr {
key,
value: value.into(),
});
self
}
}
#[derive(Debug)]
pub struct TraceBus {
sender: broadcast::Sender<Arc<TraceEvent>>,
subscriber_count: Arc<AtomicUsize>,
}
impl TraceBus {
pub fn new(capacity: usize) -> Self {
let capacity = capacity.max(1);
let (sender, _receiver) = broadcast::channel(capacity);
Self {
sender,
subscriber_count: Arc::new(AtomicUsize::new(0)),
}
}
pub fn subscriber_count(&self) -> usize {
self.subscriber_count.load(Ordering::Acquire)
}
pub fn subscribe(&self) -> TraceSubscription {
let receiver = self.sender.subscribe();
self.subscriber_count.fetch_add(1, Ordering::AcqRel);
TraceSubscription {
receiver,
subscriber_count: Arc::clone(&self.subscriber_count),
}
}
pub fn emit(&self, build: impl FnOnce() -> TraceEvent) -> bool {
if self.subscriber_count() == 0 {
return false;
}
self.sender.send(Arc::new(build())).is_ok()
}
}
impl Default for TraceBus {
fn default() -> Self {
Self::new(DEFAULT_TRACE_BUS_CAPACITY)
}
}
#[derive(Debug)]
pub struct TraceSubscription {
receiver: broadcast::Receiver<Arc<TraceEvent>>,
subscriber_count: Arc<AtomicUsize>,
}
impl TraceSubscription {
pub async fn recv(&mut self) -> Result<Arc<TraceEvent>, broadcast::error::RecvError> {
self.receiver.recv().await
}
pub fn try_recv(&mut self) -> Result<Arc<TraceEvent>, broadcast::error::TryRecvError> {
self.receiver.try_recv()
}
}
impl Drop for TraceSubscription {
fn drop(&mut self) {
self.subscriber_count.fetch_sub(1, Ordering::AcqRel);
}
}
pub fn global_trace_bus() -> &'static TraceBus {
GLOBAL_TRACE_BUS.get_or_init(TraceBus::default)
}
pub fn subscribe_trace_events() -> TraceSubscription {
global_trace_bus().subscribe()
}
pub fn trace_emit(build: impl FnOnce() -> TraceEvent) -> bool {
global_trace_bus().emit(build)
}
pub fn trace_subscriber_count() -> usize {
global_trace_bus().subscriber_count()
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::AtomicUsize;
#[test]
fn trace_emit_skips_builder_without_subscribers() {
let bus = TraceBus::new(4);
let built = AtomicUsize::new(0);
let sent = bus.emit(|| {
built.fetch_add(1, Ordering::Relaxed);
TraceEvent::new(TraceKind::Heal, TraceFunc::HealTask)
});
assert!(!sent);
assert_eq!(built.load(Ordering::Relaxed), 0);
}
#[tokio::test]
async fn trace_subscriber_receives_event() {
let bus = TraceBus::new(4);
let mut subscription = bus.subscribe();
assert!(bus.emit(|| {
TraceEvent::new(TraceKind::Heal, TraceFunc::HealObject)
.with_bucket("bucket")
.with_object("object")
.with_duration(Duration::from_millis(7))
.with_bytes(11)
.with_attr("dry", true)
}));
let event = subscription
.recv()
.await
.expect("subscriber should receive emitted trace event");
assert_eq!(event.kind, TraceKind::Heal);
assert_eq!(event.func, TraceFunc::HealObject);
assert_eq!(event.bucket.as_deref(), Some("bucket"));
assert_eq!(event.object.as_deref(), Some("object"));
assert_eq!(event.duration, Duration::from_millis(7));
assert_eq!(event.bytes, 11);
assert_eq!(
event.attrs.as_slice(),
&[TraceAttr {
key: "dry",
value: TraceVal::Bool(true)
}]
);
}
#[test]
fn trace_subscription_drop_decrements_count() {
let bus = TraceBus::new(4);
let subscription = bus.subscribe();
assert_eq!(bus.subscriber_count(), 1);
drop(subscription);
assert_eq!(bus.subscriber_count(), 0);
}
#[tokio::test]
async fn lagged_subscriber_drops_events_without_blocking_publishers() {
let bus = TraceBus::new(2);
let mut subscription = bus.subscribe();
for index in 0_u64..4 {
assert!(bus.emit(|| { TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerFolder).with_attr("index", index) }));
}
let err = subscription
.recv()
.await
.expect_err("receiver should observe lag instead of blocking publishers");
assert!(matches!(err, broadcast::error::RecvError::Lagged(_)));
}
}
+79 -4
View File
@@ -32,6 +32,7 @@ use rustfs_signer::sign_v4;
use s3s::Body;
use std::ffi::OsStr;
use std::fs as stdfs;
use std::io::ErrorKind;
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::sync::Once;
@@ -51,6 +52,11 @@ pub(crate) const FAST_DATA_USAGE_SCANNER_ENV: &[(&str, &str)] =
&[("RUSTFS_SCANNER_CYCLE", "1"), ("RUSTFS_SCANNER_START_DELAY_SECS", "0")];
pub const TEST_BUCKET: &str = "e2e-test-bucket";
const RUSTFS_FULL_FEATURE: &str = "full";
const TEST_PORT_MIN: u16 = 20_000;
const TEST_PORT_RANGE: u16 = 40_000;
const TEST_PORT_COUNTER_PATH: &str = "/tmp/rustfs_e2e_next_port";
const TEST_PORT_LOCK_DIR: &str = "/tmp/rustfs_e2e_port_allocator.lock";
const TEST_PORT_LOCK_STALE_AFTER: Duration = Duration::from_secs(30);
fn capture_log_path(log_dir: &Path, temp_dir: &str) -> Option<PathBuf> {
let temp_name = Path::new(temp_dir).file_name()?.to_string_lossy();
@@ -67,6 +73,64 @@ fn configured_capture_log_path(temp_dir: &str) -> Option<String> {
capture_log_path(Path::new(&log_dir), temp_dir).map(|path| path.to_string_lossy().into_owned())
}
struct PortAllocatorGuard;
impl PortAllocatorGuard {
async fn acquire() -> Result<Self, Box<dyn std::error::Error + Send + Sync>> {
loop {
match stdfs::create_dir(TEST_PORT_LOCK_DIR) {
Ok(()) => return Ok(Self),
Err(err) if err.kind() == ErrorKind::AlreadyExists => {
remove_stale_port_allocator_lock();
sleep(Duration::from_millis(10)).await;
}
Err(err) => return Err(err.into()),
}
}
}
}
impl Drop for PortAllocatorGuard {
fn drop(&mut self) {
let _ = stdfs::remove_dir(TEST_PORT_LOCK_DIR);
}
}
fn advance_test_port(port: u16) -> u16 {
let offset = (port - TEST_PORT_MIN + 1) % TEST_PORT_RANGE;
TEST_PORT_MIN + offset
}
fn seeded_test_port() -> u16 {
let offset = (Uuid::new_v4().as_u128() % u128::from(TEST_PORT_RANGE)) as u16;
TEST_PORT_MIN + offset
}
fn read_next_test_port() -> u16 {
stdfs::read_to_string(TEST_PORT_COUNTER_PATH)
.ok()
.and_then(|value| value.trim().parse::<u16>().ok())
.filter(|port| (TEST_PORT_MIN..TEST_PORT_MIN + TEST_PORT_RANGE).contains(port))
.unwrap_or_else(seeded_test_port)
}
fn remove_stale_port_allocator_lock() {
let Ok(metadata) = stdfs::metadata(TEST_PORT_LOCK_DIR) else {
return;
};
let Ok(modified) = metadata.modified() else {
return;
};
if modified.elapsed().is_ok_and(|elapsed| elapsed > TEST_PORT_LOCK_STALE_AFTER) {
let _ = stdfs::remove_dir(TEST_PORT_LOCK_DIR);
}
}
fn write_next_test_port(port: u16) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
stdfs::write(TEST_PORT_COUNTER_PATH, port.to_string())?;
Ok(())
}
pub(crate) fn capture_command_logs(
command: &mut Command,
log_path: Option<&str>,
@@ -508,10 +572,21 @@ impl RustFSTestEnvironment {
/// Find an available port for the test
pub async fn find_available_port() -> Result<u16, Box<dyn std::error::Error + Send + Sync>> {
use std::net::TcpListener;
let listener = TcpListener::bind("127.0.0.1:0")?;
let port = listener.local_addr()?.port();
drop(listener);
Ok(port)
let _guard = PortAllocatorGuard::acquire().await?;
let mut next_port = read_next_test_port();
for _ in 0..TEST_PORT_RANGE {
let port = next_port;
next_port = advance_test_port(next_port);
write_next_test_port(next_port)?;
if let Ok(listener) = TcpListener::bind(("127.0.0.1", port)) {
drop(listener);
return Ok(port);
}
}
Err("no available E2E test port found".into())
}
/// Kill any existing RustFS processes
+2 -1
View File
@@ -380,7 +380,7 @@ pub mod erasure {
pub mod event {
pub use crate::event::name::EventName;
pub use crate::services::event_notification::{EventArgs, register_event_dispatch_hook};
pub use crate::services::event_notification::{EventArgs, register_event_dispatch_hook, send_event};
}
pub mod global {
@@ -483,6 +483,7 @@ pub mod store_list {
}
pub mod storage {
pub use crate::core::pools::HealLifecycleExpiryContext;
pub use crate::store::HealWalkVersion;
pub use crate::store::{
ECStore, all_local_disk, all_local_disk_path, find_local_disk_by_ref, init_local_disks,
@@ -1549,8 +1549,8 @@ impl Default for PutObjectOptions {
}
}
#[allow(dead_code)]
impl PutObjectOptions {
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn set_match_etag(&mut self, etag: &str) {
if etag == "*" {
self.custom_header
@@ -1561,6 +1561,7 @@ impl PutObjectOptions {
}
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn set_match_etag_except(&mut self, etag: &str) {
if etag == "*" {
self.custom_header
@@ -1696,6 +1697,7 @@ impl PutObjectOptions {
header
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn validate(&self, _c: Arc<TargetClient>) -> Result<(), std::io::Error> {
//if self.checksum.is_set() {
/*if !self.trailing_header_support {
@@ -456,16 +456,23 @@ impl<'a> LifecycleExpiryTrace<'a> {
}
}
#[allow(dead_code)]
impl ExpiryStats {
pub fn missed_tasks(&self) -> i64 {
self.missed_expiry_tasks.load(Ordering::SeqCst)
}
#[allow(
dead_code,
reason = "asserted by this file's tests; the lib target cannot see test-only consumers (backlog#1823)"
)]
fn missed_free_vers_tasks(&self) -> i64 {
self.missed_freevers_tasks.load(Ordering::SeqCst)
}
#[allow(
dead_code,
reason = "asserted by this file's tests; the lib target cannot see test-only consumers (backlog#1823)"
)]
fn missed_tier_journal_tasks(&self) -> i64 {
self.missed_tier_journal_tasks.load(Ordering::SeqCst)
}
+1 -1
View File
@@ -19,7 +19,7 @@ pub mod core;
pub mod evaluator;
pub mod manual_transition_job;
mod metadata_boundary;
pub(crate) use metadata_boundary::get_expiry_configs;
pub(crate) use metadata_boundary::{LifecycleExpiryConfigs, get_expiry_configs};
mod object_lock_boundary;
pub use self::core as lifecycle;
mod replication_sink;
@@ -80,7 +80,10 @@ impl LastDayTierStats {
}
}
#[allow(dead_code)]
#[allow(
dead_code,
reason = "asserted by this file's tests; the lib target cannot see test-only consumers (backlog#1823)"
)]
fn merge(&self, m: LastDayTierStats) -> LastDayTierStats {
let mut cl = self.clone();
let mut cm = m;
@@ -177,9 +177,10 @@ fn should_record_remote_delete_failure(err: &std::io::Error) -> bool {
}
#[derive(Default)]
#[allow(dead_code)]
struct ObjSweeper {
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
object: String,
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
bucket: String,
version_id: Option<Uuid>,
versioned: bool,
@@ -191,9 +192,9 @@ struct ObjSweeper {
remote_object: String,
}
#[allow(dead_code)]
impl ObjSweeper {
#[allow(clippy::new_ret_no_self)]
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub async fn new(bucket: &str, object: &str) -> Result<Self, std::io::Error> {
Ok(Self {
object: object.into(),
@@ -202,17 +203,20 @@ impl ObjSweeper {
})
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub fn with_version(&mut self, vid: Option<Uuid>) -> &Self {
self.version_id = vid.clone();
self
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub fn with_versioning(&mut self, versioned: bool, suspended: bool) -> &Self {
self.versioned = versioned;
self.suspended = suspended;
self
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub fn get_opts(&self) -> lifecycle::ObjectOpts {
let mut opts = ObjectOpts {
version_id: self.version_id.clone(),
@@ -226,6 +230,7 @@ impl ObjSweeper {
opts
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub fn set_transition_state(&mut self, info: TransitionedObject) {
self.transition_tier = info.tier;
self.transition_status = info.status;
@@ -266,6 +271,7 @@ impl ObjSweeper {
None
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub async fn sweep(&self, api: Arc<ECStore>) {
let Some(je) = self.should_remove_remote_object() else {
return;
-2
View File
@@ -312,9 +312,7 @@ mod tests {
}
#[derive(Deserialize)]
struct LegacyBucketQuota {
#[allow(dead_code)]
quota: Option<u64>,
#[allow(dead_code)]
quota_type: LegacyQuotaType,
}
let legacy = serde_json::from_slice::<LegacyBucketQuota>(&json)
+2 -2
View File
@@ -95,7 +95,6 @@ impl TransitionClient {
}
#[derive(Default)]
#[allow(dead_code)]
pub struct GetRequest {
pub buffer: Vec<u8>,
pub offset: i64,
@@ -107,11 +106,12 @@ pub struct GetRequest {
pub setting_object_info: bool,
}
#[allow(dead_code)]
pub struct GetResponse {
pub size: i64,
//pub error: error,
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
pub did_read: bool,
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
pub object_info: ObjectInfo,
}
+2 -2
View File
@@ -20,6 +20,7 @@
#![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;
@@ -27,7 +28,6 @@ use tracing::warn;
use crate::client::api_error_response::err_invalid_argument;
#[derive(Default)]
#[allow(dead_code)]
pub struct AdvancedGetOptions {
pub replication_delete_marker: bool,
pub is_replication_ready_for_delete_marker: bool,
@@ -77,7 +77,7 @@ impl GetObjectOptions {
}
}
if self.checksum {
headers.insert(HeaderName::from_static("x-amz-checksum-mode"), HeaderValue::from_static("ENABLED"));
headers.insert(HeaderName::from_static(AMZ_CHECKSUM_MODE), HeaderValue::from_static("ENABLED"));
}
headers
}
-1
View File
@@ -360,7 +360,6 @@ impl TransitionClient {
}
#[derive(Default)]
#[allow(dead_code)]
pub struct ListObjectsOptions {
reverse_versions: bool,
with_versions: bool,
+3 -1
View File
@@ -137,8 +137,8 @@ impl Default for PutObjectOptions {
}
}
#[allow(dead_code)]
impl PutObjectOptions {
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn set_match_etag(&mut self, etag: &str) {
if etag == "*" {
self.custom_header.insert("If-Match", HeaderValue::from_static("*"));
@@ -149,6 +149,7 @@ impl PutObjectOptions {
}
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn set_match_etag_except(&mut self, etag: &str) {
if etag == "*" {
self.custom_header.insert("If-None-Match", HeaderValue::from_static("*"));
@@ -259,6 +260,7 @@ impl PutObjectOptions {
header
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn validate(&self, c: TransitionClient) -> Result<(), std::io::Error> {
//if self.checksum.is_set() {
/*if !self.trailing_header_support {
+2 -3
View File
@@ -55,7 +55,6 @@ pub struct RemoveBucketOptions {
const DELETE_RESPONSE_PREVIEW_LEN: usize = 1024;
#[derive(Debug)]
#[allow(dead_code)]
pub struct AdvancedRemoveOptions {
pub replication_delete_marker: bool,
pub replication_status: ReplicationStatus,
@@ -465,10 +464,10 @@ impl TransitionClient {
}
#[derive(Debug, Default)]
#[allow(dead_code)]
pub struct RemoveObjectError {
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
object_name: String,
#[allow(dead_code)]
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
version_id: String,
err: Option<std::io::Error>,
}
+3 -3
View File
@@ -372,8 +372,8 @@ pub struct Checksum {
computed: bool,
}
#[allow(dead_code)]
impl Checksum {
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn new(t: ChecksumMode, b: &[u8]) -> Checksum {
if t.is_set() && b.len() == t.raw_byte_len() {
return Checksum {
@@ -385,7 +385,7 @@ impl Checksum {
Checksum::default()
}
#[allow(dead_code)]
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn new_checksum_string(t: ChecksumMode, s: &str) -> Result<Checksum, std::io::Error> {
let b = match base64_decode(s.as_bytes()) {
Ok(b) => b,
@@ -412,7 +412,7 @@ impl Checksum {
base64_encode(&self.r)
}
#[allow(dead_code)]
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn raw(&self) -> Option<Vec<u8>> {
if !self.is_set() {
return None;
@@ -37,16 +37,17 @@ pub struct PutObjReader {
//pub sealMD5Fn: SealMD5CurrFn,
}
#[allow(dead_code)]
impl PutObjReader {
pub fn new(reader: HashReader) -> Self {
Self { reader }
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn md5_current_hex_string(&self) -> String {
self.reader.checksum().map(|v| v.encoded).unwrap_or_default()
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn with_encryption(&mut self, enc_reader: HashReader) -> Result<(), std::io::Error> {
self.reader = enc_reader;
+10 -6
View File
@@ -54,6 +54,10 @@ 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},
};
@@ -1383,12 +1387,12 @@ pub(crate) fn to_object_info_for_provider(
};
// Extract checksums
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");
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);
// Build and return the ObjectInfo struct
Ok(ObjectInfo {
-3
View File
@@ -39,7 +39,6 @@ use rustfs_config::{
};
use std::sync::LazyLock;
#[allow(dead_code)]
#[allow(clippy::declare_interior_mutable_const)]
/// Default KVS for audit webhook settings.
pub static DEFAULT_AUDIT_WEBHOOK_KVS: LazyLock<KVS> = LazyLock::new(|| {
@@ -117,7 +116,6 @@ pub static DEFAULT_AUDIT_WEBHOOK_KVS: LazyLock<KVS> = LazyLock::new(|| {
])
});
#[allow(dead_code)]
#[allow(clippy::declare_interior_mutable_const)]
/// Default KVS for audit MQTT settings.
pub static DEFAULT_AUDIT_MQTT_KVS: LazyLock<KVS> = LazyLock::new(|| {
@@ -375,7 +373,6 @@ pub static DEFAULT_AUDIT_NATS_KVS: LazyLock<KVS> = LazyLock::new(|| {
])
});
#[allow(dead_code)]
pub static DEFAULT_AUDIT_PULSAR_KVS: LazyLock<KVS> = LazyLock::new(|| {
KVS(vec![
KV {
-59
View File
@@ -12,12 +12,9 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::error::{Error, Result};
use rustfs_config::server_config::{KV, KVS};
use rustfs_config::{DEFAULT_HEAL_BITROT_CYCLE_SECS, HEAL_BITROT_CYCLE};
use rustfs_utils::string::parse_bool;
use std::sync::LazyLock;
use std::time::Duration;
pub static DEFAULT_KVS: LazyLock<KVS> = LazyLock::new(|| {
KVS(vec![KV {
@@ -26,59 +23,3 @@ pub static DEFAULT_KVS: LazyLock<KVS> = LazyLock::new(|| {
hidden_if_empty: false,
}])
});
#[derive(Debug, Default)]
pub struct Config {
pub bitrot: String,
pub sleep: Duration,
pub io_count: usize,
pub drive_workers: usize,
pub cache: Duration,
}
impl Config {
pub fn bitrot_scan_cycle(&self) -> Duration {
self.cache
}
pub fn get_workers(&self) -> usize {
self.drive_workers
}
pub fn update(&mut self, nopts: &Config) {
self.bitrot = nopts.bitrot.clone();
self.io_count = nopts.io_count;
self.sleep = nopts.sleep;
self.drive_workers = nopts.drive_workers;
}
}
const RUSTFS_BITROT_CYCLE_IN_MONTHS: u64 = 1;
fn parse_bitrot_config(s: &str) -> Result<Duration> {
match parse_bool(s) {
Ok(enabled) => {
if enabled {
Ok(Duration::from_secs_f64(0.0))
} else {
Ok(Duration::from_secs_f64(-1.0))
}
}
Err(_) => {
if !s.ends_with("m") {
return Err(Error::other("unknown format"));
}
match s.trim_end_matches('m').parse::<u64>() {
Ok(months) => {
if months < RUSTFS_BITROT_CYCLE_IN_MONTHS {
return Err(Error::other(format!("minimum bitrot cycle is {RUSTFS_BITROT_CYCLE_IN_MONTHS} month(s)")));
}
Ok(Duration::from_secs(months * 30 * 24 * 60))
}
Err(err) => Err(Error::other(err)),
}
}
}
}
-1
View File
@@ -16,7 +16,6 @@
mod audit;
pub mod com;
#[allow(dead_code)]
pub mod heal;
mod notify;
mod oidc;
+100 -3
View File
@@ -16,6 +16,7 @@ use crate::bucket::replication::replication_state_from_filemeta;
use crate::bucket::versioning_sys::BucketVersioningSys;
use crate::bucket::{
lifecycle::{
LifecycleExpiryConfigs,
bucket_lifecycle_audit::LcEventSrc,
bucket_lifecycle_ops::{
LifecycleOps, apply_expiry_on_transitioned_object, apply_expiry_rule_in, eval_action_from_lifecycle,
@@ -1996,11 +1997,11 @@ impl PoolMeta {
Ok(false)
}
#[allow(dead_code)]
pub fn validate(&self, pools: Vec<Arc<Sets>>) -> Result<bool> {
struct PoolInfo {
position: usize,
completed: bool,
#[allow(dead_code, reason = "written but never read back (backlog#1823)")]
decom_started: bool,
}
@@ -2335,6 +2336,10 @@ fn lifecycle_action_removes_data_movement_version(action: IlmAction) -> bool {
)
}
fn lifecycle_action_skips_heal_version(action: IlmAction) -> bool {
action.delete()
}
fn resolve_data_movement_lifecycle_expiry_result(action: IlmAction, apply_actions: bool, applied: bool) -> Result<bool> {
if !apply_actions || applied {
return Ok(true);
@@ -2385,7 +2390,80 @@ pub(crate) async fn should_skip_lifecycle_for_data_movement(
}
}
pub struct HealLifecycleExpiryContext {
configs: LifecycleExpiryConfigs,
}
impl ECStore {
pub async fn load_heal_lifecycle_expiry_context(&self, bucket: &str) -> Result<Option<HealLifecycleExpiryContext>> {
if bucket == RUSTFS_META_BUCKET {
return Ok(None);
}
let configs = get_expiry_configs(self, bucket).await?;
if configs.lifecycle.is_none() {
return Ok(None);
}
Ok(Some(HealLifecycleExpiryContext { configs }))
}
pub async fn enqueue_heal_lifecycle_expiry(
self: &Arc<Self>,
context: &HealLifecycleExpiryContext,
bucket: &str,
object: &str,
version_id: Option<&str>,
object_info: Option<&crate::object_api::ObjectInfo>,
) -> Result<bool> {
let Some(lifecycle_config) = context.configs.lifecycle.as_ref() else {
return Ok(false);
};
let object_info = if let Some(object_info) = object_info {
if object_info.bucket != bucket || object_info.name != object {
return Ok(false);
}
let snapshot_version_id = object_info
.version_id
.filter(|version_id| !version_id.is_nil())
.map(|version_id| version_id.to_string());
if snapshot_version_id.as_deref() != version_id {
return Ok(false);
}
object_info.clone()
} else {
match self
.get_object_info(
bucket,
object,
&ObjectOptions {
version_id: version_id.map(str::to_string),
versioned: version_id.is_some(),
expected_bucket_incarnation_id: Some(context.configs.bucket_incarnation_id),
..Default::default()
},
)
.await
{
Ok(object_info) => object_info,
Err(err) if is_err_object_not_found(&err) || is_err_version_not_found(&err) => return Ok(false),
Err(err) => return Err(err),
}
};
let event = eval_action_from_lifecycle(lifecycle_config, context.configs.object_lock.as_deref(), &object_info).await;
if !lifecycle_action_skips_heal_version(event.action) {
return Ok(false);
}
if lifecycle_delete_all_versions_blocked_by_replication(self.clone(), bucket, &object_info.name, event.action).await? {
return Ok(false);
}
Ok(apply_expiry_rule_in(self.clone(), &event, &LcEventSrc::Scanner, &object_info).await)
}
async fn save_current_pool_meta(&self) -> Result<()> {
let _save_guard = self.pool_meta_save_gate.lock().await;
let snapshot = {
@@ -4287,6 +4365,19 @@ mod tests {
));
}
#[test]
fn lifecycle_action_skips_heal_version_for_every_delete_action() {
assert!(lifecycle_action_skips_heal_version(IlmAction::DeleteAction));
assert!(lifecycle_action_skips_heal_version(IlmAction::DeleteVersionAction));
assert!(lifecycle_action_skips_heal_version(IlmAction::DeleteRestoredAction));
assert!(lifecycle_action_skips_heal_version(IlmAction::DeleteRestoredVersionAction));
assert!(lifecycle_action_skips_heal_version(IlmAction::DeleteAllVersionsAction));
assert!(lifecycle_action_skips_heal_version(IlmAction::DelMarkerDeleteAllVersionsAction));
assert!(!lifecycle_action_skips_heal_version(IlmAction::TransitionAction));
assert!(!lifecycle_action_skips_heal_version(IlmAction::TransitionVersionAction));
assert!(!lifecycle_action_skips_heal_version(IlmAction::NoneAction));
}
#[test]
fn resolve_data_movement_lifecycle_expiry_result_allows_dry_run_skip() {
let skip = resolve_data_movement_lifecycle_expiry_result(IlmAction::DeleteVersionAction, false, false)
@@ -4958,13 +5049,19 @@ fn is_disk_online_state(state: &str) -> bool {
}
#[deprecated(since = "0.1.0", note = "Use fallback_total_capacity_dedup instead")]
#[allow(dead_code)]
#[allow(
dead_code,
reason = "superseded by the replacement named in the comment at pools.rs:5071 (backlog#1823)"
)]
fn fallback_total_capacity(disks: &[rustfs_madmin::Disk]) -> usize {
fallback_total_capacity_dedup(disks)
}
#[deprecated(since = "0.1.0", note = "Use fallback_free_capacity_dedup instead")]
#[allow(dead_code)]
#[allow(
dead_code,
reason = "superseded by the replacement named in the comment at pools.rs:5071 (backlog#1823)"
)]
fn fallback_free_capacity(disks: &[rustfs_madmin::Disk]) -> usize {
fallback_free_capacity_dedup(disks)
}
+20 -9
View File
@@ -1140,11 +1140,11 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
Err(Error::DiskNotFound)
}
#[tracing::instrument(skip(self))]
async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> {
// Multipart orphan reconciliation is intentionally retained above the pool/set layers
// until there is a concrete caller and a stable lower-level contract to implement.
Err(StorageError::NotImplemented)
#[tracing::instrument(level = "debug", skip(self, opts), fields(bucket = %bucket, object = %object, dry_run = opts.dry_run))]
async fn check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> {
self.get_disks_for_heal_object(object, opts)?
.check_abandoned_parts(bucket, object, opts)
.await
}
}
@@ -1996,7 +1996,7 @@ mod tests {
}
#[tokio::test]
async fn sets_check_abandoned_parts_returns_typed_not_implemented_error() {
async fn sets_check_abandoned_parts_rejects_invalid_set_scope() {
let format = FormatV3::new(1, 1);
let sets = Sets {
id: format.id,
@@ -2021,10 +2021,21 @@ mod tests {
};
let err = sets
.check_abandoned_parts("bucket", "object", &HealOpts::default())
.check_abandoned_parts(
"bucket",
"object",
&HealOpts {
set: Some(1),
..Default::default()
},
)
.await
.expect_err("abandoned-parts ownership should stay above the pool/set storage layers");
assert!(matches!(err, StorageError::NotImplemented));
.expect_err("out-of-range abandoned-parts set scope must fail closed");
assert!(
matches!(err, StorageError::InvalidArgument(_, ref field, ref reason)
if field == "set" && reason.contains("invalid heal set index 1")),
"unexpected invalid set error: {err:?}"
);
}
// Builds a single-set `Sets` over `SET_DRIVE_COUNT` local temp-dir disks,
+1 -1
View File
@@ -6562,7 +6562,7 @@ impl LocalDisk {
Ok(f)
}
#[allow(dead_code)]
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn get_metrics(&self) -> DiskMetrics {
DiskMetrics::default()
}
@@ -132,7 +132,6 @@ impl RebalanceStopPropagationRecord {
}
}
#[allow(dead_code)]
#[derive(Debug, Clone, Default)]
pub struct DiskStat {
pub total_space: u64,
@@ -16,8 +16,16 @@ use serde::{Deserialize, Serialize};
use std::{fmt::Display, io};
use tracing::info;
#[allow(
dead_code,
reason = "tier config wire version stamped by the parity constructors below (backlog#1823)"
)]
const C_TIER_CONFIG_VER: &str = "v1";
#[allow(
dead_code,
reason = "tier-name validation message reached only from the parity constructors below (backlog#1823)"
)]
const ERR_TIER_NAME_EMPTY: &str = "remote tier name empty";
const WASABI_US_EAST_ENDPOINT: &str = "https://s3.wasabisys.com";
const WASABI_ALTERNATIVE_ENDPOINTS: &[(&str, &str)] = &[
@@ -264,7 +272,6 @@ impl Clone for TierConfig {
}
}
#[allow(dead_code)]
impl TierConfig {
pub(crate) fn clone_with_credentials(&self) -> Self {
Self {
@@ -284,6 +291,7 @@ impl TierConfig {
}
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn endpoint(&self) -> String {
match self.tier_type {
TierType::S3 => self.s3.as_ref().map(|s| s.endpoint.clone()).unwrap_or_default(),
@@ -303,6 +311,7 @@ impl TierConfig {
}
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn bucket(&self) -> String {
match self.tier_type {
TierType::S3 => self.s3.as_ref().map(|s| s.bucket.clone()).unwrap_or_default(),
@@ -322,6 +331,7 @@ impl TierConfig {
}
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn prefix(&self) -> String {
match self.tier_type {
TierType::S3 => self.s3.as_ref().map(|s| s.prefix.clone()).unwrap_or_default(),
@@ -341,6 +351,7 @@ impl TierConfig {
}
}
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn region(&self) -> String {
match self.tier_type {
TierType::S3 => self.s3.as_ref().map(|s| s.region.clone()).unwrap_or_default(),
@@ -457,7 +468,7 @@ impl TierWasabi {
}
impl TierS3 {
#[allow(dead_code)]
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn create<F>(
name: &str,
access_key: &str,
@@ -528,7 +539,7 @@ pub struct TierMinIO {
}
impl TierMinIO {
#[allow(dead_code)]
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
fn create<F>(
name: &str,
endpoint: &str,
@@ -14,7 +14,6 @@
use crate::services::tier::tier::TierConfigMgr;
#[allow(dead_code)]
impl TierConfigMgr {
pub fn msg_size(&self) -> usize {
100
@@ -4860,6 +4860,14 @@ impl SetDisks {
/// is best-effort maintenance: individual delete failures are logged and
/// skipped rather than propagated.
pub(crate) async fn reclaim_orphan_data_dirs(&self, bucket: &str, object: &str) -> disk::error::Result<usize> {
self.reclaim_orphan_data_dirs_inner(bucket, object, false).await
}
pub(crate) async fn dry_run_reclaim_orphan_data_dirs(&self, bucket: &str, object: &str) -> disk::error::Result<usize> {
self.reclaim_orphan_data_dirs_inner(bucket, object, true).await
}
async fn reclaim_orphan_data_dirs_inner(&self, bucket: &str, object: &str, dry_run: bool) -> disk::error::Result<usize> {
let disks = self.get_disks_internal().await;
// Phase 1 (read-only): build the referenced-data-dir union and record the
@@ -4967,6 +4975,20 @@ impl SetDisks {
continue;
}
let stray = format!("{object}/{dir}");
if dry_run {
removed += 1;
debug!(
target: "rustfs_ecstore::set_disk",
event = "heal_abandoned_parts",
component = "ecstore",
subsystem = "heal",
state = "dry_run_matched",
result = "matched",
bucket, object, data_dir = %dir,
"Heal abandoned parts dry-run matched orphaned data directory"
);
continue;
}
match disk
.delete(
bucket,
+105 -4
View File
@@ -6998,6 +6998,100 @@ mod tests {
assert!(object_dir.join(STORAGE_FORMAT_FILE).exists(), "metadata must be preserved");
}
async fn recv_abandoned_parts_trace(
trace: &mut rustfs_common::trace_bus::TraceSubscription,
bucket: &str,
object: &str,
state: &str,
) -> rustfs_common::trace_bus::TraceEvent {
for _ in 0..32 {
let event = tokio::time::timeout(std::time::Duration::from_secs(1), trace.recv())
.await
.expect("abandoned-parts trace event should arrive")
.expect("trace bus should stay open");
if event.kind == rustfs_common::trace_bus::TraceKind::Heal
&& event.func == rustfs_common::trace_bus::TraceFunc::HealCheckAbandonedParts
&& event.bucket.as_deref() == Some(bucket)
&& event.object.as_deref() == Some(object)
&& trace_attr_string(&event, "state").as_deref() == Some(state)
{
return (*event).clone();
}
}
panic!("expected abandoned-parts trace state {state} for {bucket}/{object}");
}
fn trace_attr_string(event: &rustfs_common::trace_bus::TraceEvent, key: &str) -> Option<String> {
event.attrs.iter().find_map(|attr| {
if attr.key != key {
return None;
}
Some(match &attr.value {
rustfs_common::trace_bus::TraceVal::Bool(value) => value.to_string(),
rustfs_common::trace_bus::TraceVal::U64(value) => value.to_string(),
rustfs_common::trace_bus::TraceVal::I64(value) => value.to_string(),
rustfs_common::trace_bus::TraceVal::Str(value) => value.to_string(),
})
})
}
#[tokio::test]
async fn check_abandoned_parts_dry_run_counts_without_deleting() {
let mut trace = rustfs_common::trace_bus::subscribe_trace_events();
let (dir, disk) = make_single_local_disk().await;
let live = Uuid::new_v4();
let orphan = Uuid::new_v4();
let object_dir = dir.path().join("bucket").join("obj");
write_object_meta_with_data_dirs(&object_dir, "bucket", "obj", &[live]).await;
fs::create_dir_all(object_dir.join(live.to_string()))
.await
.expect("live data dir should be created");
fs::create_dir_all(object_dir.join(orphan.to_string()))
.await
.expect("orphan data dir should be created");
let set = make_set_disks_with(vec![Some(disk)]).await;
set.check_abandoned_parts(
"bucket",
"obj",
&HealOpts {
dry_run: true,
no_lock: true,
..Default::default()
},
)
.await
.expect("dry-run abandoned-parts check should succeed");
let dry_run_trace = recv_abandoned_parts_trace(&mut trace, "bucket", "obj", "dry_run_matched").await;
assert_eq!(trace_attr_string(&dry_run_trace, "dry_run").as_deref(), Some("true"));
assert_eq!(trace_attr_string(&dry_run_trace, "data_dirs").as_deref(), Some("1"));
assert!(object_dir.join(live.to_string()).exists(), "referenced data dir must be preserved");
assert!(object_dir.join(orphan.to_string()).exists(), "dry-run must not remove orphaned data dir");
set.check_abandoned_parts(
"bucket",
"obj",
&HealOpts {
no_lock: true,
..Default::default()
},
)
.await
.expect("abandoned-parts check should reclaim stale data dir");
let reclaim_trace = recv_abandoned_parts_trace(&mut trace, "bucket", "obj", "reclaimed").await;
assert_eq!(trace_attr_string(&reclaim_trace, "dry_run").as_deref(), Some("false"));
assert_eq!(trace_attr_string(&reclaim_trace, "data_dirs").as_deref(), Some("1"));
assert!(
object_dir.join(live.to_string()).exists(),
"referenced data dir must remain after reclaim"
);
assert!(!object_dir.join(orphan.to_string()).exists(), "orphaned data dir must be removed");
}
#[tokio::test]
async fn reclaim_orphan_data_dirs_recovers_deferred_cleanup_after_restart() {
let (dir, disk) = make_single_local_disk().await;
@@ -12233,11 +12327,18 @@ mod tests {
.expect_err("unsupported copy_object_part should return a typed error");
assert!(matches!(copy_part_err, StorageError::NotImplemented));
let abandoned_err = set_disks
.check_abandoned_parts("bucket", "object", &HealOpts::default())
set_disks
.check_abandoned_parts(
"bucket",
"object",
&HealOpts {
dry_run: true,
no_lock: true,
..Default::default()
},
)
.await
.expect_err("abandoned-parts check should stay in the upper reconciliation layer");
assert!(matches!(abandoned_err, StorageError::NotImplemented));
.expect("abandoned-parts check should be callable on empty disk sets");
}
#[tokio::test]
+56 -5
View File
@@ -16,6 +16,7 @@ use super::super::*;
use crate::disk::disk_store::DiskStoreRenameDataExt;
use crate::io_support::bitrot::object_mmap_read_enabled;
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit};
use tracing::trace;
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
@@ -2057,11 +2058,61 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks {
Err(Error::DiskNotFound)
}
#[tracing::instrument(skip(self))]
async fn check_abandoned_parts(&self, _bucket: &str, _object: &str, _opts: &HealOpts) -> Result<()> {
// Multipart orphan reconciliation is intentionally retained above the set layer
// until there is a concrete caller and a stable lower-level contract to implement.
Err(StorageError::NotImplemented)
#[tracing::instrument(level = "debug", skip(self, opts), fields(bucket = %bucket, object = %object, dry_run = opts.dry_run))]
async fn check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> {
let started_at = std::time::Instant::now();
let _write_lock_guard = if !opts.no_lock {
let ns_lock = self.new_ns_lock(bucket, object).await?;
Some(
ns_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(|e| self.map_namespace_lock_error(bucket, object, "write", e))?,
)
} else {
None
};
let removed = if opts.dry_run {
self.dry_run_reclaim_orphan_data_dirs(bucket, object).await?
} else {
self.reclaim_orphan_data_dirs(bucket, object).await?
};
let state = if opts.dry_run && removed > 0 {
"dry_run_matched"
} else if removed > 0 {
"reclaimed"
} else {
"checked"
};
let data_dirs = u64::try_from(removed).unwrap_or(u64::MAX);
trace_emit(|| {
TraceEvent::new(TraceKind::Heal, TraceFunc::HealCheckAbandonedParts)
.with_bucket(bucket)
.with_object(object)
.with_duration(started_at.elapsed())
.with_attr("state", state)
.with_attr("dry_run", opts.dry_run)
.with_attr("data_dirs", data_dirs)
});
if removed > 0 {
trace!(
event = "heal_abandoned_parts",
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_HEAL,
state = if opts.dry_run { "dry_run_matched" } else { "reclaimed" },
result = "ok",
bucket,
object,
dry_run = opts.dry_run,
data_dirs = removed,
"Heal abandoned parts checked object data directories"
);
}
Ok(())
}
}
+46 -4
View File
@@ -23,6 +23,7 @@
//! per-version `SetDisks::heal_object`.
use super::super::*;
use crate::object_api::ObjectInfo;
use std::collections::HashSet;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
@@ -39,12 +40,16 @@ const BACKGROUND_WALKDIR_STALL_TIMEOUT: Duration = Duration::from_secs(60);
/// it must not gate healing logic — the delete-marker vs data path is chosen
/// inside `ops/heal.rs` from the resolved latest metadata. `version_id` is
/// normalized (nil/absent UUID => `None`).
#[derive(Debug, Clone, PartialEq, Eq)]
#[derive(Debug, Clone)]
pub struct HealWalkVersion {
/// object key
pub name: String,
/// normalized version id (`None` when the version is nil/absent)
pub version_id: Option<String>,
/// version modification time as Unix nanoseconds
pub mod_time_unix_nanos: Option<i128>,
/// object snapshot for lifecycle evaluation
pub lifecycle_object_info: Option<ObjectInfo>,
/// whether this version is a delete marker (observability only)
pub is_delete_marker: bool,
}
@@ -63,6 +68,7 @@ struct HealWalkCollector {
bucket: String,
batch_objects: usize,
version_budget: usize,
include_lifecycle_object_info: bool,
objects: Mutex<Vec<HealWalkObject>>,
decode_error: Mutex<Option<DiskError>>,
version_total: AtomicUsize,
@@ -116,10 +122,25 @@ impl HealWalkCollector {
let mut versions = Vec::with_capacity(fiv.versions.len() + fiv.free_versions.len());
for fi in fiv.versions.iter().chain(fiv.free_versions.iter()) {
let version_uuid = fi.version_id.filter(|version_id| !version_id.is_nil());
let lifecycle_object_info = if self.include_lifecycle_object_info {
let mut lifecycle_fi = fi.clone();
lifecycle_fi.version_id = version_uuid;
Some(ObjectInfo::from_file_info(
&lifecycle_fi,
&self.bucket,
&entry.name,
version_uuid.is_some(),
))
} else {
None
};
versions.push(HealWalkVersion {
name: entry.name.clone(),
// Normalize: nil/absent version id => None.
version_id: fi.version_id.filter(|u| !u.is_nil()).map(|u| u.to_string()),
version_id: version_uuid.map(|u| u.to_string()),
mod_time_unix_nanos: fi.mod_time.map(|mod_time| mod_time.unix_timestamp_nanos()),
lifecycle_object_info,
is_delete_marker: fi.deleted,
});
}
@@ -173,11 +194,26 @@ impl HealWalkCollector {
}
};
for fi in fiv.versions.iter().chain(fiv.free_versions.iter()) {
let vid = fi.version_id.filter(|u| !u.is_nil()).map(|u| u.to_string());
let version_uuid = fi.version_id.filter(|version_id| !version_id.is_nil());
let vid = version_uuid.map(|u| u.to_string());
if seen.insert(vid.clone()) {
let lifecycle_object_info = if self.include_lifecycle_object_info {
let mut lifecycle_fi = fi.clone();
lifecycle_fi.version_id = version_uuid;
Some(ObjectInfo::from_file_info(
&lifecycle_fi,
&self.bucket,
&entry.name,
version_uuid.is_some(),
))
} else {
None
};
versions.push(HealWalkVersion {
name: entry.name.clone(),
version_id: vid,
mod_time_unix_nanos: fi.mod_time.map(|mod_time| mod_time.unix_timestamp_nanos()),
lifecycle_object_info,
is_delete_marker: fi.deleted,
});
}
@@ -255,6 +291,7 @@ impl SetDisks {
forward_to: Option<&str>,
batch_objects: usize,
version_budget: usize,
include_lifecycle_object_info: bool,
) -> disk::error::Result<(Vec<HealWalkVersion>, Option<String>, bool)> {
assert!(batch_objects >= 2, "heal_walk_versions_page requires batch_objects >= 2");
@@ -264,6 +301,7 @@ impl SetDisks {
bucket: bucket.to_string(),
batch_objects,
version_budget: version_budget.max(1),
include_lifecycle_object_info,
objects: Mutex::new(Vec::new()),
decode_error: Mutex::new(None),
version_total: AtomicUsize::new(0),
@@ -347,6 +385,7 @@ mod tests {
bucket: "bucket".to_string(),
batch_objects: 2,
version_budget: 2,
include_lifecycle_object_info: false,
objects: Mutex::new(Vec::new()),
decode_error: Mutex::new(None),
version_total: AtomicUsize::new(0),
@@ -388,6 +427,8 @@ mod tests {
HealWalkVersion {
name: name.to_string(),
version_id: Some(id.to_string()),
mod_time_unix_nanos: None,
lifecycle_object_info: None,
is_delete_marker: dm,
}
}
@@ -491,6 +532,7 @@ mod tests {
bucket: "bucket".to_string(),
batch_objects: 1000,
version_budget: 10_000,
include_lifecycle_object_info: false,
objects: Mutex::new(Vec::new()),
version_total: AtomicUsize::new(0),
decode_error: Mutex::new(None),
@@ -567,7 +609,7 @@ mod tests {
.expect("corrupt test metadata should be written");
let error = set_disks
.heal_walk_versions_page(bucket, "", None, 2, 2)
.heal_walk_versions_page(bucket, "", None, 2, 2, false)
.await
.expect_err("semantic metadata corruption must fail the heal disk walk");
+35 -7
View File
@@ -18,6 +18,7 @@ use tracing::trace;
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
const LOG_SUBSYSTEM_HEAL: &str = "heal";
const EVENT_HEAL_ABANDONED_PARTS: &str = "heal_abandoned_parts";
const EVENT_HEAL_FORMAT_COMPLETED: &str = "heal_format_completed";
const EVENT_HEAL_OBJECT_STARTED: &str = "heal_object_started";
@@ -256,13 +257,40 @@ impl ECStore {
#[instrument(skip(self))]
pub(super) async fn handle_check_abandoned_parts(&self, bucket: &str, object: &str, opts: &HealOpts) -> Result<()> {
let _ = (bucket, object, opts);
// Stale multipart reconciliation is already owned by the lifecycle-driven
// background cleanup path in `bucket_lifecycle_ops.rs`. There is currently
// no stable object-heal contract that should fan this request out through
// pool/set storage layers, so keep the placeholder explicit at the ECStore
// boundary instead of dispatching into lower layers.
Err(StorageError::NotImplemented)
let object = encode_dir_object(object);
let pools = self.get_pools_for_heal_object(opts)?;
let mut futures = Vec::with_capacity(pools.len());
for pool in pools.iter() {
futures.push(pool.check_abandoned_parts(bucket, &object, opts));
}
let mut first_error = None;
for result in join_all(futures).await {
if let Err(err) = result
&& first_error.is_none()
{
first_error = Some(err);
}
}
if let Some(err) = first_error {
return Err(err);
}
trace!(
event = EVENT_HEAL_ABANDONED_PARTS,
component = LOG_COMPONENT_ECSTORE,
subsystem = LOG_SUBSYSTEM_HEAL,
state = "completed",
result = "ok",
bucket,
object,
dry_run = opts.dry_run,
"Heal abandoned parts completed"
);
Ok(())
}
}
+2 -1
View File
@@ -34,6 +34,7 @@ impl ECStore {
forward_to: Option<&str>,
batch_objects: usize,
version_budget: usize,
include_lifecycle_object_info: bool,
) -> Result<(Vec<HealWalkVersion>, Option<String>, bool)> {
if pool_idx >= self.pools.len() || set_idx >= self.pools[pool_idx].disk_set.len() {
return Err(Error::other(format!(
@@ -43,7 +44,7 @@ impl ECStore {
}
self.pools[pool_idx].disk_set[set_idx]
.heal_walk_versions_page(bucket, prefix, forward_to, batch_objects, version_budget)
.heal_walk_versions_page(bucket, prefix, forward_to, batch_objects, version_budget, include_lifecycle_object_info)
.await
.map_err(Error::from)
}
+1
View File
@@ -767,6 +767,7 @@ mod tests {
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> crate::Result<(Vec<crate::heal::storage::HealListItem>, Option<String>, bool)> {
Ok((vec![], None, false))
}
+317 -27
View File
@@ -23,13 +23,14 @@ use crate::heal::{
};
use crate::{Error, Result};
use futures::{StreamExt, stream::FuturesUnordered};
use metrics::gauge;
use metrics::{counter, gauge};
use rustfs_common::heal_channel::{HealOpts, HealRequestSource, HealScanMode};
use rustfs_madmin::heal_commands::HealResultItem;
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use std::time::{Duration, UNIX_EPOCH};
use tokio::sync::{RwLock, Semaphore};
use tracing::{debug, error, warn};
@@ -47,6 +48,21 @@ enum HealObjectOutcome {
Failed,
}
fn result_object_size_u64(result: &HealResultItem) -> u64 {
u64::try_from(result.object_size).unwrap_or(u64::MAX)
}
const NEW_VERSION_SKIP_GRACE_SECS: u64 = 60;
const NANOS_PER_SECOND: i128 = 1_000_000_000;
fn should_skip_new_version(mod_time_unix_nanos: Option<i128>, started_at_secs: u64) -> bool {
let Some(mod_time_unix_nanos) = mod_time_unix_nanos else {
return false;
};
let cutoff_secs = started_at_secs.saturating_add(NEW_VERSION_SKIP_GRACE_SECS);
mod_time_unix_nanos > i128::from(cutoff_secs).saturating_mul(NANOS_PER_SECOND)
}
struct PageConcurrencyGuard {
in_flight: Arc<AtomicUsize>,
set_label: String,
@@ -492,6 +508,7 @@ impl ErasureSetHealer {
&mut skipped_objects,
resume_manager,
checkpoint_manager,
state.start_time,
)
.await;
@@ -658,6 +675,7 @@ impl ErasureSetHealer {
skipped_objects: &mut u64,
resume_manager: &ResumeManager,
checkpoint_manager: &CheckpointManager,
started_at_secs: u64,
) -> Result<()> {
debug!(
target: "rustfs::heal::erasure_healer",
@@ -710,6 +728,7 @@ impl ErasureSetHealer {
// The end-of-pass summary reports the full failed/skipped counts.
let mut transient_skip_samples_logged = 0_u64;
let mut failure_samples_logged = 0_u64;
let mut bytes_processed = self.progress.read().await.bytes_processed;
// backlog#920: select the per-erasure-set DISK-WALK union enumerator when
// the scan is Deep OR the request came from AutoHeal — these are the paths
@@ -718,17 +737,25 @@ impl ErasureSetHealer {
// which stays the default.
let use_disk_walk =
matches!(self.heal_opts.scan_mode, HealScanMode::Deep) || matches!(self.source, HealRequestSource::AutoHeal);
let lifecycle_expiry_context = self.storage.load_heal_lifecycle_expiry_context(bucket).await?;
let include_lifecycle_object_info = lifecycle_expiry_context.is_some();
loop {
self.verify_replacement_identity_fence("page scan").await?;
// Get one page of object versions
let (objects, next_token, is_truncated) = if use_disk_walk {
self.storage
.list_versions_for_heal_page_disk_walk(set_disk_id, bucket, "", continuation_token.as_deref())
.list_versions_for_heal_page_disk_walk(
set_disk_id,
bucket,
"",
continuation_token.as_deref(),
include_lifecycle_object_info,
)
.await?
} else {
self.storage
.list_objects_for_heal_page(bucket, "", continuation_token.as_deref())
.list_objects_for_heal_page(bucket, "", continuation_token.as_deref(), include_lifecycle_object_info)
.await?
};
let page_is_empty = objects.is_empty();
@@ -736,6 +763,7 @@ impl ErasureSetHealer {
let page_resume_index = *current_object_index;
let semaphore = Arc::new(Semaphore::new(page_concurrency_limit));
let mut page_tasks = FuturesUnordered::new();
let mut completed_in_page = 0usize;
// Capture the last version identity of this page for the anti-loop guard.
let page_last = objects.last().map(|item| (item.name.clone(), item.version_id.clone()));
@@ -751,6 +779,75 @@ impl ErasureSetHealer {
continue;
}
if should_skip_new_version(item.mod_time_unix_nanos, started_at_secs) {
checkpoint_manager.add_processed_object(key).await?;
*processed_objects = processed_objects.saturating_add(1);
completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_new_versions_total").increment(1);
{
let mut progress = self.progress.write().await;
progress.record_skipped_new_version();
progress.set_current_object(Some(format!("skipped_new: {bucket}/{}", item.name)));
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
}
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_ERASURE_HEALER,
set_disk_id,
bucket,
object = %item.name,
version_id = ?item.version_id,
state = "skipped_new_version",
"Erasure set object version skipped because it was written after heal started"
);
if completed_in_page.is_multiple_of(100) {
checkpoint_manager.update_position(bucket_index, page_resume_index).await?;
}
continue;
}
if let Some(context) = lifecycle_expiry_context.as_ref()
&& self
.storage
.enqueue_heal_lifecycle_expiry(
context,
bucket,
&item.name,
item.version_id.as_deref(),
item.lifecycle_object_info.as_ref(),
)
.await?
{
checkpoint_manager.add_processed_object(key).await?;
*processed_objects = processed_objects.saturating_add(1);
completed_in_page = completed_in_page.saturating_add(1);
counter!("rustfs_heal_skipped_ilm_expired_total").increment(1);
{
let mut progress = self.progress.write().await;
progress.record_skipped_ilm_expired();
progress.set_current_object(Some(format!("skipped_ilm: {bucket}/{}", item.name)));
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
}
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_ERASURE_HEALER,
set_disk_id,
bucket,
object = %item.name,
version_id = ?item.version_id,
state = "skipped_ilm_expired",
"Erasure set object version skipped because lifecycle expiry was queued"
);
if completed_in_page.is_multiple_of(100) {
checkpoint_manager.update_position(bucket_index, page_resume_index).await?;
}
continue;
}
resume_manager
.set_current_item(Some(bucket.to_string()), Some(item.name.clone()))
.await?;
@@ -777,7 +874,7 @@ impl ErasureSetHealer {
let _permit = match permit {
Ok(permit) => permit,
Err(err) => return (dedup_key, object_name, version_id, Err(err)),
Err(err) => return (dedup_key, object_name, version_id, (0, Err(err))),
};
let _in_flight_guard = PageConcurrencyGuard::new(in_flight, set_label);
@@ -788,7 +885,7 @@ impl ErasureSetHealer {
// recorded as skipped-ok rather than failed. The delete-marker
// vs data path is chosen internally in ops/heal.rs.
let result = if cancel_token.is_cancelled() {
Err(Error::TaskCancelled)
(0, Err(Error::TaskCancelled))
} else {
match storage
.heal_object(&bucket_name, &object_name, version_id.as_deref(), &heal_opts)
@@ -797,8 +894,9 @@ impl ErasureSetHealer {
Ok((result, None))
if target_outcomes_complete(&result, &target_endpoints) =>
{
let object_size = result_object_size_u64(&result);
if !replacement_commit_evidence_required {
Ok(true)
(object_size, Ok(true))
} else {
match storage
.replacement_targets_have_version(
@@ -810,27 +908,42 @@ impl ErasureSetHealer {
)
.await
{
Ok(true) => Ok(true),
Ok(false) => Err(Error::transient_skip(format!(
Ok(true) => (object_size, Ok(true)),
Ok(false) => (object_size, Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} because replacement target readback did not confirm the committed version"
))),
Err(err) => Err(Error::transient_skip(format!(
)))),
Err(err) => (object_size, Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} because replacement target readback failed: {err}"
))),
)))),
}
}
}
Ok((_result, None)) if !target_endpoints.is_empty() => Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} because a replacement target was not committed"
))),
Ok((_result, None)) => Ok(true),
Ok((_, Some(err))) if is_missing_object_dir_heal_result(&object_name, &err) => Ok(false),
Ok((_, Some(err))) | Err(err) => match Self::classify_heal_object_error(&err) {
HealObjectOutcome::Absent => Ok(false),
HealObjectOutcome::Transient => Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} due to transient error: {err}"
},
Ok((result, None)) if !target_endpoints.is_empty() => (
result_object_size_u64(&result),
Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} because a replacement target was not committed"
))),
HealObjectOutcome::Failed => Err(err),
),
Ok((result, None)) => (result_object_size_u64(&result), Ok(true)),
Ok((result, Some(err))) if is_missing_object_dir_heal_result(&object_name, &err) => {
(result_object_size_u64(&result), Ok(false))
}
Ok((result, Some(err))) => {
let object_size = result_object_size_u64(&result);
match Self::classify_heal_object_error(&err) {
HealObjectOutcome::Absent => (object_size, Ok(false)),
HealObjectOutcome::Transient => (object_size, Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} due to transient error: {err}"
)))),
HealObjectOutcome::Failed => (object_size, Err(err)),
}
}
Err(err) => match Self::classify_heal_object_error(&err) {
HealObjectOutcome::Absent => (0, Ok(false)),
HealObjectOutcome::Transient => (0, Err(Error::transient_skip(format!(
"Skipped heal for {bucket_name}/{object_name} due to transient error: {err}"
)))),
HealObjectOutcome::Failed => (0, Err(err)),
},
}
};
@@ -839,11 +952,12 @@ impl ErasureSetHealer {
});
}
let mut completed_in_page = 0usize;
while let Some((key, object, version_id, result)) = page_tasks.next().await {
let (object_size, result) = result;
match result {
Ok(true) => {
*successful_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
checkpoint_manager.add_processed_object(key).await?;
debug!(
target: "rustfs::heal::erasure_healer",
@@ -861,6 +975,7 @@ impl ErasureSetHealer {
Ok(false) => {
checkpoint_manager.add_processed_object(key).await?;
*successful_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
debug!(
target: "rustfs::heal::erasure_healer",
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -877,6 +992,7 @@ impl ErasureSetHealer {
Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err),
Err(Error::TransientSkip { message }) => {
*skipped_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
checkpoint_manager.add_skipped_object(key).await?;
demote_to_debug_when!(!take_failure_log_sample(&mut transient_skip_samples_logged), warn, target: "rustfs::heal::erasure_healer", {
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -893,6 +1009,7 @@ impl ErasureSetHealer {
}
Err(err) => {
*failed_objects += 1;
bytes_processed = bytes_processed.saturating_add(object_size);
checkpoint_manager.add_failed_object(key).await?;
demote_to_debug_when!(!take_failure_log_sample(&mut failure_samples_logged), warn, target: "rustfs::heal::erasure_healer", {
event = EVENT_HEAL_ERASURE_OBJECT_STATE,
@@ -911,6 +1028,11 @@ impl ErasureSetHealer {
*processed_objects += 1;
completed_in_page += 1;
{
let mut progress = self.progress.write().await;
progress.set_current_object(Some(format!("{bucket}/{object}")));
progress.update_progress(*processed_objects, *successful_objects, *failed_objects, bytes_processed);
}
if completed_in_page.is_multiple_of(100) {
checkpoint_manager.update_position(bucket_index, page_resume_index).await?;
@@ -964,7 +1086,9 @@ impl ErasureSetHealer {
progress.objects_scanned = state.total_objects;
progress.objects_healed = state.successful_objects;
progress.objects_failed = state.failed_objects;
progress.bytes_processed = 0; // set to 0 for now, can be extended later
progress.bytes_processed = 0; // Resume state tracks object counts, not byte counters.
progress.start_time = UNIX_EPOCH.checked_add(Duration::from_secs(state.start_time));
progress.last_update_time = UNIX_EPOCH.checked_add(Duration::from_secs(state.last_update));
progress.set_current_object(state.current_object.clone());
}
}
@@ -1135,13 +1259,15 @@ mod resume_loop_tests {
//! that emits programmable multi-version pages. These exercise the real loop
//! logic (cursor seeding, per-version dedup, anti-loop guard, absence
//! handling) — not merely a mock's own output.
use super::{ErasureSetHealer, target_outcomes_complete};
use super::{
ErasureSetHealer, NANOS_PER_SECOND, NEW_VERSION_SKIP_GRACE_SECS, should_skip_new_version, target_outcomes_complete,
};
use crate::heal::progress::HealProgress;
use crate::heal::resume::{
CheckpointManager, RESUME_CHECKPOINT_FILE, ReplacementTargetIdentity, ResumeDeleteFailure, ResumeManager, ResumeUtils,
compose_key,
};
use crate::heal::storage::{DiskStatus, HealListItem, HealObjectInfo, HealStorageAPI};
use crate::heal::storage::{DiskStatus, HealLifecycleExpiryContext, HealListItem, HealObjectInfo, HealStorageAPI};
use crate::heal::storage_api::status::BucketInfo;
use crate::heal::{
BUCKET_META_PREFIX, DiskOption, DiskStore, EcstoreError, Endpoint, HealDiskExt as _, RUSTFS_META_BUCKET, new_disk,
@@ -1149,7 +1275,7 @@ mod resume_loop_tests {
use crate::{Error, Result};
use rustfs_common::heal_channel::{HealOpts, HealRequestSource};
use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos};
use std::collections::{HashMap, VecDeque};
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use tempfile::TempDir;
@@ -1160,10 +1286,37 @@ mod resume_loop_tests {
HealListItem {
name: name.to_string(),
version_id: version.map(str::to_string),
mod_time_unix_nanos: None,
lifecycle_object_info: None,
is_delete_marker: delete_marker,
}
}
fn item_with_mod_time(name: &str, version: Option<&str>, mod_time_secs: u64) -> HealListItem {
HealListItem {
name: name.to_string(),
version_id: version.map(str::to_string),
mod_time_unix_nanos: Some(i128::from(mod_time_secs).saturating_mul(NANOS_PER_SECOND)),
lifecycle_object_info: None,
is_delete_marker: false,
}
}
#[test]
fn new_version_filter_respects_grace_boundary() {
let started_at = 1_700_000_000;
assert!(!should_skip_new_version(None, started_at));
assert!(!should_skip_new_version(
Some(i128::from(started_at + NEW_VERSION_SKIP_GRACE_SECS).saturating_mul(NANOS_PER_SECOND)),
started_at,
));
assert!(should_skip_new_version(
Some(i128::from(started_at + NEW_VERSION_SKIP_GRACE_SECS + 1).saturating_mul(NANOS_PER_SECOND)),
started_at,
));
}
#[test]
fn target_outcomes_require_each_requested_endpoint_once_and_ok() {
let result = HealResultItem {
@@ -1246,8 +1399,10 @@ mod resume_loop_tests {
/// Target-specific physical readback evidence per `compose_key`; the
/// fake models a healthy backend unless a test explicitly revokes it.
replacement_commit_evidence: Mutex<HashMap<String, ReplacementCommitEvidence>>,
lifecycle_expired: Mutex<HashSet<String>>,
/// every heal_object call recorded as (name, version_id)
heal_calls: Mutex<Vec<(String, Option<String>)>>,
list_include_lifecycle_object_info: Mutex<Vec<bool>>,
replacement_target_identity_sequences: Mutex<VecDeque<Vec<ReplacementTargetIdentity>>>,
fail_listing: AtomicBool,
}
@@ -1274,9 +1429,15 @@ mod resume_loop_tests {
.unwrap()
.insert(compose_key(name, version), ReplacementCommitEvidence::Error(message.to_string()));
}
fn set_lifecycle_expired(&self, name: &str, version: Option<&str>) {
self.lifecycle_expired.lock().unwrap().insert(compose_key(name, version));
}
fn calls(&self) -> Vec<(String, Option<String>)> {
self.heal_calls.lock().unwrap().clone()
}
fn list_include_lifecycle_object_info_calls(&self) -> Vec<bool> {
self.list_include_lifecycle_object_info.lock().unwrap().clone()
}
fn fail_listing(&self) {
self.fail_listing.store(true, Ordering::SeqCst);
}
@@ -1330,6 +1491,23 @@ mod resume_loop_tests {
async fn get_object_checksum(&self, _b: &str, _o: &str) -> Result<Option<String>> {
Ok(None)
}
async fn load_heal_lifecycle_expiry_context(&self, _bucket: &str) -> Result<Option<HealLifecycleExpiryContext>> {
Ok((!self.lifecycle_expired.lock().unwrap().is_empty()).then(HealLifecycleExpiryContext::test))
}
async fn enqueue_heal_lifecycle_expiry(
&self,
_context: &HealLifecycleExpiryContext,
_bucket: &str,
object: &str,
version_id: Option<&str>,
_object_info: Option<&HealObjectInfo>,
) -> Result<bool> {
Ok(self
.lifecycle_expired
.lock()
.unwrap()
.contains(&compose_key(object, version_id)))
}
async fn heal_object(
&self,
_bucket: &str,
@@ -1386,7 +1564,12 @@ mod resume_loop_tests {
_bucket: &str,
_prefix: &str,
continuation_token: Option<&str>,
include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
self.list_include_lifecycle_object_info
.lock()
.unwrap()
.push(include_lifecycle_object_info);
if self.fail_listing.load(Ordering::SeqCst) {
return Err(Error::other("injected listing failure"));
}
@@ -1476,6 +1659,7 @@ mod resume_loop_tests {
/// Drive one bucket heal pass; returns (processed, successful, failed, skipped, result).
async fn run(env: &Env) -> (u64, u64, u64, u64, Result<()>) {
let state = env.resume.get_state().await;
let mut current_object_index = 0usize;
let mut processed = 0u64;
let mut successful = 0u64;
@@ -1494,6 +1678,7 @@ mod resume_loop_tests {
&mut skipped,
&env.resume,
&env.checkpoint,
state.start_time,
)
.await;
(processed, successful, failed, skipped, result)
@@ -1559,6 +1744,7 @@ mod resume_loop_tests {
let mut successful = 0;
let mut failed = 0;
let mut skipped = 0;
let started_at = env.resume.get_state().await.start_time;
let error = healer
.heal_bucket_with_resume(
@@ -1572,6 +1758,7 @@ mod resume_loop_tests {
&mut skipped,
&env.resume,
&env.checkpoint,
started_at,
)
.await
.expect_err("a remounted target must not begin a new page scan");
@@ -1641,6 +1828,109 @@ mod resume_loop_tests {
assert_eq!(skipped, 0);
}
#[tokio::test]
async fn erasure_set_progress_accumulates_healed_object_bytes() {
let env = make_env().await;
env.storage.set_page(
None,
Page {
items: vec![item("first", Some("v1"), false), item("second", Some("v2"), false)],
next: None,
truncated: false,
},
);
env.storage.set_result(
"first",
Some("v1"),
HealResultItem {
object_size: 1024,
..Default::default()
},
);
env.storage.set_result(
"second",
Some("v2"),
HealResultItem {
object_size: 2048,
..Default::default()
},
);
let (processed, successful, failed, skipped, result) = run(&env).await;
result.expect("page heal should succeed");
assert_eq!(processed, 2);
assert_eq!(successful, 2);
assert_eq!(failed, 0);
assert_eq!(skipped, 0);
let progress = env.healer.progress.read().await;
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_healed, 2);
assert_eq!(progress.objects_failed, 0);
assert_eq!(progress.bytes_processed, 3072);
assert!(matches!(progress.current_object.as_deref(), Some("b/first" | "b/second")));
}
#[tokio::test]
async fn erasure_set_skips_versions_written_after_heal_started() {
let env = make_env().await;
let started_at = env.resume.get_state().await.start_time;
env.storage.set_page(
None,
Page {
items: vec![
item_with_mod_time("old", Some("v1"), started_at + NEW_VERSION_SKIP_GRACE_SECS),
item_with_mod_time("new", Some("v2"), started_at + NEW_VERSION_SKIP_GRACE_SECS + 1),
],
next: None,
truncated: false,
},
);
let (processed, successful, failed, skipped, result) = run(&env).await;
result.expect("page heal should succeed");
assert_eq!(processed, 2);
assert_eq!(successful, 1);
assert_eq!(failed, 0);
assert_eq!(skipped, 0);
assert_eq!(env.storage.calls(), vec![("old".to_string(), Some("v1".to_string()))]);
let progress = env.healer.progress.read().await;
assert_eq!(progress.skipped_new_versions, 1);
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_healed, 1);
assert_eq!(progress.objects_failed, 0);
}
#[tokio::test]
async fn erasure_set_skips_versions_queued_for_lifecycle_expiry() {
let env = make_env().await;
env.storage.set_page(
None,
Page {
items: vec![item("expired", Some("v1"), false), item("kept", Some("v2"), false)],
next: None,
truncated: false,
},
);
env.storage.set_lifecycle_expired("expired", Some("v1"));
let (processed, successful, failed, skipped, result) = run(&env).await;
result.expect("page heal should succeed");
assert_eq!(processed, 2);
assert_eq!(successful, 1);
assert_eq!(failed, 0);
assert_eq!(skipped, 0);
assert_eq!(env.storage.calls(), vec![("kept".to_string(), Some("v2".to_string()))]);
assert_eq!(env.storage.list_include_lifecycle_object_info_calls(), vec![true]);
let progress = env.healer.progress.read().await;
assert_eq!(progress.skipped_ilm_expired, 1);
assert_eq!(progress.objects_scanned, 2);
assert_eq!(progress.objects_healed, 1);
assert_eq!(progress.objects_failed, 0);
}
#[tokio::test]
async fn bucket_listing_failure_does_not_mark_set_completed() {
let env = make_env().await;
+30
View File
@@ -2385,8 +2385,27 @@ impl HealManager {
snapshot.objects_scanned = snapshot.objects_scanned.saturating_add(progress.objects_scanned);
snapshot.objects_healed = snapshot.objects_healed.saturating_add(progress.objects_healed);
snapshot.objects_failed = snapshot.objects_failed.saturating_add(progress.objects_failed);
snapshot.skipped_new_versions = snapshot.skipped_new_versions.saturating_add(progress.skipped_new_versions);
snapshot.skipped_ilm_expired = snapshot.skipped_ilm_expired.saturating_add(progress.skipped_ilm_expired);
snapshot.objects_total_count = snapshot.objects_total_count.saturating_add(progress.objects_total_count);
snapshot.objects_total_size = snapshot.objects_total_size.saturating_add(progress.objects_total_size);
snapshot.bytes_processed = snapshot.bytes_processed.saturating_add(progress.bytes_processed);
snapshot.start_time = match (snapshot.start_time, progress.start_time) {
(Some(current), Some(next)) => Some(current.min(next)),
(None, next) => next,
(current, None) => current,
};
snapshot.last_update_time = match (snapshot.last_update_time, progress.last_update_time) {
(Some(current), Some(next)) => Some(current.max(next)),
(None, next) => next,
(current, None) => current,
};
if progress.current_object.is_some() {
snapshot.current_object = progress.current_object;
}
}
snapshot.refresh_progress_percentage();
snapshot.refresh_estimated_completion_time();
Some(snapshot)
}
@@ -3208,6 +3227,7 @@ impl HealManager {
} else {
completed_task.get_status().await
};
let completed_progress = completed_task.get_progress().await;
let completed_status_entry = CompletedHealStatus {
heal_type: completed_task.heal_type.clone(),
status: completed_status.clone(),
@@ -3223,6 +3243,7 @@ impl HealManager {
match completed_status {
HealTaskStatus::Completed => {
stats.update_task_completion(true);
stats.add_healed_objects(completed_progress.objects_healed, completed_progress.bytes_processed);
}
HealTaskStatus::Retrying { .. } => {}
_ => {
@@ -3749,6 +3770,7 @@ mod tests {
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> Result<(Vec<crate::heal::storage::HealListItem>, Option<String>, bool)> {
Ok((Vec::new(), None, false))
}
@@ -5396,6 +5418,8 @@ mod tests {
));
{
let mut progress = first.progress.write().await;
progress.start_time = Some(SystemTime::now() - Duration::from_secs(20));
progress.set_total_baseline(12, 8192);
progress.update_progress(7, 3, 1, 4096);
}
@@ -5405,6 +5429,8 @@ mod tests {
));
{
let mut progress = second.progress.write().await;
progress.start_time = Some(SystemTime::now() - Duration::from_secs(10));
progress.set_total_baseline(8, 4096);
progress.update_progress(11, 5, 2, 2048);
}
@@ -5419,7 +5445,11 @@ mod tests {
assert_eq!(progress.objects_scanned, 18);
assert_eq!(progress.objects_healed, 8);
assert_eq!(progress.objects_failed, 3);
assert_eq!(progress.objects_total_count, 20);
assert_eq!(progress.objects_total_size, 12288);
assert_eq!(progress.bytes_processed, 6144);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
assert!(progress.estimated_completion_time.is_some());
}
#[tokio::test]
+160 -6
View File
@@ -13,7 +13,7 @@
// limitations under the License.
use serde::{Deserialize, Serialize};
use std::time::SystemTime;
use std::time::{Duration, SystemTime};
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
@@ -24,6 +24,14 @@ pub struct HealProgress {
pub objects_healed: u64,
/// Objects failed
pub objects_failed: u64,
/// Versions skipped because they were written after this heal started
pub skipped_new_versions: u64,
/// Versions skipped because lifecycle already selected them for expiry
pub skipped_ilm_expired: u64,
/// Baseline object count from the latest complete usage snapshot
pub objects_total_count: u64,
/// Baseline object bytes from the latest complete usage snapshot
pub objects_total_size: u64,
/// Bytes processed
pub bytes_processed: u64,
/// Current object
@@ -54,10 +62,56 @@ impl HealProgress {
self.bytes_processed = bytes;
self.last_update_time = Some(SystemTime::now());
// calculate progress percentage
let total = scanned + healed + failed;
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
pub fn set_total_baseline(&mut self, objects_total_count: u64, objects_total_size: u64) {
self.objects_total_count = objects_total_count;
self.objects_total_size = objects_total_size;
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
pub fn record_skipped_new_version(&mut self) {
self.skipped_new_versions = self.skipped_new_versions.saturating_add(1);
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
pub fn record_skipped_ilm_expired(&mut self) {
self.skipped_ilm_expired = self.skipped_ilm_expired.saturating_add(1);
self.last_update_time = Some(SystemTime::now());
self.refresh_progress_percentage();
self.refresh_estimated_completion_time();
}
fn completed_for_baseline(&self) -> u64 {
self.objects_healed
.saturating_add(self.objects_failed)
.saturating_add(self.skipped_new_versions)
.saturating_add(self.skipped_ilm_expired)
}
pub(crate) fn refresh_progress_percentage(&mut self) {
if self.objects_total_size > 0 {
self.progress_percentage = ((self.bytes_processed as f64 / self.objects_total_size as f64) * 100.0).min(100.0);
return;
}
if self.objects_total_count > 0 {
let completed = self.completed_for_baseline();
self.progress_percentage = ((completed as f64 / self.objects_total_count as f64) * 100.0).min(100.0);
return;
}
let total = self
.objects_scanned
.saturating_add(self.objects_healed)
.saturating_add(self.objects_failed);
if total > 0 {
self.progress_percentage = (healed as f64 / total as f64) * 100.0;
self.progress_percentage = (self.objects_healed as f64 / total as f64) * 100.0;
}
}
@@ -66,9 +120,36 @@ impl HealProgress {
self.last_update_time = Some(SystemTime::now());
}
pub fn refresh_estimated_completion_time(&mut self) {
let Some(start_time) = self.start_time else {
self.estimated_completion_time = None;
return;
};
if self.is_completed() || !(0.0..100.0).contains(&self.progress_percentage) || self.bytes_processed == 0 {
self.estimated_completion_time = None;
return;
}
let elapsed = match SystemTime::now().duration_since(start_time) {
Ok(elapsed) if !elapsed.is_zero() => elapsed,
_ => {
self.estimated_completion_time = None;
return;
}
};
let estimated_total_secs = elapsed.as_secs_f64() * 100.0 / self.progress_percentage;
self.estimated_completion_time = start_time.checked_add(Duration::from_secs_f64(estimated_total_secs));
}
pub fn is_completed(&self) -> bool {
self.progress_percentage >= 100.0
|| self.objects_scanned > 0 && self.objects_healed + self.objects_failed >= self.objects_scanned
if self.progress_percentage >= 100.0 {
return true;
}
if self.objects_total_count > 0 || self.objects_total_size > 0 {
return false;
}
self.objects_scanned > 0 && self.objects_healed.saturating_add(self.objects_failed) >= self.objects_scanned
}
pub fn get_success_rate(&self) -> f64 {
@@ -158,6 +239,10 @@ mod tests {
assert_eq!(progress.objects_scanned, 0);
assert_eq!(progress.objects_healed, 0);
assert_eq!(progress.objects_failed, 0);
assert_eq!(progress.skipped_new_versions, 0);
assert_eq!(progress.skipped_ilm_expired, 0);
assert_eq!(progress.objects_total_count, 0);
assert_eq!(progress.objects_total_size, 0);
assert_eq!(progress.bytes_processed, 0);
assert_eq!(progress.progress_percentage, 0.0);
assert!(progress.start_time.is_some());
@@ -181,6 +266,73 @@ mod tests {
assert!(progress.last_update_time.is_some());
}
#[test]
fn test_heal_progress_estimates_completion_time_from_progress() {
let mut progress = HealProgress::new();
progress.start_time = Some(SystemTime::now() - Duration::from_secs(10));
progress.update_progress(100, 25, 0, 4096);
let eta = progress
.estimated_completion_time
.expect("partial byte progress should estimate completion");
assert!(eta > SystemTime::now());
}
#[test]
fn test_heal_progress_uses_byte_baseline_for_percentage() {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 8192);
progress.update_progress(100, 25, 0, 4096);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
}
#[test]
fn test_heal_progress_uses_object_baseline_when_bytes_unknown() {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 0);
progress.update_progress(100, 3, 2, 0);
assert!((progress.progress_percentage - 50.0).abs() < 0.001);
}
#[test]
fn test_heal_progress_counts_skipped_versions_for_object_baseline() {
let mut progress = HealProgress::new();
progress.set_total_baseline(10, 0);
progress.update_progress(100, 3, 2, 0);
progress.record_skipped_new_version();
assert_eq!(progress.skipped_new_versions, 1);
assert!((progress.progress_percentage - 60.0).abs() < 0.001);
}
#[test]
fn test_heal_progress_does_not_estimate_completion_without_bytes() {
let mut progress = HealProgress::new();
progress.start_time = Some(SystemTime::now() - Duration::from_secs(10));
progress.update_progress(100, 25, 0, 0);
assert!(progress.estimated_completion_time.is_none());
}
#[test]
fn test_heal_progress_with_baseline_is_not_completed_by_processed_count() {
let mut progress = HealProgress::new();
progress.start_time = Some(SystemTime::now() - Duration::from_secs(10));
progress.set_total_baseline(10, 8192);
progress.update_progress(1, 1, 0, 1024);
assert!(!progress.is_completed());
assert!(progress.estimated_completion_time.is_some());
}
#[test]
fn test_heal_progress_update_progress_zero_total() {
let mut progress = HealProgress::new();
@@ -251,6 +403,8 @@ mod tests {
assert_eq!(json["objectsScanned"], 10);
assert_eq!(json["objectsHealed"], 8);
assert_eq!(json["objectsFailed"], 2);
assert_eq!(json["skippedNewVersions"], 0);
assert_eq!(json["skippedIlmExpired"], 0);
assert_eq!(json["bytesProcessed"], 1024);
assert_eq!(json["currentObject"], "test-bucket/test-object");
assert!(json["progressPercentage"].is_number());
+169 -7
View File
@@ -22,6 +22,7 @@ use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tracing::{debug, error, warn};
use super::storage_api::owner::{EcstoreHealLifecycleExpiryContext, ecstore_load_admin_data_usage_from_backend_cached};
use super::storage_api::storage::{
BucketInfo, BucketOperations, DiskSetSelector, HealOperations as _, ListOperations as _, ObjectIO as _,
ObjectOperations as _, StorageAdminApi,
@@ -29,6 +30,37 @@ use super::storage_api::storage::{
use super::{DiskStore, ECStore, Endpoint, HealDiskExt as _, StorageError, resume::ReplacementTargetIdentity};
pub use super::{HealObjectInfo, HealObjectOptions, HealPutObjReader};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct HealBucketUsageBaseline {
pub objects_count: u64,
pub bytes: u64,
}
pub struct HealLifecycleExpiryContext {
inner: HealLifecycleExpiryContextInner,
}
enum HealLifecycleExpiryContextInner {
Ecstore(EcstoreHealLifecycleExpiryContext),
#[allow(dead_code)]
Test,
}
impl HealLifecycleExpiryContext {
fn ecstore(inner: EcstoreHealLifecycleExpiryContext) -> Self {
Self {
inner: HealLifecycleExpiryContextInner::Ecstore(inner),
}
}
#[cfg(test)]
pub(crate) fn test() -> Self {
Self {
inner: HealLifecycleExpiryContextInner::Test,
}
}
}
const LOG_COMPONENT_HEAL: &str = "heal";
const LOG_SUBSYSTEM_STORAGE: &str = "storage";
const EVENT_HEAL_STORAGE_OBJECT_IO: &str = "heal_storage_object_io";
@@ -272,6 +304,10 @@ pub struct HealListItem {
pub name: String,
/// normalized version id (`None` when the version is nil/absent)
pub version_id: Option<String>,
/// version modification time as Unix nanoseconds
pub mod_time_unix_nanos: Option<i128>,
/// object snapshot for lifecycle evaluation
pub lifecycle_object_info: Option<HealObjectInfo>,
/// whether this version is a delete marker (observability only)
pub is_delete_marker: bool,
}
@@ -329,6 +365,28 @@ pub trait HealStorageAPI: Send + Sync {
/// Get bucket info
async fn get_bucket_info(&self, bucket: &str) -> Result<Option<BucketInfo>>;
/// Aggregate usage-cache baselines for the requested buckets.
async fn erasure_set_usage_baseline(&self, _buckets: &[String]) -> Result<Option<HealBucketUsageBaseline>> {
Ok(None)
}
/// Load per-bucket lifecycle expiry context for heal skips.
async fn load_heal_lifecycle_expiry_context(&self, _bucket: &str) -> Result<Option<HealLifecycleExpiryContext>> {
Ok(None)
}
/// Queue lifecycle expiry for a version that heal can skip.
async fn enqueue_heal_lifecycle_expiry(
&self,
_context: &HealLifecycleExpiryContext,
_bucket: &str,
_object: &str,
_version_id: Option<&str>,
_object_info: Option<&HealObjectInfo>,
) -> Result<bool> {
Ok(false)
}
/// Fix bucket metadata
async fn heal_bucket_metadata(&self, bucket: &str) -> Result<()>;
@@ -409,6 +467,7 @@ pub trait HealStorageAPI: Send + Sync {
bucket: &str,
prefix: &str,
continuation_token: Option<&str>,
include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)>;
/// List versions for healing via a per-erasure-set DISK-WALK union enumerator
@@ -427,8 +486,10 @@ pub trait HealStorageAPI: Send + Sync {
bucket: &str,
prefix: &str,
continuation_token: Option<&str>,
include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
self.list_objects_for_heal_page(bucket, prefix, continuation_token).await
self.list_objects_for_heal_page(bucket, prefix, continuation_token, include_lifecycle_object_info)
.await
}
/// Get disk for resume functionality.
@@ -1021,6 +1082,85 @@ impl HealStorageAPI for ECStoreHealStorage {
}
}
async fn erasure_set_usage_baseline(&self, buckets: &[String]) -> Result<Option<HealBucketUsageBaseline>> {
if buckets.is_empty() {
return Ok(None);
}
let info = match ecstore_load_admin_data_usage_from_backend_cached(self.ecstore.clone()).await {
Ok(info) if info.is_complete_bucket_usage_snapshot() => info,
Ok(_) | Err(_) => return Ok(None),
};
let mut baseline = HealBucketUsageBaseline::default();
for bucket in buckets {
if let Some(usage) = info.buckets_usage.get(bucket) {
baseline.objects_count = baseline.objects_count.saturating_add(usage.objects_count);
baseline.bytes = baseline.bytes.saturating_add(usage.size);
}
}
Ok(Some(baseline))
}
async fn load_heal_lifecycle_expiry_context(&self, bucket: &str) -> Result<Option<HealLifecycleExpiryContext>> {
match self.ecstore.load_heal_lifecycle_expiry_context(bucket).await {
Ok(Some(context)) => Ok(Some(HealLifecycleExpiryContext::ecstore(context))),
Ok(None) => Ok(None),
Err(err) => {
debug!(
target: "rustfs::heal::storage",
event = EVENT_HEAL_STORAGE_ADMIN_OP,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_STORAGE,
operation = "load_heal_lifecycle_expiry_context",
bucket,
result = "failed",
error = %err,
"Heal storage lifecycle expiry context load failed"
);
Ok(None)
}
}
}
async fn enqueue_heal_lifecycle_expiry(
&self,
context: &HealLifecycleExpiryContext,
bucket: &str,
object: &str,
version_id: Option<&str>,
object_info: Option<&HealObjectInfo>,
) -> Result<bool> {
let context = match &context.inner {
HealLifecycleExpiryContextInner::Ecstore(context) => context,
HealLifecycleExpiryContextInner::Test => return Ok(false),
};
match self
.ecstore
.enqueue_heal_lifecycle_expiry(context, bucket, object, version_id, object_info)
.await
{
Ok(queued) => Ok(queued),
Err(err) => {
debug!(
target: "rustfs::heal::storage",
event = EVENT_HEAL_STORAGE_ADMIN_OP,
component = LOG_COMPONENT_HEAL,
subsystem = LOG_SUBSYSTEM_STORAGE,
operation = "enqueue_heal_lifecycle_expiry",
bucket,
object,
version_id = ?version_id,
result = "failed",
error = %err,
"Heal storage lifecycle expiry check failed"
);
Ok(false)
}
}
}
async fn heal_bucket_metadata(&self, bucket: &str) -> Result<()> {
debug!(
target: "rustfs::heal::storage",
@@ -1436,7 +1576,7 @@ impl HealStorageAPI for ECStoreHealStorage {
loop {
let (page_objects, next_token, is_truncated) = self
.list_objects_for_heal_page(bucket, prefix, continuation_token.as_deref())
.list_objects_for_heal_page(bucket, prefix, continuation_token.as_deref(), false)
.await?;
all_objects.extend(page_objects);
@@ -1471,6 +1611,7 @@ impl HealStorageAPI for ECStoreHealStorage {
bucket: &str,
prefix: &str,
continuation_token: Option<&str>,
include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
debug!(
target: "rustfs::heal::storage",
@@ -1522,10 +1663,19 @@ impl HealStorageAPI for ECStoreHealStorage {
let page_objects: Vec<HealListItem> = list_info
.objects
.into_iter()
.map(|obj| HealListItem {
name: obj.name,
version_id: obj.version_id.filter(|u| !u.is_nil()).map(|u| u.to_string()),
is_delete_marker: obj.delete_marker,
.map(|mut obj| {
obj.version_id = obj.version_id.filter(|u| !u.is_nil());
let version_id = obj.version_id.map(|u| u.to_string());
let mod_time_unix_nanos = obj.mod_time.map(|mod_time| mod_time.unix_timestamp_nanos());
let is_delete_marker = obj.delete_marker;
let lifecycle_object_info = include_lifecycle_object_info.then(|| obj.clone());
HealListItem {
name: obj.name,
version_id,
mod_time_unix_nanos,
lifecycle_object_info,
is_delete_marker,
}
})
.collect();
let page_count = page_objects.len();
@@ -1562,6 +1712,7 @@ impl HealStorageAPI for ECStoreHealStorage {
bucket: &str,
prefix: &str,
continuation_token: Option<&str>,
include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
// Per-page bounds for the disk-walk union enumerator. Objects are atomic
// (never split across pages), so version_budget only bounds how many
@@ -1590,7 +1741,16 @@ impl HealStorageAPI for ECStoreHealStorage {
let (versions, next_forward, is_truncated) = self
.ecstore
.heal_walk_versions_page(pool_idx, set_idx, bucket, prefix, forward_to.as_deref(), BATCH_OBJECTS, VERSION_BUDGET)
.heal_walk_versions_page(
pool_idx,
set_idx,
bucket,
prefix,
forward_to.as_deref(),
BATCH_OBJECTS,
VERSION_BUDGET,
include_lifecycle_object_info,
)
.await
.map_err(|e| {
error!(
@@ -1614,6 +1774,8 @@ impl HealStorageAPI for ECStoreHealStorage {
.map(|v| HealListItem {
name: v.name,
version_id: v.version_id,
mod_time_unix_nanos: v.mod_time_unix_nanos,
lifecycle_object_info: v.lifecycle_object_info,
is_delete_marker: v.is_delete_marker,
})
.collect();
+9 -4
View File
@@ -12,7 +12,10 @@
// See the License for the specific language governing permissions and
// limitations under the License.
pub(crate) use rustfs_ecstore::api::data_usage::DATA_USAGE_CACHE_NAME as ECSTORE_DATA_USAGE_CACHE_NAME;
pub(crate) use rustfs_ecstore::api::data_usage::{
DATA_USAGE_CACHE_NAME as ECSTORE_DATA_USAGE_CACHE_NAME,
load_admin_data_usage_from_backend_cached as ecstore_load_admin_data_usage_from_backend_cached,
};
pub(crate) use rustfs_ecstore::api::disk::endpoint::Endpoint as EcstoreEndpoint;
pub(crate) use rustfs_ecstore::api::disk::error::{DiskError as EcstoreDiskError, Result as EcstoreDiskResult};
pub(crate) use rustfs_ecstore::api::disk::{
@@ -25,7 +28,9 @@ pub(crate) use rustfs_ecstore::api::disk::{
pub(crate) use rustfs_ecstore::api::disk::{DiskOption as EcstoreDiskOption, new_disk as ecstore_new_disk};
pub(crate) use rustfs_ecstore::api::error::{Error as EcstoreErrorType, StorageError as EcstoreStorageError};
pub(crate) use rustfs_ecstore::api::runtime::local_disk_map_read as ecstore_local_disk_map_read;
pub(crate) use rustfs_ecstore::api::storage::ECStore as EcstoreStore;
pub(crate) use rustfs_ecstore::api::storage::{
ECStore as EcstoreStore, HealLifecycleExpiryContext as EcstoreHealLifecycleExpiryContext,
};
use rustfs_storage_api as storage_contracts;
pub(crate) mod owner {
@@ -34,8 +39,8 @@ pub(crate) mod owner {
pub(crate) use super::{
ECSTORE_BUCKET_META_PREFIX, ECSTORE_DATA_USAGE_CACHE_NAME, ECSTORE_HEALING_MARKER_PATH, ECSTORE_RUSTFS_META_BUCKET,
EcstoreConditionalFileUpdate, EcstoreDeleteOptions, EcstoreDiskAPI, EcstoreDiskBytes, EcstoreDiskError,
EcstoreDiskResult, EcstoreDiskStore, EcstoreEndpoint, EcstoreErrorType, EcstoreStorageError, EcstoreStore,
ecstore_local_disk_map_read,
EcstoreDiskResult, EcstoreDiskStore, EcstoreEndpoint, EcstoreErrorType, EcstoreHealLifecycleExpiryContext,
EcstoreStorageError, EcstoreStore, ecstore_load_admin_data_usage_from_backend_cached, ecstore_local_disk_map_read,
};
#[cfg(test)]
+235 -3
View File
@@ -19,11 +19,12 @@ use crate::heal::{
resume::{
CheckpointManager, ReplacementPhase, ReplacementTargetIdentity, ResumeManager, replacement_target_identities_match,
},
storage::{HealStorageAPI, next_heal_listing_token},
storage::{HealBucketUsageBaseline, HealStorageAPI, next_heal_listing_token},
};
use crate::{Error, Result};
use metrics::{counter, histogram};
use rustfs_common::heal_channel::{HealOpts, HealRequestSource, HealScanMode};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit};
use rustfs_madmin::heal_commands::HealResultItem;
use rustfs_utils::path::SLASH_SEPARATOR;
use serde::{Deserialize, Serialize};
@@ -178,6 +179,17 @@ pub enum HealPriority {
Urgent = 3,
}
impl HealPriority {
fn as_str(self) -> &'static str {
match self {
Self::Low => "low",
Self::Normal => "normal",
Self::High => "high",
Self::Urgent => "urgent",
}
}
}
/// Heal options
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealOptions {
@@ -498,6 +510,61 @@ impl HealTask {
}
}
fn emit_trace_task_state(&self, state: &'static str, duration: Duration, error: Option<&Error>) {
trace_emit(|| {
let mut event = TraceEvent::new(TraceKind::Heal, TraceFunc::HealTask)
.with_duration(duration)
.with_attr("task_id", self.id.as_str())
.with_attr("heal_type", self.heal_type.log_kind())
.with_attr("state", state)
.with_attr("source", self.source.as_str())
.with_attr("priority", self.priority.as_str())
.with_attr("retry_attempts", u64::from(self.retry_attempts))
.with_attr("dry_run", self.options.dry_run);
event = match &self.heal_type {
HealType::Cluster => event,
HealType::Object {
bucket,
object,
version_id,
} => {
let event = event.with_bucket(bucket.as_str()).with_object(object.as_str());
match version_id {
Some(version_id) => event.with_attr("version_id", version_id.as_str()),
None => event,
}
}
HealType::Bucket { bucket } => event.with_bucket(bucket.as_str()),
HealType::Prefix { bucket, prefix } => event.with_bucket(bucket.as_str()).with_object(prefix.as_str()),
HealType::ErasureSet { buckets, set_disk_id } => {
let bucket_count = u64::try_from(buckets.len()).unwrap_or(u64::MAX);
event
.with_attr("set_disk_id", set_disk_id.as_str())
.with_attr("bucket_count", bucket_count)
}
HealType::Metadata { bucket, object } => event.with_bucket(bucket.as_str()).with_object(object.as_str()),
HealType::ECDecode {
bucket,
object,
version_id,
} => {
let event = event.with_bucket(bucket.as_str()).with_object(object.as_str());
match version_id {
Some(version_id) => event.with_attr("version_id", version_id.as_str()),
None => event,
}
}
HealType::MRF { meta_path } => event.with_object(meta_path.as_str()),
};
match error {
Some(error) => event.with_attr("error", error.to_string()),
None => event,
}
});
}
async fn remaining_timeout(&self) -> Result<Option<Duration>> {
if let Some(total) = self.options.timeout {
let start_instant = { *self.task_start_instant.read().await };
@@ -717,6 +784,7 @@ impl HealTask {
queue_delay = ?queue_delay,
"Heal task started"
});
self.emit_trace_task_state("started", Duration::ZERO, None);
let result = match &self.heal_type {
HealType::Cluster => self.heal_cluster().await,
@@ -805,6 +873,14 @@ impl HealTask {
}
}
let terminal_state = match &result {
Ok(_) => "completed",
Err(Error::TaskCancelled) => "cancelled",
Err(Error::TaskTimeout) => "timed_out",
Err(_) => "failed",
};
self.emit_trace_task_state(terminal_state, start_instant.elapsed(), result.as_ref().err());
result
}
@@ -1535,7 +1611,7 @@ impl HealTask {
let (objects, next_token, is_truncated) = self
.await_with_control(
self.storage
.list_objects_for_heal_page(bucket, prefix, continuation_token.as_deref()),
.list_objects_for_heal_page(bucket, prefix, continuation_token.as_deref(), false),
)
.await?;
@@ -1697,6 +1773,23 @@ impl HealTask {
Ok(())
}
async fn apply_erasure_set_usage_baseline(&self, buckets: &[String]) -> Result<()> {
let baseline = match self
.await_with_control(self.storage.erasure_set_usage_baseline(buckets))
.await
{
Ok(Some(baseline)) => baseline,
Ok(None) => return Ok(()),
Err(err @ Error::TaskCancelled) | Err(err @ Error::TaskTimeout) => return Err(err),
Err(_) => return Ok(()),
};
let HealBucketUsageBaseline { objects_count, bytes } = baseline;
let mut progress = self.progress.write().await;
progress.set_total_baseline(objects_count, bytes);
Ok(())
}
async fn heal_metadata(&self, bucket: &str, object: &str) -> Result<()> {
debug!(
target: "rustfs::heal::task",
@@ -2298,6 +2391,8 @@ impl HealTask {
None
};
self.apply_erasure_set_usage_baseline(&buckets).await?;
let healing_marker = format!("{set_disk_id}:{}", self.id);
if let Some((disk, resume_manager, _)) = replacement_resume.as_ref() {
let state = resume_manager.get_state().await;
@@ -2602,7 +2697,8 @@ impl HealTask {
{
let mut progress = self.progress.write().await;
progress.update_progress(4, 4, 0, 0);
let bytes_processed = progress.bytes_processed;
progress.update_progress(4, 4, 0, bytes_processed);
}
match result {
@@ -2658,6 +2754,7 @@ mod tests {
use super::super::{DiskOption, DiskStore, Endpoint, HealDiskExt as _, new_disk};
use super::*;
use crate::heal::storage::{DiskStatus, HealListItem, HealObjectInfo};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, TraceSubscription, TraceVal, subscribe_trace_events};
use rustfs_madmin::heal_commands::{HealDriveInfo, HealResultItem, Infos};
use std::collections::{HashMap, VecDeque};
use std::sync::Mutex;
@@ -3203,6 +3300,8 @@ mod tests {
block_heal_object: Mutex<bool>,
resume_disk: Mutex<Option<DiskStore>>,
replacement_resume_disk: Mutex<Option<DiskStore>>,
usage_baseline: Mutex<Option<HealBucketUsageBaseline>>,
usage_baseline_error: Mutex<bool>,
}
#[test]
@@ -3265,11 +3364,69 @@ mod tests {
assert_eq!(samples_logged, MAX_BUCKET_FAILURE_LOG_SAMPLES);
}
#[tokio::test]
async fn execute_emits_heal_trace_task_state() {
let mut trace = subscribe_trace_events();
let storage = Arc::new(MockStorage::default());
let task = HealTask::from_request(
HealRequest::object("bucket-a".to_string(), "object-a".to_string(), Some("version-a".to_string())),
storage,
);
task.execute().await.expect("mock object heal should complete");
let started = recv_trace_task_state(&mut trace, &task.id, "started").await;
assert_eq!(started.kind, TraceKind::Heal);
assert_eq!(started.func, TraceFunc::HealTask);
assert_eq!(started.bucket.as_deref(), Some("bucket-a"));
assert_eq!(started.object.as_deref(), Some("object-a"));
assert_eq!(trace_attr_string(&started, "heal_type").as_deref(), Some("object"));
assert_eq!(trace_attr_string(&started, "source").as_deref(), Some("internal"));
assert_eq!(trace_attr_string(&started, "version_id").as_deref(), Some("version-a"));
let completed = recv_trace_task_state(&mut trace, &task.id, "completed").await;
assert_eq!(completed.kind, TraceKind::Heal);
assert_eq!(completed.func, TraceFunc::HealTask);
assert_eq!(trace_attr_string(&completed, "state").as_deref(), Some("completed"));
}
async fn recv_trace_task_state(trace: &mut TraceSubscription, task_id: &str, state: &str) -> TraceEvent {
for _ in 0..32 {
let event = tokio::time::timeout(Duration::from_secs(1), trace.recv())
.await
.expect("trace event should arrive")
.expect("trace bus should stay open");
if trace_attr_string(&event, "task_id").as_deref() == Some(task_id)
&& trace_attr_string(&event, "state").as_deref() == Some(state)
{
return (*event).clone();
}
}
panic!("expected trace state {state} for task {task_id}");
}
fn trace_attr_string(event: &TraceEvent, key: &str) -> Option<String> {
event.attrs.iter().find_map(|attr| {
if attr.key != key {
return None;
}
Some(match &attr.value {
TraceVal::Bool(value) => value.to_string(),
TraceVal::U64(value) => value.to_string(),
TraceVal::I64(value) => value.to_string(),
TraceVal::Str(value) => value.to_string(),
})
})
}
/// Build a latest, non-delete-marker heal list item with no version id.
fn heal_item(name: &str) -> HealListItem {
HealListItem {
name: name.to_string(),
version_id: None,
mod_time_unix_nanos: None,
lifecycle_object_info: None,
is_delete_marker: false,
}
}
@@ -3357,6 +3514,13 @@ mod tests {
}))
}
async fn erasure_set_usage_baseline(&self, _buckets: &[String]) -> Result<Option<HealBucketUsageBaseline>> {
if *self.usage_baseline_error.lock().unwrap() {
return Err(Error::Other("usage baseline unavailable".to_string()));
}
Ok(*self.usage_baseline.lock().unwrap())
}
async fn heal_bucket_metadata(&self, _bucket: &str) -> Result<()> {
Ok(())
}
@@ -3540,6 +3704,7 @@ mod tests {
bucket: &str,
prefix: &str,
continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool)> {
self.listed_prefixes.lock().unwrap().push(prefix.to_string());
if *self.truncate_without_token.lock().unwrap() {
@@ -4654,6 +4819,73 @@ mod tests {
assert!(storage.object_heal_opts.lock().unwrap().is_empty());
}
#[tokio::test]
async fn erasure_set_heal_applies_usage_baseline_to_progress() {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
resume_disk: Mutex::new(Some(disk)),
usage_baseline: Mutex::new(Some(HealBucketUsageBaseline {
objects_count: 10,
bytes: 8,
})),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
task.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect("erasure set heal should complete");
let progress = task.get_progress().await;
assert_eq!(progress.objects_total_count, 10);
assert_eq!(progress.objects_total_size, 8);
assert_eq!(progress.bytes_processed, 2);
assert!((progress.progress_percentage - 25.0).abs() < 0.001);
}
#[tokio::test]
async fn erasure_set_heal_ignores_usage_baseline_errors() {
let temp = TempDir::new().expect("temporary directory should be created");
let disk = make_resume_disk(&temp).await;
let storage = Arc::new(MockStorage {
resume_disk: Mutex::new(Some(disk)),
usage_baseline_error: Mutex::new(true),
..Default::default()
});
let request = HealRequest::new(
HealType::ErasureSet {
buckets: vec!["bucket-a".to_string()],
set_disk_id: "pool_0_set_0".to_string(),
},
HealOptions {
timeout: None,
..Default::default()
},
HealPriority::Normal,
);
let task = HealTask::from_request(request, storage);
task.heal_erasure_set(vec!["bucket-a".to_string()], "pool_0_set_0".to_string())
.await
.expect("usage baseline failures should not fail erasure set heal");
let progress = task.get_progress().await;
assert_eq!(progress.objects_total_count, 0);
assert_eq!(progress.objects_total_size, 0);
}
#[tokio::test]
async fn resumable_erasure_set_execution_is_cancelled_while_object_heal_is_pending() {
let temp = TempDir::new().expect("temporary directory should be created");
+1
View File
@@ -445,6 +445,7 @@ mod tests {
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> Result<(Vec<HealListItem>, Option<String>, bool), Error> {
Ok((Vec::new(), None, false))
}
@@ -176,7 +176,7 @@ async fn enumerate_all_versions(heal_storage: &Arc<ECStoreHealStorage>, bucket:
let mut token: Option<String> = None;
loop {
let (page, next, truncated) = heal_storage
.list_objects_for_heal_page(bucket, "", token.as_deref())
.list_objects_for_heal_page(bucket, "", token.as_deref(), false)
.await
.expect("list_objects_for_heal_page failed");
items.extend(page);
@@ -166,7 +166,7 @@ async fn enumerate_b5(heal_storage: &Arc<ECStoreHealStorage>, bucket: &str) -> V
let mut token: Option<String> = None;
loop {
let (page, next, truncated) = heal_storage
.list_objects_for_heal_page(bucket, "", token.as_deref())
.list_objects_for_heal_page(bucket, "", token.as_deref(), false)
.await
.expect("b5 list page failed");
items.extend(page);
@@ -187,7 +187,7 @@ async fn enumerate_disk_walk(heal_storage: &Arc<ECStoreHealStorage>, bucket: &st
let mut token: Option<String> = None;
loop {
let (page, next, truncated) = heal_storage
.list_versions_for_heal_page_disk_walk(SET_DISK_ID, bucket, "", token.as_deref())
.list_versions_for_heal_page_disk_walk(SET_DISK_ID, bucket, "", token.as_deref(), false)
.await
.expect("disk-walk list page failed");
items.extend(page);
@@ -418,7 +418,7 @@ mod serial_tests {
let mut pages = 0usize;
loop {
let (versions, next_forward, truncated) = ecstore
.heal_walk_versions_page(0, 0, bucket, "", forward.as_deref(), 2, 100_000)
.heal_walk_versions_page(0, 0, bucket, "", forward.as_deref(), 2, 100_000, false)
.await
.expect("heal_walk_versions_page failed");
pages += 1;
+2
View File
@@ -242,6 +242,7 @@ fn test_heal_task_status_atomic_update() {
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> rustfs_heal::Result<(Vec<HealListItem>, Option<String>, bool)> {
Ok((vec![], None, false))
}
@@ -385,6 +386,7 @@ async fn test_heal_task_transient_object_exists_skip_avoids_recreate() {
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> rustfs_heal::Result<(Vec<HealListItem>, Option<String>, bool)> {
Ok((Vec::new(), None, false))
}
+9 -1
View File
@@ -43,7 +43,7 @@ pub struct ServiceTraceOpts {
#[allow(dead_code)]
impl ServiceTraceOpts {
fn trace_types(&self) -> TraceType {
pub fn trace_types(&self) -> TraceType {
let mut tt = TraceType::default();
tt.set_if(self.s3, &TraceType::S3);
tt.set_if(self.internal, &TraceType::INTERNAL);
@@ -72,6 +72,14 @@ impl ServiceTraceOpts {
tt
}
pub fn only_errors(&self) -> bool {
self.only_errors
}
pub fn threshold(&self) -> Duration {
self.threshold
}
pub fn parse_params(&mut self, uri: &Uri) -> Result<(), String> {
let query_pairs: HashMap<_, _> = uri
.query()
+3
View File
@@ -108,6 +108,9 @@ temp-env = { workspace = true }
tempfile = { workspace = true }
uuid = { workspace = true, features = ["v4", "serde", "fast-rng", "macro-diagnostics"] }
tokio = { workspace = true, features = ["test-util", "fs", "rt-multi-thread"] }
# Test-only: pins the emitted scanner alert wire names against the canonical
# EventName string forms subscribers configure (rustfs/backlog#1868).
rustfs-s3-types.workspace = true
# Enables the shared MockWarmBackend / xl.meta assertion helpers exposed via
# the ecstore `api::tier::test_util` facade module (rustfs/backlog#1148 ilm-6).
rustfs-ecstore = { workspace = true, features = ["test-util"] }
+479 -4
View File
@@ -12,10 +12,10 @@
// See the License for the specific language governing permissions and
// limitations under the License.
use std::collections::HashSet;
use std::collections::{HashMap, HashSet};
use std::fs::FileType;
use std::io::ErrorKind;
use std::sync::{Arc, Once};
use std::sync::{Arc, Mutex, Once};
use std::time::{Duration, Instant, SystemTime};
use crate::ReplTargetSizeSummary;
@@ -32,6 +32,7 @@ use crate::scanner_io::{
SCANNER_SKIP_FILE_ERROR, ScannerIODisk as _, is_scanner_metadata_corrupt_error, is_scanner_metadata_transient_error,
};
use crate::sleeper::DynamicSleeper;
use crate::storage_api::owner::{EcstoreEventArgs, ecstore_send_event};
use metrics::{counter, describe_counter};
use rustfs_common::heal_channel::{
HEAL_DELETE_DANGLING, HealAdmissionDropReason, HealAdmissionResult, HealChannelPriority, HealChannelRequest,
@@ -41,6 +42,7 @@ use rustfs_common::metrics::{
CloseDiskGuard, IlmAction, Metric, Metrics, ScannerReplicationRepairKind, ScannerSourceWorkUpdate, ScannerWorkSource,
UpdateCurrentPathFn, current_path_updater, global_metrics,
};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit, trace_subscriber_count};
use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams};
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration};
@@ -97,6 +99,101 @@ const METRIC_SCANNER_EXCESS_FOLDERS_TOTAL: &str = "rustfs_scanner_excess_folders
const METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL: &str = "rustfs_scanner_pending_heal_prune_total";
const METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL: &str = "rustfs_scanner_pending_heal_malformed_total";
const MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET: usize = 128;
// --- scanner excess alerts as S3 notification events (rustfs/backlog#1868) --
//
// The excess-versions / excess-version-size / excess-folders alerts were
// metrics-and-logs only; subscribers (consoles, external auditors) had no way
// to hear them. MinIO emits s3:ObjectManyVersions / s3:ObjectLargeVersions /
// s3:PrefixManyFolders for the same conditions — RustFS carries those as
// EventName::Scanner* with the wire names below. Without a cooldown a single
// over-threshold object would re-emit on every scan cycle (~a minute), so
// emissions are edge-held per (kind, bucket, object) for 24h.
/// `s3:Scanner:ManyVersions` (MinIO `s3:ObjectManyVersions`).
pub const EVENT_SCANNER_MANY_VERSIONS: &str = "s3:Scanner:ManyVersions";
/// `s3:Scanner:LargeVersions` (MinIO `s3:ObjectLargeVersions`).
pub const EVENT_SCANNER_LARGE_VERSIONS: &str = "s3:Scanner:LargeVersions";
/// `s3:Scanner:BigPrefix` (MinIO `s3:PrefixManyFolders`).
pub const EVENT_SCANNER_BIG_PREFIX: &str = "s3:Scanner:BigPrefix";
const ENV_SCANNER_ALERT_COOLDOWN_SECS: &str = "RUSTFS_SCANNER_ALERT_COOLDOWN_SECS";
const DEFAULT_SCANNER_ALERT_COOLDOWN_SECS: u64 = 86_400;
/// Hard cap on distinct cooldown keys; a pathological number of over-threshold
/// objects clears the map wholesale instead of growing without bound (the
/// worst case is one re-emission per still-hot key per scan cycle).
const MAX_SCANNER_ALERT_COOLDOWN_KEYS: usize = 4096;
/// Distinct alert kinds sharing one cooldown map.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
enum ScannerAlertKind {
ManyVersions,
LargeVersions,
BigPrefix,
}
type ScannerAlertCooldownKey = (ScannerAlertKind, String, String);
type ScannerAlertCooldownMap = HashMap<ScannerAlertCooldownKey, Instant>;
static SCANNER_ALERT_EMISSION_COOLDOWN: Mutex<Option<ScannerAlertCooldownMap>> = Mutex::new(None);
fn scanner_alert_cooldown() -> Duration {
let raw = std::env::var(ENV_SCANNER_ALERT_COOLDOWN_SECS)
.ok()
.and_then(|v| v.parse::<u64>().ok());
Duration::from_secs(raw.unwrap_or(DEFAULT_SCANNER_ALERT_COOLDOWN_SECS))
}
/// Edge-held emission gate: returns `true` (and records the cooldown) only
/// when this (kind, bucket, object) last fired longer than the cooldown ago —
/// or never. Metrics and logs stay level-triggered every cycle; only the
/// notification events are held back.
fn scanner_alert_emission_allows(kind: ScannerAlertKind, bucket: &str, object: &str, cooldown: Duration) -> bool {
let key = (kind, bucket.to_string(), object.to_string());
let mut guard = SCANNER_ALERT_EMISSION_COOLDOWN
.lock()
.unwrap_or_else(|poison| poison.into_inner());
let guard = guard.get_or_insert_with(ScannerAlertCooldownMap::new);
let now = Instant::now();
// Expired entries leave first; the cap is still exceeded only when live
// keys alone overflow it, in which case a wholesale clear trades one
// extra emission per hot key for a hard memory bound.
if guard.len() >= MAX_SCANNER_ALERT_COOLDOWN_KEYS {
guard.retain(|_, fired_at| now.duration_since(*fired_at) < cooldown);
if guard.len() >= MAX_SCANNER_ALERT_COOLDOWN_KEYS {
guard.clear();
}
}
match guard.get(&key) {
Some(fired_at) if now.duration_since(*fired_at) < cooldown => false,
_ => {
guard.insert(key, now);
true
}
}
}
/// Emit a scanner alert as an S3 notification event through the standard
/// dispatch pipeline. Fire-and-forget: the notify layer owns delivery,
/// retry, and target filtering; the scanner never waits on it.
fn emit_scanner_alert_event(event_name: &str, bucket: &str, object: &str, size: i64, details: &[(&str, String)]) {
let mut req_params = HashMap::with_capacity(details.len());
for (key, value) in details {
req_params.insert((*key).to_string(), value.clone());
}
ecstore_send_event(EcstoreEventArgs {
event_name: event_name.to_string(),
bucket_name: bucket.to_string(),
object: crate::ScannerObjectInfo {
bucket: bucket.to_string(),
name: object.to_string(),
size,
..Default::default()
},
req_params,
user_agent: "Scanner".to_string(),
..Default::default()
});
}
const MAX_PENDING_SCANNER_HEALS_PER_BUCKET: usize = 10_000;
static SCANNER_INLINE_HEAL_WARN_ONCE: Once = Once::new();
@@ -430,6 +527,113 @@ fn non_negative_i64_to_u64(value: i64) -> u64 {
value.max(0) as u64
}
fn trace_start_instant() -> Option<Instant> {
(trace_subscriber_count() > 0).then(Instant::now)
}
fn emit_scanner_folder_trace(root: &str, folder: &str, objects: u64, started_at: Option<Instant>, state: &'static str) {
let Some(started_at) = started_at else {
return;
};
trace_emit(|| {
let (bucket, prefix) = path2_bucket_object_with_base_path(root, folder);
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerFolder)
.with_bucket(bucket)
.with_object(prefix)
.with_duration(started_at.elapsed())
.with_attr("state", state)
.with_attr("objects", objects)
});
}
fn emit_scanner_ilm_action_trace(
bucket: &str,
object: &str,
action: IlmAction,
count: u64,
queued: bool,
started_at: Option<Instant>,
) {
let Some(started_at) = started_at else {
return;
};
let state = if queued { "queued" } else { "not_queued" };
trace_emit(|| {
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerIlmAction)
.with_bucket(bucket)
.with_object(object)
.with_duration(started_at.elapsed())
.with_attr("state", state)
.with_attr("action", action.as_str())
.with_attr("count", count)
.with_attr("queued", queued)
});
}
struct ScannerHealCandidateTraceContext {
bucket: String,
object: Option<String>,
version_id: Option<String>,
scan_mode: Option<HealScanMode>,
started_at: Instant,
}
fn scanner_heal_candidate_trace_context(request: &HealChannelRequest) -> Option<ScannerHealCandidateTraceContext> {
let started_at = trace_start_instant()?;
Some(ScannerHealCandidateTraceContext {
bucket: request.bucket.clone(),
object: request.object_prefix.clone(),
version_id: request.object_version_id.clone(),
scan_mode: request.scan_mode,
started_at,
})
}
struct ScannerHealCandidateTrace<'a> {
candidate_type: &'static str,
bucket: &'a str,
object: Option<&'a str>,
version_id: Option<&'a str>,
priority: HealChannelPriority,
scan_mode: Option<HealScanMode>,
result: Result<HealAdmissionResult, &'a str>,
started_at: Instant,
}
fn emit_scanner_heal_candidate_trace(trace: ScannerHealCandidateTrace<'_>) {
trace_emit(|| {
let (state, admission, error) = match trace.result {
Ok(result) if result.is_admitted() => ("admitted", describe_heal_admission(result), None),
Ok(result) => ("not_admitted", describe_heal_admission(result), None),
Err(error) => ("submit_failed", "channel_error".to_string(), Some(error)),
};
let mut event = TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerHealCandidate)
.with_bucket(trace.bucket)
.with_duration(trace.started_at.elapsed())
.with_attr("state", state)
.with_attr("candidate_type", trace.candidate_type)
.with_attr("priority", heal_priority_label(trace.priority))
.with_attr("admission", admission);
if let Some(object) = trace.object {
event = event.with_object(object);
}
if let Some(version_id) = trace.version_id {
event = event.with_attr("version_id", version_id);
}
if let Some(scan_mode) = trace.scan_mode {
event = event.with_attr("scan_mode", scan_mode.as_str());
}
if let Some(error) = error {
event = event.with_attr("error", error);
}
event
});
}
fn apply_scanner_size_summary(into: &mut DataUsageEntry, summary: &SizeSummary) {
into.size = into.size.saturating_add(summary.total_size);
into.versions = into.versions.saturating_add(summary.versions);
@@ -677,9 +881,22 @@ async fn send_scanner_heal_request(
request: HealChannelRequest,
) -> Result<HealAdmissionResult, ScannerError> {
let priority = request.priority;
let trace_context = scanner_heal_candidate_trace_context(&request);
match send_heal_request_with_admission(request).await {
Ok(result) => {
record_heal_candidate_admission(candidate_type, priority, result);
if let Some(trace_context) = trace_context.as_ref() {
emit_scanner_heal_candidate_trace(ScannerHealCandidateTrace {
candidate_type,
bucket: &trace_context.bucket,
object: trace_context.object.as_deref(),
version_id: trace_context.version_id.as_deref(),
priority,
scan_mode: trace_context.scan_mode,
result: Ok(result),
started_at: trace_context.started_at,
});
}
Ok(result)
}
Err(err) => {
@@ -690,6 +907,18 @@ async fn send_scanner_heal_request(
"result" => "channel_error".to_string()
)
.increment(1);
if let Some(trace_context) = trace_context.as_ref() {
emit_scanner_heal_candidate_trace(ScannerHealCandidateTrace {
candidate_type,
bucket: &trace_context.bucket,
object: trace_context.object.as_deref(),
version_id: trace_context.version_id.as_deref(),
priority,
scan_mode: trace_context.scan_mode,
result: Err(err.as_str()),
started_at: trace_context.started_at,
});
}
Err(ScannerError::Other(err))
}
}
@@ -905,7 +1134,9 @@ impl ScannerItem {
"Scanner lifecycle action dispatched"
);
let done_ilm = Metrics::time_ilm(event.action);
let trace_started_at = trace_start_instant();
let queued = apply_expiry_rule(event, &LcEventSrc::Scanner, oi).await;
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
done_ilm(1)();
remaining_versions = 0;
@@ -957,7 +1188,9 @@ impl ScannerItem {
"Scanner lifecycle action dispatched"
);
let done_ilm = Metrics::time_ilm(event.action);
let trace_started_at = trace_start_instant();
let queued = apply_expiry_rule(event, &LcEventSrc::Scanner, oi).await;
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
done_ilm(1)();
if !versioning_config.prefix_enabled(&self.object_path()) && event.action == IlmAction::DeleteAction {
@@ -995,7 +1228,9 @@ impl ScannerItem {
"Scanner lifecycle action dispatched"
);
let done_ilm = Metrics::time_ilm(event.action);
let trace_started_at = trace_start_instant();
let queued = apply_transition_rule(event, &LcEventSrc::Scanner, oi).await;
emit_scanner_ilm_action_trace(&self.bucket, &oi.name, event.action, 1, queued, trace_started_at);
if record_scanner_ilm_action_if_queued(global_metrics(), event.action, 1, queued) {
done_ilm(1)();
}
@@ -1019,7 +1254,21 @@ impl ScannerItem {
let action = event.action;
let count = u64::try_from(to_delete_objs.len()).unwrap_or(u64::MAX);
let done_ilm = Metrics::time_ilm(action);
let trace_started_at = trace_start_instant();
let queued = enqueue_runtime_newer_noncurrent(&self.bucket, to_delete_objs, event, &LcEventSrc::Scanner).await;
if let Some(trace_started_at) = trace_started_at {
let state = if queued { "queued" } else { "not_queued" };
trace_emit(|| {
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerIlmAction)
.with_bucket(self.bucket.as_str())
.with_object(self.object_path())
.with_duration(trace_started_at.elapsed())
.with_attr("state", state)
.with_attr("action", action.as_str())
.with_attr("count", count)
.with_attr("queued", queued)
});
}
if record_scanner_ilm_action_if_queued(global_metrics(), action, count, queued) {
done_ilm(count)();
remaining_versions = remaining_versions.saturating_sub(noncurrent_accounting.len());
@@ -1197,6 +1446,7 @@ impl ScannerItem {
fn alert_excessive_versions(&self, remaining_versions: usize, cumulative_size: i64) {
ensure_scanner_alert_metrics_registered();
let (too_many_versions, too_large_versions) = should_alert_excessive_versions(remaining_versions, cumulative_size);
let object_path = self.object_path();
if too_many_versions {
global_metrics().record_scanner_source_executed(ScannerWorkSource::Alerts, 1);
counter!(
@@ -1204,13 +1454,26 @@ impl ScannerItem {
"bucket" => self.bucket.clone()
)
.increment(1);
if scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, &self.bucket, &object_path, scanner_alert_cooldown())
{
emit_scanner_alert_event(
EVENT_SCANNER_MANY_VERSIONS,
&self.bucket,
&object_path,
cumulative_size,
&[
("versions", remaining_versions.to_string()),
("threshold", scanner_excess_versions_threshold().to_string()),
],
);
}
warn!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_ALERT_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_FOLDER,
bucket = %self.bucket,
object = %self.object_path(),
object = %object_path,
versions = remaining_versions,
threshold = scanner_excess_versions_threshold(),
state = "excess_versions",
@@ -1224,13 +1487,31 @@ impl ScannerItem {
"bucket" => self.bucket.clone()
)
.increment(1);
if scanner_alert_emission_allows(
ScannerAlertKind::LargeVersions,
&self.bucket,
&object_path,
scanner_alert_cooldown(),
) {
emit_scanner_alert_event(
EVENT_SCANNER_LARGE_VERSIONS,
&self.bucket,
&object_path,
cumulative_size,
&[
("versions", remaining_versions.to_string()),
("cumulativeSize", cumulative_size.to_string()),
("threshold", scanner_excess_version_size_threshold().to_string()),
],
);
}
warn!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_ALERT_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_FOLDER,
bucket = %self.bucket,
object = %self.object_path(),
object = %object_path,
versions = remaining_versions,
cumulative_size,
threshold = scanner_excess_version_size_threshold(),
@@ -1611,6 +1892,15 @@ impl FolderScanner {
"root" => self.root.clone()
)
.increment(1);
if scanner_alert_emission_allows(ScannerAlertKind::BigPrefix, &self.root, folder, scanner_alert_cooldown()) {
emit_scanner_alert_event(
EVENT_SCANNER_BIG_PREFIX,
&self.root,
folder,
0,
&[("folders", total_folders.to_string()), ("threshold", threshold.to_string())],
);
}
warn!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_ALERT_STATE,
@@ -1830,6 +2120,7 @@ impl FolderScanner {
into: &mut DataUsageEntry,
) -> Result<(), ScannerError> {
let done_folder = Metrics::time(Metric::ScanFolder);
let trace_started_at = trace_start_instant();
if ctx.is_cancelled() {
return Err(ScannerError::Other("Operation cancelled".to_string()));
@@ -2895,6 +3186,8 @@ impl FolderScanner {
}
done_folder();
let scanned_objects = u64::try_from(into.objects).unwrap_or(u64::MAX);
emit_scanner_folder_trace(&self.root, &folder.name, scanned_objects, trace_started_at, "completed");
Ok(())
}
@@ -3076,6 +3369,90 @@ mod tests {
#[cfg(unix)]
use std::os::unix::fs::{PermissionsExt, symlink};
use std::sync::Mutex;
/// Reset the process-global alert cooldown map; test-only.
fn reset_alert_cooldowns() {
*SCANNER_ALERT_EMISSION_COOLDOWN
.lock()
.unwrap_or_else(|poison| poison.into_inner()) = Some(ScannerAlertCooldownMap::new());
}
/// The emitted event-name strings must be exactly what `EventName`
/// serializes, or a bucket notification subscribed to the documented name
/// would silently never match (rustfs/backlog#1868).
#[test]
fn scanner_alert_wire_names_match_canonical_event_names() {
use rustfs_s3_types::EventName;
assert_eq!(EVENT_SCANNER_MANY_VERSIONS, EventName::ScannerManyVersions.to_string());
assert_eq!(EVENT_SCANNER_LARGE_VERSIONS, EventName::ScannerLargeVersions.to_string());
assert_eq!(EVENT_SCANNER_BIG_PREFIX, EventName::ScannerBigPrefix.to_string());
}
fn cooldown_map_len() -> usize {
SCANNER_ALERT_EMISSION_COOLDOWN
.lock()
.unwrap_or_else(|poison| poison.into_inner())
.as_ref()
.map(|map| map.len())
.unwrap_or(0)
}
/// Backdate every recorded cooldown so the next check fires again.
fn expire_all_alert_cooldowns(cooldown: Duration) {
let now = Instant::now();
let mut guard = SCANNER_ALERT_EMISSION_COOLDOWN
.lock()
.unwrap_or_else(|poison| poison.into_inner());
if let Some(map) = guard.as_mut() {
for fired_at in map.values_mut() {
if let Some(expired) = now.checked_sub(cooldown + Duration::from_secs(1)) {
*fired_at = expired;
}
}
}
}
/// The emission gate is the only thing standing between an over-threshold
/// object and one S3 event per scan cycle, so its edge semantics get
/// pinned directly. All scenarios share one #[test] because the cooldown
/// map is process-global and parallel tests would read each other's
/// firings.
#[test]
fn scanner_alert_emission_is_edge_held_per_key_and_bounded() {
reset_alert_cooldowns();
let cooldown = Duration::from_secs(3600);
// First firing allows, an immediate re-check is held.
assert!(scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, "bkt", "obj", cooldown));
assert!(!scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, "bkt", "obj", cooldown));
// Different kind, object, and bucket are independent keys.
assert!(scanner_alert_emission_allows(ScannerAlertKind::LargeVersions, "bkt", "obj", cooldown));
assert!(scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, "bkt", "other", cooldown));
assert!(scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, "other", "obj", cooldown));
assert_eq!(cooldown_map_len(), 4);
// After the cooldown elapses the same key fires again.
expire_all_alert_cooldowns(cooldown);
assert!(scanner_alert_emission_allows(ScannerAlertKind::ManyVersions, "bkt", "obj", cooldown));
// A zero cooldown degenerates to always-emit (operators may want that).
assert!(scanner_alert_emission_allows(ScannerAlertKind::BigPrefix, "bkt", "dir", Duration::ZERO));
assert!(scanner_alert_emission_allows(ScannerAlertKind::BigPrefix, "bkt", "dir", Duration::ZERO));
// Hard bound: overflow the cap with zero-cooldown keys and confirm the
// map clears rather than growing past it.
reset_alert_cooldowns();
for index in 0..=(MAX_SCANNER_ALERT_COOLDOWN_KEYS + 8) {
let _ = scanner_alert_emission_allows(ScannerAlertKind::BigPrefix, "bkt", &format!("dir-{index}"), Duration::ZERO);
}
assert!(
cooldown_map_len() <= MAX_SCANNER_ALERT_COOLDOWN_KEYS,
"cooldown map must stay bounded, got {}",
cooldown_map_len()
);
}
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use temp_env::{with_var, with_var_unset};
use tracing_subscriber::fmt::MakeWriter;
@@ -4400,6 +4777,104 @@ mod tests {
);
}
#[tokio::test]
async fn scanner_trace_helpers_emit_expected_events() {
let mut trace = rustfs_common::trace_bus::subscribe_trace_events();
emit_scanner_folder_trace(
"/tmp/rustfs-scanner-trace",
"/tmp/rustfs-scanner-trace/bucket-a/folder-a",
7,
Some(Instant::now()),
"completed",
);
let folder = recv_scanner_trace_event(
&mut trace,
TraceFunc::ScannerFolder,
Some("bucket-a"),
Some("folder-a"),
Some("completed"),
)
.await;
assert_eq!(trace_attr_string(&folder, "objects").as_deref(), Some("7"));
emit_scanner_ilm_action_trace("bucket-a", "object-a", IlmAction::DeleteAction, 2, true, Some(Instant::now()));
let ilm = recv_scanner_trace_event(
&mut trace,
TraceFunc::ScannerIlmAction,
Some("bucket-a"),
Some("object-a"),
Some("queued"),
)
.await;
assert_eq!(trace_attr_string(&ilm, "action").as_deref(), Some("delete"));
assert_eq!(trace_attr_string(&ilm, "count").as_deref(), Some("2"));
assert_eq!(trace_attr_string(&ilm, "queued").as_deref(), Some("true"));
emit_scanner_heal_candidate_trace(ScannerHealCandidateTrace {
candidate_type: "object",
bucket: "bucket-a",
object: Some("object-a"),
version_id: Some("version-a"),
priority: HealChannelPriority::High,
scan_mode: Some(HealScanMode::Deep),
result: Ok(HealAdmissionResult::Merged),
started_at: Instant::now(),
});
let heal_candidate = recv_scanner_trace_event(
&mut trace,
TraceFunc::ScannerHealCandidate,
Some("bucket-a"),
Some("object-a"),
Some("admitted"),
)
.await;
assert_eq!(trace_attr_string(&heal_candidate, "candidate_type").as_deref(), Some("object"));
assert_eq!(trace_attr_string(&heal_candidate, "priority").as_deref(), Some("high"));
assert_eq!(trace_attr_string(&heal_candidate, "scan_mode").as_deref(), Some("deep"));
assert_eq!(trace_attr_string(&heal_candidate, "version_id").as_deref(), Some("version-a"));
assert_eq!(trace_attr_string(&heal_candidate, "admission").as_deref(), Some("merged"));
}
async fn recv_scanner_trace_event(
trace: &mut rustfs_common::trace_bus::TraceSubscription,
func: TraceFunc,
bucket: Option<&str>,
object: Option<&str>,
state: Option<&str>,
) -> TraceEvent {
for _ in 0..32 {
let event = tokio::time::timeout(Duration::from_secs(1), trace.recv())
.await
.expect("scanner trace event should arrive")
.expect("trace bus should stay open");
if event.kind == TraceKind::Scanner
&& event.func == func
&& event.bucket.as_deref() == bucket
&& event.object.as_deref() == object
&& state.is_none_or(|state| trace_attr_string(&event, "state").as_deref() == Some(state))
{
return (*event).clone();
}
}
panic!("expected scanner trace event {func:?} for bucket {bucket:?} object {object:?}");
}
fn trace_attr_string(event: &TraceEvent, key: &str) -> Option<String> {
event.attrs.iter().find_map(|attr| {
if attr.key != key {
return None;
}
Some(match &attr.value {
rustfs_common::trace_bus::TraceVal::Bool(value) => value.to_string(),
rustfs_common::trace_bus::TraceVal::U64(value) => value.to_string(),
rustfs_common::trace_bus::TraceVal::I64(value) => value.to_string(),
rustfs_common::trace_bus::TraceVal::Str(value) => value.to_string(),
})
})
}
#[test]
fn test_build_high_priority_heal_admission_error_contains_context() {
let err = build_high_priority_heal_admission_error(
+4 -3
View File
@@ -78,6 +78,7 @@ pub(crate) use rustfs_ecstore::api::disk::{
pub(crate) use rustfs_ecstore::api::error::{
Error as EcstoreErrorType, Result as EcstoreResultType, StorageError as EcstoreStorageError,
};
pub(crate) use rustfs_ecstore::api::event::{EventArgs as EcstoreEventArgs, send_event as ecstore_send_event};
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::layout::{
EndpointServerPools as EcstoreEndpointServerPools, Endpoints as EcstoreEndpoints, PoolEndpoints as EcstorePoolEndpoints,
@@ -110,8 +111,8 @@ pub(crate) mod owner {
ECSTORE_BUCKET_META_PREFIX, ECSTORE_RUSTFS_META_BUCKET, ECSTORE_STORAGE_FORMAT_FILE, ECSTORE_STORAGECLASS_RRS,
ECSTORE_STORAGECLASS_STANDARD, ECSTORE_TRANSITION_COMPLETE, EcstoreBucketTargetSys, EcstoreBucketVersioningSys,
EcstoreDisk, EcstoreDiskAPI, EcstoreDiskBytes, EcstoreDiskError, EcstoreDiskInfo, EcstoreDiskInfoOptions,
EcstoreDiskLocation, EcstoreDiskResult, EcstoreErrorType, EcstoreEvaluator, EcstoreEvent, EcstoreLcEventSrc,
EcstoreLifecycle, EcstoreListPathRawOptions, EcstoreNsScannerOpenRequest, EcstoreObjectOpts,
EcstoreDiskLocation, EcstoreDiskResult, EcstoreErrorType, EcstoreEvaluator, EcstoreEvent, EcstoreEventArgs,
EcstoreLcEventSrc, EcstoreLifecycle, EcstoreListPathRawOptions, EcstoreNsScannerOpenRequest, EcstoreObjectOpts,
EcstoreReplicationConfigurationExt, EcstoreReplicationScannerBridge, EcstoreResultType, EcstoreScanGuard,
EcstoreSetDisks, EcstoreStorageError, EcstoreStore, EcstoreTierConfig, EcstoreVersioningApi,
ScannerReplicationHealObject, ScannerReplicationHealResult, ScannerReplicationQueueAdmission, ecstore_apply_expiry_rule,
@@ -121,7 +122,7 @@ pub(crate) mod owner {
ecstore_is_erasure_sd, ecstore_is_reserved_or_invalid_bucket, ecstore_list_path_raw,
ecstore_object_opts_from_object_info, ecstore_path2_bucket_object, ecstore_path2_bucket_object_with_base_path,
ecstore_read_config, ecstore_replace_bucket_usage_memory_from_info, ecstore_resolve_object_store_handle,
ecstore_save_config, scanner_replication_config_for_lifecycle_eval,
ecstore_save_config, ecstore_send_event, scanner_replication_config_for_lifecycle_eval,
};
#[cfg(test)]
+37
View File
@@ -0,0 +1,37 @@
# Scanner Excess Alerts: Metrics, S3 Events, and Thresholds
> 中文版:[scanner-excess-alerts_zh.md](scanner-excess-alerts_zh.md)
Date: 2026-08-18 (rustfs/backlog#1868 / HS-04; includes the HS-15 threshold-delta notes)
The background scanner detects three classes of "excess" conditions while it walks buckets and surfaces them as alerts. This page documents each alert's trigger condition, the subscribable S3 event, the cooldown semantics, and the threshold differences versus MinIO — for operators debugging alerts and for event consumers wiring up subscriptions.
## The three alerts
| Alert | Trigger (per scan cycle) | Metric | S3 event (RustFS wire name) | MinIO event name |
|---|---|---|---|---|
| Excess versions | Retained versions of one object ≥ `scanner:alert_excess_versions` | `rustfs_scanner_excess_object_versions_total{bucket}` | `s3:Scanner:ManyVersions` | `s3:ObjectManyVersions` |
| Excess version size | Cumulative bytes of all versions of one object ≥ `scanner:alert_excess_version_size` | `rustfs_scanner_excess_object_version_size_total{bucket}` | `s3:Scanner:LargeVersions` | `s3:ObjectLargeVersions` |
| Excess folders | Direct subfolders of one directory > `scanner:alert_excess_folders` | `rustfs_scanner_excess_folders_total{root}` | `s3:Scanner:BigPrefix` | `s3:PrefixManyFolders` |
Subscribe like any bucket notification: configure a notification on the target bucket with the RustFS wire name above (or the `s3:Scanner:*` wildcard). Events carry `UserAgent: Scanner` as their origin marker, and `req_params` holds the observed value and the threshold (`versions` / `cumulativeSize` / `folders` / `threshold`), so consumers can judge severity directly.
## Metrics and events fire on different cadences
- **Metrics and structured logs are level-triggered**: as long as the object stays over the threshold, every scan cycle counts and logs it (default cycle ≈ 60s; see `scanner:speed`).
- **S3 events are edge-triggered with a cooldown**: the same (alert kind, bucket, object) emits at most once per cooldown window — 24 hours by default (`RUSTFS_SCANNER_ALERT_COOLDOWN_SECS`; set it to 0 to emit every cycle). When the window lapses and the object is still over the threshold, the event fires again. The cooldown table lives in process memory with a 4096-entry hard cap; on overflow it is cleared and rebuilt (worst case: one extra emission per still-hot key).
- A process restart resets the cooldown (every still-over-threshold object emits once more after a restart) — deliberately: restarts usually accompany incident response, and the re-emission buys visibility.
## Threshold defaults and the MinIO deltas (HS-15)
| Config key | ENV | RustFS default | MinIO default | Notes |
|---|---|---|---|---|
| `scanner:alert_excess_versions` | `RUSTFS_SCANNER_ALERT_EXCESS_VERSIONS` | 100 | 100 | Identical |
| `scanner:alert_excess_version_size` | `RUSTFS_SCANNER_ALERT_EXCESS_VERSION_SIZE` | 1 TiB | 1 TB | Same order of magnitude; different unit basis (TiB vs TB) |
| `scanner:alert_excess_folders` | `RUSTFS_SCANNER_ALERT_EXCESS_FOLDERS` | 65538 | 50000 | **Deliberate divergence**: 65538 tolerates the Proxmox Backup Server chunk layout (65536 chunks per directory plus the directory's own entries); MinIO's 50000 would fire continuously for PBS users. Set it to 50000 explicitly to match MinIO behavior |
All three keys accept both env and admin config (`PUT /rustfs/admin/v3/config`, `scanner` subsystem); hot updates take effect immediately.
## Why the event names are mapped
RustFS's event enum (`rustfs_s3_types::EventName::ScannerManyVersions/LargeVersions/BigPrefix`) keeps the repo's established `s3:Scanner:*` wire names (literally different from MinIO's `s3:ObjectManyVersions`; the enum comments preserve the mapping). Subscribers should use the RustFS wire names in this page. If you need MinIO-literal compatibility, map the names on the console/consumer side — do not change the published wire names.
@@ -0,0 +1,37 @@
# Scanner 超限告警:指标、S3 事件与阈值
> English version: [scanner-excess-alerts.md](scanner-excess-alerts.md)
日期:2026-08-18rustfs/backlog#1868 / HS-04,含 HS-15 阈值差异说明)
后台 scanner 在扫描过程中检测三类"超限"状态并对外告警。本文说明每类告警的触发条件、可订阅的 S3 事件、冷却语义,以及与 MinIO 的阈值差异,供运维排障与事件消费方对接。
## 三类告警
| 告警 | 触发条件(任一扫描周期) | 指标 | S3 事件(RustFS wire 名) | MinIO 对应事件名 |
|---|---|---|---|---|
| 版本数超限 | 单对象保留版本数 ≥ `scanner:alert_excess_versions` | `rustfs_scanner_excess_object_versions_total{bucket}` | `s3:Scanner:ManyVersions` | `s3:ObjectManyVersions` |
| 版本总大小超限 | 单对象全部版本累计字节 ≥ `scanner:alert_excess_version_size` | `rustfs_scanner_excess_object_version_size_total{bucket}` | `s3:Scanner:LargeVersions` | `s3:ObjectLargeVersions` |
| 子目录数超限 | 单目录直接子目录数 > `scanner:alert_excess_folders` | `rustfs_scanner_excess_folders_total{root}` | `s3:Scanner:BigPrefix` | `s3:PrefixManyFolders` |
订阅方式与普通桶通知一致:对目标桶配置 notification,事件名填上表 RustFS wire 名(或通配 `s3:Scanner:*`)。事件以 `UserAgent: Scanner` 标记来源,`req_params` 携带实际值与阈值(`versions` / `cumulativeSize` / `folders` / `threshold`),便于消费方直接判断严重程度。
## 指标与事件的触发节奏不同
- **指标与结构化日志是电平触发**:只要对象仍在阈值之上,每个扫描周期都会计数/打日志(默认周期约 60s,见 `scanner:speed`)。
- **S3 事件是边沿触发 + 冷却**:同一 (告警类型, 桶, 对象) 在冷却窗口内只发一次,默认 24 小时(`RUSTFS_SCANNER_ALERT_COOLDOWN_SECS`,设 0 表示每周期都发)。窗口过后对象仍超限会再次发出。冷却表在进程内有 4096 条硬顶,超限清空重建(最坏情况是每个仍超限的 key 多发一次)。
- 进程重启会重置冷却(重启后每个仍超限的对象会再发一次)——这是有意为之:重启常伴随排障,重发提供可见性。
## 阈值默认值与 MinIO 差异(HS-15
| 配置键 | ENV | RustFS 默认 | MinIO 默认 | 差异说明 |
|---|---|---|---|---|
| `scanner:alert_excess_versions` | `RUSTFS_SCANNER_ALERT_EXCESS_VERSIONS` | 100 | 100 | 一致 |
| `scanner:alert_excess_version_size` | `RUSTFS_SCANNER_ALERT_EXCESS_VERSION_SIZE` | 1 TiB | 1 TB | 语义同量级,单位口径不同(TiB vs TB) |
| `scanner:alert_excess_folders` | `RUSTFS_SCANNER_ALERT_EXCESS_FOLDERS` | 65538 | 50000 | **有意差异**65538 兼容 Proxmox Backup Server 的 chunk 布局(每目录 65536 个 chunk + 目录自身条目),按 MinIO 的 50000 会对 PBS 用户持续误报。如需与 MinIO 行为一致可显式配置为 50000 |
三个键均支持 env 与 admin config`PUT /rustfs/admin/v3/config``scanner` 子系统)双通道,热更新即时生效。
## 事件名映射的由来
RustFS 的事件枚举(`rustfs_s3_types::EventName::ScannerManyVersions/LargeVersions/BigPrefix`)沿用仓库既有 wire 名 `s3:Scanner:*`(与 MinIO 的 `s3:ObjectManyVersions` 字面不同,枚举注释中保留了映射关系)。订阅方应以本文的 RustFS wire 名为准;如需 MinIO 字面兼容,请在 console/消费侧做名称映射,不要修改已发布的 wire 名。
+397 -25
View File
@@ -23,18 +23,24 @@ use futures::{Stream, StreamExt};
use http::{HeaderMap, HeaderValue};
use hyper::{Method, StatusCode};
use matchit::Params;
use regex::Regex;
use rustfs_common::trace_bus::{TraceEvent, TraceKind, TraceVal, subscribe_trace_events};
use rustfs_madmin::service_commands::ServiceTraceOpts;
use rustfs_madmin::trace::TraceType;
use rustfs_policy::policy::action::{Action, AdminAction};
use s3s::header::CONTENT_TYPE;
use s3s::stream::{ByteStream, DynByteStream};
use s3s::{Body, S3Request, S3Response, S3Result, StdError, s3_error};
use serde::Serialize;
use std::collections::HashMap;
use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::Duration;
use std::time::{Duration, SystemTime};
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use tracing::error;
use url::form_urlencoded;
#[derive(Serialize)]
struct ProfileStatus {
@@ -206,16 +212,164 @@ impl Stream for TraceStream {
impl ByteStream for TraceStream {}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TraceKindFilter {
heal: bool,
scanner: bool,
}
impl TraceKindFilter {
const ALL_SUPPORTED: Self = Self {
heal: true,
scanner: true,
};
fn from_request(uri: &hyper::Uri, trace_types: TraceType) -> S3Result<Self> {
let mut has_kind = false;
let mut filter = Self {
heal: false,
scanner: false,
};
for (key, value) in trace_query_pairs(uri) {
if key != "kind" {
continue;
}
has_kind = true;
for item in value.split(',') {
match item.trim().to_ascii_lowercase().as_str() {
"heal" | "healing" => filter.heal = true,
"scanner" => filter.scanner = true,
"all" => return Ok(Self::ALL_SUPPORTED),
_ => return Err(s3_error!(InvalidRequest, "invalid trace kind")),
}
}
}
if has_kind {
return Ok(filter);
}
if trace_types.mask() == 0 || trace_query_flag(uri, "all") {
return Ok(Self::ALL_SUPPORTED);
}
Ok(Self {
heal: trace_types.overlaps(&TraceType::HEALING),
scanner: trace_types.overlaps(&TraceType::SCANNER),
})
}
const fn matches(self, kind: TraceKind) -> bool {
match kind {
TraceKind::Heal => self.heal,
TraceKind::Scanner => self.scanner,
}
}
}
#[derive(Debug)]
struct TraceStreamFilter {
kinds: TraceKindFilter,
regex: Option<Regex>,
threshold: Duration,
}
impl TraceStreamFilter {
fn from_request(uri: &hyper::Uri, opts: &ServiceTraceOpts) -> S3Result<Self> {
if opts.only_errors() {
return Err(s3_error!(
InvalidRequest,
"trace error-only filter is not supported for heal/scanner trace"
));
}
Ok(Self {
kinds: TraceKindFilter::from_request(uri, opts.trace_types())?,
regex: trace_regex_filter(uri)?,
threshold: opts.threshold(),
})
}
fn matches_kind(&self, kind: TraceKind) -> bool {
self.kinds.matches(kind)
}
fn matches_record(&self, record: &TraceWireRecord) -> bool {
record.duration >= self.threshold && self.regex.as_ref().is_none_or(|regex| record.matches_regex(regex))
}
}
#[derive(Serialize)]
struct TraceWireRecord {
#[serde(rename = "type")]
trace_type: u64,
#[serde(rename = "nodename")]
node_name: String,
#[serde(rename = "funcname")]
func_name: String,
#[serde(rename = "time")]
time: String,
#[serde(rename = "path")]
path: String,
#[serde(rename = "dur")]
duration: Duration,
#[serde(rename = "bytes", skip_serializing_if = "Option::is_none")]
bytes: Option<i64>,
#[serde(rename = "msg", skip_serializing_if = "Option::is_none")]
message: Option<String>,
#[serde(rename = "custom", skip_serializing_if = "Option::is_none")]
custom: Option<HashMap<String, String>>,
}
impl TraceWireRecord {
fn from_event(node_name: &str, event: &TraceEvent) -> Self {
Self {
trace_type: trace_type_mask(event.kind),
node_name: node_name.to_owned(),
func_name: event.func.as_str().to_owned(),
time: trace_time_string(event.time),
path: trace_path(event),
duration: event.duration,
bytes: trace_bytes(event.bytes),
message: None,
custom: trace_custom_attrs(event),
}
}
fn dropped(node_name: &str, dropped: u64) -> Self {
let mut custom = HashMap::new();
custom.insert("dropped_events".to_string(), dropped.to_string());
Self {
trace_type: 0,
node_name: node_name.to_owned(),
func_name: "trace.Dropped".to_string(),
time: trace_time_string(SystemTime::now()),
path: String::new(),
duration: Duration::ZERO,
bytes: None,
message: Some("trace subscriber lagged".to_string()),
custom: Some(custom),
}
}
fn matches_regex(&self, regex: &Regex) -> bool {
regex.is_match(&self.func_name)
|| regex.is_match(&self.path)
|| self.message.as_ref().is_some_and(|message| regex.is_match(message))
|| self
.custom
.as_ref()
.is_some_and(|custom| custom.iter().any(|(key, value)| regex.is_match(key) || regex.is_match(value)))
}
}
/// `GET /v3/trace` — stream real-time server trace events.
///
/// RustFS emits diagnostics through the `tracing` pipeline but does not expose
/// an in-process subscriber that can fan trace events out to an admin client
/// (there is no request-trace broadcast channel). Rather than return an opaque
/// `501` — which would make `mc admin trace` fail to connect — this honors the
/// streaming NDJSON contract: it validates the requested trace filters, opens
/// the stream, emits a single capability record explaining that live tracing is
/// not wired, then holds the connection open with keep-alives. No fabricated
/// trace records are ever sent.
/// RustFS currently publishes heal and scanner diagnostics through the common
/// trace bus. The admin endpoint exposes those events as MinIO-shaped NDJSON
/// records while keeping unsupported trace classes filtered out.
pub struct TraceHandler {}
#[async_trait::async_trait]
@@ -228,24 +382,13 @@ impl Operation for TraceHandler {
let mut opts = ServiceTraceOpts::default();
opts.parse_params(&req.uri)
.map_err(|_| s3_error!(InvalidRequest, "invalid trace parameters"))?;
let filter = TraceStreamFilter::from_request(&req.uri, &opts)?;
let node_name = sysinfo::System::host_name().unwrap_or_else(|| "rustfs".to_string());
let (tx, rx) = mpsc::channel::<Result<Bytes, StdError>>(8);
let mut subscription = subscribe_trace_events();
let (tx, rx) = mpsc::channel::<Result<Bytes, StdError>>(64);
spawn_traced(async move {
let notice = serde_json::json!({
"nodename": node_name,
"funcname": "admin.Trace",
"msg": "RustFS does not expose an in-process trace-event subscriber; live tracing is not yet available",
"err": "trace_streaming_unsupported",
});
if let Ok(mut encoded) = serde_json::to_vec(&notice) {
encoded.push(b'\n');
if tx.send(Ok(Bytes::from(encoded))).await.is_err() {
return;
}
}
let mut ticker = tokio::time::interval(Duration::from_secs(15));
ticker.tick().await;
loop {
@@ -256,6 +399,26 @@ impl Operation for TraceHandler {
break;
}
}
received = subscription.recv() => {
match received {
Ok(event) => {
if !filter.matches_kind(event.kind) {
continue;
}
let record = TraceWireRecord::from_event(&node_name, &event);
if filter.matches_record(&record) && send_trace_record(&tx, &record).await.is_err() {
break;
}
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(dropped)) => {
let record = TraceWireRecord::dropped(&node_name, dropped);
if send_trace_record(&tx, &record).await.is_err() {
break;
}
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
}
}
}
}
});
@@ -269,17 +432,115 @@ impl Operation for TraceHandler {
}
}
async fn send_trace_record(tx: &mpsc::Sender<Result<Bytes, StdError>>, record: &TraceWireRecord) -> Result<(), ()> {
let Some(encoded) = encode_ndjson(record) else {
return Ok(());
};
tx.send(Ok(encoded)).await.map_err(|_| ())
}
fn encode_ndjson(value: &impl Serialize) -> Option<Bytes> {
let mut encoded = serde_json::to_vec(value).ok()?;
encoded.push(b'\n');
Some(Bytes::from(encoded))
}
fn trace_query_pairs(uri: &hyper::Uri) -> impl Iterator<Item = (String, String)> + '_ {
uri.query()
.into_iter()
.flat_map(|query| form_urlencoded::parse(query.as_bytes()))
.map(|(key, value)| (key.into_owned(), value.into_owned()))
}
fn trace_query_flag(uri: &hyper::Uri, flag: &str) -> bool {
trace_query_pairs(uri).any(|(key, value)| key == flag && value == "true")
}
fn trace_regex_filter(uri: &hyper::Uri) -> S3Result<Option<Regex>> {
trace_query_pairs(uri)
.find_map(|(key, value)| {
if key == "filter" && !value.is_empty() {
Some(value)
} else {
None
}
})
.map(|pattern| Regex::new(&pattern).map_err(|_| s3_error!(InvalidRequest, "invalid trace filter")))
.transpose()
}
fn trace_type_mask(kind: TraceKind) -> u64 {
match kind {
TraceKind::Heal => TraceType::HEALING.mask(),
TraceKind::Scanner => TraceType::SCANNER.mask(),
}
}
fn trace_time_string(time: SystemTime) -> String {
match OffsetDateTime::from(time).format(&Rfc3339) {
Ok(value) => value,
Err(_) => "1970-01-01T00:00:00Z".to_string(),
}
}
fn trace_path(event: &TraceEvent) -> String {
match (event.bucket.as_deref(), event.object.as_deref()) {
(Some(bucket), Some(object)) if !object.is_empty() => format!("{bucket}/{object}"),
(Some(bucket), _) => bucket.to_owned(),
(None, Some(object)) => object.to_owned(),
(None, None) => String::new(),
}
}
fn trace_bytes(bytes: u64) -> Option<i64> {
if bytes == 0 {
return None;
}
match i64::try_from(bytes) {
Ok(value) => Some(value),
Err(_) => Some(i64::MAX),
}
}
fn trace_custom_attrs(event: &TraceEvent) -> Option<HashMap<String, String>> {
if event.attrs.is_empty() {
return None;
}
Some(
event
.attrs
.iter()
.map(|attr| (attr.key.to_string(), trace_value_string(&attr.value)))
.collect(),
)
}
fn trace_value_string(value: &TraceVal) -> String {
match value {
TraceVal::Bool(value) => value.to_string(),
TraceVal::U64(value) => value.to_string(),
TraceVal::I64(value) => value.to_string(),
TraceVal::Str(value) => value.to_string(),
}
}
#[cfg(test)]
mod tests {
use super::{
ProfileControlHandler, ProfileHandler, ProfileStatusHandler, ProfilingDownloadHandler, ProfilingStartHandler,
TraceHandler,
TraceHandler, TraceKindFilter, TraceStreamFilter, TraceWireRecord,
};
use crate::admin::router::Operation;
use http::{Extensions, HeaderMap, Uri};
use hyper::Method;
use matchit::Params;
use s3s::{Body, S3ErrorCode, S3Request};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind};
use rustfs_madmin::service_commands::ServiceTraceOpts;
use rustfs_madmin::trace::TraceType;
use s3s::{Body, S3ErrorCode, S3Request, S3Result};
use std::time::{Duration, UNIX_EPOCH};
fn build_profile_request(uri: &'static str) -> S3Request<Body> {
S3Request {
@@ -295,6 +556,13 @@ mod tests {
}
}
fn build_trace_stream_filter(uri: &'static str) -> S3Result<TraceStreamFilter> {
let uri = Uri::from_static(uri);
let mut opts = ServiceTraceOpts::default();
opts.parse_params(&uri).expect("test trace params should parse");
TraceStreamFilter::from_request(&uri, &opts)
}
#[tokio::test]
async fn profile_handler_rejects_missing_credentials() {
let result = ProfileHandler {}
@@ -358,4 +626,108 @@ mod tests {
.expect_err("trace must reject anonymous requests");
assert_eq!(err.code(), &S3ErrorCode::AccessDenied);
}
#[test]
fn trace_kind_filter_supports_kind_query() {
let uri = Uri::from_static("/rustfs/admin/v3/trace?kind=heal");
let filter = TraceKindFilter::from_request(&uri, TraceType::default()).expect("kind filter should parse");
assert!(filter.matches(TraceKind::Heal));
assert!(!filter.matches(TraceKind::Scanner));
}
#[test]
fn trace_kind_filter_defaults_to_supported_events_without_type_flags() {
let uri = Uri::from_static("/rustfs/admin/v3/trace");
let filter = TraceKindFilter::from_request(&uri, TraceType::default()).expect("empty filter should parse");
assert!(filter.matches(TraceKind::Heal));
assert!(filter.matches(TraceKind::Scanner));
}
#[test]
fn trace_kind_filter_rejects_unknown_kind() {
let uri = Uri::from_static("/rustfs/admin/v3/trace?kind=s3");
let err = TraceKindFilter::from_request(&uri, TraceType::default()).expect_err("unknown kind should fail");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
}
#[test]
fn trace_stream_filter_matches_regex_against_path_and_attrs() {
let filter = build_trace_stream_filter("/rustfs/admin/v3/trace?kind=heal&filter=data/.%2Bxl.meta")
.expect("regex filter should parse");
let event = TraceEvent::new(TraceKind::Heal, TraceFunc::HealObject)
.with_bucket("data")
.with_object("dir/xl.meta")
.with_attr("dry_run", true);
let record = TraceWireRecord::from_event("node-a", &event);
assert!(filter.matches_kind(event.kind));
assert!(filter.matches_record(&record));
}
#[test]
fn trace_stream_filter_rejects_invalid_regex() {
let err = build_trace_stream_filter("/rustfs/admin/v3/trace?kind=heal&filter=[").expect_err("invalid regex should fail");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
}
#[test]
fn trace_stream_filter_applies_threshold() {
let filter =
build_trace_stream_filter("/rustfs/admin/v3/trace?kind=heal&threshold=10ms").expect("threshold should parse");
let short = TraceWireRecord::from_event(
"node-a",
&TraceEvent::new(TraceKind::Heal, TraceFunc::HealObject).with_duration(Duration::from_millis(9)),
);
let long = TraceWireRecord::from_event(
"node-a",
&TraceEvent::new(TraceKind::Heal, TraceFunc::HealObject).with_duration(Duration::from_millis(10)),
);
assert!(!filter.matches_record(&short));
assert!(filter.matches_record(&long));
}
#[test]
fn trace_stream_filter_rejects_error_only_filter() {
let err = build_trace_stream_filter("/rustfs/admin/v3/trace?kind=heal&err=true").expect_err("err filter should fail");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
}
#[test]
fn trace_wire_record_contains_madmin_trace_fields() {
let event = TraceEvent::new(TraceKind::Heal, TraceFunc::HealObject)
.with_bucket("bucket")
.with_object("object")
.with_duration(Duration::from_millis(3))
.with_bytes(17)
.with_attr("dry", true);
let mut record = TraceWireRecord::from_event("node-a", &event);
record.time = "1970-01-01T00:00:00Z".to_string();
let value = serde_json::to_value(&record).expect("trace record should serialize");
assert_eq!(value["type"], TraceType::HEALING.mask());
assert_eq!(value["nodename"], "node-a");
assert_eq!(value["funcname"], "heal.Object");
assert_eq!(value["time"], "1970-01-01T00:00:00Z");
assert_eq!(value["path"], "bucket/object");
assert_eq!(value["bytes"], 17);
assert_eq!(value["custom"]["dry"], "true");
}
#[test]
fn trace_wire_record_formats_epoch_time() {
let event = TraceEvent {
time: UNIX_EPOCH,
..TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerFolder)
};
let record = TraceWireRecord::from_event("node-a", &event);
assert_eq!(record.time, "1970-01-01T00:00:00Z");
}
}
+15 -16
View File
@@ -67,6 +67,9 @@ 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,
@@ -1031,28 +1034,24 @@ 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("x-amz-checksum-crc32"), checksum_crc32.clone())?;
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_CRC32), checksum_crc32.clone())?;
}
if let Some(checksum_crc32c) = &output.checksum_crc32c {
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-crc32c"), checksum_crc32c.clone())?;
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_CRC32C), checksum_crc32c.clone())?;
}
if let Some(checksum_crc64nvme) = &output.checksum_crc64nvme {
insert_string_header(
&mut headers,
HeaderName::from_static("x-amz-checksum-crc64nvme"),
checksum_crc64nvme.clone(),
)?;
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_CRC64NVME), checksum_crc64nvme.clone())?;
}
if let Some(checksum_sha1) = &output.checksum_sha1 {
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-sha1"), checksum_sha1.clone())?;
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_SHA1), checksum_sha1.clone())?;
}
if let Some(checksum_sha256) = &output.checksum_sha256 {
insert_string_header(&mut headers, HeaderName::from_static("x-amz-checksum-sha256"), checksum_sha256.clone())?;
insert_string_header(&mut headers, HeaderName::from_static(AMZ_CHECKSUM_SHA256), checksum_sha256.clone())?;
}
if let Some(checksum_type) = &output.checksum_type {
insert_string_header(
&mut headers,
HeaderName::from_static("x-amz-checksum-type"),
HeaderName::from_static(AMZ_CHECKSUM_TYPE),
checksum_type.as_str().to_string(),
)?;
}
@@ -1114,12 +1113,12 @@ fn clear_object_lambda_variant_headers(headers: &mut HeaderMap) {
http::header::ETAG,
http::header::LAST_MODIFIED,
http::header::EXPIRES,
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(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-tagging-count"),
HeaderName::from_static("x-amz-request-route"),
HeaderName::from_static("x-amz-request-token"),
+1
View File
@@ -2420,6 +2420,7 @@ mod tests {
_bucket: &str,
_prefix: &str,
_continuation_token: Option<&str>,
_include_lifecycle_object_info: bool,
) -> rustfs_heal::Result<(Vec<rustfs_heal::heal::storage::HealListItem>, Option<String>, bool)> {
Ok((Vec::new(), None, false))
}