Skip to main content

nmbrs_metrics/queryapi/
mod.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! The metrics **query API** — the data-access *service* boundary
5//! (SRD-86 §"The metric-reader surface").
6//!
7//! `nmbrs-metrics` is the foundational data-access library; this module
8//! exposes its query surface as a service. It provides the native
9//! result shape ([`Vector`] — multiple [`Series`], each with sample
10//! points), the selector semantics ([`Matcher`]), and a small access
11//! contract ([`MetricAccess`]) whose signatures map **1:1 onto a
12//! MetricsQL parsed selector**:
13//!
14//! - a bare / instant selector → [`MetricAccess::select_instant`];
15//! - a range selector (`m[w]`) → [`MetricAccess::select_range`].
16//!
17//! Deliberately **not** here: aggregation, rollups, arithmetic — the
18//! "bells and whistles" of the query language. Those stay in the
19//! MetricsQL engine, layered *over* this access surface. Keeping the
20//! contract this thin is what lets the engine sit on any data service.
21//!
22//! ## Service location (a "data service object")
23//!
24//! Consumers (the MetricsQL engine) **locate** a service at runtime
25//! rather than binding a concrete impl:
26//!
27//! - the **live in-process** service wraps the session's `MetricsQuery`
28//!   (which is not static), so the runner [`install_live_access`]es it
29//!   per session and consumers read it via [`live_access`];
30//! - **file / external** backends (e.g. the sqlite reader) register an
31//!   [`AccessProvider`] via `inventory`, so a consumer can [`provider`]
32//!   one by scheme and open it — without the engine depending on the
33//!   backend's crate or features.
34
35pub mod catalog;
36mod hybrid;
37mod live;
38mod shapes;
39#[cfg(feature = "sqlite")]
40pub mod sqlite;
41
42pub use catalog::{
43    CachedCatalog, ExemplarPoint, LabelSet, MetricCatalog, MetricFamilyMeta, MetricType,
44};
45pub use hybrid::{HorizonAware, HybridStore, Tier};
46pub use live::MetricsQueryAccess;
47pub use shapes::{MatchOp, Matcher, Sample, Series, Vector};
48
49use std::sync::{Arc, LazyLock, OnceLock};
50
51use arc_swap::ArcSwapOption;
52
53/// Error from a metrics access backend. Backends own their taxonomy;
54/// the engine treats these as opaque from a flow-control standpoint.
55#[derive(Debug, Clone)]
56pub struct QueryError {
57    pub message: String,
58}
59
60impl QueryError {
61    pub fn new(msg: impl Into<String>) -> Self {
62        Self {
63            message: msg.into(),
64        }
65    }
66}
67
68impl std::fmt::Display for QueryError {
69    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
70        write!(f, "metrics query: {}", self.message)
71    }
72}
73
74impl std::error::Error for QueryError {}
75
76/// The metrics data-access **service**. Backends implement it; a query
77/// layer locates one at runtime and reads `Vector`s through it. See the
78/// module docs for the access/aggregation cut line.
79pub trait MetricAccess: Send + Sync {
80    /// Range selection: series matching every `matcher`, with every
81    /// sample in `[start_ms, end_ms]` (ascending). Yields a range
82    /// vector — the input a rollup (`rate(m[w])`, `*_over_time`)
83    /// consumes. This is the one required method; [`select_instant`]
84    /// derives from it.
85    ///
86    /// [`select_instant`]: MetricAccess::select_instant
87    fn select_range(
88        &self,
89        matchers: &[Matcher],
90        start_ms: i64,
91        end_ms: i64,
92    ) -> Result<Vector, QueryError>;
93
94    /// Instant selection: series matching every `matcher`, each reduced
95    /// to its latest sample within `[at_ms - lookback, at_ms]` (the
96    /// PromQL stale-tolerance window; `lookback_ms = None` means a
97    /// strict `[at_ms, at_ms]`). Yields an instant vector — one sample
98    /// per series.
99    ///
100    /// Default: `select_range` over the lookback window, then the
101    /// latest sample per series. A backend can override if it can do
102    /// the reduction more cheaply (e.g. push it into SQL).
103    fn select_instant(
104        &self,
105        matchers: &[Matcher],
106        at_ms: i64,
107        lookback_ms: Option<i64>,
108    ) -> Result<Vector, QueryError> {
109        let start_ms = at_ms - lookback_ms.unwrap_or(0);
110        let range = self.select_range(matchers, start_ms, at_ms)?;
111        let reduced = range
112            .into_series()
113            .into_iter()
114            .filter_map(|s| {
115                s.samples.last().copied().map(|last| Series {
116                    labels: s.labels,
117                    samples: vec![last],
118                })
119            })
120            .collect();
121        Ok(Vector::new(reduced))
122    }
123}
124
125// ---------------------------------------------------------------------
126// Runtime service location
127// ---------------------------------------------------------------------
128
129/// Sized holder for the trait-object service, so it can live in an
130/// `ArcSwapOption` (whose pointee must be `Sized`; a bare `dyn MetricAccess`
131/// is not).
132struct LiveHolder(Arc<dyn MetricAccess>);
133
134/// The live in-process access service for the current session. Wraps a
135/// per-session `MetricsQuery`, so it's *installed* (not static).
136///
137/// SRD-90 §M4 — an `ArcSwapOption`, not a `Mutex`: `live_access()` is on the
138/// metricsql read hot path (every settle pulse / TUI refresh resolves it), so
139/// the read is a single lock-free atomic load, never a mutex acquire that could
140/// contend under concurrent executions.
141static LIVE: LazyLock<ArcSwapOption<LiveHolder>> = LazyLock::new(ArcSwapOption::empty);
142
143/// Install the live in-process access service. Called once by the
144/// runner when the session's `MetricsQuery` is built.
145pub fn install_live_access(service: Arc<dyn MetricAccess>) {
146    LIVE.store(Some(Arc::new(LiveHolder(service))));
147}
148
149/// Drop the installed live-access service and everything it owns —
150/// crucially the HybridStore's sqlite "cold tier" reader connection on
151/// `metrics.db`. The runner calls this at session shutdown BEFORE the
152/// SQLite reporter consolidates the WAL: `PRAGMA journal_mode=DELETE` needs
153/// an EXCLUSIVE lock on the db, which a still-open reader connection on the
154/// same file blocks ("database is locked"). `swap` + explicit `drop` so the
155/// holder's Arc chain is released here rather than deferred to the next
156/// `store`. Idempotent; a no-op if nothing was installed.
157pub fn uninstall_live_access() {
158    drop(LIVE.swap(None));
159}
160
161/// The live in-process access service, if a session has installed one.
162/// Lock-free: one atomic `ArcSwap` load.
163pub fn live_access() -> Option<Arc<dyn MetricAccess>> {
164    LIVE.load_full().map(|h| h.0.clone())
165}
166
167/// Resolves the **reading execution's** `exec_id`, so a live metric read
168/// can scope itself to its own execution's series instead of every
169/// execution sharing the session (SRD-88 encapsulation — without this, an
170/// optimizer's `sum(rate(errors_total[…]))` would sum a concurrent
171/// neighbour's errors too). `exec_id` lives in nmbrs-runtime's task-local
172/// `ExecutionContext`, a layer above this crate, so the runtime installs
173/// a small resolver hook here. `None` ⇒ no scope (single-run / outside any
174/// execution — read everything, A1).
175static READ_EXEC_ID_HOOK: OnceLock<fn() -> Option<u64>> = OnceLock::new();
176
177/// Install the reading-execution `exec_id` resolver (idempotent — first
178/// wins). The runtime calls this once with a fn that reads its task-local
179/// execution context.
180pub fn install_read_exec_id_hook(hook: fn() -> Option<u64>) {
181    let _ = READ_EXEC_ID_HOOK.set(hook);
182}
183
184/// The reading execution's `exec_id`, if a hook is installed and a scope
185/// is active. `None` ⇒ live reads are unscoped (A1 single-run).
186pub fn current_read_exec_id() -> Option<u64> {
187    READ_EXEC_ID_HOOK.get().and_then(|h| h())
188}
189
190/// SRD-89 §3b / SRD-90 §M6 — scope every read to its **reading execution** by
191/// injecting `exec_id` as a uniform **dimensional-label matcher**, so each
192/// interior backend applies it wherever `exec_id` lives (the in-memory tier's
193/// label set, the sqlite tier's `exec_id` column) with the same value. This
194/// replaces the per-backend special-casing (the live tier's bespoke post-filter,
195/// a sqlite selection mode) with one mechanism: `exec_id` is just a label.
196///
197/// `None` (single-run / outside any execution scope) ⇒ no injection — the read
198/// is unscoped and sees the sole execution's data, byte-identical to before
199/// (axiom A1). If the caller already constrained `exec_id`, nothing is injected.
200pub struct ExecScopedAccess {
201    inner: Arc<dyn MetricAccess>,
202}
203
204impl ExecScopedAccess {
205    pub fn new(inner: Arc<dyn MetricAccess>) -> Self {
206        Self { inner }
207    }
208
209    /// The matcher set with the reading execution's `exec_id` injected (when a
210    /// scope is active and the caller hasn't already constrained it).
211    fn scoped(&self, matchers: &[Matcher]) -> Option<Vec<Matcher>> {
212        let id = current_read_exec_id()?;
213        if matchers.iter().any(|m| m.label == "exec_id") {
214            return None; // caller already scoped — don't double-inject
215        }
216        let mut v = matchers.to_vec();
217        v.push(Matcher::eq("exec_id", &id.to_string()));
218        Some(v)
219    }
220}
221
222impl MetricAccess for ExecScopedAccess {
223    fn select_range(
224        &self,
225        matchers: &[Matcher],
226        start_ms: i64,
227        end_ms: i64,
228    ) -> Result<Vector, QueryError> {
229        match self.scoped(matchers) {
230            Some(scoped) => self.inner.select_range(&scoped, start_ms, end_ms),
231            None => self.inner.select_range(matchers, start_ms, end_ms),
232        }
233    }
234
235    fn select_instant(
236        &self,
237        matchers: &[Matcher],
238        at_ms: i64,
239        lookback_ms: Option<i64>,
240    ) -> Result<Vector, QueryError> {
241        match self.scoped(matchers) {
242            Some(scoped) => self.inner.select_instant(&scoped, at_ms, lookback_ms),
243            None => self.inner.select_instant(matchers, at_ms, lookback_ms),
244        }
245    }
246}
247
248/// A pluggable access backend, discovered at runtime via `inventory`.
249/// A consumer opens a service for a scheme-specific `target` (e.g. a
250/// db path). The sqlite reader registers one; future backends can too,
251/// without the query engine depending on them.
252pub struct AccessProvider {
253    /// The scheme this provider answers for (e.g. `"sqlite"`).
254    pub scheme: &'static str,
255    /// Open an access service for `target` (scheme-specific).
256    pub open: fn(target: &str) -> Result<Box<dyn MetricAccess>, QueryError>,
257}
258
259inventory::collect!(AccessProvider);
260
261/// Locate a registered [`AccessProvider`] by scheme.
262pub fn provider(scheme: &str) -> Option<&'static AccessProvider> {
263    inventory::iter::<AccessProvider>
264        .into_iter()
265        .find(|p| p.scheme == scheme)
266}
267
268#[cfg(test)]
269mod tests {
270    use super::*;
271
272    #[test]
273    fn live_access_installs_and_reads_back() {
274        struct Stub;
275        impl MetricAccess for Stub {
276            fn select_instant(
277                &self,
278                _: &[Matcher],
279                _: i64,
280                _: Option<i64>,
281            ) -> Result<Vector, QueryError> {
282                Ok(Vector::default())
283            }
284            fn select_range(&self, _: &[Matcher], _: i64, _: i64) -> Result<Vector, QueryError> {
285                Ok(Vector::default())
286            }
287        }
288        install_live_access(Arc::new(Stub));
289        assert!(live_access().is_some());
290    }
291
292    #[test]
293    fn unknown_provider_scheme_is_none() {
294        assert!(provider("no-such-scheme").is_none());
295    }
296}