Skip to main content

isb_apps/monitor/
store.rs

1//! Check history, per org, in SQLite (`<state>/orgs/<org>/monitors/
2//! monitors.db`), the bundled rusqlite the metrics history uses.
3//!
4//! | Table | Rows | Kept |
5//! |---|---|---|
6//! | `checks` | every check: ok, pending, latency, HTTP status, error | 7 days |
7//! | `hourly` | per monitor and hour: checks, successes, pending checks, latency p50 and p95 | 90 days |
8//! | `incidents` | each time a monitor went down, and when it came back | 90 days after it ended |
9//! | `state` | each monitor's [`super::state::State`] and last check | while the monitor exists |
10//!
11//! A *pending* check is a failure before the monitor's first success: kept
12//! (grey in the history) but never counted as a check, so it is not in an
13//! uptime percentage.
14//!
15//! A row is about 60 bytes: a monitor checked every 30 s keeps about
16//! 20,000 raw rows and 2,160 hourly ones.
17
18use std::path::Path;
19
20use rusqlite::{Connection, OptionalExtension, params};
21use serde::{Deserialize, Serialize};
22
23use crate::error::{Error, Result};
24
25pub const RAW_KEEP_MS: u64 = 7 * 86_400_000;
26pub const HOURLY_KEEP_MS: u64 = 90 * 86_400_000;
27pub const HOUR_MS: u64 = 3_600_000;
28
29/// One check, as kept.
30#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
31pub struct Check {
32    /// Unix milliseconds.
33    pub at: u64,
34    pub ok: bool,
35    #[serde(skip_serializing_if = "Option::is_none")]
36    pub latency_ms: Option<u64>,
37    #[serde(skip_serializing_if = "Option::is_none")]
38    pub status: Option<u16>,
39    #[serde(skip_serializing_if = "Option::is_none")]
40    pub error: Option<String>,
41    /// A failure while the monitor waited for its first success.
42    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
43    pub pending: bool,
44}
45
46/// A time bucket of checks.
47#[derive(Debug, Clone, PartialEq, Serialize)]
48pub struct Bucket {
49    /// Unix milliseconds the bucket starts.
50    pub at: u64,
51    /// Counted checks (pending ones are not).
52    pub checks: u64,
53    pub ok: u64,
54    /// Checks that failed while the monitor waited for its first success.
55    pub pending: u64,
56    /// Percent of checks that succeeded; none without checks.
57    pub uptime: Option<f64>,
58    pub p50: Option<u64>,
59    pub p95: Option<u64>,
60}
61
62/// One outage.
63#[derive(Debug, Clone, PartialEq, Serialize)]
64pub struct Incident {
65    pub id: i64,
66    pub monitor: String,
67    /// Unix milliseconds: the first failed check.
68    pub started: u64,
69    #[serde(skip_serializing_if = "Option::is_none")]
70    pub ended: Option<u64>,
71    pub duration_ms: u64,
72    #[serde(skip_serializing_if = "Option::is_none")]
73    pub error: Option<String>,
74}
75
76/// An org's check history.
77pub struct Db {
78    conn: Connection,
79}
80
81fn db_err(e: rusqlite::Error) -> Error {
82    Error::invalid(format!("monitor history: {e}"))
83}
84
85const SCHEMA: &str = "
86CREATE TABLE IF NOT EXISTS checks (
87    monitor TEXT NOT NULL, ts INTEGER NOT NULL, ok INTEGER NOT NULL,
88    latency INTEGER, status INTEGER, error TEXT);
89CREATE INDEX IF NOT EXISTS checks_monitor_ts ON checks (monitor, ts);
90CREATE INDEX IF NOT EXISTS checks_ts ON checks (ts);
91CREATE TABLE IF NOT EXISTS hourly (
92    monitor TEXT NOT NULL, hour INTEGER NOT NULL, total INTEGER NOT NULL,
93    ok INTEGER NOT NULL, p50 INTEGER, p95 INTEGER,
94    PRIMARY KEY (monitor, hour));
95CREATE TABLE IF NOT EXISTS incidents (
96    id INTEGER PRIMARY KEY AUTOINCREMENT, monitor TEXT NOT NULL,
97    started INTEGER NOT NULL, ended INTEGER, error TEXT);
98CREATE INDEX IF NOT EXISTS incidents_monitor ON incidents (monitor, started);
99CREATE TABLE IF NOT EXISTS state (monitor TEXT PRIMARY KEY, json TEXT NOT NULL);
100";
101
102/// The `p`th percentile (0..=100) of sorted values, nearest rank.
103pub fn percentile(sorted: &[u64], p: f64) -> Option<u64> {
104    if sorted.is_empty() {
105        return None;
106    }
107    let rank = ((p / 100.0) * sorted.len() as f64).ceil() as usize;
108    Some(sorted[rank.clamp(1, sorted.len()) - 1])
109}
110
111fn i(v: u64) -> i64 {
112    i64::try_from(v).unwrap_or(i64::MAX)
113}
114
115impl Db {
116    pub fn open(path: &Path) -> Result<Db> {
117        if let Some(d) = path.parent() {
118            std::fs::create_dir_all(d)?;
119        }
120        let conn = Connection::open(path).map_err(db_err)?;
121        conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL;")
122            .map_err(db_err)?;
123        conn.execute_batch(SCHEMA).map_err(db_err)?;
124        let db = Db { conn };
125        db.migrate()?;
126        Ok(db)
127    }
128
129    /// Add the pending columns to older files, and once drop the downtime
130    /// recorded before a monitor's first success (see [`Db::heal_pending`]).
131    fn migrate(&self) -> Result<()> {
132        let mut added = false;
133        for t in ["checks", "hourly"] {
134            let has = self
135                .conn
136                .query_row(
137                    &format!(
138                        "SELECT COUNT(*) FROM pragma_table_info('{t}') WHERE name = 'pending'"
139                    ),
140                    [],
141                    |r| r.get::<_, i64>(0),
142                )
143                .map_err(db_err)?
144                > 0;
145            if !has {
146                self.conn
147                    .execute_batch(&format!(
148                        "ALTER TABLE {t} ADD COLUMN pending INTEGER NOT NULL DEFAULT 0"
149                    ))
150                    .map_err(db_err)?;
151                added = true;
152            }
153        }
154        if added {
155            self.heal_pending()?;
156        }
157        Ok(())
158    }
159
160    /// Failures before a monitor's first success were never downtime:
161    /// mark those checks pending, drop the incidents that began before it,
162    /// and put a monitor that never succeeded back to pending.
163    fn heal_pending(&self) -> Result<()> {
164        let monitors: Vec<String> = {
165            let mut st = self
166                .conn
167                .prepare("SELECT monitor FROM state UNION SELECT monitor FROM checks")
168                .map_err(db_err)?;
169            let rows = st.query_map([], |r| r.get(0)).map_err(db_err)?;
170            rows.collect::<std::result::Result<_, _>>()
171                .map_err(db_err)?
172        };
173        for m in monitors {
174            let first_ok: Option<i64> = self
175                .conn
176                .query_row(
177                    "SELECT MIN(t) FROM (SELECT MIN(ts) AS t FROM checks WHERE monitor = ?1 AND ok = 1 UNION ALL SELECT MIN(hour) FROM hourly WHERE monitor = ?1 AND ok > 0)",
178                    [&m],
179                    |r| r.get(0),
180                )
181                .map_err(db_err)?;
182            let until = first_ok.unwrap_or(i64::MAX);
183            self.conn
184                .execute(
185                    "UPDATE checks SET pending = 1 WHERE monitor = ?1 AND ok = 0 AND ts < ?2",
186                    params![m, until],
187                )
188                .map_err(db_err)?;
189            self.conn
190                .execute(
191                    "UPDATE hourly SET pending = total, total = 0 WHERE monitor = ?1 AND ok = 0 AND hour < ?2",
192                    params![m, until],
193                )
194                .map_err(db_err)?;
195            self.conn
196                .execute(
197                    "DELETE FROM incidents WHERE monitor = ?1 AND started < ?2",
198                    params![m, until],
199                )
200                .map_err(db_err)?;
201            if first_ok.is_none() {
202                self.reset_down_state(&m)?;
203            }
204        }
205        Ok(())
206    }
207
208    /// A monitor that never succeeded is not down: back to pending.
209    fn reset_down_state(&self, monitor: &str) -> Result<()> {
210        let Some(mut v) = self.load_state::<serde_json::Value>(monitor)? else {
211            return Ok(());
212        };
213        let Some(st) = v["state"].as_object_mut().filter(|o| o["status"] == "down") else {
214            return Ok(());
215        };
216        st.insert("status".into(), "pending".into());
217        st.insert("fails".into(), 0.into());
218        st.insert("downs".into(), serde_json::json!([]));
219        st.insert("flapping".into(), false.into());
220        st.remove("failing_since");
221        self.save_state(monitor, &v)
222    }
223
224    pub fn memory() -> Result<Db> {
225        let conn = Connection::open_in_memory().map_err(db_err)?;
226        conn.execute_batch(SCHEMA).map_err(db_err)?;
227        let db = Db { conn };
228        db.migrate()?;
229        Ok(db)
230    }
231
232    pub fn insert(&self, monitor: &str, c: &Check) -> Result<()> {
233        self.conn
234            .execute(
235                "INSERT INTO checks (monitor, ts, ok, latency, status, error, pending) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
236                params![monitor, i(c.at), c.ok, c.latency_ms.map(i), c.status, c.error, c.pending],
237            )
238            .map_err(db_err)?;
239        Ok(())
240    }
241
242    /// The newest checks first.
243    pub fn recent(&self, monitor: &str, limit: usize) -> Result<Vec<Check>> {
244        let mut st = self
245            .conn
246            .prepare(
247                "SELECT ts, ok, latency, status, error, pending FROM checks WHERE monitor = ?1 ORDER BY ts DESC LIMIT ?2",
248            )
249            .map_err(db_err)?;
250        let rows = st
251            .query_map(params![monitor, limit as i64], |r| {
252                Ok(Check {
253                    at: r.get::<_, i64>(0)? as u64,
254                    ok: r.get(1)?,
255                    latency_ms: r.get::<_, Option<i64>>(2)?.map(|v| v as u64),
256                    status: r.get(3)?,
257                    error: r.get(4)?,
258                    pending: r.get(5)?,
259                })
260            })
261            .map_err(db_err)?;
262        rows.collect::<std::result::Result<_, _>>().map_err(db_err)
263    }
264
265    /// Where the hourly rollup reaches: hours before this are rolled up.
266    fn rolled_until(&self) -> Result<Option<u64>> {
267        let h: Option<i64> = self
268            .conn
269            .query_row("SELECT MAX(hour) FROM hourly", [], |r| r.get(0))
270            .map_err(db_err)?;
271        Ok(h.map(|h| h as u64 + HOUR_MS))
272    }
273
274    /// Roll every complete hour not rolled up yet into `hourly`, then drop
275    /// what is past its keep.
276    pub fn rollup(&mut self, now: u64) -> Result<()> {
277        let current = now - now % HOUR_MS;
278        let from = match self.rolled_until()? {
279            Some(h) => h,
280            None => {
281                let first: Option<i64> = self
282                    .conn
283                    .query_row("SELECT MIN(ts) FROM checks", [], |r| r.get(0))
284                    .map_err(db_err)?;
285                match first {
286                    Some(t) => t as u64 - t as u64 % HOUR_MS,
287                    None => current,
288                }
289            }
290        };
291        let tx = self.conn.transaction().map_err(db_err)?;
292        let mut h = from.max(current.saturating_sub(RAW_KEEP_MS));
293        while h < current {
294            roll_hour(&tx, h)?;
295            h += HOUR_MS;
296        }
297        tx.execute(
298            "DELETE FROM checks WHERE ts < ?1",
299            [i(now.saturating_sub(RAW_KEEP_MS))],
300        )
301        .map_err(db_err)?;
302        tx.execute(
303            "DELETE FROM hourly WHERE hour < ?1",
304            [i(now.saturating_sub(HOURLY_KEEP_MS))],
305        )
306        .map_err(db_err)?;
307        tx.execute(
308            "DELETE FROM incidents WHERE ended IS NOT NULL AND ended < ?1",
309            [i(now.saturating_sub(HOURLY_KEEP_MS))],
310        )
311        .map_err(db_err)?;
312        tx.commit().map_err(db_err)
313    }
314
315    /// Checks and successes in `[from, to)`: whole hours from the rollup,
316    /// the rest from the raw checks.
317    pub fn counts(&self, monitor: &str, from: u64, to: u64) -> Result<(u64, u64)> {
318        let split = self.rolled_until()?.unwrap_or(0).clamp(from, to);
319        let (mut total, mut ok): (i64, i64) = self
320            .conn
321            .query_row(
322                "SELECT COALESCE(SUM(total), 0), COALESCE(SUM(ok), 0) FROM hourly WHERE monitor = ?1 AND hour >= ?2 AND hour < ?3",
323                params![monitor, i(from - from % HOUR_MS), i(split)],
324                |r| Ok((r.get(0)?, r.get(1)?)),
325            )
326            .map_err(db_err)?;
327        let (t2, o2): (i64, i64) = self
328            .conn
329            .query_row(
330                "SELECT COUNT(*), COALESCE(SUM(ok), 0) FROM checks WHERE monitor = ?1 AND ts >= ?2 AND ts < ?3 AND pending = 0",
331                params![monitor, i(split), i(to)],
332                |r| Ok((r.get(0)?, r.get(1)?)),
333            )
334            .map_err(db_err)?;
335        total += t2;
336        ok += o2;
337        Ok((total as u64, ok as u64))
338    }
339
340    /// Percent up over `[from, to)`; none without checks.
341    pub fn uptime(&self, monitor: &str, from: u64, to: u64) -> Result<Option<f64>> {
342        let (t, o) = self.counts(monitor, from, to)?;
343        Ok((t > 0).then(|| o as f64 * 100.0 / t as f64))
344    }
345
346    /// Latency percentiles (p50, p95) of successful checks since `from`.
347    pub fn latency(&self, monitor: &str, from: u64) -> Result<(Option<u64>, Option<u64>)> {
348        let v = self.latencies(monitor, from, u64::MAX)?;
349        Ok((percentile(&v, 50.0), percentile(&v, 95.0)))
350    }
351
352    fn latencies(&self, monitor: &str, from: u64, to: u64) -> Result<Vec<u64>> {
353        let mut st = self
354            .conn
355            .prepare_cached(
356                "SELECT latency FROM checks WHERE monitor = ?1 AND ts >= ?2 AND ts < ?3 AND ok = 1 AND latency IS NOT NULL",
357            )
358            .map_err(db_err)?;
359        let mut v: Vec<u64> = st
360            .query_map(params![monitor, i(from), i(to)], |r| r.get::<_, i64>(0))
361            .map_err(db_err)?
362            .filter_map(|x| x.ok().map(|x| x as u64))
363            .collect();
364        v.sort_unstable();
365        Ok(v)
366    }
367
368    /// Buckets of `step` over `[from, to)`. Steps under an hour come from
369    /// the raw checks (the last 7 days); an hour or more from the rollup
370    /// (p50 and p95 then the median of the hours').
371    pub fn buckets(&self, monitor: &str, from: u64, to: u64, step: u64) -> Result<Vec<Bucket>> {
372        let step = step.max(60_000);
373        let from = from - from % step;
374        let mut out: Vec<Bucket> = (0..(to.saturating_sub(from)).div_ceil(step))
375            .map(|k| Bucket {
376                at: from + k * step,
377                checks: 0,
378                ok: 0,
379                pending: 0,
380                uptime: None,
381                p50: None,
382                p95: None,
383            })
384            .collect();
385        if step < HOUR_MS {
386            self.fill_raw(monitor, (from, to), from, step, &mut out)?;
387        } else {
388            self.fill_hourly(monitor, from, to, step, &mut out)?;
389        }
390        for b in &mut out {
391            b.uptime = (b.checks > 0).then(|| b.ok as f64 * 100.0 / b.checks as f64);
392        }
393        Ok(out)
394    }
395
396    /// Count the raw checks in `[lo, hi)` into buckets of `step` from `base`.
397    fn fill_raw(
398        &self,
399        monitor: &str,
400        (lo, hi): (u64, u64),
401        base: u64,
402        step: u64,
403        out: &mut [Bucket],
404    ) -> Result<()> {
405        let mut st = self
406            .conn
407            .prepare(
408                "SELECT ts, ok, latency, pending FROM checks WHERE monitor = ?1 AND ts >= ?2 AND ts < ?3",
409            )
410            .map_err(db_err)?;
411        let rows = st
412            .query_map(params![monitor, i(lo), i(hi)], |r| {
413                Ok((
414                    r.get::<_, i64>(0)? as u64,
415                    r.get::<_, bool>(1)?,
416                    r.get::<_, Option<i64>>(2)?,
417                    r.get::<_, bool>(3)?,
418                ))
419            })
420            .map_err(db_err)?;
421        let mut lat: Vec<Vec<u64>> = vec![Vec::new(); out.len()];
422        for row in rows {
423            let (ts, ok, l, pending) = row.map_err(db_err)?;
424            let k = (ts.saturating_sub(base) / step) as usize;
425            let Some(b) = out.get_mut(k) else { continue };
426            if pending {
427                b.pending += 1;
428                continue;
429            }
430            b.checks += 1;
431            if ok {
432                b.ok += 1;
433                if let Some(l) = l {
434                    lat[k].push(l as u64);
435                }
436            }
437        }
438        for (b, mut l) in out.iter_mut().zip(lat) {
439            l.sort_unstable();
440            b.p50 = percentile(&l, 50.0);
441            b.p95 = percentile(&l, 95.0);
442        }
443        Ok(())
444    }
445
446    fn fill_hourly(
447        &self,
448        monitor: &str,
449        from: u64,
450        to: u64,
451        step: u64,
452        out: &mut [Bucket],
453    ) -> Result<()> {
454        let mut st = self
455            .conn
456            .prepare("SELECT hour, total, ok, p50, p95, pending FROM hourly WHERE monitor = ?1 AND hour >= ?2 AND hour < ?3")
457            .map_err(db_err)?;
458        type Row = (i64, i64, i64, Option<i64>, Option<i64>, i64);
459        let rows = st
460            .query_map(
461                params![monitor, i(from), i(to)],
462                |r| -> rusqlite::Result<Row> {
463                    Ok((
464                        r.get(0)?,
465                        r.get(1)?,
466                        r.get(2)?,
467                        r.get(3)?,
468                        r.get(4)?,
469                        r.get(5)?,
470                    ))
471                },
472            )
473            .map_err(db_err)?;
474        let mut p: Vec<(Vec<u64>, Vec<u64>)> = vec![(Vec::new(), Vec::new()); out.len()];
475        for row in rows {
476            let (h, total, ok, p50, p95, pending) = row.map_err(db_err)?;
477            let k = ((h as u64 - from) / step) as usize;
478            let Some(b) = out.get_mut(k) else { continue };
479            b.pending += pending as u64;
480            b.checks += total as u64;
481            b.ok += ok as u64;
482            p[k].0.extend(p50.map(|x| x as u64));
483            p[k].1.extend(p95.map(|x| x as u64));
484        }
485        // The hours not rolled up yet (the current one, mostly).
486        let split = self.rolled_until()?.unwrap_or(from).max(from);
487        let mut tail: Vec<Bucket> = out.to_vec();
488        for b in &mut tail {
489            b.checks = 0;
490            b.ok = 0;
491            b.pending = 0;
492        }
493        if split < to {
494            self.fill_raw(monitor, (split, to), from, step, &mut tail)?;
495        }
496        for ((b, t), (mut a, mut c)) in out.iter_mut().zip(tail).zip(p) {
497            b.checks += t.checks;
498            b.ok += t.ok;
499            b.pending += t.pending;
500            a.sort_unstable();
501            c.sort_unstable();
502            b.p50 = percentile(&a, 50.0).or(t.p50);
503            b.p95 = percentile(&c, 50.0).or(t.p95);
504        }
505        Ok(())
506    }
507
508    /// Open an incident; its id.
509    pub fn open_incident(&self, monitor: &str, started: u64, error: Option<&str>) -> Result<i64> {
510        self.conn
511            .execute(
512                "INSERT INTO incidents (monitor, started, error) VALUES (?1, ?2, ?3)",
513                params![monitor, i(started), error],
514            )
515            .map_err(db_err)?;
516        Ok(self.conn.last_insert_rowid())
517    }
518
519    /// End the monitor's open incident, if any.
520    pub fn close_incident(&self, monitor: &str, ended: u64) -> Result<()> {
521        self.conn
522            .execute(
523                "UPDATE incidents SET ended = ?2 WHERE monitor = ?1 AND ended IS NULL",
524                params![monitor, i(ended)],
525            )
526            .map_err(db_err)?;
527        Ok(())
528    }
529
530    /// Incidents, newest first: one monitor's, or every one's.
531    pub fn incidents(
532        &self,
533        monitor: Option<&str>,
534        limit: usize,
535        now: u64,
536    ) -> Result<Vec<Incident>> {
537        let mut st = self
538            .conn
539            .prepare(
540                "SELECT id, monitor, started, ended, error FROM incidents WHERE ?1 IS NULL OR monitor = ?1 ORDER BY started DESC LIMIT ?2",
541            )
542            .map_err(db_err)?;
543        let rows = st
544            .query_map(params![monitor, limit as i64], |r| {
545                let started = r.get::<_, i64>(2)? as u64;
546                let ended = r.get::<_, Option<i64>>(3)?.map(|v| v as u64);
547                Ok(Incident {
548                    id: r.get(0)?,
549                    monitor: r.get(1)?,
550                    started,
551                    ended,
552                    duration_ms: ended.unwrap_or(now).saturating_sub(started),
553                    error: r.get(4)?,
554                })
555            })
556            .map_err(db_err)?;
557        rows.collect::<std::result::Result<_, _>>().map_err(db_err)
558    }
559
560    pub fn load_state<T: for<'de> Deserialize<'de>>(&self, monitor: &str) -> Result<Option<T>> {
561        let j: Option<String> = self
562            .conn
563            .query_row(
564                "SELECT json FROM state WHERE monitor = ?1",
565                [monitor],
566                |r| r.get(0),
567            )
568            .optional()
569            .map_err(db_err)?;
570        Ok(j.and_then(|j| serde_json::from_str(&j).ok()))
571    }
572
573    pub fn save_state<T: Serialize>(&self, monitor: &str, s: &T) -> Result<()> {
574        self.conn
575            .execute(
576                "INSERT INTO state (monitor, json) VALUES (?1, ?2) ON CONFLICT (monitor) DO UPDATE SET json = ?2",
577                params![monitor, serde_json::to_string(s)?],
578            )
579            .map_err(db_err)?;
580        Ok(())
581    }
582
583    /// Forget a monitor: its checks, rollups, incidents and state.
584    pub fn forget(&self, monitor: &str) -> Result<()> {
585        for t in ["checks", "hourly", "incidents", "state"] {
586            self.conn
587                .execute(&format!("DELETE FROM {t} WHERE monitor = ?1"), [monitor])
588                .map_err(db_err)?;
589        }
590        Ok(())
591    }
592}
593
594/// Roll one hour of raw checks into `hourly`.
595fn roll_hour(tx: &rusqlite::Transaction, h: u64) -> Result<()> {
596    let mut st = tx
597        .prepare_cached(
598            "SELECT monitor, ok, latency, pending FROM checks WHERE ts >= ?1 AND ts < ?2 ORDER BY monitor",
599        )
600        .map_err(db_err)?;
601    let rows = st
602        .query_map([i(h), i(h + HOUR_MS)], |r| {
603            Ok((
604                r.get::<_, String>(0)?,
605                r.get::<_, bool>(1)?,
606                r.get::<_, Option<i64>>(2)?,
607                r.get::<_, bool>(3)?,
608            ))
609        })
610        .map_err(db_err)?;
611    let mut per: std::collections::BTreeMap<String, (u64, u64, Vec<u64>, u64)> = Default::default();
612    for row in rows {
613        let (m, ok, l, pending) = row.map_err(db_err)?;
614        let e = per.entry(m).or_default();
615        if pending {
616            e.3 += 1;
617            continue;
618        }
619        e.0 += 1;
620        if ok {
621            e.1 += 1;
622            e.2.extend(l.map(|l| l as u64));
623        }
624    }
625    for (m, (total, ok, mut l, pending)) in per {
626        l.sort_unstable();
627        tx.execute(
628            "INSERT OR REPLACE INTO hourly (monitor, hour, total, ok, p50, p95, pending) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
629            params![m, i(h), i(total), i(ok), percentile(&l, 50.0).map(i), percentile(&l, 95.0).map(i), i(pending)],
630        )
631        .map_err(db_err)?;
632    }
633    Ok(())
634}
635
636#[cfg(test)]
637mod tests {
638    use super::*;
639
640    fn check(at: u64, ok: bool, l: u64) -> Check {
641        Check {
642            at,
643            ok,
644            latency_ms: Some(l),
645            status: Some(if ok { 200 } else { 503 }),
646            error: (!ok).then(|| "HTTP 503".into()),
647            pending: false,
648        }
649    }
650
651    #[test]
652    fn percentiles() {
653        assert_eq!(percentile(&[], 50.0), None);
654        let v: Vec<u64> = (1..=100).collect();
655        assert_eq!(percentile(&v, 50.0), Some(50));
656        assert_eq!(percentile(&v, 95.0), Some(95));
657        assert_eq!(percentile(&[7], 95.0), Some(7));
658    }
659
660    #[test]
661    fn rollups_keep_uptime_and_latency() {
662        let mut db = Db::memory().unwrap();
663        let day0 = 100 * 86_400_000;
664        // Two days, a check a minute; the 10th hour of each day is down.
665        for m in 0..(2 * 24 * 60) {
666            let at = day0 + m * 60_000;
667            let ok = (at / HOUR_MS) % 24 != 10;
668            db.insert("web", &check(at, ok, 10 + m % 10)).unwrap();
669        }
670        let end = day0 + 2 * 86_400_000;
671        let before = db.uptime("web", day0, end).unwrap().unwrap();
672        assert!((before - 100.0 * 23.0 / 24.0).abs() < 1e-9, "{before}");
673        db.rollup(end + 1000).unwrap();
674        // Hours are rolled up; the totals do not change.
675        assert_eq!(db.uptime("web", day0, end).unwrap().unwrap(), before);
676        let (t, _) = db.counts("web", day0, end).unwrap();
677        assert_eq!(t, 2 * 24 * 60);
678        // A second rollup adds nothing twice.
679        db.rollup(end + 2000).unwrap();
680        assert_eq!(db.counts("web", day0, end).unwrap().0, t);
681        // Daily buckets from the rollup, minute buckets from the raw rows.
682        let days = db.buckets("web", day0, end, 86_400_000).unwrap();
683        assert_eq!(days.len(), 2);
684        assert_eq!(days[0].checks, 1440);
685        assert_eq!(days[0].ok, 1380);
686        assert!(days[0].p50.is_some());
687        let mins = db.buckets("web", day0, day0 + HOUR_MS, 600_000).unwrap();
688        assert_eq!(mins.len(), 6);
689        assert!(
690            mins.iter()
691                .all(|b| b.checks == 10 && b.uptime == Some(100.0))
692        );
693        let (p50, p95) = db.latency("web", day0).unwrap();
694        assert_eq!((p50, p95), (Some(14), Some(19)));
695        // Past the keeps, everything goes.
696        db.rollup(end + 91 * 86_400_000).unwrap();
697        assert_eq!(db.counts("web", day0, end).unwrap().0, 0);
698    }
699
700    #[test]
701    fn pending_checks_are_not_counted() {
702        let mut db = Db::memory().unwrap();
703        let t0 = 100 * 86_400_000;
704        for m in 0..30 {
705            let mut c = check(t0 + m * 60_000, false, 0);
706            c.pending = true;
707            db.insert("web", &c).unwrap();
708        }
709        for m in 30..40 {
710            db.insert("web", &check(t0 + m * 60_000, true, 5)).unwrap();
711        }
712        let end = t0 + HOUR_MS;
713        assert_eq!(db.counts("web", t0, end).unwrap(), (10, 10));
714        assert_eq!(db.uptime("web", t0, end).unwrap(), Some(100.0));
715        let b = db.buckets("web", t0, end, 20 * 60_000).unwrap();
716        assert_eq!((b[0].checks, b[0].pending, b[0].uptime), (0, 20, None));
717        assert_eq!((b[1].checks, b[1].pending), (10, 10));
718        let recent = db.recent("web", 50).unwrap();
719        assert_eq!(recent.iter().filter(|c| c.pending).count(), 30);
720        // The rollup keeps them out of the totals too.
721        db.rollup(end + 1000).unwrap();
722        assert_eq!(db.counts("web", t0, end).unwrap(), (10, 10));
723        let h = db.buckets("web", t0, end, HOUR_MS).unwrap();
724        assert_eq!((h[0].checks, h[0].ok, h[0].pending), (10, 10, 30));
725    }
726
727    #[test]
728    fn old_downtime_before_the_first_success_heals() {
729        let dir = tempfile::tempdir().unwrap();
730        let path = dir.path().join("m.db");
731        {
732            // An old file: no pending columns.
733            let conn = Connection::open(&path).unwrap();
734            conn.execute_batch(
735                "CREATE TABLE checks (monitor TEXT NOT NULL, ts INTEGER NOT NULL, ok INTEGER NOT NULL, latency INTEGER, status INTEGER, error TEXT);
736                 CREATE TABLE hourly (monitor TEXT NOT NULL, hour INTEGER NOT NULL, total INTEGER NOT NULL, ok INTEGER NOT NULL, p50 INTEGER, p95 INTEGER, PRIMARY KEY (monitor, hour));
737                 CREATE TABLE incidents (id INTEGER PRIMARY KEY AUTOINCREMENT, monitor TEXT NOT NULL, started INTEGER NOT NULL, ended INTEGER, error TEXT);
738                 CREATE TABLE state (monitor TEXT PRIMARY KEY, json TEXT NOT NULL);
739                 INSERT INTO checks VALUES ('umami', 1000, 0, NULL, NULL, 'refused'), ('umami', 2000, 0, NULL, NULL, 'refused');
740                 INSERT INTO incidents (monitor, started, error) VALUES ('umami', 1000, 'refused');
741                 INSERT INTO state VALUES ('umami', '{\"state\":{\"status\":\"down\",\"fails\":2,\"notified\":\"down\"}}');
742                 INSERT INTO checks VALUES ('shop', 1000, 0, NULL, NULL, 'x'), ('shop', 5000, 1, 3, 200, NULL), ('shop', 6000, 0, NULL, NULL, 'x'), ('shop', 7000, 0, NULL, NULL, 'x');
743                 INSERT INTO incidents (monitor, started, error) VALUES ('shop', 1000, 'x'), ('shop', 6000, 'x');",
744            )
745            .unwrap();
746        }
747        let db = Db::open(&path).unwrap();
748        assert!(db.incidents(Some("umami"), 10, 9000).unwrap().is_empty());
749        let s: serde_json::Value = db.load_state("umami").unwrap().unwrap();
750        assert_eq!(s["state"]["status"], "pending");
751        assert!(db.recent("umami", 10).unwrap().iter().all(|c| c.pending));
752        assert_eq!(db.counts("umami", 0, 9000).unwrap(), (0, 0));
753        // A real outage after the first success stays.
754        let inc = db.incidents(Some("shop"), 10, 9000).unwrap();
755        assert_eq!(inc.len(), 1);
756        assert_eq!(inc[0].started, 6000);
757        assert_eq!(db.counts("shop", 0, 9000).unwrap(), (3, 1));
758        // A second open changes nothing.
759        drop(db);
760        let db = Db::open(&path).unwrap();
761        assert_eq!(db.incidents(Some("shop"), 10, 9000).unwrap().len(), 1);
762    }
763
764    #[test]
765    fn incidents_and_state() {
766        let db = Db::memory().unwrap();
767        let id = db.open_incident("web", 1000, Some("HTTP 503")).unwrap();
768        assert_eq!(db.incidents(None, 10, 5000).unwrap()[0].duration_ms, 4000);
769        db.close_incident("web", 3000).unwrap();
770        let i = &db.incidents(Some("web"), 10, 9000).unwrap()[0];
771        assert_eq!((i.id, i.ended, i.duration_ms), (id, Some(3000), 2000));
772        assert!(db.incidents(Some("api"), 10, 0).unwrap().is_empty());
773        db.save_state("web", &serde_json::json!({"a": 1})).unwrap();
774        db.save_state("web", &serde_json::json!({"a": 2})).unwrap();
775        let s: serde_json::Value = db.load_state("web").unwrap().unwrap();
776        assert_eq!(s["a"], 2);
777        db.insert("web", &check(1, true, 5)).unwrap();
778        assert_eq!(db.recent("web", 5).unwrap().len(), 1);
779        db.forget("web").unwrap();
780        assert!(db.load_state::<serde_json::Value>("web").unwrap().is_none());
781        assert!(db.recent("web", 5).unwrap().is_empty());
782    }
783}