Skip to main content

isb_core/
metrics_history.rs

1//! Metrics history: the sampler's per-instance CPU, memory, network and disk
2//! numbers, kept for a month in downsampling tiers, per org, across daemon
3//! restarts.
4//!
5//! Storage is SQLite (`<state>/orgs/<org>/metrics.db`, the bundled rusqlite
6//! the identity store already uses) rather than a custom file format: tiers
7//! are a `GROUP BY` away, retention is a `DELETE`, a crash mid-write leaves
8//! the last committed transaction, and queries aggregate in SQL.
9//!
10//! | Tier | Step | Kept |
11//! |---|---|---|
12//! | 0 | 10 s | 24 h |
13//! | 1 | 1 min | 7 d |
14//! | 2 | 10 min | 30 d |
15//!
16//! Tier 0 is written as each 10 s bucket closes (the average of the samples
17//! in it); tiers 1 and 2 are rolled up from the tier below once a minute.
18//! Disk use is bounded by instances × (8640 + 10080 + 4320) rows of about 60
19//! bytes; memory by one open bucket per instance.
20
21use 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/// A downsampling tier: rows every `step` seconds, kept for `keep` seconds.
35#[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
56/// The metrics kept, in column order.
57pub const METRICS: [&str; 6] = [
58    "cpu",
59    "memory",
60    "net_rx",
61    "net_tx",
62    "disk_read",
63    "disk_write",
64];
65
66/// One instance's numbers over a bucket: CPU in percent of one core, memory
67/// in bytes, the rest in bytes per second. `None`: not measured.
68#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize)]
69pub struct Values(pub [Option<f64>; 6]);
70
71/// One closed tier-0 bucket of one instance.
72#[derive(Debug, Clone, PartialEq)]
73pub struct Row {
74    pub org: OrgId,
75    pub instance: String,
76    pub stack: String,
77    pub service: String,
78    /// Bucket start, unix seconds.
79    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/// Turns samples into closed buckets: counters into rates, samples into
97/// bucket averages.
98#[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        // The first sample, or a counter reset (a restart).
111        _ => (None, None),
112    }
113}
114
115impl Recorder {
116    /// Add one sample taken at `at_ms`; returns the buckets it closed.
117    /// Instances outside isb's orgs, and stopped ones, are skipped.
118    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        // Gone instances: close what they had.
189        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
227/// The history database of one org.
228pub 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        // auto_vacuum only takes on a new database, before any table.
261        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    /// Write closed tier-0 buckets in one transaction.
271    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    /// Roll closed buckets up into the coarser tiers and apply retention, as
303    /// of `now` (unix seconds).
304    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            // A bucket is complete once the tier below has closed past it
309            // (one more of its steps for the buckets still being written).
310            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    /// Series of one or more instances over `[from, to)` at `step` (see
361    /// [`Query`]).
362    pub fn query(&self, q: &Query, now: u64) -> Result<Answer> {
363        let span = q.to.saturating_sub(q.from).max(1);
364        // The finest tier still holding `from`.
365        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        // At most MAX_POINTS buckets.
371        step = step.max(span.div_ceil(MAX_POINTS));
372        step = step.div_ceil(TIERS[t].step) * TIERS[t].step;
373        // The coarser tier lags its roll-up: take the newest buckets from the
374        // tier below.
375        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
462/// A transaction's connection, for helpers that take `&OrgDb`.
463struct 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
486/// The most buckets a query answers with; a longer range gets a wider step.
487pub const MAX_POINTS: u64 = 2000;
488
489/// What to read.
490#[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    /// Unix seconds, `[from, to)`.
497    pub from: u64,
498    pub to: u64,
499    /// The bucket width asked for; widened to the tier's step at least.
500    pub step: u64,
501    /// `sum`, `avg` or `max` over the instances, bucket by bucket; `None`
502    /// returns one series per instance.
503    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/// One line on a chart: `[bucket start, value]` points, oldest first.
519#[derive(Debug, Clone, Serialize, PartialEq)]
520pub struct Series {
521    /// The instance, or what was aggregated.
522    pub name: String,
523    pub stack: String,
524    pub service: String,
525    /// Aggregates: how many instances had a value, per point.
526    #[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    /// The resolution of the tier read.
538    pub tier_step: u64,
539    pub series: Vec<Series>,
540}
541
542/// Combine per-instance series bucket by bucket.
543pub 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/// The history of every org under a state directory, behind one lock.
581#[derive(Clone)]
582pub struct History {
583    state: PathBuf,
584    dbs: Arc<Mutex<BTreeMap<OrgId, OrgDb>>>,
585    /// Orgs deleted this run, and when: buckets their instances close on
586    /// the way out must not make their database again, nor an org list
587    /// read before the deletion.
588    gone: Arc<Mutex<BTreeMap<OrgId, std::time::Instant>>>,
589}
590
591mod orgs;
592pub use orgs::{OrgSet, Sample};
593
594/// How often roll-ups and retention run.
595const 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    /// Answer a query in one org; an org with no history answers empty.
620    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    /// Start the writer: samples go in through the returned sender, which
635    /// never blocks the sampler (a full queue drops the sample).
636    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
646/// Hand a sample to the writer without waiting.
647pub fn offer(tx: &SyncSender<Sample>, s: Sample) {
648    // Full: the writer is behind (a slow disk); drop rather than stall.
649    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; // a multiple of 10 s
691        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        // Next bucket: the first closes.
704        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        // rx: 1000 B/s then 2000 B/s; the first sample has no rate.
711        assert_eq!(v[2], Some(1500.0));
712        assert_eq!(v[3], Some(750.0));
713        // Disk had one reading in the bucket: no rate yet.
714        assert_eq!(v[4], None);
715        assert_eq!(rows[0].stack, "shop");
716        assert_eq!(rows[0].service, "web");
717        // A counter reset gives no rate rather than a negative one.
718        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        // 6000 -> 6000 is 0 B/s; 6000 -> 10 (a reset) counts for nothing,
722        // rather than a huge or negative rate.
723        assert_eq!(rows[0].values.0[2], Some(0.0));
724        // The disk rate came from two readings 10 s apart.
725        assert_eq!(rows[0].values.0[4], Some(1000.0));
726        // Instances outside isb's orgs are not kept.
727        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        // 2 hours of two replicas, every 10 s: cpu 10 and 30.
754        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        // And one ancient row, past every tier's retention.
760        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        // Tier 1: complete minutes up to now - 20 s rounded down.
772        assert_eq!(count(&db, 1), 2 * 119);
773        // Tier 2: complete 10-minute buckets up to now - 120 s.
774        assert_eq!(count(&db, 2), 2 * 11);
775        // Maintenance is idempotent.
776        db.maintain(now).unwrap();
777        assert_eq!(count(&db, 1), 2 * 119);
778
779        // Per instance, at 1 min over the last hour.
780        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        // Summed over the replicas.
795        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        // Two days back reads tier 1 (raw is gone past 24 h), plus the
820        // newest raw buckets the roll-up has not reached.
821        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        // An unknown metric or aggregate is refused.
837        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        // One instance only.
858        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        // An org without history answers empty, creating nothing.
897        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    /// Disk use and write cost of 20 instances over 30 days (run by hand:
915    /// `cargo test --lib metrics_history -- --ignored --nocapture`).
916    #[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            // Hourly here (the daemon: every minute), to keep the run short.
943            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}