1use std::collections::BTreeMap;
22use std::path::{Path, PathBuf};
23use std::sync::mpsc::SyncSender;
24use std::sync::{Arc, Mutex};
25use std::time::Duration;
26
27use rusqlite::{Connection, OptionalExtension, params};
28use serde::Serialize;
29
30use crate::error::{Error, Result};
31use crate::metrics::InstanceSample;
32use crate::org::OrgId;
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq)]
36pub struct Tier {
37 pub step: u64,
38 pub keep: u64,
39}
40
41pub const TIERS: [Tier; 3] = [
42 Tier {
43 step: 10,
44 keep: 86_400,
45 },
46 Tier {
47 step: 60,
48 keep: 7 * 86_400,
49 },
50 Tier {
51 step: 600,
52 keep: 30 * 86_400,
53 },
54];
55
56pub const METRICS: [&str; 6] = [
58 "cpu",
59 "memory",
60 "net_rx",
61 "net_tx",
62 "disk_read",
63 "disk_write",
64];
65
66#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize)]
69pub struct Values(pub [Option<f64>; 6]);
70
71#[derive(Debug, Clone, PartialEq)]
73pub struct Row {
74 pub org: OrgId,
75 pub instance: String,
76 pub stack: String,
77 pub service: String,
78 pub ts: u64,
80 pub values: Values,
81}
82
83#[derive(Debug, Default)]
84struct Acc {
85 ts: u64,
86 sum: [f64; 6],
87 n: [u32; 6],
88}
89
90#[derive(Debug, Default)]
91struct Prev {
92 net: Option<(u64, u64, u64)>,
93 disk: Option<(u64, u64, u64)>,
94}
95
96#[derive(Debug, Default)]
99pub struct Recorder {
100 prev: BTreeMap<(String, String), Prev>,
101 acc: BTreeMap<(String, String), (Acc, String, String)>,
102}
103
104fn rate(prev: Option<(u64, u64, u64)>, now: (u64, u64), at: u64) -> (Option<f64>, Option<f64>) {
105 match prev {
106 Some((a, b, t)) if at > t && now.0 >= a && now.1 >= b => {
107 let dt = (at - t) as f64 / 1000.0;
108 (Some((now.0 - a) as f64 / dt), Some((now.1 - b) as f64 / dt))
109 }
110 _ => (None, None),
112 }
113}
114
115impl Recorder {
116 pub fn add(&mut self, at_ms: u64, samples: &[InstanceSample]) -> Vec<Row> {
119 let step = TIERS[0].step;
120 let bucket = at_ms / 1000 / step * step;
121 let mut out = Vec::new();
122 let mut seen = Vec::new();
123 for i in samples {
124 let Some(org) = OrgId::from_incus_project(&i.project) else {
125 continue;
126 };
127 if !i.running() {
128 continue;
129 }
130 let key = (org.to_string(), i.name.clone());
131 seen.push(key.clone());
132 let p = self.prev.entry(key.clone()).or_default();
133 let (rx, tx) = match (i.net_rx_bytes, i.net_tx_bytes) {
134 (Some(r), Some(t)) => {
135 let v = rate(p.net, (r, t), at_ms);
136 p.net = Some((r, t, at_ms));
137 v
138 }
139 _ => (None, None),
140 };
141 let (rd, wr) = match (i.disk_read_bytes, i.disk_write_bytes) {
142 (Some(r), Some(w)) => {
143 let v = rate(p.disk, (r, w), at_ms);
144 p.disk = Some((r, w, at_ms));
145 v
146 }
147 _ => (None, None),
148 };
149 let vals = [
150 i.cpu_pct.map(f64::from),
151 i.mem_bytes.map(|m| m as f64),
152 rx,
153 tx,
154 rd,
155 wr,
156 ];
157 let labels = (
158 i.labels.get("isb.stack").cloned().unwrap_or_default(),
159 i.labels.get("isb.service").cloned().unwrap_or_default(),
160 );
161 let e = self.acc.entry(key.clone()).or_insert_with(|| {
162 (
163 Acc {
164 ts: bucket,
165 ..Default::default()
166 },
167 labels.0.clone(),
168 labels.1.clone(),
169 )
170 });
171 if e.0.ts != bucket {
172 if let Some(r) = close(&key, e) {
173 out.push(r);
174 }
175 e.0 = Acc {
176 ts: bucket,
177 ..Default::default()
178 };
179 }
180 (e.1, e.2) = labels;
181 for (k, v) in vals.iter().enumerate() {
182 if let Some(v) = v {
183 e.0.sum[k] += v;
184 e.0.n[k] += 1;
185 }
186 }
187 }
188 let gone: Vec<_> = self
190 .acc
191 .keys()
192 .filter(|k| !seen.contains(k))
193 .cloned()
194 .collect();
195 for k in gone {
196 if let Some(e) = self.acc.remove(&k) {
197 if let Some(r) = close(&k, &e) {
198 out.push(r);
199 }
200 }
201 self.prev.remove(&k);
202 }
203 out
204 }
205}
206
207fn close(key: &(String, String), e: &(Acc, String, String)) -> Option<Row> {
208 let mut v = Values::default();
209 for k in 0..6 {
210 if e.0.n[k] > 0 {
211 v.0[k] = Some(e.0.sum[k] / f64::from(e.0.n[k]));
212 }
213 }
214 if v.0.iter().all(Option::is_none) {
215 return None;
216 }
217 Some(Row {
218 org: OrgId::new(key.0.clone()).ok()?,
219 instance: key.1.clone(),
220 stack: e.1.clone(),
221 service: e.2.clone(),
222 ts: e.0.ts,
223 values: v,
224 })
225}
226
227pub struct OrgDb {
229 conn: Connection,
230}
231
232const SCHEMA: &str = "
233CREATE TABLE IF NOT EXISTS series(
234 id INTEGER PRIMARY KEY,
235 instance TEXT NOT NULL UNIQUE,
236 stack TEXT NOT NULL,
237 service TEXT NOT NULL,
238 last_seen INTEGER NOT NULL);
239CREATE TABLE IF NOT EXISTS samples(
240 tier INTEGER NOT NULL,
241 sid INTEGER NOT NULL,
242 ts INTEGER NOT NULL,
243 cpu REAL, mem REAL, net_rx REAL, net_tx REAL, disk_read REAL, disk_write REAL,
244 PRIMARY KEY(tier, sid, ts)) WITHOUT ROWID;
245CREATE TABLE IF NOT EXISTS meta(k TEXT PRIMARY KEY, v INTEGER NOT NULL);
246";
247
248fn db_err(step: &str, e: rusqlite::Error) -> Error {
249 Error::invalid(format!("metrics history: {step}: {e}"))
250}
251
252impl OrgDb {
253 pub fn open(path: &Path) -> Result<OrgDb> {
254 if let Some(d) = path.parent() {
255 std::fs::create_dir_all(d)?;
256 }
257 let conn = Connection::open(path).map_err(|e| db_err("open", e))?;
258 conn.busy_timeout(Duration::from_secs(5))
259 .map_err(|e| db_err("busy timeout", e))?;
260 conn.execute_batch(
262 "PRAGMA auto_vacuum=INCREMENTAL; PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL;",
263 )
264 .map_err(|e| db_err("pragmas", e))?;
265 conn.execute_batch(SCHEMA)
266 .map_err(|e| db_err("schema", e))?;
267 Ok(OrgDb { conn })
268 }
269
270 pub fn insert(&mut self, rows: &[Row]) -> Result<()> {
272 let tx = self.conn.transaction().map_err(|e| db_err("begin", e))?;
273 {
274 let me = OrgDb::borrow(&tx);
275 for r in rows {
276 let sid = me.series_id(r)?;
277 let v = r.values.0;
278 tx.execute(
279 "INSERT OR REPLACE INTO samples VALUES(0, ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
280 params![sid, r.ts as i64, v[0], v[1], v[2], v[3], v[4], v[5]],
281 )
282 .map_err(|e| db_err("insert", e))?;
283 }
284 }
285 tx.commit().map_err(|e| db_err("commit", e))
286 }
287
288 fn borrow(c: &Connection) -> Borrowed<'_> {
289 Borrowed { conn: c }
290 }
291
292 fn meta(&self, k: &str) -> Result<Option<u64>> {
293 self.conn
294 .query_row("SELECT v FROM meta WHERE k=?1", params![k], |r| {
295 r.get::<_, i64>(0)
296 })
297 .optional()
298 .map(|v| v.map(|v| v as u64))
299 .map_err(|e| db_err("meta", e))
300 }
301
302 pub fn maintain(&mut self, now: u64) -> Result<()> {
305 let tx = self.conn.transaction().map_err(|e| db_err("begin", e))?;
306 for t in 1..TIERS.len() {
307 let (step, below) = (TIERS[t].step, TIERS[t - 1].step);
308 let to = now.saturating_sub(2 * below) / step * step;
311 let key = format!("rolled_{t}");
312 let from: u64 = tx
313 .query_row("SELECT v FROM meta WHERE k=?1", params![key], |r| {
314 r.get::<_, i64>(0)
315 })
316 .optional()
317 .map_err(|e| db_err("meta", e))?
318 .map(|v| v as u64)
319 .unwrap_or_else(|| now.saturating_sub(TIERS[t - 1].keep) / step * step);
320 if to > from {
321 tx.execute(
322 "INSERT OR REPLACE INTO samples
323 SELECT ?1, sid, (ts / ?2) * ?2, avg(cpu), avg(mem), avg(net_rx), avg(net_tx),
324 avg(disk_read), avg(disk_write)
325 FROM samples WHERE tier = ?3 AND ts >= ?4 AND ts < ?5
326 GROUP BY sid, ts / ?2",
327 params![
328 t as i64,
329 step as i64,
330 (t - 1) as i64,
331 from as i64,
332 to as i64
333 ],
334 )
335 .map_err(|e| db_err("roll up", e))?;
336 tx.execute(
337 "INSERT OR REPLACE INTO meta VALUES(?1, ?2)",
338 params![key, to as i64],
339 )
340 .map_err(|e| db_err("meta", e))?;
341 }
342 }
343 for (t, tier) in TIERS.iter().enumerate() {
344 tx.execute(
345 "DELETE FROM samples WHERE tier = ?1 AND ts < ?2",
346 params![t as i64, now.saturating_sub(tier.keep) as i64],
347 )
348 .map_err(|e| db_err("retention", e))?;
349 }
350 tx.execute(
351 "DELETE FROM series WHERE NOT EXISTS (SELECT 1 FROM samples WHERE sid = series.id)",
352 [],
353 )
354 .map_err(|e| db_err("prune series", e))?;
355 tx.commit().map_err(|e| db_err("commit", e))?;
356 let _ = self.conn.execute_batch("PRAGMA incremental_vacuum;");
357 Ok(())
358 }
359
360 pub fn query(&self, q: &Query, now: u64) -> Result<Answer> {
363 let span = q.to.saturating_sub(q.from).max(1);
364 let t = TIERS
366 .iter()
367 .position(|t| q.from + t.keep >= now)
368 .unwrap_or(TIERS.len() - 1);
369 let mut step = q.step.max(TIERS[t].step);
370 step = step.max(span.div_ceil(MAX_POINTS));
372 step = step.div_ceil(TIERS[t].step) * TIERS[t].step;
373 let rolled = if t > 0 {
376 self.meta(&format!("rolled_{t}"))?.unwrap_or(0)
377 } else {
378 0
379 };
380 let mut where_series = String::new();
381 let mut args: Vec<rusqlite::types::Value> = Vec::new();
382 if let Some(s) = &q.stack {
383 where_series.push_str(" AND s.stack = ?");
384 args.push(s.clone().into());
385 }
386 if let Some(s) = &q.service {
387 where_series.push_str(" AND s.service = ?");
388 args.push(s.clone().into());
389 }
390 if let Some(i) = &q.instance {
391 where_series.push_str(" AND s.instance = ?");
392 args.push(i.clone().into());
393 }
394 let col = match q.metric.as_str() {
395 "cpu" => "cpu",
396 "memory" => "mem",
397 "net_rx" => "net_rx",
398 "net_tx" => "net_tx",
399 "disk_read" => "disk_read",
400 "disk_write" => "disk_write",
401 m => {
402 return Err(Error::invalid(format!(
403 "metric {m:?}: one of {}",
404 METRICS.join(", ")
405 )));
406 }
407 };
408 let sql = format!(
409 "SELECT s.instance, s.stack, s.service, (x.ts / {step}) * {step} AS b, avg(x.{col})
410 FROM samples x JOIN series s ON s.id = x.sid
411 WHERE ((x.tier = {t} AND x.ts >= {from} AND x.ts < {to})
412 OR (x.tier = {below} AND x.ts >= {lag} AND x.ts < {to})){where_series}
413 AND x.{col} IS NOT NULL
414 GROUP BY x.sid, b ORDER BY s.instance, b",
415 from = q.from,
416 to = q.to,
417 below = if t > 0 { t - 1 } else { t },
418 lag = if t > 0 { rolled.max(q.from) } else { q.to },
419 );
420 let mut stmt = self.conn.prepare(&sql).map_err(|e| db_err("query", e))?;
421 let rows = stmt
422 .query_map(rusqlite::params_from_iter(args), |r| {
423 Ok((
424 r.get::<_, String>(0)?,
425 r.get::<_, String>(1)?,
426 r.get::<_, String>(2)?,
427 r.get::<_, i64>(3)? as u64,
428 r.get::<_, f64>(4)?,
429 ))
430 })
431 .map_err(|e| db_err("query", e))?;
432 let mut per: BTreeMap<String, Series> = BTreeMap::new();
433 for r in rows {
434 let (inst, stack, service, b, v) = r.map_err(|e| db_err("query row", e))?;
435 per.entry(inst.clone())
436 .or_insert_with(|| Series {
437 name: inst,
438 stack,
439 service,
440 replicas: None,
441 points: Vec::new(),
442 })
443 .points
444 .push((b, v));
445 }
446 let series: Vec<Series> = per.into_values().collect();
447 let series = match q.aggregate.as_deref() {
448 None => series,
449 Some(a) => vec![aggregate(&series, a, q.label())?],
450 };
451 Ok(Answer {
452 metric: q.metric.clone(),
453 from: q.from,
454 to: q.to,
455 step,
456 tier_step: TIERS[t].step,
457 series,
458 })
459 }
460}
461
462struct Borrowed<'a> {
464 conn: &'a Connection,
465}
466
467impl Borrowed<'_> {
468 fn series_id(&self, r: &Row) -> Result<i64> {
469 self.conn
470 .execute(
471 "INSERT INTO series(instance, stack, service, last_seen) VALUES(?1, ?2, ?3, ?4)
472 ON CONFLICT(instance) DO UPDATE SET stack=?2, service=?3, last_seen=?4",
473 params![r.instance, r.stack, r.service, r.ts as i64],
474 )
475 .map_err(|e| db_err("series", e))?;
476 self.conn
477 .query_row(
478 "SELECT id FROM series WHERE instance=?1",
479 params![r.instance],
480 |x| x.get(0),
481 )
482 .map_err(|e| db_err("series id", e))
483 }
484}
485
486pub const MAX_POINTS: u64 = 2000;
488
489#[derive(Debug, Clone, Default)]
491pub struct Query {
492 pub metric: String,
493 pub stack: Option<String>,
494 pub service: Option<String>,
495 pub instance: Option<String>,
496 pub from: u64,
498 pub to: u64,
499 pub step: u64,
501 pub aggregate: Option<String>,
504}
505
506impl Query {
507 fn label(&self) -> String {
508 match (&self.stack, &self.service, &self.instance) {
509 (_, _, Some(i)) => i.clone(),
510 (Some(st), Some(sv), None) => format!("{st}/{sv}"),
511 (Some(st), None, None) => st.clone(),
512 (None, Some(sv), None) => sv.clone(),
513 (None, None, None) => "org".into(),
514 }
515 }
516}
517
518#[derive(Debug, Clone, Serialize, PartialEq)]
520pub struct Series {
521 pub name: String,
523 pub stack: String,
524 pub service: String,
525 #[serde(skip_serializing_if = "Option::is_none")]
527 pub replicas: Option<Vec<u32>>,
528 pub points: Vec<(u64, f64)>,
529}
530
531#[derive(Debug, Clone, Serialize)]
532pub struct Answer {
533 pub metric: String,
534 pub from: u64,
535 pub to: u64,
536 pub step: u64,
537 pub tier_step: u64,
539 pub series: Vec<Series>,
540}
541
542pub fn aggregate(series: &[Series], how: &str, name: String) -> Result<Series> {
544 let mut by: BTreeMap<u64, Vec<f64>> = BTreeMap::new();
545 for s in series {
546 for (b, v) in &s.points {
547 by.entry(*b).or_default().push(*v);
548 }
549 }
550 let f: fn(&[f64]) -> f64 = match how {
551 "sum" => |v| v.iter().sum(),
552 "avg" => |v| v.iter().sum::<f64>() / v.len() as f64,
553 "max" => |v| v.iter().copied().fold(f64::MIN, f64::max),
554 "min" => |v| v.iter().copied().fold(f64::MAX, f64::min),
555 a => {
556 return Err(Error::invalid(format!(
557 "aggregate {a:?}: sum, avg, max or min"
558 )));
559 }
560 };
561 let one = |pick: fn(&Series) -> &String| {
562 let mut v: Vec<&String> = series.iter().map(pick).collect();
563 v.sort();
564 v.dedup();
565 if v.len() == 1 {
566 v[0].clone()
567 } else {
568 String::new()
569 }
570 };
571 Ok(Series {
572 name,
573 stack: one(|s| &s.stack),
574 service: one(|s| &s.service),
575 replicas: Some(by.values().map(|v| v.len() as u32).collect()),
576 points: by.iter().map(|(b, v)| (*b, f(v))).collect(),
577 })
578}
579
580#[derive(Clone)]
582pub struct History {
583 state: PathBuf,
584 dbs: Arc<Mutex<BTreeMap<OrgId, OrgDb>>>,
585 gone: Arc<Mutex<BTreeMap<OrgId, std::time::Instant>>>,
589}
590
591mod orgs;
592pub use orgs::{OrgSet, Sample};
593
594const MAINTAIN_EVERY: Duration = Duration::from_secs(60);
596
597impl History {
598 pub fn new(state: &Path) -> History {
599 History {
600 state: state.to_path_buf(),
601 dbs: Arc::new(Mutex::new(BTreeMap::new())),
602 gone: Default::default(),
603 }
604 }
605
606 pub fn path(&self, org: &OrgId) -> PathBuf {
607 org.dir(&self.state).join("metrics.db")
608 }
609
610 fn with<T>(&self, org: &OrgId, f: impl FnOnce(&mut OrgDb) -> Result<T>) -> Result<T> {
611 let mut dbs = self.dbs.lock().unwrap();
612 if !dbs.contains_key(org) {
613 let db = OrgDb::open(&self.path(org))?;
614 dbs.insert(org.clone(), db);
615 }
616 f(dbs.get_mut(org).expect("just opened"))
617 }
618
619 pub fn query(&self, org: &OrgId, q: &Query) -> Result<Answer> {
621 if !self.path(org).exists() {
622 return Ok(Answer {
623 metric: q.metric.clone(),
624 from: q.from,
625 to: q.to,
626 step: q.step.max(TIERS[0].step),
627 tier_step: TIERS[0].step,
628 series: Vec::new(),
629 });
630 }
631 self.with(org, |db| db.query(q, now_secs()))
632 }
633
634 pub fn start(&self) -> SyncSender<Sample> {
637 let (tx, rx) = std::sync::mpsc::sync_channel::<Sample>(8);
638 let me = self.clone();
639 let _ = std::thread::Builder::new()
640 .name("isb-metrics-history".into())
641 .spawn(move || me.writer(rx));
642 tx
643 }
644}
645
646pub fn offer(tx: &SyncSender<Sample>, s: Sample) {
648 let _ = tx.try_send(s);
650}
651
652fn now_secs() -> u64 {
653 crate::stack::controller::now_ms() / 1000
654}
655
656#[cfg(test)]
657mod tests {
658 use super::*;
659
660 fn inst(
661 name: &str,
662 slot: &str,
663 cpu: f32,
664 mem: u64,
665 rx: u64,
666 disk: Option<u64>,
667 ) -> InstanceSample {
668 let mut labels = BTreeMap::new();
669 labels.insert("isb.stack".to_string(), "shop".to_string());
670 labels.insert("isb.service".to_string(), "web".to_string());
671 labels.insert("isb.slot".to_string(), slot.to_string());
672 InstanceSample {
673 name: name.into(),
674 status: "Running".into(),
675 project: "isb-acme".into(),
676 cpu_pct: Some(cpu),
677 mem_bytes: Some(mem),
678 net_rx_bytes: Some(rx),
679 net_tx_bytes: Some(rx / 2),
680 disk_read_bytes: disk,
681 disk_write_bytes: disk,
682 labels,
683 ..Default::default()
684 }
685 }
686
687 #[test]
688 fn buckets_average_and_rates_follow_counters() {
689 let mut r = Recorder::default();
690 let t0 = 1_000_000_000_000; assert!(
692 r.add(t0, &[inst("a", "1", 10.0, 100, 0, Some(0))])
693 .is_empty()
694 );
695 assert!(
696 r.add(t0 + 2000, &[inst("a", "1", 30.0, 300, 2000, None)])
697 .is_empty()
698 );
699 assert!(
700 r.add(t0 + 4000, &[inst("a", "1", 20.0, 200, 6000, None)])
701 .is_empty()
702 );
703 let rows = r.add(t0 + 10_000, &[inst("a", "1", 0.0, 0, 6000, Some(10_000))]);
705 assert_eq!(rows.len(), 1);
706 let v = rows[0].values.0;
707 assert_eq!(rows[0].ts, t0 / 1000);
708 assert_eq!(v[0], Some(20.0));
709 assert_eq!(v[1], Some(200.0));
710 assert_eq!(v[2], Some(1500.0));
712 assert_eq!(v[3], Some(750.0));
713 assert_eq!(v[4], None);
715 assert_eq!(rows[0].stack, "shop");
716 assert_eq!(rows[0].service, "web");
717 r.add(t0 + 12_000, &[inst("a", "1", 0.0, 0, 10, None)]);
719 let rows = r.add(t0 + 20_000, &[]);
720 assert_eq!(rows.len(), 1, "a gone instance closes its bucket");
721 assert_eq!(rows[0].values.0[2], Some(0.0));
724 assert_eq!(rows[0].values.0[4], Some(1000.0));
726 let mut other = inst("x", "1", 1.0, 1, 1, None);
728 other.project = "titan-foo".into();
729 assert!(r.add(t0 + 30_000, &[other.clone()]).is_empty());
730 assert!(r.add(t0 + 40_000, &[other]).is_empty());
731 }
732
733 fn row(inst: &str, ts: u64, cpu: f64, mem: f64) -> Row {
734 Row {
735 org: OrgId::new("acme").unwrap(),
736 instance: inst.into(),
737 stack: "shop".into(),
738 service: "web".into(),
739 ts,
740 values: Values([Some(cpu), Some(mem), None, None, None, None]),
741 }
742 }
743
744 #[test]
745 #[expect(
746 clippy::too_many_lines,
747 reason = "predates the lint ratchet; split it when next changed"
748 )]
749 fn rollups_retention_and_queries() {
750 let dir = tempfile::tempdir().unwrap();
751 let mut db = OrgDb::open(&dir.path().join("m.db")).unwrap();
752 let now: u64 = 1_800_000_000 / 600 * 600;
753 let mut rows = Vec::new();
755 for ts in (now - 7200..now).step_by(10) {
756 rows.push(row("web-1", ts, 10.0, 100.0));
757 rows.push(row("web-2", ts, 30.0, 300.0));
758 }
759 rows.push(row("web-1", now - 40 * 86_400, 99.0, 1.0));
761 db.insert(&rows).unwrap();
762 db.maintain(now).unwrap();
763 let count = |db: &OrgDb, t: i64| -> i64 {
764 db.conn
765 .query_row("SELECT count(*) FROM samples WHERE tier=?1", [t], |r| {
766 r.get(0)
767 })
768 .unwrap()
769 };
770 assert_eq!(count(&db, 0), 2 * 720, "the ancient row is gone");
771 assert_eq!(count(&db, 1), 2 * 119);
773 assert_eq!(count(&db, 2), 2 * 11);
775 db.maintain(now).unwrap();
777 assert_eq!(count(&db, 1), 2 * 119);
778
779 let q = Query {
781 metric: "cpu".into(),
782 stack: Some("shop".into()),
783 service: Some("web".into()),
784 from: now - 3600,
785 to: now,
786 step: 60,
787 ..Default::default()
788 };
789 let a = db.query(&q, now).unwrap();
790 assert_eq!((a.step, a.tier_step), (60, 10));
791 assert_eq!(a.series.len(), 2);
792 assert_eq!(a.series[0].points.len(), 60);
793 assert!(a.series[0].points.iter().all(|(_, v)| *v == 10.0));
794 let a = db
796 .query(
797 &Query {
798 aggregate: Some("sum".into()),
799 ..q.clone()
800 },
801 now,
802 )
803 .unwrap();
804 assert_eq!(a.series.len(), 1);
805 assert_eq!(a.series[0].name, "shop/web");
806 assert!(a.series[0].points.iter().all(|(_, v)| *v == 40.0));
807 assert_eq!(a.series[0].replicas.as_ref().unwrap()[0], 2);
808 let a = db
809 .query(
810 &Query {
811 aggregate: Some("avg".into()),
812 metric: "memory".into(),
813 ..q.clone()
814 },
815 now,
816 )
817 .unwrap();
818 assert!(a.series[0].points.iter().all(|(_, v)| *v == 200.0));
819 let q2 = Query {
822 from: now - 2 * 86_400,
823 step: 0,
824 aggregate: Some("max".into()),
825 ..q.clone()
826 };
827 let a = db.query(&q2, now).unwrap();
828 assert_eq!(a.tier_step, 60);
829 assert!(
830 a.step >= 2 * 86_400 / MAX_POINTS && a.step % 60 == 0,
831 "{}",
832 a.step
833 );
834 let last = a.series[0].points.last().unwrap().0;
835 assert!(last + a.step >= now - 60, "{last} {now}");
836 assert!(
838 db.query(
839 &Query {
840 metric: "bogus".into(),
841 ..q.clone()
842 },
843 now
844 )
845 .is_err()
846 );
847 assert!(
848 db.query(
849 &Query {
850 aggregate: Some("median".into()),
851 ..q.clone()
852 },
853 now
854 )
855 .is_err()
856 );
857 let a = db
859 .query(
860 &Query {
861 instance: Some("web-2".into()),
862 stack: None,
863 service: None,
864 ..q
865 },
866 now,
867 )
868 .unwrap();
869 assert_eq!(a.series.len(), 1);
870 assert_eq!(a.series[0].name, "web-2");
871 }
872
873 #[test]
874 fn history_survives_reopening() {
875 let dir = tempfile::tempdir().unwrap();
876 let h = History::new(dir.path());
877 let org = OrgId::new("acme").unwrap();
878 let now = now_secs() / 10 * 10;
879 h.with(&org, |db| db.insert(&[row("web-1", now - 30, 5.0, 1.0)]))
880 .unwrap();
881 drop(h);
882 let h = History::new(dir.path());
883 let a = h
884 .query(
885 &org,
886 &Query {
887 metric: "cpu".into(),
888 from: now - 600,
889 to: now + 10,
890 step: 10,
891 ..Default::default()
892 },
893 )
894 .unwrap();
895 assert_eq!(a.series[0].points, vec![(now - 30, 5.0)]);
896 let other = OrgId::new("beta").unwrap();
898 let a = h
899 .query(
900 &other,
901 &Query {
902 metric: "cpu".into(),
903 from: 0,
904 to: now,
905 step: 10,
906 ..Default::default()
907 },
908 )
909 .unwrap();
910 assert!(a.series.is_empty());
911 assert!(!h.path(&other).exists());
912 }
913
914 #[test]
917 #[ignore]
918 fn measure_disk_use() {
919 let dir = tempfile::tempdir().unwrap();
920 let path = dir.path().join("m.db");
921 let mut db = OrgDb::open(&path).unwrap();
922 let start: u64 = 1_800_000_000 / 600 * 600;
923 let days = 31u64;
924 let t = std::time::Instant::now();
925 let mut writes = 0u64;
926 for ts in (start..start + days * 86_400).step_by(10) {
927 let rows: Vec<Row> = (0..20)
928 .map(|i| Row {
929 values: Values([
930 Some(i as f64 * 1.37),
931 Some(1e8 + i as f64),
932 Some(1234.5),
933 Some(99.0),
934 Some(4096.0),
935 Some(0.0),
936 ]),
937 ..row(&format!("app-web-{i}-abcd"), ts, 0.0, 0.0)
938 })
939 .collect();
940 db.insert(&rows).unwrap();
941 writes += 1;
942 if ts % 3600 == 0 {
944 db.maintain(ts).unwrap();
945 }
946 }
947 let el = t.elapsed();
948 let end = start + days * 86_400;
949 let m = std::time::Instant::now();
950 db.maintain(end).unwrap();
951 db.maintain(end + 60).unwrap();
952 let maintain = m.elapsed() / 2;
953 let q = std::time::Instant::now();
954 let a = db
955 .query(
956 &Query {
957 metric: "cpu".into(),
958 stack: Some("shop".into()),
959 service: Some("web".into()),
960 from: end - 86_400,
961 to: end,
962 aggregate: Some("sum".into()),
963 ..Default::default()
964 },
965 end,
966 )
967 .unwrap();
968 println!(
969 "steady-state maintain {maintain:?}; a 24 h service query ({} points) {:?}",
970 a.series[0].points.len(),
971 q.elapsed()
972 );
973 db.conn
974 .execute_batch("PRAGMA wal_checkpoint(TRUNCATE);")
975 .unwrap();
976 let size = std::fs::metadata(&path).unwrap().len();
977 let rows: i64 = db
978 .conn
979 .query_row("SELECT count(*) FROM samples", [], |r| r.get(0))
980 .unwrap();
981 println!(
982 "20 instances, {days} days: {rows} rows, {:.1} MiB on disk; {writes} transactions in {el:?} ({:.0} us each)",
983 size as f64 / 1048576.0,
984 el.as_micros() as f64 / writes as f64
985 );
986 }
987}