nmbrs_metrics/queryapi/catalog.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Catalog: enumerable backend introspection for metrics.
5//!
6//! [`MetricAccess`](super::MetricAccess) is for **fetching
7//! values** given a selector. This module is for **enumerating
8//! what's available**: which metric families exist, what
9//! labels they carry, what values those labels take. The split
10//! mirrors VictoriaMetrics' own surface (`/api/v1/query` and
11//! `/api/v1/labels` endpoints); a remote backend implements
12//! both traits against the same wire protocol.
13//!
14//! ## Why not put enumeration on `DataSource`?
15//!
16//! Two reasons:
17//!
18//! 1. **Separable lifetimes.** Catalog data is small (≤ a few
19//! thousand entries even on big sessions) and slow-changing.
20//! Sample data is large and time-windowed. Cache strategies
21//! diverge — see [`CachedCatalog`].
22//! 2. **Implementation effort.** A backend may serve catalog
23//! enumeration from cheaply-aggregable indices but require
24//! significant query work for sample fetching. Letting an
25//! impl provide one without the other is the right composition.
26//!
27//! ## OpenMetrics fidelity
28//!
29//! [`MetricFamilyMeta`] mirrors the OpenMetrics 1.0 metadata
30//! envelope: name + type + unit + help. [`MetricType`]
31//! enumerates every type the OpenMetrics specification
32//! recognises; backends that use a smaller subset map their
33//! types into [`MetricType::Unknown`] when no faithful match
34//! exists.
35//!
36//! ## Autocompletion building block
37//!
38//! `nmbrs::completion` consumes this trait via the
39//! [`MetricCatalog`] object trait to surface metric / label
40//! suggestions for `--metric`, `--over`, `--by`, `--where`,
41//! and (planned) for autocompletion *inside metricsql
42//! expressions*. The four enumeration methods cover everything
43//! a metricsql-aware completion engine needs:
44//!
45//! | What user types | Catalog method called |
46//! |---|---|
47//! | bare metric name | [`MetricCatalog::metric_families`] |
48//! | inside `{...}`, key | [`MetricCatalog::label_keys`] |
49//! | inside `{...}`, value | [`MetricCatalog::label_values`] |
50//! | series-existence check | [`MetricCatalog::series`] |
51
52use std::sync::{Arc, Mutex};
53use std::time::{Duration, Instant};
54
55use super::{MatchOp as MatcherOp, Matcher, QueryError as DataSourceError};
56
57/// One metric family's OpenMetrics metadata. Identifies the
58/// family by name and surfaces its type / unit / help.
59///
60/// The OpenMetrics envelope:
61///
62/// ```text
63/// # TYPE process_cpu_seconds_total counter
64/// # UNIT process_cpu_seconds_total seconds
65/// # HELP process_cpu_seconds_total Total user and system CPU time spent in seconds.
66/// ```
67///
68/// becomes:
69///
70/// ```rust,ignore
71/// MetricFamilyMeta {
72/// name: "process_cpu_seconds_total".into(),
73/// ty: MetricType::Counter,
74/// unit: Some("seconds".into()),
75/// help: Some("Total user and system CPU time spent in seconds.".into()),
76/// }
77/// ```
78#[derive(Debug, Clone, PartialEq, Eq)]
79pub struct MetricFamilyMeta {
80 /// Family name (e.g. `http_requests_total`,
81 /// `nmbrs_phase_duration_seconds`). For histograms /
82 /// summaries this is the **family** name; the time-series
83 /// names emitted on the wire (`<name>_bucket`,
84 /// `<name>_sum`, `<name>_count`) are derived siblings, not
85 /// distinct families.
86 pub name: String,
87 /// Metric type per OpenMetrics 1.0.
88 pub ty: MetricType,
89 /// Optional unit (`seconds`, `bytes`, etc.). Used by
90 /// completion to surface a hint on the metric name, and
91 /// by the metricsql evaluator to validate
92 /// type-incompatible binary ops.
93 pub unit: Option<String>,
94 /// Optional help string, lifted verbatim from the
95 /// OpenMetrics `# HELP` metadata.
96 pub help: Option<String>,
97}
98
99/// Metric types per the OpenMetrics 1.0 specification.
100///
101/// Implementations whose underlying schema only knows the
102/// Prometheus subset (counter / gauge / histogram / summary /
103/// untyped) map their types into the corresponding variant
104/// here; the OpenMetrics 1.0 additions ([`Self::GaugeHistogram`],
105/// [`Self::Info`], [`Self::StateSet`]) become [`Self::Unknown`]
106/// in those backends.
107///
108/// The string form (`as_str` / `parse`) matches the
109/// OpenMetrics text-format keyword exactly. Round-trip
110/// parse/format is identity for every named variant.
111#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
112pub enum MetricType {
113 /// Monotonically-increasing counter; resets on process
114 /// restart. `rate()` and `increase()` operate on this.
115 Counter,
116 /// Instantaneous reading that can rise or fall.
117 Gauge,
118 /// Cumulative bucketed histogram with `_bucket{le=...}`,
119 /// `_sum`, `_count` siblings.
120 Histogram,
121 /// OpenMetrics 1.0: histogram whose buckets can decrease.
122 GaugeHistogram,
123 /// φ-quantile summary with `{quantile=...}` siblings plus
124 /// `_sum` and `_count`.
125 Summary,
126 /// Always-1 metric whose label set carries descriptive
127 /// info. Joins onto other metrics via `* on (labels)`.
128 Info,
129 /// Set of named states represented as 0/1 indicators.
130 /// OpenMetrics 1.0.
131 StateSet,
132 /// Type unspecified or unmappable — Prometheus's "untyped".
133 Unknown,
134}
135
136impl MetricType {
137 /// Canonical lower-case keyword (matches OpenMetrics
138 /// text format).
139 pub fn as_str(&self) -> &'static str {
140 match self {
141 MetricType::Counter => "counter",
142 MetricType::Gauge => "gauge",
143 MetricType::Histogram => "histogram",
144 MetricType::GaugeHistogram => "gaugehistogram",
145 MetricType::Summary => "summary",
146 MetricType::Info => "info",
147 MetricType::StateSet => "stateset",
148 MetricType::Unknown => "unknown",
149 }
150 }
151
152 /// Parse an OpenMetrics-style type keyword. Unknown
153 /// keywords map to [`Self::Unknown`] rather than failing —
154 /// a metric whose backend uses a non-standard type stays
155 /// queryable (the metricsql evaluator treats it as a
156 /// generic series) instead of disappearing.
157 pub fn parse(s: &str) -> MetricType {
158 match s.trim().to_ascii_lowercase().as_str() {
159 "counter" => MetricType::Counter,
160 "gauge" => MetricType::Gauge,
161 "histogram" => MetricType::Histogram,
162 "gaugehistogram" => MetricType::GaugeHistogram,
163 "summary" => MetricType::Summary,
164 "info" => MetricType::Info,
165 "stateset" => MetricType::StateSet,
166 _ => MetricType::Unknown,
167 }
168 }
169
170 /// True if this metric type implies derived time-series
171 /// names with the family-name prefix. For histogram +
172 /// summary the catalog's enumeration of *time series* will
173 /// surface `<name>_bucket`, `<name>_sum`, `<name>_count`
174 /// (and for summary: the family name itself). Selectors
175 /// can target any of these directly. Used by completion
176 /// to expand a family-name suggestion into its OpenMetrics
177 /// siblings when the user asks.
178 pub fn has_derived_series(&self) -> bool {
179 matches!(
180 self,
181 MetricType::Histogram | MetricType::GaugeHistogram | MetricType::Summary,
182 )
183 }
184
185 /// The label key OpenMetrics implies for this type's
186 /// derived series (`le` for histograms, `quantile` for
187 /// summaries). Returns `None` for types without an implied
188 /// label.
189 pub fn implied_label(&self) -> Option<&'static str> {
190 match self {
191 MetricType::Histogram | MetricType::GaugeHistogram => Some("le"),
192 MetricType::Summary => Some("quantile"),
193 _ => None,
194 }
195 }
196}
197
198/// One series identifier — the label set that distinguishes a
199/// time series within its metric family. Returned by
200/// [`MetricCatalog::series`]. Pairs are in declaration order
201/// the backend chooses; consumers that need a stable order
202/// sort themselves.
203pub type LabelSet = Vec<(String, String)>;
204
205/// One exemplar per OpenMetrics §4.6.1. Anchored to a
206/// specific sample observation by the (series identity +
207/// sample timestamp) pair, with its own value, optional
208/// timestamp, and label set (trace ids, span ids, …).
209///
210/// Returned by [`MetricCatalog::exemplars`].
211#[derive(Debug, Clone, PartialEq)]
212pub struct ExemplarPoint {
213 /// The series this exemplar attaches to. Includes the
214 /// synthetic `__name__` label so callers can rebuild
215 /// the full selector.
216 pub series: LabelSet,
217 /// Timestamp of the sample observation the exemplar
218 /// pairs with. The catalog reader's join key.
219 pub sample_timestamp_ms: i64,
220 /// The exemplar's observed value — typically the raw
221 /// sample value the exemplar describes (the bucket
222 /// boundary that captured it for histograms; the
223 /// counter increment value for counters).
224 pub value: f64,
225 /// Optional exemplar timestamp, distinct from
226 /// `sample_timestamp_ms`. Spec §4.6.1 leaves this
227 /// optional so producers that don't track timestamps
228 /// can omit it.
229 pub timestamp_ms: Option<i64>,
230 /// Exemplar's own labels — trace context, span ids,
231 /// arbitrary diagnostic envelope.
232 pub labels: Vec<(String, String)>,
233}
234
235/// Backend catalog trait.
236///
237/// All four methods may be called concurrently from multiple
238/// threads. Implementations either provide their own
239/// concurrency story (mutex, connection pool) or compose with
240/// [`CachedCatalog`] which serialises through a single
241/// per-key lock.
242///
243/// Errors surface via [`DataSourceError`] — same channel
244/// [`super::MetricAccess`] uses, so callers that already
245/// handle one can handle the other uniformly.
246pub trait MetricCatalog: Send + Sync {
247 /// Every metric family known to the backend, with its
248 /// OpenMetrics metadata. Order is backend-chosen
249 /// (sqlite-backed impls return alphabetical; remote
250 /// backends typically don't promise an order).
251 fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError>;
252
253 /// Distinct label keys observed across the catalog,
254 /// optionally restricted to one metric family by exact
255 /// name. `family_filter = None` returns every key the
256 /// backend has seen anywhere.
257 ///
258 /// Excludes synthetic keys: `__name__` is implicit
259 /// (it's the family name) and is never returned.
260 fn label_keys(&self, family_filter: Option<&str>) -> Result<Vec<String>, DataSourceError>;
261
262 /// Distinct values observed for `key`, optionally
263 /// restricted to one metric family. Returns an empty list
264 /// if the key has no values observed (or doesn't exist) —
265 /// no error, since "no values yet" is a normal state for
266 /// a fresh session.
267 ///
268 /// For implied labels (`le` on histograms, `quantile` on
269 /// summaries) the values are the bucket boundaries /
270 /// quantile φ values, in the canonical numeric order the
271 /// underlying schema records.
272 fn label_values(
273 &self,
274 key: &str,
275 family_filter: Option<&str>,
276 ) -> Result<Vec<String>, DataSourceError>;
277
278 /// Distinct label sets in the catalog matching every
279 /// [`Matcher`]. The result is the *series identity*
280 /// surface — one entry per `(family + label set)` tuple
281 /// that satisfies the selector. Used by callers that need
282 /// to know "which concrete series exist that satisfy this
283 /// shape" without fetching their samples.
284 ///
285 /// Each returned [`LabelSet`] **includes** the synthetic
286 /// `__name__` label so the caller can reconstruct the
287 /// full selector. Backends that natively store
288 /// per-instance labels project the family name on the
289 /// way out.
290 ///
291 /// Empty matcher list returns every series. Backends MAY
292 /// reject excessively wide queries with
293 /// [`DataSourceError`] — local sqlite tolerates it,
294 /// remote backends often don't.
295 fn series(&self, matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError>;
296
297 /// OpenMetrics §4.6.1 exemplars matching `matchers`,
298 /// optionally restricted to a `[start_ms, end_ms]`
299 /// window on `sample_timestamp_ms`. Returns one
300 /// [`ExemplarPoint`] per stored exemplar; series with
301 /// no exemplars contribute nothing.
302 ///
303 /// **Default implementation returns empty.** Backends
304 /// without exemplar storage (e.g. a remote VM endpoint
305 /// that doesn't expose them) leave this as the default
306 /// and don't have to fail. Backends that DO store them
307 /// override the method.
308 fn exemplars(
309 &self,
310 matchers: &[Matcher],
311 time_range: Option<(i64, i64)>,
312 ) -> Result<Vec<ExemplarPoint>, DataSourceError> {
313 let _ = (matchers, time_range);
314 Ok(Vec::new())
315 }
316}
317
318// =====================================================================
319// Caching wrapper
320// =====================================================================
321
322/// Cache layer for any [`MetricCatalog`] implementation.
323///
324/// Catalog data is small and slow-changing (per-session: it
325/// can only **grow**, never shrink). The cache absorbs the
326/// repeated queries that completion fires on every keystroke
327/// and amortises the backend round-trip across them. For
328/// sqlite-backed sessions the speedup is modest (the queries
329/// are already cheap); for remote backends it's the
330/// difference between "completion feels instant" and "every
331/// tab spawns an HTTP roundtrip."
332///
333/// ## Invalidation
334///
335/// Three layers, evaluated in order:
336///
337/// 1. **Time-based TTL** — every cached entry expires after
338/// [`CachedCatalog::ttl`] (default: 1 second). Aggressive
339/// by design: completion users tap repeatedly, so the
340/// cache absorbs bursts but yields to fresh data quickly.
341/// 2. **Generation counter** — [`CachedCatalog::invalidate`]
342/// bumps a process-local generation to force the next
343/// read. Used by integration tests + by callers that
344/// *know* the underlying state changed (e.g. a writer
345/// just landed a new metric family).
346/// 3. **Backend-supplied mtime** — when the
347/// [`CachedCatalog`] holds a `mtime_fn`, every read
348/// consults it; if the timestamp moved past the cached
349/// snapshot's read time, the entry expires immediately.
350/// This is the path the sqlite adapter uses to detect
351/// on-disk db changes (writer flushed).
352///
353/// The cache is a soft hint — if any layer says "stale,"
354/// the next call refetches and rebuilds. There's no
355/// invariant that two concurrent readers of the same
356/// just-invalidated key see the same fresh result; both may
357/// run the underlying query. That's acceptable because
358/// catalog reads are idempotent and side-effect-free.
359pub struct CachedCatalog<C: MetricCatalog + ?Sized> {
360 inner: Arc<C>,
361 ttl: Duration,
362 /// `mtime_fn` returns the latest backend-mtime as a
363 /// monotonic instant. The default is `None` (no mtime
364 /// invalidation; rely on TTL + generation only). The
365 /// sqlite adapter wires this to the db file's
366 /// `metadata().modified()`.
367 mtime_fn: Option<Box<dyn Fn() -> Option<Instant> + Send + Sync>>,
368 state: Mutex<CacheState>,
369}
370
371#[derive(Default)]
372struct CacheState {
373 generation: u64,
374 families: Option<CacheEntry<Vec<MetricFamilyMeta>>>,
375 label_keys: std::collections::HashMap<Option<String>, CacheEntry<Vec<String>>>,
376 label_values: std::collections::HashMap<(String, Option<String>), CacheEntry<Vec<String>>>,
377 /// `series` is keyed by a stable encoding of the matcher
378 /// list (sorted on label, op, value). Computing it
379 /// per-call is cheap relative to the underlying query.
380 series: std::collections::HashMap<String, CacheEntry<Vec<LabelSet>>>,
381 /// `exemplars` keyed by `(encoded matchers, time-range)`.
382 /// The time range is part of the identity since two
383 /// queries on the same selector but different windows
384 /// can resolve to different exemplar sets.
385 exemplars: std::collections::HashMap<String, CacheEntry<Vec<ExemplarPoint>>>,
386}
387
388struct CacheEntry<T> {
389 /// Generation at which this entry was recorded.
390 generation: u64,
391 /// Wall-clock `Instant` when the entry was filled.
392 filled_at: Instant,
393 /// Backend mtime at fill time (when known). On read, if
394 /// the current mtime is newer, the entry is stale.
395 backend_mtime: Option<Instant>,
396 value: T,
397}
398
399impl<C: MetricCatalog + ?Sized + 'static> CachedCatalog<C> {
400 /// Wrap an inner catalog with a default 1-second TTL and
401 /// no backend-mtime hook.
402 pub fn new(inner: Arc<C>) -> Self {
403 Self {
404 inner,
405 ttl: Duration::from_secs(1),
406 mtime_fn: None,
407 state: Mutex::new(CacheState::default()),
408 }
409 }
410
411 /// Builder: set the TTL window. Setting `Duration::ZERO`
412 /// disables the TTL layer (entries still expire on
413 /// generation bump or mtime change).
414 pub fn with_ttl(mut self, ttl: Duration) -> Self {
415 self.ttl = ttl;
416 self
417 }
418
419 /// Builder: install a backend-mtime hook. The closure
420 /// returns the latest backend-side mtime as a monotonic
421 /// `Instant`; entries cached *before* the latest mtime
422 /// are considered stale on the next read.
423 ///
424 /// `None` from the closure (e.g. the underlying file
425 /// disappeared) keeps the existing entries — the
426 /// behaviour matches "mtime unknown ⇒ trust the TTL."
427 pub fn with_mtime_fn<F>(mut self, mtime_fn: F) -> Self
428 where
429 F: Fn() -> Option<Instant> + Send + Sync + 'static,
430 {
431 self.mtime_fn = Some(Box::new(mtime_fn));
432 self
433 }
434
435 /// Force the next read of every key to refetch.
436 /// Internally bumps the generation counter; existing
437 /// entries become stale at next access.
438 pub fn invalidate(&self) {
439 if let Ok(mut s) = self.state.lock() {
440 s.generation = s.generation.wrapping_add(1);
441 }
442 }
443
444 /// True if `entry` is still fresh according to all three
445 /// invalidation layers. Caller holds the state lock.
446 fn is_fresh<T>(&self, entry: &CacheEntry<T>, gen_now: u64) -> bool {
447 if entry.generation != gen_now {
448 return false;
449 }
450 if !self.ttl.is_zero() && entry.filled_at.elapsed() >= self.ttl {
451 return false;
452 }
453 if let Some(f) = &self.mtime_fn
454 && let Some(now_mtime) = f()
455 {
456 match entry.backend_mtime {
457 Some(prev) if now_mtime > prev => return false,
458 None => return false,
459 _ => {}
460 }
461 }
462 true
463 }
464
465 fn current_mtime(&self) -> Option<Instant> {
466 self.mtime_fn.as_ref().and_then(|f| f())
467 }
468}
469
470impl<C: MetricCatalog + ?Sized + 'static> MetricCatalog for CachedCatalog<C> {
471 fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError> {
472 let gen_now;
473 {
474 let state = self
475 .state
476 .lock()
477 .map_err(|_| DataSourceError::new("cache poisoned"))?;
478 gen_now = state.generation;
479 if let Some(entry) = &state.families
480 && self.is_fresh(entry, gen_now)
481 {
482 return Ok(entry.value.clone());
483 }
484 }
485 let value = self.inner.metric_families()?;
486 let entry = CacheEntry {
487 generation: gen_now,
488 filled_at: Instant::now(),
489 backend_mtime: self.current_mtime(),
490 value: value.clone(),
491 };
492 if let Ok(mut state) = self.state.lock() {
493 state.families = Some(entry);
494 }
495 Ok(value)
496 }
497
498 fn label_keys(&self, family_filter: Option<&str>) -> Result<Vec<String>, DataSourceError> {
499 let key = family_filter.map(|s| s.to_string());
500 let gen_now;
501 {
502 let state = self
503 .state
504 .lock()
505 .map_err(|_| DataSourceError::new("cache poisoned"))?;
506 gen_now = state.generation;
507 if let Some(entry) = state.label_keys.get(&key)
508 && self.is_fresh(entry, gen_now)
509 {
510 return Ok(entry.value.clone());
511 }
512 }
513 let value = self.inner.label_keys(family_filter)?;
514 let entry = CacheEntry {
515 generation: gen_now,
516 filled_at: Instant::now(),
517 backend_mtime: self.current_mtime(),
518 value: value.clone(),
519 };
520 if let Ok(mut state) = self.state.lock() {
521 state.label_keys.insert(key, entry);
522 }
523 Ok(value)
524 }
525
526 fn label_values(
527 &self,
528 key: &str,
529 family_filter: Option<&str>,
530 ) -> Result<Vec<String>, DataSourceError> {
531 let cache_key = (key.to_string(), family_filter.map(|s| s.to_string()));
532 let gen_now;
533 {
534 let state = self
535 .state
536 .lock()
537 .map_err(|_| DataSourceError::new("cache poisoned"))?;
538 gen_now = state.generation;
539 if let Some(entry) = state.label_values.get(&cache_key)
540 && self.is_fresh(entry, gen_now)
541 {
542 return Ok(entry.value.clone());
543 }
544 }
545 let value = self.inner.label_values(key, family_filter)?;
546 let entry = CacheEntry {
547 generation: gen_now,
548 filled_at: Instant::now(),
549 backend_mtime: self.current_mtime(),
550 value: value.clone(),
551 };
552 if let Ok(mut state) = self.state.lock() {
553 state.label_values.insert(cache_key, entry);
554 }
555 Ok(value)
556 }
557
558 fn series(&self, matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError> {
559 let cache_key = encode_matchers(matchers);
560 let gen_now;
561 {
562 let state = self
563 .state
564 .lock()
565 .map_err(|_| DataSourceError::new("cache poisoned"))?;
566 gen_now = state.generation;
567 if let Some(entry) = state.series.get(&cache_key)
568 && self.is_fresh(entry, gen_now)
569 {
570 return Ok(entry.value.clone());
571 }
572 }
573 let value = self.inner.series(matchers)?;
574 let entry = CacheEntry {
575 generation: gen_now,
576 filled_at: Instant::now(),
577 backend_mtime: self.current_mtime(),
578 value: value.clone(),
579 };
580 if let Ok(mut state) = self.state.lock() {
581 state.series.insert(cache_key, entry);
582 }
583 Ok(value)
584 }
585
586 fn exemplars(
587 &self,
588 matchers: &[Matcher],
589 time_range: Option<(i64, i64)>,
590 ) -> Result<Vec<ExemplarPoint>, DataSourceError> {
591 // Cache key = matcher encoding ⊕ time range. The time
592 // range is part of the identity because two callers
593 // asking for different windows on the same selector
594 // need different cached results.
595 let mut cache_key = encode_matchers(matchers);
596 if let Some((s, e)) = time_range {
597 cache_key.push_str(&format!("@{s}..{e}"));
598 }
599 let gen_now;
600 {
601 let state = self
602 .state
603 .lock()
604 .map_err(|_| DataSourceError::new("cache poisoned"))?;
605 gen_now = state.generation;
606 if let Some(entry) = state.exemplars.get(&cache_key)
607 && self.is_fresh(entry, gen_now)
608 {
609 return Ok(entry.value.clone());
610 }
611 }
612 let value = self.inner.exemplars(matchers, time_range)?;
613 let entry = CacheEntry {
614 generation: gen_now,
615 filled_at: Instant::now(),
616 backend_mtime: self.current_mtime(),
617 value: value.clone(),
618 };
619 if let Ok(mut state) = self.state.lock() {
620 state.exemplars.insert(cache_key, entry);
621 }
622 Ok(value)
623 }
624}
625
626/// Stable string encoding of a matcher list — used as a
627/// cache key for `series` queries. Sorted by label so two
628/// matcher lists differing only in order map to the same key.
629fn encode_matchers(matchers: &[Matcher]) -> String {
630 let mut sorted: Vec<&Matcher> = matchers.iter().collect();
631 sorted.sort_by(|a, b| a.label.cmp(&b.label));
632 let mut out = String::new();
633 for m in sorted {
634 let op = match m.op {
635 MatcherOp::Eq => "=",
636 MatcherOp::Ne => "!=",
637 MatcherOp::EqRegex => "=~",
638 MatcherOp::NeRegex => "!~",
639 };
640 out.push_str(&m.label);
641 out.push_str(op);
642 out.push_str(&m.value);
643 out.push('\x1f'); // unit separator — won't appear in legit matchers
644 }
645 out
646}
647
648#[cfg(test)]
649mod tests {
650 use super::*;
651 use MatcherOp;
652 use std::sync::atomic::{AtomicUsize, Ordering};
653
654 /// Simple in-memory catalog used to drive the cache tests.
655 /// Counts how many times each method is called so the
656 /// tests can assert "this call hit the inner catalog" or
657 /// "this one was served from cache."
658 struct MockCatalog {
659 families: Vec<MetricFamilyMeta>,
660 keys: Vec<(Option<String>, Vec<String>)>,
661 values: Vec<(String, Option<String>, Vec<String>)>,
662 series: Vec<LabelSet>,
663 family_calls: AtomicUsize,
664 key_calls: AtomicUsize,
665 value_calls: AtomicUsize,
666 series_calls: AtomicUsize,
667 }
668
669 impl MockCatalog {
670 fn new() -> Self {
671 Self {
672 families: vec![
673 MetricFamilyMeta {
674 name: "ops_total".into(),
675 ty: MetricType::Counter,
676 unit: None,
677 help: Some("ops".into()),
678 },
679 MetricFamilyMeta {
680 name: "latency".into(),
681 ty: MetricType::Histogram,
682 unit: Some("seconds".into()),
683 help: None,
684 },
685 ],
686 keys: vec![
687 (None, vec!["phase".into(), "scenario".into()]),
688 (Some("ops_total".into()), vec!["phase".into()]),
689 ],
690 values: vec![("phase".into(), None, vec!["setup".into(), "run".into()])],
691 series: vec![vec![
692 ("__name__".into(), "ops_total".into()),
693 ("phase".into(), "setup".into()),
694 ]],
695 family_calls: AtomicUsize::new(0),
696 key_calls: AtomicUsize::new(0),
697 value_calls: AtomicUsize::new(0),
698 series_calls: AtomicUsize::new(0),
699 }
700 }
701 }
702
703 impl MetricCatalog for MockCatalog {
704 fn metric_families(&self) -> Result<Vec<MetricFamilyMeta>, DataSourceError> {
705 self.family_calls.fetch_add(1, Ordering::SeqCst);
706 Ok(self.families.clone())
707 }
708 fn label_keys(&self, filter: Option<&str>) -> Result<Vec<String>, DataSourceError> {
709 self.key_calls.fetch_add(1, Ordering::SeqCst);
710 for (f, v) in &self.keys {
711 if f.as_deref() == filter {
712 return Ok(v.clone());
713 }
714 }
715 Ok(Vec::new())
716 }
717 fn label_values(
718 &self,
719 key: &str,
720 filter: Option<&str>,
721 ) -> Result<Vec<String>, DataSourceError> {
722 self.value_calls.fetch_add(1, Ordering::SeqCst);
723 for (k, f, v) in &self.values {
724 if k == key && f.as_deref() == filter {
725 return Ok(v.clone());
726 }
727 }
728 Ok(Vec::new())
729 }
730 fn series(&self, _matchers: &[Matcher]) -> Result<Vec<LabelSet>, DataSourceError> {
731 self.series_calls.fetch_add(1, Ordering::SeqCst);
732 Ok(self.series.clone())
733 }
734 }
735
736 #[test]
737 fn metric_type_round_trips_via_str() {
738 for ty in [
739 MetricType::Counter,
740 MetricType::Gauge,
741 MetricType::Histogram,
742 MetricType::GaugeHistogram,
743 MetricType::Summary,
744 MetricType::Info,
745 MetricType::StateSet,
746 MetricType::Unknown,
747 ] {
748 assert_eq!(MetricType::parse(ty.as_str()), ty);
749 }
750 }
751
752 #[test]
753 fn metric_type_parse_unknown_keyword_yields_unknown() {
754 assert_eq!(MetricType::parse("garbage"), MetricType::Unknown);
755 assert_eq!(MetricType::parse(""), MetricType::Unknown);
756 }
757
758 #[test]
759 fn metric_type_implied_labels_match_openmetrics() {
760 assert_eq!(MetricType::Histogram.implied_label(), Some("le"));
761 assert_eq!(MetricType::GaugeHistogram.implied_label(), Some("le"));
762 assert_eq!(MetricType::Summary.implied_label(), Some("quantile"));
763 assert_eq!(MetricType::Counter.implied_label(), None);
764 assert_eq!(MetricType::Gauge.implied_label(), None);
765 }
766
767 #[test]
768 fn metric_type_has_derived_series_only_for_histograms_and_summary() {
769 assert!(MetricType::Histogram.has_derived_series());
770 assert!(MetricType::GaugeHistogram.has_derived_series());
771 assert!(MetricType::Summary.has_derived_series());
772 assert!(!MetricType::Counter.has_derived_series());
773 assert!(!MetricType::Gauge.has_derived_series());
774 assert!(!MetricType::Info.has_derived_series());
775 }
776
777 #[test]
778 fn cache_serves_repeat_calls_from_cache() {
779 let inner = Arc::new(MockCatalog::new());
780 let cache = CachedCatalog::new(inner.clone());
781 let _ = cache.metric_families().unwrap();
782 let _ = cache.metric_families().unwrap();
783 let _ = cache.metric_families().unwrap();
784 assert_eq!(
785 inner.family_calls.load(Ordering::SeqCst),
786 1,
787 "cache should have served repeated calls"
788 );
789 }
790
791 #[test]
792 fn cache_invalidate_forces_refetch() {
793 let inner = Arc::new(MockCatalog::new());
794 let cache = CachedCatalog::new(inner.clone());
795 let _ = cache.metric_families().unwrap();
796 cache.invalidate();
797 let _ = cache.metric_families().unwrap();
798 assert_eq!(
799 inner.family_calls.load(Ordering::SeqCst),
800 2,
801 "invalidate should have forced a refetch"
802 );
803 }
804
805 #[test]
806 fn cache_zero_ttl_still_caches_via_generation() {
807 let inner = Arc::new(MockCatalog::new());
808 // TTL = 0 is "TTL disabled, fall through to generation
809 // / mtime layers." Without an mtime hook, generation
810 // alone keeps things cached.
811 let cache = CachedCatalog::new(inner.clone()).with_ttl(Duration::ZERO);
812 let _ = cache.metric_families().unwrap();
813 let _ = cache.metric_families().unwrap();
814 assert_eq!(
815 inner.family_calls.load(Ordering::SeqCst),
816 1,
817 "TTL=0 alone still caches via generation"
818 );
819 }
820
821 #[test]
822 fn cache_distinguishes_label_keys_by_family_filter() {
823 let inner = Arc::new(MockCatalog::new());
824 let cache = CachedCatalog::new(inner.clone());
825 let _ = cache.label_keys(None).unwrap();
826 let _ = cache.label_keys(Some("ops_total")).unwrap();
827 let _ = cache.label_keys(None).unwrap();
828 let _ = cache.label_keys(Some("ops_total")).unwrap();
829 assert_eq!(
830 inner.key_calls.load(Ordering::SeqCst),
831 2,
832 "different family filters should be cached separately"
833 );
834 }
835
836 #[test]
837 fn cache_distinguishes_label_values_by_key_and_filter() {
838 let inner = Arc::new(MockCatalog::new());
839 let cache = CachedCatalog::new(inner.clone());
840 let _ = cache.label_values("phase", None).unwrap();
841 let _ = cache.label_values("phase", Some("ops_total")).unwrap();
842 let _ = cache.label_values("phase", None).unwrap();
843 assert_eq!(inner.value_calls.load(Ordering::SeqCst), 2);
844 }
845
846 #[test]
847 fn cache_series_keyed_by_matchers_independent_of_order() {
848 let inner = Arc::new(MockCatalog::new());
849 let cache = CachedCatalog::new(inner.clone());
850 let m1 = vec![
851 Matcher {
852 label: "a".into(),
853 op: MatcherOp::Eq,
854 value: "1".into(),
855 },
856 Matcher {
857 label: "b".into(),
858 op: MatcherOp::Eq,
859 value: "2".into(),
860 },
861 ];
862 let m2 = vec![
863 Matcher {
864 label: "b".into(),
865 op: MatcherOp::Eq,
866 value: "2".into(),
867 },
868 Matcher {
869 label: "a".into(),
870 op: MatcherOp::Eq,
871 value: "1".into(),
872 },
873 ];
874 let _ = cache.series(&m1).unwrap();
875 let _ = cache.series(&m2).unwrap();
876 assert_eq!(
877 inner.series_calls.load(Ordering::SeqCst),
878 1,
879 "matcher-order shouldn't fragment the cache"
880 );
881 }
882
883 #[test]
884 fn cache_mtime_hook_invalidates_on_advance() {
885 let inner = Arc::new(MockCatalog::new());
886 let mtime = Arc::new(Mutex::new(Instant::now()));
887 let mtime_clone = mtime.clone();
888 let cache = CachedCatalog::new(inner.clone())
889 .with_mtime_fn(move || Some(*mtime_clone.lock().unwrap()));
890 let _ = cache.metric_families().unwrap();
891 // First repeat is fresh.
892 let _ = cache.metric_families().unwrap();
893 assert_eq!(inner.family_calls.load(Ordering::SeqCst), 1);
894 // Advance the mtime — next call should refetch.
895 *mtime.lock().unwrap() += Duration::from_millis(1);
896 let _ = cache.metric_families().unwrap();
897 assert_eq!(
898 inner.family_calls.load(Ordering::SeqCst),
899 2,
900 "advanced mtime should have invalidated the cache"
901 );
902 }
903
904 #[test]
905 fn encode_matchers_sorts_by_label() {
906 let m1 = vec![
907 Matcher {
908 label: "a".into(),
909 op: MatcherOp::Eq,
910 value: "1".into(),
911 },
912 Matcher {
913 label: "b".into(),
914 op: MatcherOp::EqRegex,
915 value: ".*".into(),
916 },
917 ];
918 let m2 = vec![
919 Matcher {
920 label: "b".into(),
921 op: MatcherOp::EqRegex,
922 value: ".*".into(),
923 },
924 Matcher {
925 label: "a".into(),
926 op: MatcherOp::Eq,
927 value: "1".into(),
928 },
929 ];
930 assert_eq!(encode_matchers(&m1), encode_matchers(&m2));
931 }
932}