1use std::collections::HashMap;
29use std::sync::Arc;
30
31use super::{Matcher, MetricAccess, QueryError, Sample, Series, Vector};
32
33pub trait HorizonAware: Send + Sync {
37 fn earliest_ms(&self) -> Option<i64>;
38}
39
40pub struct Tier {
42 access: Arc<dyn MetricAccess>,
43 earliest_ms: Arc<dyn Fn() -> Option<i64> + Send + Sync>,
45}
46
47impl Tier {
48 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 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
74pub struct HybridStore {
76 tiers: Vec<Tier>,
77}
78
79impl HybridStore {
80 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 let mut chosen: Vec<&Tier> = Vec::new();
95 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; }
102 let earliest = tier.earliest();
103 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, };
115 }
116
117 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 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
150fn 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 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
184fn 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 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 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 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 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}