Skip to main content

nmbrs_metrics/queryapi/
live.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! The live in-process [`MetricAccess`] backend: reads the current
5//! session's cadence ring through [`MetricsQuery`].
6//!
7//! ## Time model
8//!
9//! Cadence windows carry a monotonic [`std::time::Instant`]
10//! ([`crate::snapshot::MetricSet::captured_at`]), not a wall-clock
11//! stamp; the query API works in Unix-ms. Each window is converted as
12//! `now_ms − captured_at().elapsed()`, with `now_ms` read **once per
13//! call** so every window is stamped against one clock — the same
14//! approximation the sqlite writer makes (`reporters::sqlite` stamps
15//! `SystemTime::now()` at write time). A window is captured before any
16//! query, so `window_ms ≤ now`, and the caller's `end_ms` is also
17//! ~now — the freshest window is never excluded.
18//!
19//! ## Coverage
20//!
21//! Series are drawn from the **smallest declared cadence**'s ring +
22//! latest + in-flight prebuffer (deduped by capture instant); the
23//! queryable lookback is bounded by the ring depth. `__name__` is the
24//! metric family name, with the component's scope tags merged in.
25//! `Counter`/`Gauge`/`Histogram` project to f64; `Info`/`StateSet` are
26//! skipped (proper histogram `_count`/`_sum`/`_bucket` decomposition is
27//! a follow-up).
28
29use std::sync::Arc;
30use std::time::SystemTime;
31
32/// A window's Unix-ms stamp against `now_ms`, read once per call. The
33/// elapsed time is rounded **up**: `now_ms` is already floored, and
34/// subtracting a floored elapsed could stamp a window up to 1 ms after
35/// its capture — past a caller's `end_ms` read in that same millisecond,
36/// which would drop the freshest window. Rounding up keeps
37/// `window_ms ≤ capture`.
38fn stamp_ms(now_ms: i64, captured_at: std::time::Instant) -> i64 {
39    now_ms - captured_at.elapsed().as_micros().div_ceil(1000) as i64
40}
41
42use super::{Matcher, MetricAccess, QueryError, Sample, Series, Vector};
43use crate::metrics_query::MetricsQuery;
44use crate::snapshot::MetricValue;
45
46/// The live in-process access backend. Cheap to construct (wraps the
47/// shared `Arc`).
48pub struct MetricsQueryAccess {
49    query: Arc<MetricsQuery>,
50}
51
52impl MetricsQueryAccess {
53    /// Wrap the session's live metrics query.
54    pub fn new(query: Arc<MetricsQuery>) -> Self {
55        Self { query }
56    }
57
58    /// SRD-90 §M5 — the oldest sample-time this in-memory tier can answer
59    /// (Unix-ms), for the [`super::HybridStore`] coverage walk: the minimum
60    /// `captured_at` across every component's smallest-cadence retention ring,
61    /// converted on the same clock as [`Self::select_range`]. `None` when
62    /// nothing is retained yet (empty store / no cadence) — the hybrid then
63    /// never short-circuits the colder tail.
64    pub fn earliest_ms(&self) -> Option<i64> {
65        let reporter = self.query.reporter();
66        let cadence = reporter.declared_cadences().smallest();
67        if cadence.is_zero() {
68            return None;
69        }
70        let now_ms = SystemTime::now()
71            .duration_since(SystemTime::UNIX_EPOCH)
72            .map(|d| d.as_millis() as i64)
73            .unwrap_or(0);
74        let mut earliest: Option<i64> = None;
75        for component in reporter.component_labels() {
76            // The ring is oldest-first; its first entry is the oldest retained.
77            // One window-view lookup, then read `ring.first()` directly — no
78            // cloning the whole ring just to peek the head.
79            if let Some(oldest) = reporter
80                .window_view(&component, cadence)
81                .and_then(|v| v.ring.first().cloned())
82            {
83                let ms = stamp_ms(now_ms, oldest.captured_at());
84                earliest = Some(earliest.map_or(ms, |e: i64| e.min(ms)));
85            }
86        }
87        earliest
88    }
89}
90
91impl super::HorizonAware for MetricsQueryAccess {
92    fn earliest_ms(&self) -> Option<i64> {
93        MetricsQueryAccess::earliest_ms(self)
94    }
95}
96
97impl MetricAccess for MetricsQueryAccess {
98    fn select_range(
99        &self,
100        matchers: &[Matcher],
101        start_ms: i64,
102        end_ms: i64,
103    ) -> Result<Vector, QueryError> {
104        let reporter = self.query.reporter();
105        let cadence = reporter.declared_cadences().smallest();
106        if cadence.is_zero() {
107            return Ok(Vector::default());
108        }
109        let now_ms = SystemTime::now()
110            .duration_since(SystemTime::UNIX_EPOCH)
111            .map(|d| d.as_millis() as i64)
112            .unwrap_or(0);
113
114        // SRD-89 §3b / SRD-90 §M6 — execution scoping is a **uniform dimensional
115        // label**: `exec_id` arrives as a matcher (injected by
116        // `ExecScopedAccess` for the reading execution) and is applied by the
117        // matcher check below against each series' `exec_id` scope tag — there
118        // is no separate post-filter. A `None`-scope read (single-run) injects
119        // nothing, so it sees the sole execution's series (A1). Each series here
120        // carries `exec_id` merged in from the component scope tags, which is
121        // what the matcher matches against.
122        let mut out: Vec<Series> = Vec::new();
123        for component in reporter.component_labels() {
124            // Closed-window history + freshest closed window + in-flight
125            // partial, deduped by capture instant. ONE window-view lookup per
126            // component (was three: ring + latest + prebuffer each rebuilt the
127            // component-path string and re-hashed it); the partial is an `Arc`
128            // clone off the view, not a deep `MetricSet` copy.
129            let Some(view) = reporter.window_view(&component, cadence) else {
130                continue;
131            };
132            let mut windows: Vec<Arc<crate::snapshot::MetricSet>> =
133                Vec::with_capacity(view.ring.len() + 2);
134            windows.extend(view.ring.iter().cloned());
135            if let Some(l) = &view.latest {
136                windows.push(l.clone());
137            }
138            if let Some(p) = &view.prebuffer {
139                windows.push(p.clone());
140            }
141            windows.sort_by_key(|w| w.captured_at());
142            windows.dedup_by_key(|w| w.captured_at());
143
144            for window in &windows {
145                let window_ms = stamp_ms(now_ms, window.captured_at());
146                if window_ms < start_ms || window_ms > end_ms {
147                    continue;
148                }
149                for family in window.families() {
150                    let name = family.name();
151                    for metric in family.metrics() {
152                        let labels = series_labels(name, &component, metric.labels());
153                        if !matchers.iter().all(|m| m.matches(&labels)) {
154                            continue;
155                        }
156                        let Some(value) = metric.point().and_then(|p| value_to_f64(p.value()))
157                        else {
158                            continue;
159                        };
160                        push_sample(
161                            &mut out,
162                            labels,
163                            Sample {
164                                timestamp_ms: window_ms,
165                                value,
166                            },
167                        );
168                    }
169                }
170            }
171        }
172
173        for s in &mut out {
174            s.samples.sort_by_key(|x| x.timestamp_ms);
175        }
176        Ok(Vector::new(out))
177    }
178
179    // `select_instant` uses the trait default (range + latest-per-series).
180}
181
182/// Build a series label set: `__name__` (front), then the union of the
183/// instrument's `metric` labels and the `component` (ring-key) scope
184/// tags. Instrument labels win on key collision; both normally agree.
185fn series_labels(
186    name: &str,
187    component: &crate::labels::Labels,
188    metric: &crate::labels::Labels,
189) -> Vec<(String, String)> {
190    let mut out = vec![("__name__".to_string(), name.to_string())];
191    let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
192    seen.insert("__name__".to_string());
193    for (k, v) in metric.iter().chain(component.iter()) {
194        if seen.insert(k.to_string()) {
195            out.push((k.to_string(), v.to_string()));
196        }
197    }
198    out
199}
200
201/// Project a [`MetricValue`] to the f64 the query API exposes. `None`
202/// for variants with no single scalar projection (skipped by callers).
203fn value_to_f64(v: &MetricValue) -> Option<f64> {
204    match v {
205        // Counters expose the cumulative running total (Prometheus/VM
206        // semantics) so MetricsQL rate()/increase/*_over_time compute
207        // Δ/Δt correctly — see the cumulative-counter note.
208        MetricValue::Counter(c) => Some(c.cumulative as f64),
209        MetricValue::Gauge(g) => Some(g.value),
210        // Histogram count is exposed cumulative too (its lifetime total),
211        // so rate()/increase over a histogram count is PromQL-correct.
212        MetricValue::Histogram(h) => Some(h.cumulative_count as f64),
213        MetricValue::BucketedHistogram(b) => Some(b.cumulative_count as f64),
214        _ => None,
215    }
216}
217
218/// Append a sample to the series with this exact label set, creating it
219/// if new.
220fn push_sample(out: &mut Vec<Series>, labels: Vec<(String, String)>, sample: Sample) {
221    if let Some(existing) = out.iter_mut().find(|s| s.labels == labels) {
222        existing.samples.push(sample);
223    } else {
224        out.push(Series {
225            labels,
226            samples: vec![sample],
227        });
228    }
229}
230
231#[cfg(test)]
232mod tests {
233    use super::*;
234
235    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
236    async fn value_projection_per_variant() {
237        use crate::snapshot::{CounterValue, GaugeValue, HistogramValue};
238        use hdrhistogram::Histogram as HdrHistogram;
239        assert_eq!(
240            value_to_f64(&MetricValue::Counter(CounterValue::new(7))),
241            Some(7.0)
242        );
243        assert_eq!(
244            value_to_f64(&MetricValue::Gauge(GaugeValue::new(3.5))),
245            Some(3.5)
246        );
247        // A histogram projects its CUMULATIVE count (lifetime total), not the
248        // per-window reservoir count — here a 2-sample window over a lifetime
249        // of 42.
250        let mut hdr = HdrHistogram::<u64>::new(3).unwrap();
251        hdr.record(5).unwrap();
252        hdr.record(6).unwrap();
253        let hv = HistogramValue::from_hdr(hdr).with_cumulative_count(42);
254        assert_eq!(value_to_f64(&MetricValue::Histogram(hv)), Some(42.0));
255    }
256
257    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
258    async fn series_labels_promote_name_and_merge_scope() {
259        let component = crate::labels::Labels::empty().with("phase", "saturate");
260        let metric = crate::labels::Labels::empty().with("__name__", "stale");
261        let out = series_labels("errors_total", &component, &metric);
262        assert_eq!(out[0], ("__name__".to_string(), "errors_total".to_string()));
263        assert_eq!(out.iter().filter(|(k, _)| k == "__name__").count(), 1);
264        assert!(out.iter().any(|(k, v)| k == "phase" && v == "saturate"));
265    }
266
267    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
268    async fn select_instant_returns_a_live_counter() {
269        use crate::cadence::{CadenceTree, Cadences};
270        use crate::cadence_reporter::CadenceReporter;
271        use crate::labels::Labels;
272        use crate::snapshot::MetricSet;
273        use std::collections::HashMap;
274        use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
275
276        // Populate the ring synchronously.
277        let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
278        let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
279        let comp = Labels::of("phase", "saturate");
280        let mut delta = MetricSet::new(Duration::from_millis(100));
281        delta.insert_counter("ops", Labels::default(), 7, Instant::now());
282        reporter.scope_close(&comp, delta);
283        reporter.flush_for_tests();
284
285        let root = crate::component::Component::root(Labels::of("session", "s1"), HashMap::new());
286        let access = MetricsQueryAccess::new(Arc::new(MetricsQuery::new(reporter.clone(), root)));
287
288        let now_ms = SystemTime::now()
289            .duration_since(UNIX_EPOCH)
290            .unwrap()
291            .as_millis() as i64;
292        let v = access
293            .select_instant(&[Matcher::eq("__name__", "ops")], now_ms, Some(60_000))
294            .expect("select");
295
296        assert_eq!(v.len(), 1, "one series, got {v:?}");
297        let s = &v.series()[0];
298        assert!(
299            s.labels
300                .iter()
301                .any(|(k, val)| k == "__name__" && val == "ops")
302        );
303        assert!(
304            s.labels
305                .iter()
306                .any(|(k, val)| k == "phase" && val == "saturate"),
307            "merged scope tag: {:?}",
308            s.labels
309        );
310        assert_eq!(s.samples.len(), 1, "instant vector: one sample per series");
311        assert_eq!(s.samples[0].value, 7.0);
312    }
313}