feat(scanner): expose pacing pressure status (#3319)

* feat(scanner): expose pacing pressure status

* fix(scanner): preserve merged pause pressure

* fix(scanner): default missing primary pressure

---------

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-06-10 15:33:41 +08:00
committed by GitHub
parent bb5d9565a6
commit 66fd55a8e0
3 changed files with 311 additions and 2 deletions
+131
View File
@@ -748,6 +748,33 @@ pub struct ScannerSourceCycleSnapshot {
pub cycles: u64, pub cycles: u64,
} }
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
pub struct ScannerPacingPressureSnapshot {
pub primary_pressure: String,
pub current_queued_scans: u64,
pub current_active_scans: u64,
pub last_cycle_budget_limited: bool,
pub last_cycle_pause_observed: bool,
pub last_cycle_throttle_sleep_ratio: f64,
pub last_cycle_yield_ratio: f64,
pub last_cycle_total_pause_ratio: f64,
}
impl Default for ScannerPacingPressureSnapshot {
fn default() -> Self {
Self {
primary_pressure: "none".to_string(),
current_queued_scans: 0,
current_active_scans: 0,
last_cycle_budget_limited: false,
last_cycle_pause_observed: false,
last_cycle_throttle_sleep_ratio: 0.0,
last_cycle_yield_ratio: 0.0,
last_cycle_total_pause_ratio: 0.0,
}
}
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)] #[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct ScannerLastMinute { pub struct ScannerLastMinute {
pub actions: HashMap<String, ScannerTimedAction>, pub actions: HashMap<String, ScannerTimedAction>,
@@ -853,6 +880,8 @@ pub struct ScannerMetricsReport {
#[serde(default)] #[serde(default)]
pub partial_cycles_by_source: Vec<ScannerSourceCycleSnapshot>, pub partial_cycles_by_source: Vec<ScannerSourceCycleSnapshot>,
#[serde(default)] #[serde(default)]
pub pacing_pressure: ScannerPacingPressureSnapshot,
#[serde(default)]
pub throttle_idle_mode_enabled: bool, pub throttle_idle_mode_enabled: bool,
#[serde(default)] #[serde(default)]
pub throttle_sleep_factor: f64, pub throttle_sleep_factor: f64,
@@ -933,6 +962,76 @@ fn scan_cycle_result_label(result: u8) -> &'static str {
} }
} }
fn scanner_ratio(numerator: f64, denominator: f64) -> f64 {
if denominator <= 0.0 {
0.0
} else {
(numerator / denominator).clamp(0.0, 1.0)
}
}
fn scanner_last_cycle_budget_limited(result_code: u64, partial_reason: &str) -> bool {
result_code == u64::from(SCAN_CYCLE_RESULT_PARTIAL) && matches!(partial_reason, "runtime" | "objects" | "directories")
}
fn scanner_primary_pressure(
current_queued_scans: u64,
current_active_scans: u64,
last_cycle_budget_limited: bool,
last_cycle_throttle_sleep_events: u64,
last_cycle_yield_events: u64,
) -> &'static str {
if current_queued_scans > 0 {
"queued_scans"
} else if last_cycle_budget_limited {
"cycle_budget"
} else if last_cycle_throttle_sleep_events > 0 || last_cycle_yield_events > 0 {
"throttle_pause"
} else if current_active_scans > 0 {
"active_scans"
} else {
"none"
}
}
fn scanner_pacing_pressure(metrics: &ScannerMetricsReport) -> ScannerPacingPressureSnapshot {
let current_queued_scans = metrics
.current_set_scans_queued
.saturating_add(metrics.current_disk_bucket_scans_queued);
let current_active_scans = metrics
.current_set_scans_active
.saturating_add(metrics.current_disk_bucket_scans_active)
.max(usize_to_u64_saturated(metrics.active_scan_paths));
let last_cycle_budget_limited =
scanner_last_cycle_budget_limited(metrics.last_cycle_result_code, metrics.last_cycle_partial_reason.as_str());
let last_cycle_pause_observed = metrics.last_cycle_throttle_sleep_events > 0 || metrics.last_cycle_yield_events > 0;
let last_cycle_throttle_sleep_ratio =
scanner_ratio(metrics.last_cycle_throttle_sleep_duration_seconds, metrics.last_cycle_duration_seconds);
let last_cycle_yield_ratio = scanner_ratio(metrics.last_cycle_yield_duration_seconds, metrics.last_cycle_duration_seconds);
let last_cycle_total_pause_ratio = scanner_ratio(
metrics.last_cycle_throttle_sleep_duration_seconds.max(0.0) + metrics.last_cycle_yield_duration_seconds.max(0.0),
metrics.last_cycle_duration_seconds,
);
ScannerPacingPressureSnapshot {
primary_pressure: scanner_primary_pressure(
current_queued_scans,
current_active_scans,
last_cycle_budget_limited,
metrics.last_cycle_throttle_sleep_events,
metrics.last_cycle_yield_events,
)
.to_string(),
current_queued_scans,
current_active_scans,
last_cycle_budget_limited,
last_cycle_pause_observed,
last_cycle_throttle_sleep_ratio,
last_cycle_yield_ratio,
last_cycle_total_pause_ratio,
}
}
fn duration_millis_saturated(duration: Duration) -> u64 { fn duration_millis_saturated(duration: Duration) -> u64 {
duration.as_millis().min(u64::MAX as u128) as u64 duration.as_millis().min(u64::MAX as u128) as u64
} }
@@ -1904,6 +2003,8 @@ impl Metrics {
} }
} }
m.pacing_pressure = scanner_pacing_pressure(&m);
m m
} }
} }
@@ -2060,6 +2161,36 @@ mod tests {
assert_eq!(report.current_disk_bucket_scans_active, 0); assert_eq!(report.current_disk_bucket_scans_active, 0);
} }
#[tokio::test]
async fn report_derives_scanner_pacing_pressure() {
let metrics = Metrics::new();
metrics.record_scanner_set_scan_state(Some(2), Some(3), Some(1));
metrics.record_scanner_disk_bucket_scan_state("0", "0", Some(1), Some(4), Some(1));
metrics.record_scan_cycle_partial_with_source(
Duration::from_secs(10),
ScanCyclePartialReason::Runtime,
Some(ScannerWorkSource::Usage),
);
metrics.record_scan_cycle_work(ScanCycleWorkSnapshot {
yield_events: 2,
yield_duration_millis: 500,
throttle_sleep_events: 4,
throttle_sleep_duration_millis: 2500,
..Default::default()
});
let report = metrics.report().await;
assert_eq!(report.pacing_pressure.primary_pressure, "queued_scans");
assert_eq!(report.pacing_pressure.current_queued_scans, 7);
assert_eq!(report.pacing_pressure.current_active_scans, 2);
assert!(report.pacing_pressure.last_cycle_budget_limited);
assert!(report.pacing_pressure.last_cycle_pause_observed);
assert_eq!(report.pacing_pressure.last_cycle_throttle_sleep_ratio, 0.25);
assert_eq!(report.pacing_pressure.last_cycle_yield_ratio, 0.05);
assert_eq!(report.pacing_pressure.last_cycle_total_pause_ratio, 0.3);
}
#[tokio::test] #[tokio::test]
async fn report_includes_scanner_checkpoint_status() { async fn report_includes_scanner_checkpoint_status() {
let metrics = Metrics::new(); let metrics = Metrics::new();
+38 -2
View File
@@ -18,8 +18,8 @@ use rustfs_common::{GLOBAL_LOCAL_NODE_NAME, GLOBAL_RUSTFS_ADDR, heal_channel::Dr
use rustfs_io_metrics::internode_metrics::global_internode_metrics; use rustfs_io_metrics::internode_metrics::global_internode_metrics;
use rustfs_madmin::metrics::{ use rustfs_madmin::metrics::{
DiskIOStats, DiskMetric, LastMinute as MadminLastMinute, NetDevLine, NetMetrics, RPCMetrics, RealtimeMetrics, DiskIOStats, DiskMetric, LastMinute as MadminLastMinute, NetDevLine, NetMetrics, RPCMetrics, RealtimeMetrics,
ScannerMetrics as MadminScannerMetrics, ScannerSourceCycleSnapshot as MadminScannerSourceCycleSnapshot, ScannerMetrics as MadminScannerMetrics, ScannerPacingPressureSnapshot as MadminScannerPacingPressureSnapshot,
TimedAction as MadminTimedAction, ScannerSourceCycleSnapshot as MadminScannerSourceCycleSnapshot, TimedAction as MadminTimedAction,
}; };
use rustfs_utils::os::get_drive_stats; use rustfs_utils::os::get_drive_stats;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
@@ -102,6 +102,16 @@ fn to_madmin_scanner_metrics(metrics: rustfs_common::metrics::ScannerMetricsRepo
active_paths: metrics.active_paths, active_paths: metrics.active_paths,
last_cycle_partial_source: metrics.last_cycle_partial_source, last_cycle_partial_source: metrics.last_cycle_partial_source,
last_cycle_partial_source_code: metrics.last_cycle_partial_source_code, last_cycle_partial_source_code: metrics.last_cycle_partial_source_code,
pacing_pressure: MadminScannerPacingPressureSnapshot {
primary_pressure: metrics.pacing_pressure.primary_pressure,
current_queued_scans: metrics.pacing_pressure.current_queued_scans,
current_active_scans: metrics.pacing_pressure.current_active_scans,
last_cycle_budget_limited: metrics.pacing_pressure.last_cycle_budget_limited,
last_cycle_pause_observed: metrics.pacing_pressure.last_cycle_pause_observed,
last_cycle_throttle_sleep_ratio: metrics.pacing_pressure.last_cycle_throttle_sleep_ratio,
last_cycle_yield_ratio: metrics.pacing_pressure.last_cycle_yield_ratio,
last_cycle_total_pause_ratio: metrics.pacing_pressure.last_cycle_total_pause_ratio,
},
partial_cycles_by_source: metrics partial_cycles_by_source: metrics
.partial_cycles_by_source .partial_cycles_by_source
.into_iter() .into_iter()
@@ -376,4 +386,30 @@ mod test {
.expect("usage partial source should be mapped"); .expect("usage partial source should be mapped");
assert_eq!(usage.cycles, 2); assert_eq!(usage.cycles, 2);
} }
#[test]
fn scanner_metrics_mapping_preserves_pacing_pressure() {
let scanner = to_madmin_scanner_metrics(rustfs_common::metrics::ScannerMetricsReport {
pacing_pressure: rustfs_common::metrics::ScannerPacingPressureSnapshot {
primary_pressure: "cycle_budget".to_string(),
current_queued_scans: 4,
current_active_scans: 2,
last_cycle_budget_limited: true,
last_cycle_pause_observed: true,
last_cycle_throttle_sleep_ratio: 0.25,
last_cycle_yield_ratio: 0.05,
last_cycle_total_pause_ratio: 0.3,
},
..Default::default()
});
assert_eq!(scanner.pacing_pressure.primary_pressure, "cycle_budget");
assert_eq!(scanner.pacing_pressure.current_queued_scans, 4);
assert_eq!(scanner.pacing_pressure.current_active_scans, 2);
assert!(scanner.pacing_pressure.last_cycle_budget_limited);
assert!(scanner.pacing_pressure.last_cycle_pause_observed);
assert_eq!(scanner.pacing_pressure.last_cycle_throttle_sleep_ratio, 0.25);
assert_eq!(scanner.pacing_pressure.last_cycle_yield_ratio, 0.05);
assert_eq!(scanner.pacing_pressure.last_cycle_total_pause_ratio, 0.3);
}
} }
+142
View File
@@ -128,6 +128,81 @@ pub struct ScannerSourceCycleSnapshot {
pub cycles: u64, pub cycles: u64,
} }
const SCANNER_PRIMARY_PRESSURE_NONE: &str = "none";
fn default_scanner_primary_pressure() -> String {
SCANNER_PRIMARY_PRESSURE_NONE.to_string()
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
pub struct ScannerPacingPressureSnapshot {
#[serde(rename = "primary_pressure", default = "default_scanner_primary_pressure")]
pub primary_pressure: String,
#[serde(rename = "current_queued_scans", default)]
pub current_queued_scans: u64,
#[serde(rename = "current_active_scans", default)]
pub current_active_scans: u64,
#[serde(rename = "last_cycle_budget_limited", default)]
pub last_cycle_budget_limited: bool,
#[serde(rename = "last_cycle_pause_observed", default)]
pub last_cycle_pause_observed: bool,
#[serde(rename = "last_cycle_throttle_sleep_ratio", default)]
pub last_cycle_throttle_sleep_ratio: f64,
#[serde(rename = "last_cycle_yield_ratio", default)]
pub last_cycle_yield_ratio: f64,
#[serde(rename = "last_cycle_total_pause_ratio", default)]
pub last_cycle_total_pause_ratio: f64,
}
impl Default for ScannerPacingPressureSnapshot {
fn default() -> Self {
Self {
primary_pressure: default_scanner_primary_pressure(),
current_queued_scans: 0,
current_active_scans: 0,
last_cycle_budget_limited: false,
last_cycle_pause_observed: false,
last_cycle_throttle_sleep_ratio: 0.0,
last_cycle_yield_ratio: 0.0,
last_cycle_total_pause_ratio: 0.0,
}
}
}
impl ScannerPacingPressureSnapshot {
fn refresh_primary_pressure(&mut self) {
self.primary_pressure = if self.current_queued_scans > 0 {
"queued_scans"
} else if self.last_cycle_budget_limited {
"cycle_budget"
} else if self.last_cycle_pause_observed {
"throttle_pause"
} else if self.current_active_scans > 0 {
"active_scans"
} else {
SCANNER_PRIMARY_PRESSURE_NONE
}
.to_string();
}
fn merge(&mut self, other: &Self) {
let pause_observed = self.last_cycle_pause_observed
|| self.primary_pressure == "throttle_pause"
|| other.last_cycle_pause_observed
|| other.primary_pressure == "throttle_pause";
self.current_queued_scans = self.current_queued_scans.saturating_add(other.current_queued_scans);
self.current_active_scans = self.current_active_scans.saturating_add(other.current_active_scans);
self.last_cycle_budget_limited |= other.last_cycle_budget_limited;
self.last_cycle_pause_observed = pause_observed;
self.last_cycle_throttle_sleep_ratio = self
.last_cycle_throttle_sleep_ratio
.max(other.last_cycle_throttle_sleep_ratio);
self.last_cycle_yield_ratio = self.last_cycle_yield_ratio.max(other.last_cycle_yield_ratio);
self.last_cycle_total_pause_ratio = self.last_cycle_total_pause_ratio.max(other.last_cycle_total_pause_ratio);
self.refresh_primary_pressure();
}
}
#[derive(Clone, Debug, Default, Serialize, Deserialize)] #[derive(Clone, Debug, Default, Serialize, Deserialize)]
pub struct ScannerMetrics { pub struct ScannerMetrics {
#[serde(rename = "collected")] #[serde(rename = "collected")]
@@ -154,6 +229,8 @@ pub struct ScannerMetrics {
pub last_cycle_partial_source_code: u64, pub last_cycle_partial_source_code: u64,
#[serde(rename = "partial_cycles_by_source", default)] #[serde(rename = "partial_cycles_by_source", default)]
pub partial_cycles_by_source: Vec<ScannerSourceCycleSnapshot>, pub partial_cycles_by_source: Vec<ScannerSourceCycleSnapshot>,
#[serde(rename = "pacing_pressure", default)]
pub pacing_pressure: ScannerPacingPressureSnapshot,
} }
impl ScannerMetrics { impl ScannerMetrics {
@@ -165,6 +242,8 @@ impl ScannerMetrics {
self.last_cycle_partial_source_code = other.last_cycle_partial_source_code; self.last_cycle_partial_source_code = other.last_cycle_partial_source_code;
} }
self.pacing_pressure.merge(&other.pacing_pressure);
if self.ongoing_buckets < other.ongoing_buckets { if self.ongoing_buckets < other.ongoing_buckets {
self.ongoing_buckets = other.ongoing_buckets; self.ongoing_buckets = other.ongoing_buckets;
} }
@@ -711,6 +790,16 @@ mod tests {
collected_at, collected_at,
last_cycle_partial_source: "usage".to_string(), last_cycle_partial_source: "usage".to_string(),
last_cycle_partial_source_code: 1, last_cycle_partial_source_code: 1,
pacing_pressure: ScannerPacingPressureSnapshot {
primary_pressure: "cycle_budget".to_string(),
current_active_scans: 1,
last_cycle_budget_limited: true,
last_cycle_pause_observed: true,
last_cycle_throttle_sleep_ratio: 0.1,
last_cycle_yield_ratio: 0.2,
last_cycle_total_pause_ratio: 0.3,
..Default::default()
},
partial_cycles_by_source: vec![ScannerSourceCycleSnapshot { partial_cycles_by_source: vec![ScannerSourceCycleSnapshot {
source: "usage".to_string(), source: "usage".to_string(),
cycles: 1, cycles: 1,
@@ -722,6 +811,15 @@ mod tests {
collected_at: collected_at + chrono::Duration::seconds(1), collected_at: collected_at + chrono::Duration::seconds(1),
last_cycle_partial_source: "lifecycle".to_string(), last_cycle_partial_source: "lifecycle".to_string(),
last_cycle_partial_source_code: 2, last_cycle_partial_source_code: 2,
pacing_pressure: ScannerPacingPressureSnapshot {
primary_pressure: "queued_scans".to_string(),
current_queued_scans: 4,
current_active_scans: 2,
last_cycle_throttle_sleep_ratio: 0.4,
last_cycle_yield_ratio: 0.1,
last_cycle_total_pause_ratio: 0.5,
..Default::default()
},
partial_cycles_by_source: vec![ partial_cycles_by_source: vec![
ScannerSourceCycleSnapshot { ScannerSourceCycleSnapshot {
source: "usage".to_string(), source: "usage".to_string(),
@@ -737,6 +835,14 @@ mod tests {
assert_eq!(scanner.last_cycle_partial_source, "lifecycle"); assert_eq!(scanner.last_cycle_partial_source, "lifecycle");
assert_eq!(scanner.last_cycle_partial_source_code, 2); assert_eq!(scanner.last_cycle_partial_source_code, 2);
assert_eq!(scanner.pacing_pressure.primary_pressure, "queued_scans");
assert_eq!(scanner.pacing_pressure.current_queued_scans, 4);
assert_eq!(scanner.pacing_pressure.current_active_scans, 3);
assert!(scanner.pacing_pressure.last_cycle_budget_limited);
assert!(scanner.pacing_pressure.last_cycle_pause_observed);
assert_eq!(scanner.pacing_pressure.last_cycle_throttle_sleep_ratio, 0.4);
assert_eq!(scanner.pacing_pressure.last_cycle_yield_ratio, 0.2);
assert_eq!(scanner.pacing_pressure.last_cycle_total_pause_ratio, 0.5);
assert_eq!( assert_eq!(
scanner.partial_cycles_by_source, scanner.partial_cycles_by_source,
vec![ vec![
@@ -751,4 +857,40 @@ mod tests {
] ]
); );
} }
#[test]
fn scanner_metrics_merge_preserves_pause_pressure_without_duration() {
let collected_at = Utc::now();
let mut scanner = ScannerMetrics {
collected_at,
pacing_pressure: ScannerPacingPressureSnapshot {
primary_pressure: "throttle_pause".to_string(),
..Default::default()
},
..Default::default()
};
scanner.merge(&ScannerMetrics {
collected_at: collected_at + chrono::Duration::seconds(1),
pacing_pressure: ScannerPacingPressureSnapshot::default(),
..Default::default()
});
assert_eq!(scanner.pacing_pressure.primary_pressure, "throttle_pause");
assert!(scanner.pacing_pressure.last_cycle_pause_observed);
assert_eq!(scanner.pacing_pressure.last_cycle_total_pause_ratio, 0.0);
}
#[test]
fn scanner_metrics_deserializes_missing_primary_pressure_as_none() {
let pacing_pressure: ScannerPacingPressureSnapshot = serde_json::from_str(
r#"{
"current_active_scans": 1
}"#,
)
.expect("deserialize partial scanner pacing pressure");
assert_eq!(pacing_pressure.primary_pressure, "none");
assert_eq!(pacing_pressure.current_active_scans, 1);
}
} }