Skip to main content

nmbrs_metrics/
metrics_query.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Unified metrics read API (SRD-42 §"MetricsQuery").
5//!
6//! Every consumer (TUI, summary report, SQLite emitter, GK
7//! `metric()`/`metric_window()` nodes, programmatic callers) reads
8//! through this single interface. There is no per-consumer access
9//! layer — the query speaks the metrics system's native types
10//! ([`MetricSet`] / [`MetricFamily`] / [`Metric`] / [`MetricPoint`]),
11//! and exposes these query modes — the method name says whether a value is a
12//! running **total**, a span **increase** (delta), or a **distribution**:
13//!
14//! - [`MetricsQuery::now`] — running-total snapshot at the live cadence.
15//! - [`MetricsQuery::cadence_window`] — the last full closed window's running
16//!   totals for a declared cadence.
17//! - [`MetricsQuery::session_lifetime`] — the session's running totals,
18//!   walking the cascade down at read time so no in-flight data is missed.
19//! - [`MetricsQuery::increase_over`] — the counter **increase** over the last
20//!   `span` (PromQL `increase`; a rate is `increase / span`), differenced
21//!   from the retained finest ring at the finest covering resolution.
22//! - [`MetricsQuery::distribution_over`] — the merged latency/value
23//!   **distribution** (histogram reservoir) over the last `span`.
24//!
25//! ## Selection
26//!
27//! Every query takes a [`Selection`] — a label-based filter applied
28//! to each `(component_labels, metric_labels)` pair as the query
29//! walks the store and the live tree. Identity for combine /
30//! deduplication is `(family.name, label_set)` per OpenMetrics
31//! §4.5.1.
32
33use std::sync::{Arc, RwLock};
34use std::time::{Duration, Instant};
35
36use crate::cadence_reporter::CadenceReporter;
37use crate::component::Component;
38use crate::labels::Labels;
39use crate::snapshot::{CounterValue, Metric, MetricFamily, MetricPoint, MetricSet, MetricValue};
40
41/// A label-based filter for selecting which metrics a query operates
42/// on. Composes by AND — every constraint must match.
43#[derive(Clone, Debug, Default)]
44pub struct Selection {
45    /// Required metric family name. `None` matches any family.
46    family: Option<String>,
47    /// Required label `(key, value)` pairs on the metric's `LabelSet`.
48    /// All pairs must match.
49    label_eq: Vec<(String, String)>,
50    /// Required label `(key, value_substring)` pairs. The value
51    /// must contain the substring.
52    label_contains: Vec<(String, String)>,
53}
54
55impl Selection {
56    pub fn all() -> Self {
57        Self::default()
58    }
59
60    /// Match any series in the named family.
61    pub fn family(name: impl Into<String>) -> Self {
62        Self {
63            family: Some(name.into()),
64            ..Default::default()
65        }
66    }
67
68    /// Builder form: constrain this selection to the named family.
69    pub fn with_family(mut self, name: impl Into<String>) -> Self {
70        self.family = Some(name.into());
71        self
72    }
73
74    /// Restrict to series whose `LabelSet` contains `key=value`.
75    pub fn with_label(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
76        self.label_eq.push((key.into(), value.into()));
77        self
78    }
79
80    /// Restrict to series whose label value at `key` contains
81    /// `substring` (operator-friendly for path-style labels).
82    pub fn with_label_containing(
83        mut self,
84        key: impl Into<String>,
85        substring: impl Into<String>,
86    ) -> Self {
87        self.label_contains.push((key.into(), substring.into()));
88        self
89    }
90
91    /// True when the selection's family constraint matches the
92    /// candidate, or there is no family constraint.
93    pub fn matches_family(&self, family_name: &str) -> bool {
94        self.family
95            .as_deref()
96            .map(|f| f == family_name)
97            .unwrap_or(true)
98    }
99
100    /// True when every label constraint is satisfied by `labels`.
101    pub fn matches_labels(&self, labels: &Labels) -> bool {
102        for (k, v) in &self.label_eq {
103            if labels.get(k) != Some(v.as_str()) {
104                return false;
105            }
106        }
107        for (k, sub) in &self.label_contains {
108            match labels.get(k) {
109                Some(value) if value.contains(sub.as_str()) => {}
110                _ => return false,
111            }
112        }
113        true
114    }
115}
116
117/// Errors returned by selection-required queries.
118#[derive(Debug, PartialEq, Eq)]
119pub enum SelectError {
120    /// The selection matched no metric instance.
121    NoMatch,
122    /// The selection matched more than one instance — caller
123    /// requested exactly one.
124    MultipleMatches(usize),
125}
126
127impl std::fmt::Display for SelectError {
128    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
129        match self {
130            Self::NoMatch => write!(f, "selection matched no metric instance"),
131            Self::MultipleMatches(n) => {
132                write!(f, "selection matched {n} instances, expected exactly one")
133            }
134        }
135    }
136}
137
138impl std::error::Error for SelectError {}
139
140/// The unified metrics read interface. Constructed once at session
141/// start with references to the cadence reporter (for closed
142/// windows + cascade peeks) and the component tree root (for the
143/// `now` mode's live instrument walk).
144pub struct MetricsQuery {
145    reporter: Arc<CadenceReporter>,
146    /// Session-scope component root. Held for component-tree
147    /// structural queries (e.g. "how many phases are running?")
148    /// performed by display code like the TUI's Focus-LOD
149    /// placeholder logic.
150    component_root: Arc<RwLock<Component>>,
151}
152
153impl MetricsQuery {
154    pub fn new(reporter: Arc<CadenceReporter>, component_root: Arc<RwLock<Component>>) -> Self {
155        Self {
156            reporter,
157            component_root,
158        }
159    }
160
161    /// Reference to the cadence reporter — exposed so consumers that
162    /// need to enumerate declared cadences (e.g., per-cadence
163    /// columns) can ask it directly.
164    pub fn reporter(&self) -> &Arc<CadenceReporter> {
165        &self.reporter
166    }
167
168    /// Count of phases currently in `Running` state anywhere in the
169    /// session's component tree. A pure structural query — no
170    /// metric data involved. Used by display code that needs to
171    /// decide "live vs waiting vs done" without re-implementing
172    /// that logic over its own state mirror.
173    pub fn running_phase_count(&self) -> usize {
174        self.component_root
175            .read()
176            .map(|c| c.running_descendant_count())
177            .unwrap_or(0)
178    }
179
180    // ---- now: live instrument peek -------------------------------------
181
182    /// Recent snapshot at the smallest declared cadence, filtered
183    /// by `selection`.
184    ///
185    /// Reads [`Self::cadence_window`] at the smallest declared
186    /// cadence (1 s for default configurations) — the last fully-
187    /// closed window of that cadence. Does NOT pass through to
188    /// the live instruments: counter values, gauge values, and
189    /// histogram reservoirs all come from the cadence-reporter
190    /// store, which the scheduler populates via its per-tick
191    /// coalesce.
192    ///
193    /// Why not a live-instrument peek? Counters are absolute
194    /// atomics (peek is free), but histogram peeks return "samples
195    /// accumulated since the scheduler's last drain" — a partial,
196    /// drifting sub-interval window. The 1 s cadence window is a
197    /// stable, sample-weighted view that matches what every other
198    /// reader (summary, SQLite, cadence subscribers) sees for the
199    /// same time slice.
200    ///
201    /// Returns an empty `MetricSet` before the first window of the
202    /// smallest declared cadence closes (i.e., during the first
203    /// `1 s` of a run). Callers that need true sub-second live data
204    /// for a specific Timer should use
205    /// [`crate::instruments::timer::Timer::peek_live_window`].
206    pub fn now(&self, selection: &Selection) -> MetricSet {
207        let smallest = self.reporter.declared_cadences().smallest();
208        if smallest.is_zero() {
209            return MetricSet::at(Instant::now(), Duration::ZERO);
210        }
211        self.cadence_window(smallest, selection)
212    }
213
214    /// Build a [`MetricHandle`] that caches the `(selection,
215    /// cadence)` pair for repeated cheap reads. The handle reads
216    /// the smallest declared cadence's last closed window on every
217    /// `read_now` — no component-tree walk, no live-instrument
218    /// access.
219    ///
220    /// Callers that want a specific cadence (not the smallest) can
221    /// use [`Self::resolve_at`].
222    pub fn resolve(&self, selection: Selection) -> MetricHandle {
223        let cadence = self.reporter.declared_cadences().smallest();
224        MetricHandle {
225            reporter: self.reporter.clone(),
226            selection,
227            cadence,
228        }
229    }
230
231    /// Same as [`Self::resolve`] but pins the handle to a specific
232    /// cadence — use for per-cadence columns in summary reports
233    /// or for explicit longer-horizon readers.
234    pub fn resolve_at(&self, selection: Selection, cadence: Duration) -> MetricHandle {
235        MetricHandle {
236            reporter: self.reporter.clone(),
237            selection,
238            cadence,
239        }
240    }
241
242    // ---- cadence_window: last closed snapshot --------------------------
243
244    /// Latest fully-closed snapshot for the named cadence, filtered
245    /// by `selection`. Returns an empty snapshot when no closed
246    /// window has been published yet (early in a run).
247    ///
248    /// Walks every component tracked by the cadence reporter,
249    /// merging matching metrics into one result. Identity follows
250    /// OpenMetrics §4.5.1 — same `(family.name, label_set)` combines.
251    pub fn cadence_window(&self, cadence: Duration, selection: &Selection) -> MetricSet {
252        let mut out = MetricSet::at(Instant::now(), cadence);
253        for component in self.reporter.component_labels() {
254            let Some(snap) = self.reporter.latest(&component, cadence) else {
255                continue;
256            };
257            for family in snap.families() {
258                if !selection.matches_family(family.name()) {
259                    continue;
260                }
261                for metric in family.metrics() {
262                    if !selection.matches_labels(metric.labels()) {
263                        continue;
264                    }
265                    insert_metric_into(&mut out, family, metric);
266                }
267            }
268        }
269        out
270    }
271
272    // ---- increase_over / distribution_over: sliding span derivations ----
273
274    /// The finest declared cadence whose retained ring (`HISTORY_RING_CAP`
275    /// windows) covers `span`, so a sliding lookback stays at the finest
276    /// available resolution. `None` if no cadence is declared.
277    fn finest_cadence_covering(&self, span: Duration) -> Option<Duration> {
278        let cap = crate::cadence_reporter::HISTORY_RING_CAP as u128;
279        let mut chosen: Option<Duration> = None;
280        for layer in self.reporter.layers() {
281            if layer.hidden {
282                continue;
283            }
284            chosen = Some(layer.interval);
285            if layer.interval.as_nanos().saturating_mul(cap) >= span.as_nanos() {
286                break;
287            }
288        }
289        chosen
290    }
291
292    /// The counter **increase** over the trailing `span` (PromQL `increase`):
293    /// for each matched counter, `cum[now] − cum[now−span]` differenced from
294    /// the retained finest-cadence ring at the finest covering resolution — a
295    /// continuous/sliding window (e.g. "the last 10 s" off the 1 s ring).
296    /// Emits **counters only**, and their values are DELTAS — contrast the
297    /// running-total readers [`Self::now`] / [`Self::cadence_window`] /
298    /// [`Self::session_lifetime`]. A rate is `increase / span`. For the recent
299    /// latency/value distribution use [`Self::distribution_over`].
300    pub fn increase_over(&self, span: Duration, selection: &Selection) -> MetricSet {
301        let mut out = MetricSet::at(Instant::now(), span);
302        let Some(cadence) = self.finest_cadence_covering(span) else {
303            return out;
304        };
305        let windows = ((span.as_nanos().max(1)) / (cadence.as_nanos().max(1))).max(1) as usize;
306        let now = Instant::now();
307        for component in self.reporter.component_labels() {
308            let ring = self.reporter.ring(&component, cadence);
309            if ring.is_empty() {
310                continue;
311            }
312            let end = ring.len();
313            let start = end.saturating_sub(windows);
314            let per = coalesce_component_windows(
315                ring[start..end].iter().map(|a| a.as_ref()),
316                selection,
317                now,
318                span,
319            );
320            // Baseline = the running total just BEFORE the span (window at
321            // start-1); subtracting it turns a counter's window-end cumulative
322            // into the span increase. `None` (→ 0) at the ring's start.
323            let baseline = start.checked_sub(1).map(|i| ring[i].clone());
324            for family in per.families() {
325                for metric in family.metrics() {
326                    insert_counter_increase_into(&mut out, family, metric, baseline.as_deref());
327                }
328            }
329        }
330        out
331    }
332
333    /// The merged latency/value **distribution** over the trailing `span`:
334    /// for each matched histogram, the HDR reservoir merged across the windows
335    /// in the span (read `p50`/`p99`/`mean` from it — the PromQL `*_over_time`
336    /// quantile family). Same finest-covering-resolution sliding window as
337    /// [`Self::increase_over`]. Emits **histograms only** — the recent
338    /// distribution, never a counter increase.
339    pub fn distribution_over(&self, span: Duration, selection: &Selection) -> MetricSet {
340        let mut out = MetricSet::at(Instant::now(), span);
341        let Some(cadence) = self.finest_cadence_covering(span) else {
342            return out;
343        };
344        let windows = ((span.as_nanos().max(1)) / (cadence.as_nanos().max(1))).max(1) as usize;
345        let now = Instant::now();
346        for component in self.reporter.component_labels() {
347            let ring = self.reporter.ring(&component, cadence);
348            if ring.is_empty() {
349                continue;
350            }
351            let end = ring.len();
352            let start = end.saturating_sub(windows);
353            let per = coalesce_component_windows(
354                ring[start..end].iter().map(|a| a.as_ref()),
355                selection,
356                now,
357                span,
358            );
359            for family in per.families() {
360                for metric in family.metrics() {
361                    if matches!(
362                        metric.point().map(|p| p.value()),
363                        Some(MetricValue::Histogram(_)) | Some(MetricValue::BucketedHistogram(_))
364                    ) {
365                        insert_metric_into(&mut out, family, metric);
366                    }
367                }
368            }
369        }
370        out
371    }
372
373    // ---- session_lifetime: full canonical span ------------------------
374
375    /// Full canonical session span as of *now*, filtered by
376    /// `selection`. Walks the cascade *down* at read time:
377    ///
378    /// Per component, COALESCEs the cascade's disjoint time-slices — every
379    /// layer's in-flight prebuffer plus the largest cadence's retained
380    /// accumulator (the lifetime buffer) and its last-closed window — then
381    /// AGGREGATEs the per-component results across components. Coalescing
382    /// the time dimension keeps the latest `cumulative` and merges
383    /// reservoirs, so a session-cumulative counter is the latest running
384    /// total — not multiplied by the number of cascade sources it appears in.
385    ///
386    /// Per SRD-42 §"Cost rule for recent_window", only matched metric
387    /// instances combine — same shape as `increase_over` / `distribution_over`.
388    pub fn session_lifetime(&self, selection: &Selection) -> MetricSet {
389        let session_age = self.reporter.started_at().elapsed();
390        let now = Instant::now();
391        let mut out = MetricSet::at(now, session_age);
392        let largest = self.reporter.layers().last().map(|l| l.interval);
393
394        for component in self.reporter.component_labels() {
395            // Gather this component's disjoint cascade sources.
396            let mut sources: Vec<MetricSet> = Vec::new();
397            for layer in self.reporter.layers() {
398                if let Some(pre) = self.reporter.prebuffer(&component, layer.interval) {
399                    sources.push(pre);
400                }
401                // Only the LARGEST cadence's last-closed window is read: a
402                // smaller layer's closed window already folded into the next
403                // layer's prebuffer (reading it too would double-count); the
404                // largest cadence has no parent to fold into. (The earlier
405                // `now`/closed-window top-up re-read a promoted window — the
406                // source of the session-cumulative overcount.)
407                if Some(layer.interval) == largest
408                    && let Some(latest) = self.reporter.latest(&component, layer.interval)
409                {
410                    sources.push((*latest).clone());
411                }
412            }
413            if sources.is_empty() {
414                continue;
415            }
416            // Coalesce the time dimension, then aggregate across components.
417            let per = coalesce_component_windows(sources.iter(), selection, now, session_age);
418            for family in per.families() {
419                for metric in family.metrics() {
420                    insert_metric_into(&mut out, family, metric);
421                }
422            }
423        }
424
425        out
426    }
427
428    // ---- expect-exactly-one helpers ------------------------------------
429
430    /// Run a query mode and assert exactly one matching `Metric` per
431    /// the SRD's "specific metric" semantics. Returns `Err` if 0 or
432    /// >1 matches.
433    pub fn select_one<F>(&self, mode: F) -> Result<MetricSet, SelectError>
434    where
435        F: FnOnce(&Self) -> MetricSet,
436    {
437        let snap = mode(self);
438        let total: usize = snap.families().map(|f| f.len()).sum();
439        match total {
440            0 => Err(SelectError::NoMatch),
441            1 => Ok(snap),
442            n => Err(SelectError::MultipleMatches(n)),
443        }
444    }
445}
446
447/// Memoized pull handle — resolved once via [`MetricsQuery::resolve`],
448/// then reused for cheap per-draw / per-frame reads.
449///
450/// The handle caches the `(selection, cadence)` pair the caller
451/// asked for. Each `read_now` issues a `cadence_window` query
452/// against the reporter — O(components) per call, no component-tree
453/// walk, no instrument access.
454///
455/// Per SRD-42 (revised), "recent info" always routes through the
456/// cadence-reporter store, never through live instruments. The
457/// handle's `read_now` reads the smallest declared cadence's last
458/// closed window; callers who need true sub-second live data for a
459/// specific Timer should use
460/// [`crate::instruments::timer::Timer::peek_live_window`] directly.
461pub struct MetricHandle {
462    reporter: Arc<CadenceReporter>,
463    selection: Selection,
464    cadence: Duration,
465}
466
467impl MetricHandle {
468    /// Read the last closed window at this handle's cadence,
469    /// filtered by the handle's selection. Non-mutating, safe to
470    /// call arbitrarily often.
471    pub fn read_now(&self) -> MetricSet {
472        let mut out = MetricSet::at(Instant::now(), self.cadence);
473        for component in self.reporter.component_labels() {
474            let Some(snap) = self.reporter.latest(&component, self.cadence) else {
475                continue;
476            };
477            for family in snap.families() {
478                if !self.selection.matches_family(family.name()) {
479                    continue;
480                }
481                for metric in family.metrics() {
482                    if !self.selection.matches_labels(metric.labels()) {
483                        continue;
484                    }
485                    insert_metric_into(&mut out, family, metric);
486                }
487            }
488        }
489        out
490    }
491
492    /// No-op retained for API compatibility with callers that
493    /// expect to "refresh" after phase transitions. The new handle
494    /// reads through the cadence reporter's store, which already
495    /// reflects the current set of tracked components — no resync
496    /// needed.
497    pub fn refresh(&mut self) {}
498
499    /// The selection this handle was resolved against.
500    pub fn selection(&self) -> &Selection {
501        &self.selection
502    }
503
504    /// Cadence this handle reads from (the smallest declared
505    /// cadence at resolve time).
506    pub fn cadence(&self) -> Duration {
507        self.cadence
508    }
509
510    /// Number of components currently tracked by the reporter —
511    /// informational. Not cached; queried fresh each call.
512    pub fn source_count(&self) -> usize {
513        self.reporter.component_labels().len()
514    }
515}
516
517/// Insert one `(family, metric)` pair into `out`, merging an existing
518/// same-identity entry (OpenMetrics §4.5.1) as a **cross-component
519/// aggregate** ([`CombineMode::Aggregate`]): counter `cumulative` and
520/// histogram `cumulative_count` SUM across the matching series from
521/// different components.
522fn insert_metric_into(out: &mut MetricSet, family: &MetricFamily, metric: &Metric) {
523    insert_metric_with_mode(out, family, metric, crate::snapshot::CombineMode::Aggregate);
524}
525
526/// Insert one `(family, metric)` pair into `out`, merging an existing
527/// same-identity entry under `mode`. Use [`CombineMode::Coalesce`] to
528/// fold the **time dimension** (consecutive windows / cascade layers of
529/// one series: keep the latest `cumulative`, merge reservoirs) and
530/// [`CombineMode::Aggregate`] to fold **across components** (sum). Mixing
531/// them up double-counts a cumulative value — see `session_lifetime` /
532/// `increase_over`.
533fn insert_metric_with_mode(
534    out: &mut MetricSet,
535    family: &MetricFamily,
536    metric: &Metric,
537    mode: crate::snapshot::CombineMode,
538) {
539    let Some(point) = metric.point() else { return };
540    let existing = out
541        .family(family.name())
542        .and_then(|f| f.metric_with_labels(metric.labels()))
543        .is_some();
544    if existing {
545        // Combine into existing — requires owned access via a
546        // rebuild. Simpler pattern: drop the existing family and
547        // rebuild with a coalesced replacement. For Phase 7 v1 we
548        // accept the cost and use `MetricSet::coalesce` over a
549        // two-element slice.
550        let mut tmp = MetricSet::at(out.captured_at(), out.interval());
551        tmp.insert_metric(
552            family.name().to_string(),
553            family.r#type(),
554            metric.labels().clone(),
555            point.value().clone(),
556            point.timestamp().unwrap_or(out.captured_at()),
557        );
558        let merged = MetricSet::coalesce_with_mode(
559            std::slice::from_ref(out)
560                .iter()
561                .chain(std::slice::from_ref(&tmp).iter())
562                .cloned()
563                .collect::<Vec<_>>()
564                .as_slice(),
565            mode,
566        );
567        *out = merged;
568    } else {
569        out.insert_metric(
570            family.name().to_string(),
571            family.r#type(),
572            metric.labels().clone(),
573            point.value().clone(),
574            point.timestamp().unwrap_or(out.captured_at()),
575        );
576    }
577}
578
579/// Fold one component's time-ordered window `sources` into a single
580/// per-component snapshot, COALESCING the time dimension (keep the latest
581/// `cumulative`, merge reservoirs) over the matched selection only. The
582/// caller then AGGREGATEs the result across components. This split keeps a
583/// cumulative counter at its latest running total rather than multiplied
584/// by the number of cascade sources it appears in.
585fn coalesce_component_windows<'a>(
586    sources: impl IntoIterator<Item = &'a MetricSet>,
587    selection: &Selection,
588    captured_at: Instant,
589    interval: Duration,
590) -> MetricSet {
591    let mut per = MetricSet::at(captured_at, interval);
592    for src in sources {
593        for family in src.families() {
594            if !selection.matches_family(family.name()) {
595                continue;
596            }
597            for metric in family.metrics() {
598                if !selection.matches_labels(metric.labels()) {
599                    continue;
600                }
601                insert_metric_with_mode(
602                    &mut per,
603                    family,
604                    metric,
605                    crate::snapshot::CombineMode::Coalesce,
606                );
607            }
608        }
609    }
610    per
611}
612
613/// Insert a counter's span **increase** into `out`, converting its
614/// window-end cumulative into `cum[end] − cum[before span]` by subtracting
615/// the `baseline` running total (the window just before the span), then
616/// aggregating (summing) across components. **Counters only** — non-counter
617/// points are skipped (the distribution lives in `distribution_over`). This
618/// is the consumer-side derivation `increase_over` applies.
619fn insert_counter_increase_into(
620    out: &mut MetricSet,
621    family: &MetricFamily,
622    metric: &Metric,
623    baseline: Option<&MetricSet>,
624) {
625    let Some(point) = metric.point() else { return };
626    let MetricValue::Counter(c) = point.value() else {
627        return;
628    };
629    let base = baseline
630        .and_then(|b| b.family(family.name()))
631        .and_then(|f| f.metric_with_labels(metric.labels()))
632        .and_then(|m| m.point())
633        .and_then(|p| match p.value() {
634            MetricValue::Counter(bc) => Some(bc.cumulative),
635            _ => None,
636        })
637        .unwrap_or(0);
638    let increase = c.cumulative.saturating_sub(base);
639    let inc = Metric::single(
640        metric.labels().clone(),
641        MetricPoint::new(
642            MetricValue::Counter(CounterValue::new(increase)),
643            point.timestamp().unwrap_or(out.captured_at()),
644        ),
645    );
646    insert_metric_into(out, family, &inc);
647}
648
649// =========================================================================
650// Tests
651// =========================================================================
652
653#[cfg(test)]
654mod tests {
655    use super::*;
656    use crate::cadence::{CadenceTree, Cadences};
657    use crate::component::{Component, ComponentState, InstrumentRef, attach};
658    use crate::instruments::counter::Counter;
659    use crate::snapshot::MetricValue;
660    use std::collections::HashMap;
661
662    fn build_one_component_query() -> (Arc<RwLock<Component>>, Arc<CadenceReporter>, MetricsQuery) {
663        let root = Component::root(Labels::of("session", "s1"), HashMap::new());
664        let phase = Arc::new(RwLock::new(Component::new(
665            Labels::of("phase", "load"),
666            HashMap::new(),
667        )));
668        attach(&root, &phase);
669        {
670            let mut p = phase.write().unwrap();
671            p.set_state(ComponentState::Running);
672            let counter = Arc::new(Counter::new(Labels::of("name", "ops")));
673            counter.inc_by(7);
674            p.register_instrument("ops", InstrumentRef::Counter(counter))
675                .unwrap();
676        }
677
678        let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
679        let tree = CadenceTree::plan_default(cadences);
680        let reporter = Arc::new(CadenceReporter::new(tree));
681        let query = MetricsQuery::new(reporter.clone(), root.clone());
682        (root, reporter, query)
683    }
684
685    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
686    async fn now_reads_smallest_cadence_window() {
687        // `now` no longer walks the live tree — it reads
688        // `cadence_window(smallest_declared)`. So we must ingest a
689        // closed window first; before any close, `now` is empty.
690        let (_root, reporter, query) = build_one_component_query();
691        assert!(
692            query.now(&Selection::family("ops")).is_empty(),
693            "pre-close now should be empty"
694        );
695
696        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
697        let mut s = MetricSet::new(Duration::from_millis(100));
698        s.insert_counter("ops", Labels::default(), 42, Instant::now());
699        reporter.ingest(&labels, s);
700        reporter.flush_for_tests();
701
702        let snap = query.now(&Selection::family("ops"));
703        let total = match snap
704            .family("ops")
705            .unwrap()
706            .metrics()
707            .next()
708            .unwrap()
709            .point()
710            .unwrap()
711            .value()
712        {
713            MetricValue::Counter(c) => c.cumulative,
714            _ => panic!("not a counter"),
715        };
716        assert_eq!(total, 42);
717    }
718
719    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
720    async fn cadence_window_returns_latest_closed_snapshot() {
721        let (_root, reporter, query) = build_one_component_query();
722        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
723        // Inject a closed window via the reporter.
724        let mut s = MetricSet::new(Duration::from_millis(100));
725        s.insert_counter("ops", Labels::default(), 99, Instant::now());
726        reporter.ingest(&labels, s);
727        reporter.flush_for_tests();
728
729        let snap = query.cadence_window(Duration::from_millis(100), &Selection::family("ops"));
730        let f = snap
731            .family("ops")
732            .expect("ops family in cadence_window result");
733        match f.metrics().next().unwrap().point().unwrap().value() {
734            MetricValue::Counter(c) => assert_eq!(c.cumulative, 99),
735            _ => panic!("not a counter"),
736        }
737    }
738
739    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
740    async fn session_lifetime_does_not_overcount_a_single_counter() {
741        // session_lifetime walks the cascade down (every layer's prebuffer
742        // + the largest's latest) and combines. With ONE ingested counter
743        // the canonical value is its cumulative — not a multiple from the
744        // same value appearing in several cascade sources.
745        let (_root, reporter, query) = build_one_component_query();
746        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
747        let mut s = MetricSet::new(Duration::from_millis(100));
748        s.insert_counter("ops", Labels::default(), 42, Instant::now());
749        reporter.ingest(&labels, s);
750        reporter.flush_for_tests();
751
752        let snap = query.session_lifetime(&Selection::family("ops"));
753        let cumulative = match snap
754            .family("ops")
755            .unwrap()
756            .metrics()
757            .next()
758            .unwrap()
759            .point()
760            .unwrap()
761            .value()
762        {
763            MetricValue::Counter(c) => c.cumulative,
764            _ => panic!("not a counter"),
765        };
766        assert_eq!(
767            cumulative, 42,
768            "session_lifetime cumulative overcounted (got {cumulative})"
769        );
770    }
771
772    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
773    async fn increase_over_gives_the_span_increment() {
774        // Two windows of a counter — cumulative 10 then 20, no prior data.
775        // The recent-span value is the increment over the span (cum[end] −
776        // cum[before span] = 20 − 0 = 20), derived from the running totals.
777        let (_root, reporter, query) = build_one_component_query();
778        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
779
780        let mut s1 = MetricSet::new(Duration::from_millis(100));
781        s1.insert_counter("ops", Labels::default(), 10, Instant::now());
782        reporter.ingest(&labels, s1);
783        reporter.flush_for_tests();
784        std::thread::sleep(Duration::from_millis(2)); // distinct window timestamps
785        let mut s2 = MetricSet::new(Duration::from_millis(100));
786        s2.insert_counter("ops", Labels::default(), 20, Instant::now());
787        reporter.ingest(&labels, s2);
788        reporter.flush_for_tests();
789
790        let snap = query.increase_over(Duration::from_millis(250), &Selection::family("ops"));
791        let value = match snap
792            .family("ops")
793            .unwrap()
794            .metrics()
795            .next()
796            .unwrap()
797            .point()
798            .unwrap()
799            .value()
800        {
801            MetricValue::Counter(c) => c.cumulative,
802            _ => panic!("not a counter"),
803        };
804        assert_eq!(
805            value, 20,
806            "span increment over the recent window (got {value})"
807        );
808    }
809
810    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
811    async fn increase_over_subtracts_prior_cumulative() {
812        // A counter climbing 100→110→120 over three windows. The recent
813        // window's value is the INCREASE over the span — the running total
814        // BEFORE the span is subtracted — not the latest cumulative.
815        let (_root, reporter, query) = build_one_component_query();
816        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
817        for v in [100u64, 110, 120] {
818            let mut s = MetricSet::new(Duration::from_millis(100));
819            s.insert_counter("ops", Labels::default(), v, Instant::now());
820            reporter.ingest(&labels, s);
821            reporter.flush_for_tests();
822            std::thread::sleep(Duration::from_millis(2));
823        }
824        let read = |span| match query
825            .increase_over(span, &Selection::family("ops"))
826            .family("ops")
827            .unwrap()
828            .metrics()
829            .next()
830            .unwrap()
831            .point()
832            .unwrap()
833            .value()
834        {
835            MetricValue::Counter(c) => c.cumulative,
836            _ => panic!("not a counter"),
837        };
838        assert_eq!(
839            read(Duration::from_millis(100)),
840            10,
841            "last window increase = 120−110"
842        );
843        assert_eq!(
844            read(Duration::from_millis(200)),
845            20,
846            "last two windows increase = 120−100"
847        );
848    }
849
850    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
851    async fn distribution_over_merges_histogram_windows() {
852        use hdrhistogram::Histogram as HdrHistogram;
853        // Two windows of a 2-sample histogram; the recent distribution over a
854        // span covering both is the MERGED reservoir (4 samples).
855        let (_root, reporter, query) = build_one_component_query();
856        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
857        for _ in 0..2 {
858            let mut s = MetricSet::new(Duration::from_millis(100));
859            let mut h = HdrHistogram::<u64>::new_with_bounds(1, 3_600_000_000_000, 3).unwrap();
860            h.record(1_000_000).unwrap();
861            h.record(2_000_000).unwrap();
862            s.insert_histogram("latency", Labels::default(), h, Instant::now());
863            reporter.ingest(&labels, s);
864            reporter.flush_for_tests();
865            std::thread::sleep(Duration::from_millis(2));
866        }
867        let snap =
868            query.distribution_over(Duration::from_millis(250), &Selection::family("latency"));
869        let count = match snap
870            .family("latency")
871            .unwrap()
872            .metrics()
873            .next()
874            .unwrap()
875            .point()
876            .unwrap()
877            .value()
878        {
879            MetricValue::Histogram(h) => h.count,
880            _ => panic!("not a histogram"),
881        };
882        assert_eq!(count, 4, "two windows of 2 samples merge to 4");
883    }
884
885    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
886    async fn selection_filter_excludes_non_matching_labels() {
887        let (_root, reporter, query) = build_one_component_query();
888        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
889        let mut s = MetricSet::new(Duration::from_millis(100));
890        s.insert_counter("ops", Labels::of("kind", "a"), 5, Instant::now());
891        s.insert_counter("ops", Labels::of("kind", "b"), 9, Instant::now());
892        reporter.ingest(&labels, s);
893        reporter.flush_for_tests();
894
895        let snap = query.cadence_window(
896            Duration::from_millis(100),
897            &Selection::family("ops").with_label("kind", "b"),
898        );
899        let f = snap.family("ops").expect("ops family");
900        assert_eq!(f.len(), 1);
901        match f.metrics().next().unwrap().point().unwrap().value() {
902            MetricValue::Counter(c) => assert_eq!(c.cumulative, 9),
903            _ => panic!("not a counter"),
904        }
905    }
906
907    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
908    async fn select_one_errors_on_zero_matches() {
909        let (_root, _reporter, query) = build_one_component_query();
910        let result = query.select_one(|q| {
911            q.cadence_window(
912                Duration::from_millis(100),
913                &Selection::family("nonexistent"),
914            )
915        });
916        assert_eq!(result.unwrap_err(), SelectError::NoMatch);
917    }
918
919    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
920    async fn select_one_succeeds_on_exact_match() {
921        let (_root, reporter, query) = build_one_component_query();
922        let labels = Labels::of("session", "s1").extend(&Labels::of("phase", "load"));
923        let mut s = MetricSet::new(Duration::from_millis(100));
924        s.insert_counter("ops", Labels::default(), 1, Instant::now());
925        reporter.ingest(&labels, s);
926        reporter.flush_for_tests();
927
928        let result = query.select_one(|q| {
929            q.cadence_window(Duration::from_millis(100), &Selection::family("ops"))
930        });
931        assert!(result.is_ok());
932    }
933}