use std::sync::Arc;
use std::time::SystemTime;
fn stamp_ms(now_ms: i64, captured_at: std::time::Instant) -> i64 {
now_ms - captured_at.elapsed().as_micros().div_ceil(1000) as i64
}
use super::{Matcher, MetricAccess, QueryError, Sample, Series, Vector};
use crate::metrics_query::MetricsQuery;
use crate::snapshot::MetricValue;
pub struct MetricsQueryAccess {
query: Arc<MetricsQuery>,
}
impl MetricsQueryAccess {
pub fn new(query: Arc<MetricsQuery>) -> Self {
Self { query }
}
pub fn earliest_ms(&self) -> Option<i64> {
let reporter = self.query.reporter();
let cadence = reporter.declared_cadences().smallest();
if cadence.is_zero() {
return None;
}
let now_ms = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
let mut earliest: Option<i64> = None;
for component in reporter.component_labels() {
if let Some(oldest) = reporter
.window_view(&component, cadence)
.and_then(|v| v.ring.first().cloned())
{
let ms = stamp_ms(now_ms, oldest.captured_at());
earliest = Some(earliest.map_or(ms, |e: i64| e.min(ms)));
}
}
earliest
}
}
impl super::HorizonAware for MetricsQueryAccess {
fn earliest_ms(&self) -> Option<i64> {
MetricsQueryAccess::earliest_ms(self)
}
}
impl MetricAccess for MetricsQueryAccess {
fn select_range(
&self,
matchers: &[Matcher],
start_ms: i64,
end_ms: i64,
) -> Result<Vector, QueryError> {
let reporter = self.query.reporter();
let cadence = reporter.declared_cadences().smallest();
if cadence.is_zero() {
return Ok(Vector::default());
}
let now_ms = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
let mut out: Vec<Series> = Vec::new();
for component in reporter.component_labels() {
let Some(view) = reporter.window_view(&component, cadence) else {
continue;
};
let mut windows: Vec<Arc<crate::snapshot::MetricSet>> =
Vec::with_capacity(view.ring.len() + 2);
windows.extend(view.ring.iter().cloned());
if let Some(l) = &view.latest {
windows.push(l.clone());
}
if let Some(p) = &view.prebuffer {
windows.push(p.clone());
}
windows.sort_by_key(|w| w.captured_at());
windows.dedup_by_key(|w| w.captured_at());
for window in &windows {
let window_ms = stamp_ms(now_ms, window.captured_at());
if window_ms < start_ms || window_ms > end_ms {
continue;
}
for family in window.families() {
let name = family.name();
for metric in family.metrics() {
let labels = series_labels(name, &component, metric.labels());
if !matchers.iter().all(|m| m.matches(&labels)) {
continue;
}
let Some(value) = metric.point().and_then(|p| value_to_f64(p.value()))
else {
continue;
};
push_sample(
&mut out,
labels,
Sample {
timestamp_ms: window_ms,
value,
},
);
}
}
}
}
for s in &mut out {
s.samples.sort_by_key(|x| x.timestamp_ms);
}
Ok(Vector::new(out))
}
}
fn series_labels(
name: &str,
component: &crate::labels::Labels,
metric: &crate::labels::Labels,
) -> Vec<(String, String)> {
let mut out = vec![("__name__".to_string(), name.to_string())];
let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
seen.insert("__name__".to_string());
for (k, v) in metric.iter().chain(component.iter()) {
if seen.insert(k.to_string()) {
out.push((k.to_string(), v.to_string()));
}
}
out
}
fn value_to_f64(v: &MetricValue) -> Option<f64> {
match v {
MetricValue::Counter(c) => Some(c.cumulative as f64),
MetricValue::Gauge(g) => Some(g.value),
MetricValue::Histogram(h) => Some(h.cumulative_count as f64),
MetricValue::BucketedHistogram(b) => Some(b.cumulative_count as f64),
_ => None,
}
}
fn push_sample(out: &mut Vec<Series>, labels: Vec<(String, String)>, sample: Sample) {
if let Some(existing) = out.iter_mut().find(|s| s.labels == labels) {
existing.samples.push(sample);
} else {
out.push(Series {
labels,
samples: vec![sample],
});
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn value_projection_per_variant() {
use crate::snapshot::{CounterValue, GaugeValue, HistogramValue};
use hdrhistogram::Histogram as HdrHistogram;
assert_eq!(
value_to_f64(&MetricValue::Counter(CounterValue::new(7))),
Some(7.0)
);
assert_eq!(
value_to_f64(&MetricValue::Gauge(GaugeValue::new(3.5))),
Some(3.5)
);
let mut hdr = HdrHistogram::<u64>::new(3).unwrap();
hdr.record(5).unwrap();
hdr.record(6).unwrap();
let hv = HistogramValue::from_hdr(hdr).with_cumulative_count(42);
assert_eq!(value_to_f64(&MetricValue::Histogram(hv)), Some(42.0));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn series_labels_promote_name_and_merge_scope() {
let component = crate::labels::Labels::empty().with("phase", "saturate");
let metric = crate::labels::Labels::empty().with("__name__", "stale");
let out = series_labels("errors_total", &component, &metric);
assert_eq!(out[0], ("__name__".to_string(), "errors_total".to_string()));
assert_eq!(out.iter().filter(|(k, _)| k == "__name__").count(), 1);
assert!(out.iter().any(|(k, v)| k == "phase" && v == "saturate"));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn select_instant_returns_a_live_counter() {
use crate::cadence::{CadenceTree, Cadences};
use crate::cadence_reporter::CadenceReporter;
use crate::labels::Labels;
use crate::snapshot::MetricSet;
use std::collections::HashMap;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let comp = Labels::of("phase", "saturate");
let mut delta = MetricSet::new(Duration::from_millis(100));
delta.insert_counter("ops", Labels::default(), 7, Instant::now());
reporter.scope_close(&comp, delta);
reporter.flush_for_tests();
let root = crate::component::Component::root(Labels::of("session", "s1"), HashMap::new());
let access = MetricsQueryAccess::new(Arc::new(MetricsQuery::new(reporter.clone(), root)));
let now_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let v = access
.select_instant(&[Matcher::eq("__name__", "ops")], now_ms, Some(60_000))
.expect("select");
assert_eq!(v.len(), 1, "one series, got {v:?}");
let s = &v.series()[0];
assert!(
s.labels
.iter()
.any(|(k, val)| k == "__name__" && val == "ops")
);
assert!(
s.labels
.iter()
.any(|(k, val)| k == "phase" && val == "saturate"),
"merged scope tag: {:?}",
s.labels
);
assert_eq!(s.samples.len(), 1, "instant vector: one sample per series");
assert_eq!(s.samples[0].value, 7.0);
}
}