1use 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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
31pub struct Check {
32 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 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
43 pub pending: bool,
44}
45
46#[derive(Debug, Clone, PartialEq, Serialize)]
48pub struct Bucket {
49 pub at: u64,
51 pub checks: u64,
53 pub ok: u64,
54 pub pending: u64,
56 pub uptime: Option<f64>,
58 pub p50: Option<u64>,
59 pub p95: Option<u64>,
60}
61
62#[derive(Debug, Clone, PartialEq, Serialize)]
64pub struct Incident {
65 pub id: i64,
66 pub monitor: String,
67 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
76pub 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
102pub 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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
594fn 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 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 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 db.rollup(end + 2000).unwrap();
680 assert_eq!(db.counts("web", day0, end).unwrap().0, t);
681 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 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 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 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 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 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}