feat(scanner): expose distributed metrics (#3452)

* feat(scanner): expose distributed metrics

* docs(scanner): clarify distributed metrics collection

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
Co-authored-by: houseme <housemecn@gmail.com>
This commit is contained in:
Henry Guo
2026-06-15 07:05:43 +08:00
committed by GitHub
parent e4c1ee9941
commit 2c5615d2ea
4 changed files with 834 additions and 1 deletions
+273 -1
View File
@@ -18,9 +18,11 @@ use rustfs_common::{GLOBAL_LOCAL_NODE_NAME, GLOBAL_RUSTFS_ADDR, heal_channel::Dr
use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_madmin::metrics::{
DiskIOStats, DiskMetric, LastMinute as MadminLastMinute, NetDevLine, NetMetrics, RPCMetrics, RealtimeMetrics,
ScannerCheckpointReport as MadminScannerCheckpointReport,
ScannerLifecycleTransitionSnapshot as MadminScannerLifecycleTransitionSnapshot, ScannerMetrics as MadminScannerMetrics,
ScannerPacingPressureSnapshot as MadminScannerPacingPressureSnapshot,
ScannerSourceCycleSnapshot as MadminScannerSourceCycleSnapshot, TimedAction as MadminTimedAction,
ScannerSourceCycleSnapshot as MadminScannerSourceCycleSnapshot, ScannerSourceWorkSnapshot as MadminScannerSourceWorkSnapshot,
TimedAction as MadminTimedAction,
};
use rustfs_storage_api::StorageAdminApi;
use rustfs_utils::os::get_drive_stats;
@@ -67,6 +69,8 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
current_started: metrics.current_started,
cycles_completed_at: metrics.cycles_completed_at,
ongoing_buckets: metrics.ongoing_buckets,
active_scan_paths: metrics.active_scan_paths,
oldest_active_path_age_seconds: metrics.oldest_active_path_age_seconds,
life_time_ops: metrics.life_time_ops,
life_time_ilm: metrics.life_time_ilm,
last_minute: MadminLastMinute {
@@ -102,12 +106,53 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
.collect(),
},
active_paths: metrics.active_paths,
current_scan_mode: metrics.current_scan_mode,
current_set_scan_concurrency_limit: metrics.current_set_scan_concurrency_limit,
current_set_scans_queued: metrics.current_set_scans_queued,
current_set_scans_active: metrics.current_set_scans_active,
current_disk_scan_concurrency_limit: metrics.current_disk_scan_concurrency_limit,
current_disk_bucket_scans_queued: metrics.current_disk_bucket_scans_queued,
current_disk_bucket_scans_active: metrics.current_disk_bucket_scans_active,
current_cycle_objects_scanned: metrics.current_cycle_objects_scanned,
current_cycle_directories_scanned: metrics.current_cycle_directories_scanned,
current_cycle_bucket_drive_scans: metrics.current_cycle_bucket_drive_scans,
current_cycle_bucket_drive_failures: metrics.current_cycle_bucket_drive_failures,
current_cycle_yield_events: metrics.current_cycle_yield_events,
current_cycle_yield_duration_seconds: metrics.current_cycle_yield_duration_seconds,
current_cycle_throttle_sleep_events: metrics.current_cycle_throttle_sleep_events,
current_cycle_throttle_sleep_duration_seconds: metrics.current_cycle_throttle_sleep_duration_seconds,
current_cycle_ilm_actions: metrics.current_cycle_ilm_actions,
current_cycle_lifecycle_expiry_actions: metrics.current_cycle_lifecycle_expiry_actions,
current_cycle_lifecycle_transition_actions: metrics.current_cycle_lifecycle_transition_actions,
current_cycle_heal_objects: metrics.current_cycle_heal_objects,
current_cycle_replication_checks: metrics.current_cycle_replication_checks,
current_cycle_usage_saves: metrics.current_cycle_usage_saves,
last_cycle_result: metrics.last_cycle_result,
last_cycle_result_code: metrics.last_cycle_result_code,
last_cycle_partial_reason: metrics.last_cycle_partial_reason,
last_cycle_partial_reason_code: metrics.last_cycle_partial_reason_code,
last_cycle_lifecycle_expiry_actions: metrics.last_cycle_lifecycle_expiry_actions,
last_cycle_lifecycle_transition_actions: metrics.last_cycle_lifecycle_transition_actions,
last_cycle_partial_source: metrics.last_cycle_partial_source,
last_cycle_partial_source_code: metrics.last_cycle_partial_source_code,
last_cycle_duration_seconds: metrics.last_cycle_duration_seconds,
last_cycle_objects_scanned: metrics.last_cycle_objects_scanned,
last_cycle_directories_scanned: metrics.last_cycle_directories_scanned,
last_cycle_bucket_drive_scans: metrics.last_cycle_bucket_drive_scans,
last_cycle_bucket_drive_failures: metrics.last_cycle_bucket_drive_failures,
last_cycle_yield_events: metrics.last_cycle_yield_events,
last_cycle_yield_duration_seconds: metrics.last_cycle_yield_duration_seconds,
last_cycle_throttle_sleep_events: metrics.last_cycle_throttle_sleep_events,
last_cycle_throttle_sleep_duration_seconds: metrics.last_cycle_throttle_sleep_duration_seconds,
last_cycle_ilm_actions: metrics.last_cycle_ilm_actions,
last_cycle_heal_objects: metrics.last_cycle_heal_objects,
last_cycle_replication_checks: metrics.last_cycle_replication_checks,
last_cycle_usage_saves: metrics.last_cycle_usage_saves,
failed_cycles: metrics.failed_cycles,
partial_cycles_unknown: metrics.partial_cycles_unknown,
partial_cycles_runtime: metrics.partial_cycles_runtime,
partial_cycles_objects: metrics.partial_cycles_objects,
partial_cycles_directories: metrics.partial_cycles_directories,
pacing_pressure: MadminScannerPacingPressureSnapshot {
primary_pressure: metrics.pacing_pressure.primary_pressure,
current_queued_scans: metrics.pacing_pressure.current_queued_scans,
@@ -140,6 +185,66 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
cycles: source.cycles,
})
.collect(),
throttle_idle_mode_enabled: metrics.throttle_idle_mode_enabled,
throttle_sleep_factor: metrics.throttle_sleep_factor,
throttle_max_sleep_seconds: metrics.throttle_max_sleep_seconds,
yield_every_n_objects: metrics.yield_every_n_objects,
cycle_interval_seconds: metrics.cycle_interval_seconds,
cycle_max_duration_seconds: metrics.cycle_max_duration_seconds,
cycle_max_objects: metrics.cycle_max_objects,
cycle_max_directories: metrics.cycle_max_directories,
bitrot_cycle_enabled: metrics.bitrot_cycle_enabled,
bitrot_cycle_seconds: metrics.bitrot_cycle_seconds,
scan_checkpoint: metrics.scan_checkpoint.map(|checkpoint| MadminScannerCheckpointReport {
version: checkpoint.version,
resume_after: checkpoint.resume_after,
reason: checkpoint.reason,
last_event: checkpoint.last_event,
}),
scan_checkpoint_used: metrics.scan_checkpoint_used,
scan_checkpoint_cleared: metrics.scan_checkpoint_cleared,
scan_checkpoint_ignored: metrics.scan_checkpoint_ignored,
scan_checkpoint_stale: metrics.scan_checkpoint_stale,
source_work: metrics
.source_work
.into_iter()
.map(|work| MadminScannerSourceWorkSnapshot {
source: work.source,
checked: work.checked,
queued: work.queued,
executed: work.executed,
failed: work.failed,
skipped: work.skipped,
missed: work.missed,
})
.collect(),
current_cycle_source_work: metrics
.current_cycle_source_work
.into_iter()
.map(|work| MadminScannerSourceWorkSnapshot {
source: work.source,
checked: work.checked,
queued: work.queued,
executed: work.executed,
failed: work.failed,
skipped: work.skipped,
missed: work.missed,
})
.collect(),
last_cycle_source_work: metrics
.last_cycle_source_work
.into_iter()
.map(|work| MadminScannerSourceWorkSnapshot {
source: work.source,
checked: work.checked,
queued: work.queued,
executed: work.executed,
failed: work.failed,
skipped: work.skipped,
missed: work.missed,
})
.collect(),
partial_cycles: metrics.partial_cycles,
}
}
@@ -474,4 +579,171 @@ mod test {
assert_eq!(scanner.last_cycle_lifecycle_expiry_actions, 5);
assert_eq!(scanner.last_cycle_lifecycle_transition_actions, 7);
}
#[test]
fn scanner_metrics_mapping_preserves_distributed_status_fields() {
let scanner = to_madmin_scanner_metrics(rustfs_common::metrics::ScannerMetricsReport {
active_scan_paths: 2,
oldest_active_path_age_seconds: 45,
active_paths: vec!["disk-a/bucket-a".to_string(), "disk-b/bucket-b".to_string()],
current_scan_mode: "normal".to_string(),
current_set_scan_concurrency_limit: 3,
current_set_scans_queued: 4,
current_set_scans_active: 5,
current_disk_scan_concurrency_limit: 6,
current_disk_bucket_scans_queued: 7,
current_disk_bucket_scans_active: 8,
current_cycle_objects_scanned: 9,
current_cycle_directories_scanned: 10,
current_cycle_bucket_drive_scans: 11,
current_cycle_bucket_drive_failures: 12,
current_cycle_yield_events: 13,
current_cycle_yield_duration_seconds: 1.5,
current_cycle_throttle_sleep_events: 14,
current_cycle_throttle_sleep_duration_seconds: 2.5,
current_cycle_ilm_actions: 15,
current_cycle_heal_objects: 16,
current_cycle_replication_checks: 17,
current_cycle_usage_saves: 18,
last_cycle_result: "partial".to_string(),
last_cycle_result_code: 2,
last_cycle_partial_reason: "objects".to_string(),
last_cycle_partial_reason_code: 3,
last_cycle_duration_seconds: 4.5,
last_cycle_objects_scanned: 19,
last_cycle_directories_scanned: 20,
last_cycle_bucket_drive_scans: 21,
last_cycle_bucket_drive_failures: 22,
last_cycle_yield_events: 23,
last_cycle_yield_duration_seconds: 5.5,
last_cycle_throttle_sleep_events: 24,
last_cycle_throttle_sleep_duration_seconds: 6.5,
last_cycle_ilm_actions: 25,
last_cycle_heal_objects: 26,
last_cycle_replication_checks: 27,
last_cycle_usage_saves: 28,
failed_cycles: 29,
partial_cycles_unknown: 30,
partial_cycles_runtime: 31,
partial_cycles_objects: 32,
partial_cycles_directories: 33,
throttle_idle_mode_enabled: true,
throttle_sleep_factor: 2.0,
throttle_max_sleep_seconds: 3.0,
yield_every_n_objects: 34,
cycle_interval_seconds: 35.0,
cycle_max_duration_seconds: 36.0,
cycle_max_objects: 37,
cycle_max_directories: 38,
bitrot_cycle_enabled: true,
bitrot_cycle_seconds: 39.0,
scan_checkpoint: Some(rustfs_common::metrics::ScannerCheckpointReport {
version: 1,
resume_after: "bucket-a/prefix-a".to_string(),
reason: "directories".to_string(),
last_event: "used".to_string(),
}),
scan_checkpoint_used: 40,
scan_checkpoint_cleared: 41,
scan_checkpoint_ignored: 42,
scan_checkpoint_stale: 43,
source_work: vec![rustfs_common::metrics::ScannerSourceWorkSnapshot {
source: "usage".to_string(),
checked: 44,
queued: 45,
executed: 46,
failed: 47,
skipped: 48,
missed: 49,
}],
current_cycle_source_work: vec![rustfs_common::metrics::ScannerSourceWorkSnapshot {
source: "lifecycle".to_string(),
checked: 50,
queued: 51,
executed: 52,
failed: 53,
skipped: 54,
missed: 55,
}],
last_cycle_source_work: vec![rustfs_common::metrics::ScannerSourceWorkSnapshot {
source: "heal".to_string(),
checked: 56,
queued: 57,
executed: 58,
failed: 59,
skipped: 60,
missed: 61,
}],
partial_cycles: 62,
..Default::default()
});
assert_eq!(scanner.active_scan_paths, 2);
assert_eq!(scanner.oldest_active_path_age_seconds, 45);
assert_eq!(scanner.active_paths, vec!["disk-a/bucket-a".to_string(), "disk-b/bucket-b".to_string()]);
assert_eq!(scanner.current_scan_mode, "normal");
assert_eq!(scanner.current_set_scan_concurrency_limit, 3);
assert_eq!(scanner.current_set_scans_queued, 4);
assert_eq!(scanner.current_set_scans_active, 5);
assert_eq!(scanner.current_disk_scan_concurrency_limit, 6);
assert_eq!(scanner.current_disk_bucket_scans_queued, 7);
assert_eq!(scanner.current_disk_bucket_scans_active, 8);
assert_eq!(scanner.current_cycle_objects_scanned, 9);
assert_eq!(scanner.current_cycle_directories_scanned, 10);
assert_eq!(scanner.current_cycle_bucket_drive_scans, 11);
assert_eq!(scanner.current_cycle_bucket_drive_failures, 12);
assert_eq!(scanner.current_cycle_yield_events, 13);
assert_eq!(scanner.current_cycle_yield_duration_seconds, 1.5);
assert_eq!(scanner.current_cycle_throttle_sleep_events, 14);
assert_eq!(scanner.current_cycle_throttle_sleep_duration_seconds, 2.5);
assert_eq!(scanner.current_cycle_ilm_actions, 15);
assert_eq!(scanner.current_cycle_heal_objects, 16);
assert_eq!(scanner.current_cycle_replication_checks, 17);
assert_eq!(scanner.current_cycle_usage_saves, 18);
assert_eq!(scanner.last_cycle_result, "partial");
assert_eq!(scanner.last_cycle_result_code, 2);
assert_eq!(scanner.last_cycle_partial_reason, "objects");
assert_eq!(scanner.last_cycle_partial_reason_code, 3);
assert_eq!(scanner.last_cycle_duration_seconds, 4.5);
assert_eq!(scanner.last_cycle_objects_scanned, 19);
assert_eq!(scanner.last_cycle_directories_scanned, 20);
assert_eq!(scanner.last_cycle_bucket_drive_scans, 21);
assert_eq!(scanner.last_cycle_bucket_drive_failures, 22);
assert_eq!(scanner.last_cycle_yield_events, 23);
assert_eq!(scanner.last_cycle_yield_duration_seconds, 5.5);
assert_eq!(scanner.last_cycle_throttle_sleep_events, 24);
assert_eq!(scanner.last_cycle_throttle_sleep_duration_seconds, 6.5);
assert_eq!(scanner.last_cycle_ilm_actions, 25);
assert_eq!(scanner.last_cycle_heal_objects, 26);
assert_eq!(scanner.last_cycle_replication_checks, 27);
assert_eq!(scanner.last_cycle_usage_saves, 28);
assert_eq!(scanner.failed_cycles, 29);
assert_eq!(scanner.partial_cycles_unknown, 30);
assert_eq!(scanner.partial_cycles_runtime, 31);
assert_eq!(scanner.partial_cycles_objects, 32);
assert_eq!(scanner.partial_cycles_directories, 33);
assert!(scanner.throttle_idle_mode_enabled);
assert_eq!(scanner.throttle_sleep_factor, 2.0);
assert_eq!(scanner.throttle_max_sleep_seconds, 3.0);
assert_eq!(scanner.yield_every_n_objects, 34);
assert_eq!(scanner.cycle_interval_seconds, 35.0);
assert_eq!(scanner.cycle_max_duration_seconds, 36.0);
assert_eq!(scanner.cycle_max_objects, 37);
assert_eq!(scanner.cycle_max_directories, 38);
assert!(scanner.bitrot_cycle_enabled);
assert_eq!(scanner.bitrot_cycle_seconds, 39.0);
let checkpoint = scanner.scan_checkpoint.expect("checkpoint should be mapped");
assert_eq!(checkpoint.resume_after, "bucket-a/prefix-a");
assert_eq!(scanner.scan_checkpoint_used, 40);
assert_eq!(scanner.scan_checkpoint_cleared, 41);
assert_eq!(scanner.scan_checkpoint_ignored, 42);
assert_eq!(scanner.scan_checkpoint_stale, 43);
assert_eq!(scanner.source_work[0].source, "usage");
assert_eq!(scanner.source_work[0].missed, 49);
assert_eq!(scanner.current_cycle_source_work[0].source, "lifecycle");
assert_eq!(scanner.current_cycle_source_work[0].queued, 51);
assert_eq!(scanner.last_cycle_source_work[0].source, "heal");
assert_eq!(scanner.last_cycle_source_work[0].skipped, 60);
assert_eq!(scanner.partial_cycles, 62);
}
}
+511
View File
@@ -128,6 +128,58 @@ pub struct ScannerSourceCycleSnapshot {
pub cycles: u64,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct ScannerCheckpointReport {
#[serde(rename = "version", default)]
pub version: u16,
#[serde(rename = "resume_after", default)]
pub resume_after: String,
#[serde(rename = "reason", default)]
pub reason: String,
#[serde(rename = "last_event", default)]
pub last_event: String,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct ScannerSourceWorkSnapshot {
#[serde(rename = "source", default)]
pub source: String,
#[serde(rename = "checked", default)]
pub checked: u64,
#[serde(rename = "queued", default)]
pub queued: u64,
#[serde(rename = "executed", default)]
pub executed: u64,
#[serde(rename = "failed", default)]
pub failed: u64,
#[serde(rename = "skipped", default)]
pub skipped: u64,
#[serde(rename = "missed", default)]
pub missed: u64,
}
impl ScannerSourceWorkSnapshot {
fn merge(&mut self, other: &Self) {
self.checked = self.checked.saturating_add(other.checked);
self.queued = self.queued.saturating_add(other.queued);
self.executed = self.executed.saturating_add(other.executed);
self.failed = self.failed.saturating_add(other.failed);
self.skipped = self.skipped.saturating_add(other.skipped);
self.missed = self.missed.saturating_add(other.missed);
}
}
fn merge_source_work_snapshots(target: &mut Vec<ScannerSourceWorkSnapshot>, source: &[ScannerSourceWorkSnapshot]) {
for work in source {
if let Some(existing) = target.iter_mut().find(|existing| existing.source == work.source) {
existing.merge(work);
} else {
target.push(work.clone());
}
}
target.sort_by(|left, right| left.source.cmp(&right.source));
}
const SCANNER_PRIMARY_PRESSURE_NONE: &str = "none";
fn default_scanner_primary_pressure() -> String {
@@ -260,6 +312,10 @@ pub struct ScannerMetrics {
pub cycles_completed_at: Vec<DateTime<Utc>>,
#[serde(rename = "ongoing_buckets")]
pub ongoing_buckets: usize,
#[serde(rename = "active_scan_paths", default)]
pub active_scan_paths: usize,
#[serde(rename = "oldest_active_path_age_seconds", default)]
pub oldest_active_path_age_seconds: u64,
#[serde(rename = "life_time_ops")]
pub life_time_ops: HashMap<String, u64>,
#[serde(rename = "ilm_ops")]
@@ -268,10 +324,56 @@ pub struct ScannerMetrics {
pub last_minute: LastMinute,
#[serde(rename = "active")]
pub active_paths: Vec<String>,
#[serde(rename = "current_scan_mode", default)]
pub current_scan_mode: String,
#[serde(rename = "current_set_scan_concurrency_limit", default)]
pub current_set_scan_concurrency_limit: u64,
#[serde(rename = "current_set_scans_queued", default)]
pub current_set_scans_queued: u64,
#[serde(rename = "current_set_scans_active", default)]
pub current_set_scans_active: u64,
#[serde(rename = "current_disk_scan_concurrency_limit", default)]
pub current_disk_scan_concurrency_limit: u64,
#[serde(rename = "current_disk_bucket_scans_queued", default)]
pub current_disk_bucket_scans_queued: u64,
#[serde(rename = "current_disk_bucket_scans_active", default)]
pub current_disk_bucket_scans_active: u64,
#[serde(rename = "current_cycle_objects_scanned", default)]
pub current_cycle_objects_scanned: u64,
#[serde(rename = "current_cycle_directories_scanned", default)]
pub current_cycle_directories_scanned: u64,
#[serde(rename = "current_cycle_bucket_drive_scans", default)]
pub current_cycle_bucket_drive_scans: u64,
#[serde(rename = "current_cycle_bucket_drive_failures", default)]
pub current_cycle_bucket_drive_failures: u64,
#[serde(rename = "current_cycle_yield_events", default)]
pub current_cycle_yield_events: u64,
#[serde(rename = "current_cycle_yield_duration_seconds", default)]
pub current_cycle_yield_duration_seconds: f64,
#[serde(rename = "current_cycle_throttle_sleep_events", default)]
pub current_cycle_throttle_sleep_events: u64,
#[serde(rename = "current_cycle_throttle_sleep_duration_seconds", default)]
pub current_cycle_throttle_sleep_duration_seconds: f64,
#[serde(rename = "current_cycle_ilm_actions", default)]
pub current_cycle_ilm_actions: u64,
#[serde(rename = "current_cycle_lifecycle_expiry_actions", default)]
pub current_cycle_lifecycle_expiry_actions: u64,
#[serde(rename = "current_cycle_lifecycle_transition_actions", default)]
pub current_cycle_lifecycle_transition_actions: u64,
#[serde(rename = "current_cycle_heal_objects", default)]
pub current_cycle_heal_objects: u64,
#[serde(rename = "current_cycle_replication_checks", default)]
pub current_cycle_replication_checks: u64,
#[serde(rename = "current_cycle_usage_saves", default)]
pub current_cycle_usage_saves: u64,
#[serde(rename = "last_cycle_result", default)]
pub last_cycle_result: String,
#[serde(rename = "last_cycle_result_code", default)]
pub last_cycle_result_code: u64,
#[serde(rename = "last_cycle_partial_reason", default)]
pub last_cycle_partial_reason: String,
#[serde(rename = "last_cycle_partial_reason_code", default)]
pub last_cycle_partial_reason_code: u64,
#[serde(rename = "last_cycle_lifecycle_expiry_actions", default)]
pub last_cycle_lifecycle_expiry_actions: u64,
#[serde(rename = "last_cycle_lifecycle_transition_actions", default)]
@@ -280,12 +382,86 @@ pub struct ScannerMetrics {
pub last_cycle_partial_source: String,
#[serde(rename = "last_cycle_partial_source_code", default)]
pub last_cycle_partial_source_code: u64,
#[serde(rename = "last_cycle_duration_seconds", default)]
pub last_cycle_duration_seconds: f64,
#[serde(rename = "last_cycle_objects_scanned", default)]
pub last_cycle_objects_scanned: u64,
#[serde(rename = "last_cycle_directories_scanned", default)]
pub last_cycle_directories_scanned: u64,
#[serde(rename = "last_cycle_bucket_drive_scans", default)]
pub last_cycle_bucket_drive_scans: u64,
#[serde(rename = "last_cycle_bucket_drive_failures", default)]
pub last_cycle_bucket_drive_failures: u64,
#[serde(rename = "last_cycle_yield_events", default)]
pub last_cycle_yield_events: u64,
#[serde(rename = "last_cycle_yield_duration_seconds", default)]
pub last_cycle_yield_duration_seconds: f64,
#[serde(rename = "last_cycle_throttle_sleep_events", default)]
pub last_cycle_throttle_sleep_events: u64,
#[serde(rename = "last_cycle_throttle_sleep_duration_seconds", default)]
pub last_cycle_throttle_sleep_duration_seconds: f64,
#[serde(rename = "last_cycle_ilm_actions", default)]
pub last_cycle_ilm_actions: u64,
#[serde(rename = "last_cycle_heal_objects", default)]
pub last_cycle_heal_objects: u64,
#[serde(rename = "last_cycle_replication_checks", default)]
pub last_cycle_replication_checks: u64,
#[serde(rename = "last_cycle_usage_saves", default)]
pub last_cycle_usage_saves: u64,
#[serde(rename = "failed_cycles", default)]
pub failed_cycles: u64,
#[serde(rename = "partial_cycles_unknown", default)]
pub partial_cycles_unknown: u64,
#[serde(rename = "partial_cycles_runtime", default)]
pub partial_cycles_runtime: u64,
#[serde(rename = "partial_cycles_objects", default)]
pub partial_cycles_objects: u64,
#[serde(rename = "partial_cycles_directories", default)]
pub partial_cycles_directories: u64,
#[serde(rename = "partial_cycles_by_source", default)]
pub partial_cycles_by_source: Vec<ScannerSourceCycleSnapshot>,
#[serde(rename = "pacing_pressure", default)]
pub pacing_pressure: ScannerPacingPressureSnapshot,
#[serde(rename = "lifecycle_transition", default)]
pub lifecycle_transition: ScannerLifecycleTransitionSnapshot,
#[serde(rename = "throttle_idle_mode_enabled", default)]
pub throttle_idle_mode_enabled: bool,
#[serde(rename = "throttle_sleep_factor", default)]
pub throttle_sleep_factor: f64,
#[serde(rename = "throttle_max_sleep_seconds", default)]
pub throttle_max_sleep_seconds: f64,
#[serde(rename = "yield_every_n_objects", default)]
pub yield_every_n_objects: u64,
#[serde(rename = "cycle_interval_seconds", default)]
pub cycle_interval_seconds: f64,
#[serde(rename = "cycle_max_duration_seconds", default)]
pub cycle_max_duration_seconds: f64,
#[serde(rename = "cycle_max_objects", default)]
pub cycle_max_objects: u64,
#[serde(rename = "cycle_max_directories", default)]
pub cycle_max_directories: u64,
#[serde(rename = "bitrot_cycle_enabled", default)]
pub bitrot_cycle_enabled: bool,
#[serde(rename = "bitrot_cycle_seconds", default)]
pub bitrot_cycle_seconds: f64,
#[serde(rename = "scan_checkpoint", default, skip_serializing_if = "Option::is_none")]
pub scan_checkpoint: Option<ScannerCheckpointReport>,
#[serde(rename = "scan_checkpoint_used", default)]
pub scan_checkpoint_used: u64,
#[serde(rename = "scan_checkpoint_cleared", default)]
pub scan_checkpoint_cleared: u64,
#[serde(rename = "scan_checkpoint_ignored", default)]
pub scan_checkpoint_ignored: u64,
#[serde(rename = "scan_checkpoint_stale", default)]
pub scan_checkpoint_stale: u64,
#[serde(rename = "source_work", default)]
pub source_work: Vec<ScannerSourceWorkSnapshot>,
#[serde(rename = "current_cycle_source_work", default)]
pub current_cycle_source_work: Vec<ScannerSourceWorkSnapshot>,
#[serde(rename = "last_cycle_source_work", default)]
pub last_cycle_source_work: Vec<ScannerSourceWorkSnapshot>,
#[serde(rename = "partial_cycles", default)]
pub partial_cycles: u64,
}
impl ScannerMetrics {
@@ -293,24 +469,126 @@ impl ScannerMetrics {
let other_is_newer = self.collected_at < other.collected_at;
if other_is_newer {
self.collected_at = other.collected_at;
self.current_scan_mode = other.current_scan_mode.clone();
self.last_cycle_result = other.last_cycle_result.clone();
self.last_cycle_result_code = other.last_cycle_result_code;
self.last_cycle_partial_reason = other.last_cycle_partial_reason.clone();
self.last_cycle_partial_reason_code = other.last_cycle_partial_reason_code;
self.last_cycle_partial_source = other.last_cycle_partial_source.clone();
self.last_cycle_partial_source_code = other.last_cycle_partial_source_code;
self.last_cycle_duration_seconds = other.last_cycle_duration_seconds;
self.throttle_idle_mode_enabled = other.throttle_idle_mode_enabled;
self.throttle_sleep_factor = other.throttle_sleep_factor;
self.throttle_max_sleep_seconds = other.throttle_max_sleep_seconds;
self.yield_every_n_objects = other.yield_every_n_objects;
self.cycle_interval_seconds = other.cycle_interval_seconds;
self.cycle_max_duration_seconds = other.cycle_max_duration_seconds;
self.cycle_max_objects = other.cycle_max_objects;
self.cycle_max_directories = other.cycle_max_directories;
self.bitrot_cycle_enabled = other.bitrot_cycle_enabled;
self.bitrot_cycle_seconds = other.bitrot_cycle_seconds;
}
if other.scan_checkpoint.is_some() && (other_is_newer || self.scan_checkpoint.is_none()) {
self.scan_checkpoint = other.scan_checkpoint.clone();
}
self.active_scan_paths = self.active_scan_paths.saturating_add(other.active_scan_paths);
self.oldest_active_path_age_seconds = self.oldest_active_path_age_seconds.max(other.oldest_active_path_age_seconds);
self.pacing_pressure.merge(&other.pacing_pressure);
self.lifecycle_transition.merge(&other.lifecycle_transition);
self.current_set_scan_concurrency_limit = self
.current_set_scan_concurrency_limit
.saturating_add(other.current_set_scan_concurrency_limit);
self.current_set_scans_queued = self.current_set_scans_queued.saturating_add(other.current_set_scans_queued);
self.current_set_scans_active = self.current_set_scans_active.saturating_add(other.current_set_scans_active);
self.current_disk_scan_concurrency_limit = self
.current_disk_scan_concurrency_limit
.saturating_add(other.current_disk_scan_concurrency_limit);
self.current_disk_bucket_scans_queued = self
.current_disk_bucket_scans_queued
.saturating_add(other.current_disk_bucket_scans_queued);
self.current_disk_bucket_scans_active = self
.current_disk_bucket_scans_active
.saturating_add(other.current_disk_bucket_scans_active);
self.current_cycle_objects_scanned = self
.current_cycle_objects_scanned
.saturating_add(other.current_cycle_objects_scanned);
self.current_cycle_directories_scanned = self
.current_cycle_directories_scanned
.saturating_add(other.current_cycle_directories_scanned);
self.current_cycle_bucket_drive_scans = self
.current_cycle_bucket_drive_scans
.saturating_add(other.current_cycle_bucket_drive_scans);
self.current_cycle_bucket_drive_failures = self
.current_cycle_bucket_drive_failures
.saturating_add(other.current_cycle_bucket_drive_failures);
self.current_cycle_yield_events = self
.current_cycle_yield_events
.saturating_add(other.current_cycle_yield_events);
self.current_cycle_yield_duration_seconds += other.current_cycle_yield_duration_seconds;
self.current_cycle_throttle_sleep_events = self
.current_cycle_throttle_sleep_events
.saturating_add(other.current_cycle_throttle_sleep_events);
self.current_cycle_throttle_sleep_duration_seconds += other.current_cycle_throttle_sleep_duration_seconds;
self.current_cycle_ilm_actions = self.current_cycle_ilm_actions.saturating_add(other.current_cycle_ilm_actions);
self.current_cycle_lifecycle_expiry_actions = self
.current_cycle_lifecycle_expiry_actions
.saturating_add(other.current_cycle_lifecycle_expiry_actions);
self.current_cycle_lifecycle_transition_actions = self
.current_cycle_lifecycle_transition_actions
.saturating_add(other.current_cycle_lifecycle_transition_actions);
self.current_cycle_heal_objects = self
.current_cycle_heal_objects
.saturating_add(other.current_cycle_heal_objects);
self.current_cycle_replication_checks = self
.current_cycle_replication_checks
.saturating_add(other.current_cycle_replication_checks);
self.current_cycle_usage_saves = self.current_cycle_usage_saves.saturating_add(other.current_cycle_usage_saves);
self.last_cycle_objects_scanned = self
.last_cycle_objects_scanned
.saturating_add(other.last_cycle_objects_scanned);
self.last_cycle_directories_scanned = self
.last_cycle_directories_scanned
.saturating_add(other.last_cycle_directories_scanned);
self.last_cycle_bucket_drive_scans = self
.last_cycle_bucket_drive_scans
.saturating_add(other.last_cycle_bucket_drive_scans);
self.last_cycle_bucket_drive_failures = self
.last_cycle_bucket_drive_failures
.saturating_add(other.last_cycle_bucket_drive_failures);
self.last_cycle_yield_events = self.last_cycle_yield_events.saturating_add(other.last_cycle_yield_events);
self.last_cycle_yield_duration_seconds += other.last_cycle_yield_duration_seconds;
self.last_cycle_throttle_sleep_events = self
.last_cycle_throttle_sleep_events
.saturating_add(other.last_cycle_throttle_sleep_events);
self.last_cycle_throttle_sleep_duration_seconds += other.last_cycle_throttle_sleep_duration_seconds;
self.last_cycle_ilm_actions = self.last_cycle_ilm_actions.saturating_add(other.last_cycle_ilm_actions);
self.last_cycle_lifecycle_expiry_actions = self
.last_cycle_lifecycle_expiry_actions
.saturating_add(other.last_cycle_lifecycle_expiry_actions);
self.last_cycle_lifecycle_transition_actions = self
.last_cycle_lifecycle_transition_actions
.saturating_add(other.last_cycle_lifecycle_transition_actions);
self.last_cycle_heal_objects = self.last_cycle_heal_objects.saturating_add(other.last_cycle_heal_objects);
self.last_cycle_replication_checks = self
.last_cycle_replication_checks
.saturating_add(other.last_cycle_replication_checks);
self.last_cycle_usage_saves = self.last_cycle_usage_saves.saturating_add(other.last_cycle_usage_saves);
self.failed_cycles = self.failed_cycles.saturating_add(other.failed_cycles);
self.partial_cycles_unknown = self.partial_cycles_unknown.saturating_add(other.partial_cycles_unknown);
self.partial_cycles_runtime = self.partial_cycles_runtime.saturating_add(other.partial_cycles_runtime);
self.partial_cycles_objects = self.partial_cycles_objects.saturating_add(other.partial_cycles_objects);
self.partial_cycles_directories = self
.partial_cycles_directories
.saturating_add(other.partial_cycles_directories);
self.scan_checkpoint_used = self.scan_checkpoint_used.saturating_add(other.scan_checkpoint_used);
self.scan_checkpoint_cleared = self.scan_checkpoint_cleared.saturating_add(other.scan_checkpoint_cleared);
self.scan_checkpoint_ignored = self.scan_checkpoint_ignored.saturating_add(other.scan_checkpoint_ignored);
self.scan_checkpoint_stale = self.scan_checkpoint_stale.saturating_add(other.scan_checkpoint_stale);
self.partial_cycles = self.partial_cycles.saturating_add(other.partial_cycles);
merge_source_work_snapshots(&mut self.source_work, &other.source_work);
merge_source_work_snapshots(&mut self.current_cycle_source_work, &other.current_cycle_source_work);
merge_source_work_snapshots(&mut self.last_cycle_source_work, &other.last_cycle_source_work);
if self.ongoing_buckets < other.ongoing_buckets {
self.ongoing_buckets = other.ongoing_buckets;
@@ -1028,4 +1306,237 @@ mod tests {
assert_eq!(scanner.last_cycle_lifecycle_expiry_actions, 22);
assert_eq!(scanner.last_cycle_lifecycle_transition_actions, 26);
}
#[test]
fn scanner_metrics_merge_preserves_distributed_status_fields() {
let collected_at = Utc::now();
let mut scanner = ScannerMetrics {
collected_at,
active_scan_paths: 1,
oldest_active_path_age_seconds: 15,
active_paths: vec!["node-a/disk-a/bucket-a".to_string()],
current_set_scan_concurrency_limit: 1,
current_set_scans_queued: 2,
current_set_scans_active: 3,
current_disk_scan_concurrency_limit: 4,
current_disk_bucket_scans_queued: 5,
current_disk_bucket_scans_active: 6,
current_cycle_objects_scanned: 7,
current_cycle_directories_scanned: 8,
current_cycle_bucket_drive_scans: 9,
current_cycle_bucket_drive_failures: 1,
current_cycle_yield_events: 2,
current_cycle_yield_duration_seconds: 0.5,
current_cycle_throttle_sleep_events: 3,
current_cycle_throttle_sleep_duration_seconds: 1.5,
current_cycle_ilm_actions: 4,
current_cycle_heal_objects: 5,
current_cycle_replication_checks: 6,
current_cycle_usage_saves: 7,
last_cycle_objects_scanned: 8,
last_cycle_directories_scanned: 9,
last_cycle_bucket_drive_scans: 10,
last_cycle_bucket_drive_failures: 2,
last_cycle_yield_events: 3,
last_cycle_yield_duration_seconds: 0.75,
last_cycle_throttle_sleep_events: 4,
last_cycle_throttle_sleep_duration_seconds: 1.75,
last_cycle_ilm_actions: 5,
last_cycle_heal_objects: 6,
last_cycle_replication_checks: 7,
last_cycle_usage_saves: 8,
failed_cycles: 1,
partial_cycles_unknown: 2,
partial_cycles_runtime: 3,
partial_cycles_objects: 4,
partial_cycles_directories: 5,
scan_checkpoint: Some(ScannerCheckpointReport {
version: 1,
resume_after: "bucket-a/prefix-a".to_string(),
reason: "directories".to_string(),
last_event: "used".to_string(),
}),
scan_checkpoint_used: 1,
scan_checkpoint_cleared: 2,
scan_checkpoint_ignored: 3,
scan_checkpoint_stale: 4,
source_work: vec![ScannerSourceWorkSnapshot {
source: "usage".to_string(),
checked: 10,
executed: 1,
..Default::default()
}],
current_cycle_source_work: vec![ScannerSourceWorkSnapshot {
source: "lifecycle".to_string(),
queued: 2,
missed: 1,
..Default::default()
}],
last_cycle_source_work: vec![ScannerSourceWorkSnapshot {
source: "heal".to_string(),
skipped: 3,
..Default::default()
}],
partial_cycles: 6,
..Default::default()
};
scanner.merge(&ScannerMetrics {
collected_at: collected_at + chrono::Duration::seconds(1),
active_scan_paths: 2,
oldest_active_path_age_seconds: 45,
active_paths: vec!["node-b/disk-b/bucket-b".to_string()],
current_set_scan_concurrency_limit: 10,
current_set_scans_queued: 20,
current_set_scans_active: 30,
current_disk_scan_concurrency_limit: 40,
current_disk_bucket_scans_queued: 50,
current_disk_bucket_scans_active: 60,
current_cycle_objects_scanned: 70,
current_cycle_directories_scanned: 80,
current_cycle_bucket_drive_scans: 90,
current_cycle_bucket_drive_failures: 10,
current_cycle_yield_events: 20,
current_cycle_yield_duration_seconds: 2.5,
current_cycle_throttle_sleep_events: 30,
current_cycle_throttle_sleep_duration_seconds: 3.5,
current_cycle_ilm_actions: 40,
current_cycle_heal_objects: 50,
current_cycle_replication_checks: 60,
current_cycle_usage_saves: 70,
last_cycle_objects_scanned: 80,
last_cycle_directories_scanned: 90,
last_cycle_bucket_drive_scans: 100,
last_cycle_bucket_drive_failures: 20,
last_cycle_yield_events: 30,
last_cycle_yield_duration_seconds: 2.75,
last_cycle_throttle_sleep_events: 40,
last_cycle_throttle_sleep_duration_seconds: 3.75,
last_cycle_ilm_actions: 50,
last_cycle_heal_objects: 60,
last_cycle_replication_checks: 70,
last_cycle_usage_saves: 80,
failed_cycles: 10,
partial_cycles_unknown: 20,
partial_cycles_runtime: 30,
partial_cycles_objects: 40,
partial_cycles_directories: 50,
scan_checkpoint: Some(ScannerCheckpointReport {
version: 1,
resume_after: "bucket-b/prefix-b".to_string(),
reason: "objects".to_string(),
last_event: "set".to_string(),
}),
scan_checkpoint_used: 10,
scan_checkpoint_cleared: 20,
scan_checkpoint_ignored: 30,
scan_checkpoint_stale: 40,
source_work: vec![
ScannerSourceWorkSnapshot {
source: "usage".to_string(),
checked: 5,
executed: 2,
..Default::default()
},
ScannerSourceWorkSnapshot {
source: "replication".to_string(),
queued: 7,
missed: 2,
..Default::default()
},
],
current_cycle_source_work: vec![ScannerSourceWorkSnapshot {
source: "lifecycle".to_string(),
queued: 3,
missed: 4,
..Default::default()
}],
last_cycle_source_work: vec![ScannerSourceWorkSnapshot {
source: "heal".to_string(),
skipped: 4,
missed: 5,
..Default::default()
}],
partial_cycles: 60,
..Default::default()
});
assert_eq!(scanner.active_scan_paths, 3);
assert_eq!(scanner.oldest_active_path_age_seconds, 45);
assert_eq!(
scanner.active_paths,
vec!["node-a/disk-a/bucket-a".to_string(), "node-b/disk-b/bucket-b".to_string()]
);
assert_eq!(scanner.current_set_scan_concurrency_limit, 11);
assert_eq!(scanner.current_set_scans_queued, 22);
assert_eq!(scanner.current_set_scans_active, 33);
assert_eq!(scanner.current_disk_scan_concurrency_limit, 44);
assert_eq!(scanner.current_disk_bucket_scans_queued, 55);
assert_eq!(scanner.current_disk_bucket_scans_active, 66);
assert_eq!(scanner.current_cycle_objects_scanned, 77);
assert_eq!(scanner.current_cycle_directories_scanned, 88);
assert_eq!(scanner.current_cycle_bucket_drive_scans, 99);
assert_eq!(scanner.current_cycle_bucket_drive_failures, 11);
assert_eq!(scanner.current_cycle_yield_events, 22);
assert_eq!(scanner.current_cycle_yield_duration_seconds, 3.0);
assert_eq!(scanner.current_cycle_throttle_sleep_events, 33);
assert_eq!(scanner.current_cycle_throttle_sleep_duration_seconds, 5.0);
assert_eq!(scanner.current_cycle_ilm_actions, 44);
assert_eq!(scanner.current_cycle_heal_objects, 55);
assert_eq!(scanner.current_cycle_replication_checks, 66);
assert_eq!(scanner.current_cycle_usage_saves, 77);
assert_eq!(scanner.last_cycle_objects_scanned, 88);
assert_eq!(scanner.last_cycle_directories_scanned, 99);
assert_eq!(scanner.last_cycle_bucket_drive_scans, 110);
assert_eq!(scanner.last_cycle_bucket_drive_failures, 22);
assert_eq!(scanner.last_cycle_yield_events, 33);
assert_eq!(scanner.last_cycle_yield_duration_seconds, 3.5);
assert_eq!(scanner.last_cycle_throttle_sleep_events, 44);
assert_eq!(scanner.last_cycle_throttle_sleep_duration_seconds, 5.5);
assert_eq!(scanner.last_cycle_ilm_actions, 55);
assert_eq!(scanner.last_cycle_heal_objects, 66);
assert_eq!(scanner.last_cycle_replication_checks, 77);
assert_eq!(scanner.last_cycle_usage_saves, 88);
assert_eq!(scanner.failed_cycles, 11);
assert_eq!(scanner.partial_cycles_unknown, 22);
assert_eq!(scanner.partial_cycles_runtime, 33);
assert_eq!(scanner.partial_cycles_objects, 44);
assert_eq!(scanner.partial_cycles_directories, 55);
let checkpoint = scanner.scan_checkpoint.expect("newest checkpoint should be preserved");
assert_eq!(checkpoint.resume_after, "bucket-b/prefix-b");
assert_eq!(scanner.scan_checkpoint_used, 11);
assert_eq!(scanner.scan_checkpoint_cleared, 22);
assert_eq!(scanner.scan_checkpoint_ignored, 33);
assert_eq!(scanner.scan_checkpoint_stale, 44);
assert_eq!(scanner.partial_cycles, 66);
let usage = scanner
.source_work
.iter()
.find(|work| work.source == "usage")
.expect("usage work should be merged");
assert_eq!(usage.checked, 15);
assert_eq!(usage.executed, 3);
let replication = scanner
.source_work
.iter()
.find(|work| work.source == "replication")
.expect("replication work should be preserved");
assert_eq!(replication.queued, 7);
assert_eq!(replication.missed, 2);
let lifecycle = scanner
.current_cycle_source_work
.iter()
.find(|work| work.source == "lifecycle")
.expect("current cycle lifecycle work should be merged");
assert_eq!(lifecycle.queued, 5);
assert_eq!(lifecycle.missed, 5);
let heal = scanner
.last_cycle_source_work
.iter()
.find(|work| work.source == "heal")
.expect("last cycle heal work should be merged");
assert_eq!(heal.skipped, 7);
assert_eq!(heal.missed, 5);
}
}
@@ -156,6 +156,27 @@ scripts/run_scanner_validation_harness.sh \
The harness writes scanner/heal config snapshots, scanner status samples, host
telemetry when available, run metadata, and `scanner-summary.csv`.
For distributed runs, capture scanner admin metrics from every node with
`by-host=true`. The metrics endpoint reports the node that handles the request;
`by-host=true` preserves that node's host view but does not collect peer nodes.
These per-node artifacts include active path age, checkpoint state, pacing
pressure, source work, and queued/skipped/missed downstream admission counters.
```bash
for endpoint in http://node-a:9000 http://node-b:9000 http://node-c:9000; do
node="${endpoint#http://}"
node="${node%%:*}"
awscurl \
--service s3 \
--region us-east-1 \
--access_key "$RUSTFS_ACCESS_KEY" \
--secret_key "$RUSTFS_SECRET_KEY" \
--request GET \
"${endpoint}/rustfs/admin/v3/metrics?types=1&by-host=true&n=1" \
> "artifacts/scanner-metrics.${node}.$(date -u +%Y%m%dT%H%M%SZ).ndjson"
done
```
Example status request:
```bash
@@ -146,6 +146,35 @@ Use these counters to decide whether scan progress is limited by scanner pacing
or by a downstream subsystem such as lifecycle transition, replication repair,
or heal admission.
## Reading Distributed Metrics
`/rustfs/admin/v3/scanner/status` and `/rustfs/admin/v3/metrics` report the
node that handles the HTTP request. The metrics endpoint does not fan out to
peer nodes. In distributed deployments, query every node explicitly and keep
`by-host=true` enabled so each response includes that node's host view:
```bash
for endpoint in http://node-a:9000 http://node-b:9000 http://node-c:9000; do
node="${endpoint#http://}"
node="${node%%:*}"
awscurl \
--service s3 \
--region us-east-1 \
--access_key "$RUSTFS_ACCESS_KEY" \
--secret_key "$RUSTFS_SECRET_KEY" \
--request GET \
"${endpoint}/rustfs/admin/v3/metrics?types=1&by-host=true&n=1" \
> "artifacts/scanner-metrics.${node}.$(date -u +%Y%m%dT%H%M%SZ).ndjson"
done
```
The `aggregated.scanner` payload preserves the same scanner progress,
checkpoint, pacing, source work, and lifecycle transition fields used by the
local scanner status, but only for the node that returned the response. The
`by_host.*.scanner` payload keeps that node's host view. Compare the per-node
artifacts externally to find old active paths, partial checkpoints, pacing
pressure, or downstream queue admission problems across the deployment.
## Reading Lifecycle Transition Status
`metrics.lifecycle_transition` focuses on scanner-driven lifecycle transition