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}