use std::collections::HashMap;
use std::sync::Mutex;
use crate::metrics::registry::get_registry;
use crate::proto::grpc::{Metric, MetricType};
pub struct ClientMetricsReporter {
last_counter: Mutex<HashMap<String, i64>>,
}
impl Default for ClientMetricsReporter {
fn default() -> Self {
Self {
last_counter: Mutex::new(HashMap::new()),
}
}
}
impl ClientMetricsReporter {
pub fn new() -> Self {
Self {
last_counter: Mutex::new(HashMap::new()),
}
}
pub fn snapshot(&self) -> Vec<Metric> {
let mut out = Vec::new();
let mut last = self.last_counter.lock().unwrap();
let registry = get_registry();
for entry in registry.counters.iter() {
let name = entry.key().clone();
let cur = entry.value().get();
let diff = match last.get(&name) {
Some(&prev) => cur - prev, None => cur, };
last.insert(name.clone(), cur);
if diff != 0 {
out.push(Metric {
instance: Some(INSTANCE_CLIENT.to_string()),
source: None,
name: Some(strip_instance_prefix(&name).to_string()),
value: Some(diff as f64),
metric_type: MetricType::Counter as i32,
tags: Default::default(),
});
}
}
for entry in registry.gauges.iter() {
out.push(Metric {
instance: Some(INSTANCE_CLIENT.to_string()),
name: Some(strip_instance_prefix(entry.key()).to_string()),
value: Some(entry.value().get() as f64),
metric_type: MetricType::Gauge as i32,
..Default::default()
});
}
out
}
}
const INSTANCE_CLIENT: &str = "Client";
fn strip_instance_prefix(full: &str) -> &str {
match full.split_once('.') {
Some((_instance, rest)) => rest,
None => full,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metrics::registry::{counter, gauge, name};
#[test]
fn counter_first_snapshot_reports_full_value() {
let reporter = ClientMetricsReporter::default();
let c = counter("test_reporter_first_snap");
c.inc(100);
let snap = reporter.snapshot();
let found = snap
.iter()
.find(|m| m.name.as_deref() == Some("test_reporter_first_snap"));
assert!(found.is_some(), "counter should appear in first snapshot");
assert_eq!(found.unwrap().value, Some(100.0));
assert_eq!(found.unwrap().metric_type, MetricType::Counter as i32);
}
#[test]
fn counter_diff_correct_between_snapshots() {
let reporter = ClientMetricsReporter::default();
let c = counter("test_reporter_diff");
c.inc(50);
let snap1 = reporter.snapshot();
let v1 = snap1
.iter()
.find(|m| m.name.as_deref() == Some("test_reporter_diff"))
.and_then(|m| m.value);
assert_eq!(v1, Some(50.0));
c.inc(30);
let snap2 = reporter.snapshot();
let v2 = snap2
.iter()
.find(|m| m.name.as_deref() == Some("test_reporter_diff"))
.and_then(|m| m.value);
assert_eq!(v2, Some(30.0));
}
#[test]
fn counter_unchanged_not_in_snapshot() {
let reporter = ClientMetricsReporter::default();
let c = counter("test_reporter_unchanged");
c.inc(10);
let snap1 = reporter.snapshot();
assert!(snap1
.iter()
.any(|m| m.name.as_deref() == Some("test_reporter_unchanged")));
let snap2 = reporter.snapshot();
assert!(
!snap2
.iter()
.any(|m| m.name.as_deref() == Some("test_reporter_unchanged")),
"unchanged counter must not appear in snapshot"
);
}
#[test]
fn counter_zero_from_start_not_in_snapshot() {
let reporter = ClientMetricsReporter::default();
let _c = counter("test_reporter_zero_start");
let snap = reporter.snapshot();
assert!(
!snap
.iter()
.any(|m| m.name.as_deref() == Some("test_reporter_zero_start")),
"zero-value counter must not appear in snapshot"
);
}
#[test]
fn counter_baseline_written_even_when_diff_zero() {
let reporter = ClientMetricsReporter::default();
let c = counter("test_reporter_baseline");
let snap1 = reporter.snapshot();
assert!(!snap1
.iter()
.any(|m| m.name.as_deref() == Some("test_reporter_baseline")));
c.inc(75);
let snap2 = reporter.snapshot();
let v = snap2
.iter()
.find(|m| m.name.as_deref() == Some("test_reporter_baseline"))
.and_then(|m| m.value);
assert_eq!(
v,
Some(75.0),
"diff after baseline write must be incremental"
);
}
#[test]
fn gauge_returns_current_value() {
let reporter = ClientMetricsReporter::default();
let g = gauge("test_reporter_gauge");
g.set(42);
let snap = reporter.snapshot();
let found = snap
.iter()
.find(|m| m.name.as_deref() == Some("test_reporter_gauge"));
assert!(found.is_some(), "gauge must appear in snapshot");
assert_eq!(found.unwrap().value, Some(42.0));
assert_eq!(found.unwrap().metric_type, MetricType::Gauge as i32);
}
#[test]
fn gauge_no_diff_always_reported() {
let reporter = ClientMetricsReporter::default();
let g = gauge("test_reporter_gauge_stable");
g.set(7);
let snap1 = reporter.snapshot();
let snap2 = reporter.snapshot();
let in_snap1 = snap1
.iter()
.any(|m| m.name.as_deref() == Some("test_reporter_gauge_stable"));
let in_snap2 = snap2
.iter()
.any(|m| m.name.as_deref() == Some("test_reporter_gauge_stable"));
assert!(in_snap1, "gauge must appear in first snapshot");
assert!(
in_snap2,
"gauge must appear in second snapshot even if unchanged"
);
}
#[test]
fn name_constants_round_trip() {
let reporter = ClientMetricsReporter::default();
counter(name::CLIENT_BYTES_READ_LOCAL).inc(1024);
counter(name::CLIENT_BYTES_WRITTEN_LOCAL).inc(2048);
counter(name::CLIENT_BYTES_WRITTEN_UFS).inc(512);
let snap = reporter.snapshot();
for bare in &["BytesReadLocal", "BytesWrittenLocal", "BytesWrittenUfs"] {
let m = snap
.iter()
.find(|m| m.name.as_deref() == Some(*bare))
.unwrap_or_else(|| panic!("metric {} missing from snapshot", bare));
assert_eq!(
m.instance.as_deref(),
Some(INSTANCE_CLIENT),
"metric {} must carry instance=\"Client\"",
bare
);
assert_eq!(m.metric_type, MetricType::Counter as i32);
}
}
#[test]
fn strip_instance_prefix_behaviour() {
assert_eq!(
strip_instance_prefix("Client.BytesReadLocal"),
"BytesReadLocal"
);
assert_eq!(
strip_instance_prefix("Worker.BytesReadAlluxio"),
"BytesReadAlluxio"
);
assert_eq!(
strip_instance_prefix("Client.BytesReadPerUfs.s3a"),
"BytesReadPerUfs.s3a"
);
assert_eq!(strip_instance_prefix("BytesReadLocal"), "BytesReadLocal");
assert_eq!(strip_instance_prefix(""), "");
}
}