feat(scanner): track scan cycle source work (#3240)

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-06-06 20:53:48 +08:00
committed by GitHub
parent fade13d950
commit 83a4e5712e
+171 -11
View File
@@ -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<ScannerSourceWorkCounters>,
current_scan_cycle_source_work_start: Vec<ScannerSourceWorkCounters>,
last_scan_cycle_source_work: Vec<ScannerSourceWorkCounters>,
partial_scan_cycles: AtomicU64,
}
@@ -737,6 +782,10 @@ pub struct ScannerMetricsReport {
#[serde(default)]
pub source_work: Vec<ScannerSourceWorkSnapshot>,
#[serde(default)]
pub current_cycle_source_work: Vec<ScannerSourceWorkSnapshot>,
#[serde(default)]
pub last_cycle_source_work: Vec<ScannerSourceWorkSnapshot>,
#[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<ScannerSourceWorkValues> {
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<ScannerSourceWorkValues> {
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<ScannerSourceWorkValues> {
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<ScannerSourceWorkSnapshot> {
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<ScannerSourceWorkSnapshot> {
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<String> {
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(&current_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;