Skip to main content

nmbrs_metrics/queryapi/
hybrid.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-90 §M5 — the **hybrid** metric read backend.
5//!
6//! A composite [`MetricAccess`] over an **ordered list of tiers**, finest/
7//! freshest first (today: the in-memory cadence tier, then the per-session
8//! sqlite store; extensible to a remote tail, etc.). The read model:
9//!
10//! 1. Walk the tiers finest-first. A tier advertises the oldest time it can
11//!    answer ([`Tier::earliest_ms`]). Query a tier only while the query is **not
12//!    yet covered back to `start`** by the finer tiers already chosen — and skip
13//!    a tier whose data begins after the window (no intersection). Stop as soon
14//!    as a tier reaches back to (or before) `start`, or advertises an unbounded
15//!    horizon.
16//! 2. Issue the **same** `[start_ms, end_ms]` bounds to every chosen tier,
17//!    **concurrently** when more than one is needed.
18//! 3. Fold the results finest-first by **union-minus-overlap**: a finer tier
19//!    wins every time-span it covers, so "high-centered" recent data comes from
20//!    memory at full sub-interval resolution and only the older tail is served,
21//!    coarser, from the durable store. The union/overlap rule is **per series**
22//!    (boundary = wherever the finer tier's own samples start), not a single
23//!    global cut.
24//!
25//! The common case — a recent windowed read served entirely from the in-memory
26//! horizon — chooses exactly one tier and never opens sqlite.
27
28use std::collections::HashMap;
29use std::sync::Arc;
30
31use super::{Matcher, MetricAccess, QueryError, Sample, Series, Vector};
32
33/// A horizon-advertising backend: the oldest sample-time it still holds, in
34/// Unix-ms. `None` ⇒ unbounded (covers back as far as the query asks) — the
35/// natural answer for a durable tail tier.
36pub trait HorizonAware: Send + Sync {
37    fn earliest_ms(&self) -> Option<i64>;
38}
39
40/// One tier of a [`HybridStore`]: a backend plus its horizon advertiser.
41pub struct Tier {
42    access: Arc<dyn MetricAccess>,
43    /// Oldest answerable time (Unix-ms). `None` ⇒ unbounded tail.
44    earliest_ms: Arc<dyn Fn() -> Option<i64> + Send + Sync>,
45}
46
47impl Tier {
48    /// A tier whose horizon is computed on demand (e.g. the in-memory tier's
49    /// retention edge).
50    pub fn new(
51        access: Arc<dyn MetricAccess>,
52        earliest_ms: Arc<dyn Fn() -> Option<i64> + Send + Sync>,
53    ) -> Self {
54        Self {
55            access,
56            earliest_ms,
57        }
58    }
59
60    /// A tier that covers back as far as any query asks (a durable tail such as
61    /// sqlite) — always serves the older remainder.
62    pub fn unbounded(access: Arc<dyn MetricAccess>) -> Self {
63        Self {
64            access,
65            earliest_ms: Arc::new(|| None),
66        }
67    }
68
69    fn earliest(&self) -> Option<i64> {
70        (self.earliest_ms)()
71    }
72}
73
74/// Composite read backend over an ordered tier list (finest first).
75pub struct HybridStore {
76    tiers: Vec<Tier>,
77}
78
79impl HybridStore {
80    /// Compose `tiers`, finest/freshest first.
81    pub fn new(tiers: Vec<Tier>) -> Self {
82        Self { tiers }
83    }
84}
85
86impl MetricAccess for HybridStore {
87    fn select_range(
88        &self,
89        matchers: &[Matcher],
90        start_ms: i64,
91        end_ms: i64,
92    ) -> Result<Vector, QueryError> {
93        // 1. Coverage-aware tier selection (finest first).
94        let mut chosen: Vec<&Tier> = Vec::new();
95        // `covered_back_to` = oldest time covered so far (exclusive lower edge);
96        // start above the window so the first intersecting tier is always taken.
97        let mut covered_back_to = end_ms.saturating_add(1);
98        for tier in &self.tiers {
99            if covered_back_to <= start_ms {
100                break; // query already covered back to `start`
101            }
102            let earliest = tier.earliest();
103            // Skip a tier whose data begins after the window — it intersects
104            // nothing in `[start, end]`.
105            if let Some(e) = earliest
106                && e > end_ms
107            {
108                continue;
109            }
110            chosen.push(tier);
111            covered_back_to = match earliest {
112                Some(e) => covered_back_to.min(e),
113                None => start_ms, // unbounded tail closes the coverage
114            };
115        }
116
117        // 2. Issue the same bounds to every chosen tier — concurrently when >1.
118        let results: Vec<Vector> = match chosen.as_slice() {
119            [] => return Ok(Vector::default()),
120            [only] => vec![only.access.select_range(matchers, start_ms, end_ms)?],
121            many => {
122                let collected: Vec<Result<Vector, QueryError>> = std::thread::scope(|scope| {
123                    let handles: Vec<_> = many
124                        .iter()
125                        .map(|tier| {
126                            scope
127                                .spawn(move || tier.access.select_range(matchers, start_ms, end_ms))
128                        })
129                        .collect();
130                    handles
131                        .into_iter()
132                        .map(|h| h.join().expect("tier query panicked"))
133                        .collect()
134                });
135                collected.into_iter().collect::<Result<Vec<_>, _>>()?
136            }
137        };
138
139        // 3. Fold finest-first; each finer accumulation wins overlaps with the
140        //    next, coarser tier.
141        let mut iter = results.into_iter();
142        let mut acc = iter.next().unwrap_or_default();
143        for coarser in iter {
144            acc = union_minus_overlap(acc, coarser);
145        }
146        Ok(acc)
147    }
148}
149
150/// Merge two range-vectors, preferring `fine` in every overlap.
151///
152/// Per series (matched by label *set*, order-independent): keep all of `fine`'s
153/// samples, and from `coarse` keep only the samples **strictly older** than
154/// `fine`'s earliest sample for that series — the non-overlapping tail. A series
155/// present only in `coarse` passes through whole. Output samples are time-ordered.
156fn union_minus_overlap(fine: Vector, coarse: Vector) -> Vector {
157    let mut out: HashMap<Vec<(String, String)>, Series> = HashMap::new();
158    for s in fine.into_series() {
159        out.insert(canonical_labels(&s.labels), s);
160    }
161    for c in coarse.into_series() {
162        let key = canonical_labels(&c.labels);
163        match out.get_mut(&key) {
164            Some(fine_series) => {
165                // `fine`'s samples are ascending; its earliest is the overlap edge.
166                let edge = fine_series.samples.first().map(|s| s.timestamp_ms);
167                let mut merged: Vec<Sample> = c
168                    .samples
169                    .into_iter()
170                    .filter(|s| edge.map_or(true, |e| s.timestamp_ms < e))
171                    .collect();
172                merged.append(&mut fine_series.samples);
173                merged.sort_by_key(|s| s.timestamp_ms);
174                fine_series.samples = merged;
175            }
176            None => {
177                out.insert(key, c);
178            }
179        }
180    }
181    Vector::new(out.into_values().collect())
182}
183
184/// A canonical, order-independent key for a series' label set.
185fn canonical_labels(labels: &[(String, String)]) -> Vec<(String, String)> {
186    let mut v = labels.to_vec();
187    v.sort();
188    v
189}
190
191#[cfg(test)]
192mod tests {
193    use super::*;
194    use std::sync::atomic::{AtomicUsize, Ordering};
195
196    fn series(name: &str, pts: &[(i64, f64)]) -> Series {
197        Series {
198            labels: vec![("__name__".to_string(), name.to_string())],
199            samples: pts
200                .iter()
201                .map(|&(t, v)| Sample {
202                    timestamp_ms: t,
203                    value: v,
204                })
205                .collect(),
206        }
207    }
208
209    struct Stub {
210        out: Vector,
211        calls: Arc<AtomicUsize>,
212    }
213    impl MetricAccess for Stub {
214        fn select_range(&self, _: &[Matcher], _: i64, _: i64) -> Result<Vector, QueryError> {
215            self.calls.fetch_add(1, Ordering::SeqCst);
216            Ok(self.out.clone())
217        }
218    }
219
220    fn tier(out: Vector, earliest: Option<i64>) -> (Tier, Arc<AtomicUsize>) {
221        let calls = Arc::new(AtomicUsize::new(0));
222        let t = Tier::new(
223            Arc::new(Stub {
224                out,
225                calls: calls.clone(),
226            }),
227            Arc::new(move || earliest),
228        );
229        (t, calls)
230    }
231
232    #[test]
233    fn one_tier_covers_query_no_lower_tier_queried() {
234        // mem reaches back to t=0; query [100,200] is covered by mem alone.
235        let (mem, mem_calls) = tier(
236            Vector::new(vec![series("ops", &[(100, 1.0), (150, 2.0), (200, 3.0)])]),
237            Some(0),
238        );
239        let (cold, cold_calls) = tier(Vector::new(vec![series("ops", &[(100, 1.0)])]), None);
240        let store = HybridStore::new(vec![mem, cold]);
241
242        let v = store.select_range(&[], 100, 200).unwrap();
243        assert_eq!(mem_calls.load(Ordering::SeqCst), 1);
244        assert_eq!(
245            cold_calls.load(Ordering::SeqCst),
246            0,
247            "cold not consulted when mem covers the query"
248        );
249        assert_eq!(v.series()[0].samples.len(), 3);
250    }
251
252    #[test]
253    fn spill_stitches_cold_tail_under_mem_recent() {
254        // mem only holds from t=150; query [0,300] needs the older tail from cold.
255        let (mem, _) = tier(
256            Vector::new(vec![series("ops", &[(150, 5.0), (300, 7.0)])]),
257            Some(150),
258        );
259        let (cold, cold_calls) = tier(
260            Vector::new(vec![series(
261                "ops",
262                &[(0, 1.0), (100, 2.0), (150, 5.0), (300, 7.0)],
263            )]),
264            None,
265        );
266        let store = HybridStore::new(vec![mem, cold]);
267
268        let v = store.select_range(&[], 0, 300).unwrap();
269        assert_eq!(
270            cold_calls.load(Ordering::SeqCst),
271            1,
272            "cold consulted for the older tail"
273        );
274        let s = &v.series()[0];
275        let ts: Vec<i64> = s.samples.iter().map(|x| x.timestamp_ms).collect();
276        // cold's 150 & 300 (>= mem edge 150) dropped as overlap; 0 & 100 kept;
277        // then mem's 150 & 300 — one smooth timeline, no double-count at 150.
278        assert_eq!(ts, vec![0, 100, 150, 300]);
279        assert_eq!(
280            s.samples.iter().filter(|x| x.timestamp_ms == 150).count(),
281            1
282        );
283        assert_eq!(
284            s.samples
285                .iter()
286                .find(|x| x.timestamp_ms == 150)
287                .unwrap()
288                .value,
289            5.0
290        );
291    }
292
293    #[test]
294    fn tier_whose_data_starts_after_the_window_is_skipped() {
295        // Query an OLD window [0,100]; mem only holds from t=500 (after the
296        // window) → mem skipped, cold serves it.
297        let (mem, mem_calls) = tier(Vector::new(vec![series("ops", &[(500, 9.0)])]), Some(500));
298        let (cold, cold_calls) = tier(
299            Vector::new(vec![series("ops", &[(0, 1.0), (100, 2.0)])]),
300            None,
301        );
302        let store = HybridStore::new(vec![mem, cold]);
303
304        let v = store.select_range(&[], 0, 100).unwrap();
305        assert_eq!(
306            mem_calls.load(Ordering::SeqCst),
307            0,
308            "mem has no data in the window — skipped"
309        );
310        assert_eq!(cold_calls.load(Ordering::SeqCst), 1);
311        assert_eq!(v.series()[0].samples.len(), 2);
312    }
313
314    #[test]
315    fn series_only_in_cold_survives_the_union() {
316        let (mem, _) = tier(Vector::new(vec![series("ops", &[(150, 5.0)])]), Some(150));
317        let (cold, _) = tier(
318            Vector::new(vec![
319                series("ops", &[(150, 5.0)]),
320                series("errors", &[(0, 9.0), (100, 9.0)]),
321            ]),
322            None,
323        );
324        let store = HybridStore::new(vec![mem, cold]);
325
326        let v = store.select_range(&[], 0, 200).unwrap();
327        assert_eq!(v.len(), 2);
328        let errors = v
329            .series()
330            .iter()
331            .find(|s| s.labels.iter().any(|(_, v)| v == "errors"))
332            .unwrap();
333        assert_eq!(errors.samples.len(), 2);
334    }
335}