diff --git a/crates/common/src/metrics.rs b/crates/common/src/metrics.rs index e38fb583c..bbd447857 100644 --- a/crates/common/src/metrics.rs +++ b/crates/common/src/metrics.rs @@ -302,9 +302,16 @@ impl ScannerSourceWorkCounters { self.skipped.fetch_add(skipped, Ordering::Relaxed); } - fn snapshot(&self, source: ScannerWorkSource) -> ScannerSourceWorkSnapshot { - ScannerSourceWorkSnapshot { - source: source.as_str().to_string(), + fn store(&self, values: ScannerSourceWorkValues) { + self.checked.store(values.checked, Ordering::Relaxed); + self.queued.store(values.queued, Ordering::Relaxed); + self.executed.store(values.executed, Ordering::Relaxed); + self.failed.store(values.failed, Ordering::Relaxed); + self.skipped.store(values.skipped, Ordering::Relaxed); + } + + fn values(&self) -> ScannerSourceWorkValues { + ScannerSourceWorkValues { checked: self.checked.load(Ordering::Relaxed), queued: self.queued.load(Ordering::Relaxed), executed: self.executed.load(Ordering::Relaxed), @@ -312,6 +319,42 @@ impl ScannerSourceWorkCounters { skipped: self.skipped.load(Ordering::Relaxed), } } + + fn snapshot(&self, source: ScannerWorkSource) -> ScannerSourceWorkSnapshot { + self.values().snapshot(source) + } +} + +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +struct ScannerSourceWorkValues { + checked: u64, + queued: u64, + executed: u64, + failed: u64, + skipped: u64, +} + +impl ScannerSourceWorkValues { + fn saturating_sub(self, start: Self) -> Self { + Self { + checked: self.checked.saturating_sub(start.checked), + queued: self.queued.saturating_sub(start.queued), + executed: self.executed.saturating_sub(start.executed), + failed: self.failed.saturating_sub(start.failed), + skipped: self.skipped.saturating_sub(start.skipped), + } + } + + fn snapshot(self, source: ScannerWorkSource) -> ScannerSourceWorkSnapshot { + ScannerSourceWorkSnapshot { + source: source.as_str().to_string(), + checked: self.checked, + queued: self.queued, + executed: self.executed, + failed: self.failed, + skipped: self.skipped, + } + } } // --------------------------------------------------------------------------- @@ -515,6 +558,8 @@ pub struct Metrics { scanner_checkpoint_ignored: AtomicU64, scanner_checkpoint_stale: AtomicU64, scanner_source_work: Vec, + current_scan_cycle_source_work_start: Vec, + last_scan_cycle_source_work: Vec, partial_scan_cycles: AtomicU64, } @@ -737,6 +782,10 @@ pub struct ScannerMetricsReport { #[serde(default)] pub source_work: Vec, #[serde(default)] + pub current_cycle_source_work: Vec, + #[serde(default)] + pub last_cycle_source_work: Vec, + #[serde(default)] pub partial_cycles: u64, } @@ -922,6 +971,14 @@ impl Metrics { .iter() .map(|_| ScannerSourceWorkCounters::default()) .collect(), + current_scan_cycle_source_work_start: ScannerWorkSource::all() + .iter() + .map(|_| ScannerSourceWorkCounters::default()) + .collect(), + last_scan_cycle_source_work: ScannerWorkSource::all() + .iter() + .map(|_| ScannerSourceWorkCounters::default()) + .collect(), partial_scan_cycles: AtomicU64::new(0), } } @@ -1323,6 +1380,7 @@ impl Metrics { pub fn start_scan_cycle_work(&self) -> ScanCycleWorkSnapshot { let snapshot = self.scan_cycle_work_snapshot(); + let source_snapshot = self.scanner_source_work_values(); self.current_scan_cycle_objects_start .store(snapshot.objects_scanned, Ordering::Relaxed); self.current_scan_cycle_directories_start @@ -1347,13 +1405,16 @@ impl Metrics { .store(snapshot.replication_checks, Ordering::Relaxed); self.current_scan_cycle_usage_saves_start .store(snapshot.usage_saves, Ordering::Relaxed); + self.store_scanner_source_work_values(&self.current_scan_cycle_source_work_start, &source_snapshot); self.current_scan_cycle_work_active.store(true, Ordering::Relaxed); snapshot } pub fn finish_scan_cycle_work(&self, start: ScanCycleWorkSnapshot) { let work = self.scan_cycle_work_since(start); + let source_work = self.scanner_source_work_since(&self.current_scan_cycle_source_work_start_values()); self.record_scan_cycle_work(work); + self.record_scan_cycle_source_work(&source_work); self.current_scan_cycle_work_active.store(false, Ordering::Relaxed); } @@ -1413,6 +1474,58 @@ impl Metrics { } } + fn scanner_source_work_values(&self) -> Vec { + ScannerWorkSource::all() + .iter() + .filter_map(|source| { + self.scanner_source_work + .get(source.index()) + .map(ScannerSourceWorkCounters::values) + }) + .collect() + } + + fn current_scan_cycle_source_work_start_values(&self) -> Vec { + ScannerWorkSource::all() + .iter() + .filter_map(|source| { + self.current_scan_cycle_source_work_start + .get(source.index()) + .map(ScannerSourceWorkCounters::values) + }) + .collect() + } + + fn scanner_source_work_since(&self, start: &[ScannerSourceWorkValues]) -> Vec { + self.scanner_source_work_values() + .into_iter() + .enumerate() + .map(|(index, current)| current.saturating_sub(start.get(index).copied().unwrap_or_default())) + .collect() + } + + fn scanner_source_work_snapshots(&self, values: &[ScannerSourceWorkValues]) -> Vec { + ScannerWorkSource::all() + .iter() + .filter_map(|source| values.get(source.index()).map(|values| values.snapshot(*source))) + .collect() + } + + fn scanner_source_work_counter_snapshots(&self, counters: &[ScannerSourceWorkCounters]) -> Vec { + ScannerWorkSource::all() + .iter() + .filter_map(|source| counters.get(source.index()).map(|counters| counters.snapshot(*source))) + .collect() + } + + fn store_scanner_source_work_values(&self, counters: &[ScannerSourceWorkCounters], values: &[ScannerSourceWorkValues]) { + for source in ScannerWorkSource::all() { + if let (Some(counter), Some(values)) = (counters.get(source.index()), values.get(source.index())) { + counter.store(*values); + } + } + } + pub fn record_scan_cycle_work(&self, work: ScanCycleWorkSnapshot) { // Telemetry-only gauges: readers may observe a transient mixed snapshot // while these independent atomic fields are updated. @@ -1438,6 +1551,10 @@ impl Metrics { self.last_scan_cycle_usage_saves.store(work.usage_saves, Ordering::Relaxed); } + fn record_scan_cycle_source_work(&self, work: &[ScannerSourceWorkValues]) { + self.store_scanner_source_work_values(&self.last_scan_cycle_source_work, work); + } + /// Snapshot of every path currently being scanned. pub async fn get_current_paths(&self) -> Vec { self.current_path_snapshots() @@ -1509,6 +1626,7 @@ impl Metrics { m.current_disk_bucket_scans_active = disk_bucket_scans_active; if self.current_scan_cycle_work_active.load(Ordering::Relaxed) { let current_work = self.scan_cycle_work_since(self.current_scan_cycle_work_start()); + let current_source_work = self.scanner_source_work_since(&self.current_scan_cycle_source_work_start_values()); m.current_cycle_objects_scanned = current_work.objects_scanned; m.current_cycle_directories_scanned = current_work.directories_scanned; m.current_cycle_bucket_drive_scans = current_work.bucket_drive_scans; @@ -1521,6 +1639,7 @@ impl Metrics { m.current_cycle_heal_objects = current_work.heal_objects; m.current_cycle_replication_checks = current_work.replication_checks; m.current_cycle_usage_saves = current_work.usage_saves; + m.current_cycle_source_work = self.scanner_source_work_snapshots(¤t_source_work); } let last_cycle_result = self.last_scan_cycle_result.load(Ordering::Relaxed); m.last_cycle_result = scan_cycle_result_label(last_cycle_result).to_string(); @@ -1542,6 +1661,7 @@ impl Metrics { m.last_cycle_heal_objects = self.last_scan_cycle_heal_objects.load(Ordering::Relaxed); m.last_cycle_replication_checks = self.last_scan_cycle_replication_checks.load(Ordering::Relaxed); m.last_cycle_usage_saves = self.last_scan_cycle_usage_saves.load(Ordering::Relaxed); + m.last_cycle_source_work = self.scanner_source_work_counter_snapshots(&self.last_scan_cycle_source_work); m.failed_cycles = self.failed_scan_cycles.load(Ordering::Relaxed); m.partial_cycles_unknown = self.partial_scan_cycles_unknown.load(Ordering::Relaxed); m.partial_cycles_runtime = self.partial_scan_cycles_runtime.load(Ordering::Relaxed); @@ -1565,14 +1685,7 @@ impl Metrics { m.scan_checkpoint_cleared = self.scanner_checkpoint_cleared.load(Ordering::Relaxed); m.scan_checkpoint_ignored = self.scanner_checkpoint_ignored.load(Ordering::Relaxed); m.scan_checkpoint_stale = self.scanner_checkpoint_stale.load(Ordering::Relaxed); - m.source_work = ScannerWorkSource::all() - .iter() - .filter_map(|source| { - self.scanner_source_work - .get(source.index()) - .map(|counters| counters.snapshot(*source)) - }) - .collect(); + m.source_work = self.scanner_source_work_counter_snapshots(&self.scanner_source_work); m.partial_cycles = self.partial_scan_cycles.load(Ordering::Relaxed); // Lifetime operation counts @@ -1873,6 +1986,53 @@ mod tests { assert_eq!(heal.executed, 6); } + #[tokio::test] + async fn report_includes_scan_cycle_source_work() { + let metrics = Metrics::new(); + metrics.record_scanner_source_work(ScannerWorkSource::Usage, 10, 0, 1, 0, 0); + + let start = metrics.start_scan_cycle_work(); + metrics.record_scanner_source_work(ScannerWorkSource::Usage, 3, 0, 2, 0, 0); + metrics.record_scanner_source_work(ScannerWorkSource::Lifecycle, 0, 1, 1, 0, 0); + + let report = metrics.report().await; + let usage = report + .current_cycle_source_work + .iter() + .find(|work| work.source == ScannerWorkSource::Usage.as_str()) + .expect("current cycle usage source work should be visible"); + assert_eq!(usage.checked, 3); + assert_eq!(usage.executed, 2); + + let lifecycle = report + .current_cycle_source_work + .iter() + .find(|work| work.source == ScannerWorkSource::Lifecycle.as_str()) + .expect("current cycle lifecycle source work should be visible"); + assert_eq!(lifecycle.queued, 1); + assert_eq!(lifecycle.executed, 1); + + metrics.finish_scan_cycle_work(start); + let report = metrics.report().await; + assert!(report.current_cycle_source_work.is_empty()); + + let usage = report + .last_cycle_source_work + .iter() + .find(|work| work.source == ScannerWorkSource::Usage.as_str()) + .expect("last cycle usage source work should be visible"); + assert_eq!(usage.checked, 3); + assert_eq!(usage.executed, 2); + + let lifecycle = report + .last_cycle_source_work + .iter() + .find(|work| work.source == ScannerWorkSource::Lifecycle.as_str()) + .expect("last cycle lifecycle source work should be visible"); + assert_eq!(lifecycle.queued, 1); + assert_eq!(lifecycle.executed, 1); + } + #[tokio::test] async fn report_preserves_current_cycle_started_time() { let previous_init_time = *crate::globals::GLOBAL_INIT_TIME.read().await;