nodedb_lite/engine/timeseries/
query_routing.rs1use std::collections::HashMap;
18
19use serde::{Deserialize, Serialize};
20
21use nodedb_types::timeseries::{SeriesId, TimeRange};
22
23#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
25pub enum QueryScope {
26 #[default]
28 Local,
29 Cloud,
31 Hybrid,
34}
35
36#[derive(Debug, Clone, Serialize, Deserialize)]
42pub struct TimeseriesShape {
43 pub shape_id: String,
45 pub collection: String,
47 pub metric: String,
49 pub group_by: Vec<String>,
51 pub aggregate: String,
53 pub interval_ms: u64,
55}
56
57#[derive(Debug, Clone)]
61pub struct CachedShapeData {
62 pub shape: TimeseriesShape,
63 pub buckets: Vec<(i64, f64)>,
65 pub last_updated_ms: u64,
67 pub max_ts: i64,
69}
70
71impl CachedShapeData {
72 pub fn staleness_ms(&self, now_ms: u64) -> u64 {
74 now_ms.saturating_sub(self.last_updated_ms)
75 }
76}
77
78#[derive(Debug, Default)]
80pub struct TimeseriesShapeManager {
81 shapes: HashMap<String, CachedShapeData>,
83}
84
85impl TimeseriesShapeManager {
86 pub fn new() -> Self {
87 Self::default()
88 }
89
90 pub fn subscribe(&mut self, shape: TimeseriesShape) {
92 let shape_id = shape.shape_id.clone();
93 self.shapes.insert(
94 shape_id,
95 CachedShapeData {
96 shape,
97 buckets: Vec::new(),
98 last_updated_ms: 0,
99 max_ts: 0,
100 },
101 );
102 }
103
104 pub fn unsubscribe(&mut self, shape_id: &str) {
106 self.shapes.remove(shape_id);
107 }
108
109 pub fn update(&mut self, shape_id: &str, buckets: Vec<(i64, f64)>, now_ms: u64) {
111 if let Some(cached) = self.shapes.get_mut(shape_id) {
112 cached.max_ts = buckets.iter().map(|(ts, _)| *ts).max().unwrap_or(0);
113 cached.buckets = buckets;
114 cached.last_updated_ms = now_ms;
115 }
116 }
117
118 pub fn query(&self, shape_id: &str, range: &TimeRange, now_ms: u64) -> (Vec<(i64, f64)>, u64) {
123 match self.shapes.get(shape_id) {
124 Some(cached) => {
125 let filtered: Vec<(i64, f64)> = cached
126 .buckets
127 .iter()
128 .filter(|(ts, _)| range.contains(*ts))
129 .copied()
130 .collect();
131 (filtered, cached.staleness_ms(now_ms))
132 }
133 None => (Vec::new(), u64::MAX),
134 }
135 }
136
137 pub fn active_shapes(&self) -> Vec<&TimeseriesShape> {
139 self.shapes.values().map(|c| &c.shape).collect()
140 }
141
142 pub fn len(&self) -> usize {
144 self.shapes.len()
145 }
146
147 pub fn is_empty(&self) -> bool {
148 self.shapes.is_empty()
149 }
150
151 pub fn export(&self) -> Vec<(String, CachedShapeData)> {
153 self.shapes
154 .iter()
155 .map(|(k, v)| (k.clone(), v.clone()))
156 .collect()
157 }
158
159 pub fn import(&mut self, entries: Vec<(String, CachedShapeData)>) {
161 for (k, v) in entries {
162 self.shapes.insert(k, v);
163 }
164 }
165}
166
167#[derive(Debug)]
169pub struct HybridQueryResult {
170 pub local: Vec<(i64, f64, SeriesId)>,
172 pub fleet: Vec<(i64, f64)>,
174 pub shape_staleness_ms: u64,
177 pub local_available: bool,
179}
180
181pub struct RoutedQueryParams<'a> {
183 pub scope: QueryScope,
184 pub collection: &'a str,
185 pub range: &'a TimeRange,
186 pub bucket_ms: Option<i64>,
187 pub shape_id: Option<&'a str>,
188 pub now_ms: u64,
189}
190
191pub fn execute_routed_query(
200 params: &RoutedQueryParams<'_>,
201 engine: &super::engine::TimeseriesEngine,
202 shape_mgr: &TimeseriesShapeManager,
203) -> HybridQueryResult {
204 let RoutedQueryParams {
205 scope,
206 collection,
207 range,
208 bucket_ms,
209 shape_id,
210 now_ms,
211 } = params;
212 let _ = bucket_ms; let local = match scope {
215 QueryScope::Local | QueryScope::Hybrid => engine.scan(collection, range),
216 QueryScope::Cloud => {
217 Vec::new()
220 }
221 };
222
223 let (fleet, shape_staleness_ms) = match scope {
224 QueryScope::Hybrid => {
225 if let Some(sid) = shape_id {
226 shape_mgr.query(sid, range, *now_ms)
227 } else {
228 (Vec::new(), u64::MAX)
229 }
230 }
231 _ => (Vec::new(), u64::MAX),
232 };
233
234 HybridQueryResult {
235 local_available: !matches!(scope, QueryScope::Cloud),
236 local,
237 fleet,
238 shape_staleness_ms,
239 }
240}
241
242#[cfg(test)]
243mod tests {
244 use super::*;
245
246 fn make_shape(id: &str) -> TimeseriesShape {
247 TimeseriesShape {
248 shape_id: id.into(),
249 collection: "metrics".into(),
250 metric: "cpu_usage".into(),
251 group_by: vec![],
252 aggregate: "avg".into(),
253 interval_ms: 300_000,
254 }
255 }
256
257 #[test]
258 fn subscribe_and_query() {
259 let mut mgr = TimeseriesShapeManager::new();
260 mgr.subscribe(make_shape("fleet_cpu"));
261 assert_eq!(mgr.len(), 1);
262
263 let (buckets, staleness) = mgr.query("fleet_cpu", &TimeRange::new(0, 1_000_000), 1000);
265 assert!(buckets.is_empty());
266 assert_eq!(staleness, 1000); mgr.update(
270 "fleet_cpu",
271 vec![(300_000, 45.0), (600_000, 52.0), (900_000, 48.0)],
272 1_000_000,
273 );
274
275 let (buckets, staleness) = mgr.query("fleet_cpu", &TimeRange::new(0, 1_000_000), 1_000_000);
276 assert_eq!(buckets.len(), 3);
277 assert_eq!(staleness, 0); }
279
280 #[test]
281 fn unsubscribe_removes() {
282 let mut mgr = TimeseriesShapeManager::new();
283 mgr.subscribe(make_shape("s1"));
284 assert_eq!(mgr.len(), 1);
285 mgr.unsubscribe("s1");
286 assert_eq!(mgr.len(), 0);
287 }
288
289 #[test]
290 fn query_range_filtering() {
291 let mut mgr = TimeseriesShapeManager::new();
292 mgr.subscribe(make_shape("s1"));
293 mgr.update(
294 "s1",
295 vec![(100, 1.0), (200, 2.0), (300, 3.0), (400, 4.0)],
296 500,
297 );
298
299 let (buckets, _) = mgr.query("s1", &TimeRange::new(200, 300), 500);
301 assert_eq!(buckets.len(), 2);
302 assert_eq!(buckets[0], (200, 2.0));
303 assert_eq!(buckets[1], (300, 3.0));
304 }
305
306 #[test]
307 fn missing_shape_returns_max_staleness() {
308 let mgr = TimeseriesShapeManager::new();
309 let (buckets, staleness) = mgr.query("nonexistent", &TimeRange::new(0, 1000), 1000);
310 assert!(buckets.is_empty());
311 assert_eq!(staleness, u64::MAX);
312 }
313
314 #[test]
315 fn export_import_roundtrip() {
316 let mut mgr = TimeseriesShapeManager::new();
317 mgr.subscribe(make_shape("s1"));
318 mgr.update("s1", vec![(100, 42.0)], 200);
319
320 let exported = mgr.export();
321 let mut mgr2 = TimeseriesShapeManager::new();
322 mgr2.import(exported);
323 assert_eq!(mgr2.len(), 1);
324
325 let (buckets, _) = mgr2.query("s1", &TimeRange::new(0, 1000), 200);
326 assert_eq!(buckets.len(), 1);
327 }
328
329 #[test]
330 fn query_scope_default_is_local() {
331 assert_eq!(QueryScope::default(), QueryScope::Local);
332 }
333}