Skip to main content

isb_apps/monitor/
service.rs

1//! The monitoring service of a daemon: definitions, the scheduler, the
2//! workers that run checks, and the events they raise.
3//!
4//! One scheduler thread wakes every second and hands due checks to a fixed
5//! pool of workers through a bounded queue ([`WORKERS`], [`QUEUE`]), so a
6//! hundred slow targets never start a hundred threads, and a check never
7//! runs twice at once. A monitor's first check comes at a phase spread by
8//! its name, later ones every interval plus a little jitter, so monitors
9//! created together do not fire together.
10
11use std::collections::{BTreeMap, VecDeque};
12use std::path::{Path, PathBuf};
13use std::sync::atomic::{AtomicBool, Ordering};
14use std::sync::mpsc::{Receiver, SyncSender, TrySendError};
15use std::sync::{Arc, Mutex};
16use std::time::Duration;
17
18use serde::{Deserialize, Serialize};
19use serde_json::{Value, json};
20
21use super::state::{NEVER_UP_MS, Notify, State, Status, Thresholds};
22use super::store::{Check, Db};
23use super::target::{self, Ctx, Outcome};
24use super::{AUTO_PREFIX, Kind, MAX_PER_ORG, Monitor, Settings, dir, human, read_json, write_json};
25use crate::app::Apps;
26use crate::error::{Error, Result};
27use crate::org::OrgId;
28use crate::secrets::Secrets;
29use crate::stack::controller::now_ms;
30
31/// Checks running at once, at most.
32pub const WORKERS: usize = 8;
33/// Due checks waiting for a worker, at most; more wait for the next tick.
34pub const QUEUE: usize = 256;
35/// How often apps are looked at for their own monitors.
36const AUTO_SYNC_MS: u64 = 60_000;
37/// How often the history is rolled up and pruned.
38const ROLLUP_MS: u64 = 300_000;
39/// Event details kept for the notifier to pick up.
40const DETAILS_KEPT: usize = 256;
41
42/// Why an app's monitor has not checked yet.
43pub const WAITING_FOR_APP: &str = "waiting for the app's first live deployment";
44
45/// Is the address policy relaxed (the platform's private-targets setting)?
46pub type AllowPrivate = Arc<dyn Fn() -> bool + Send + Sync>;
47
48/// A monitor's state as kept: the state machine and its last check.
49#[derive(Debug, Clone, Default, Serialize, Deserialize)]
50#[serde(default)]
51pub struct Stored {
52    pub state: State,
53    pub last: Option<Outcome>,
54}
55
56#[derive(Debug, Clone, Copy, Default)]
57struct Slot {
58    next: u64,
59    running: bool,
60}
61
62struct Job {
63    org: OrgId,
64    monitor: Monitor,
65}
66
67type DetailKey = (OrgId, String, String);
68
69struct Inner {
70    state: PathBuf,
71    apps: Apps,
72    secrets: Arc<Secrets>,
73    allow_private: AllowPrivate,
74    public_url: Option<String>,
75    /// Held across read-modify-write of an org's files.
76    edit: Mutex<()>,
77    dbs: Mutex<BTreeMap<OrgId, Arc<Mutex<Db>>>>,
78    slots: Mutex<BTreeMap<(OrgId, String), Slot>>,
79    details: Mutex<VecDeque<(DetailKey, Value)>>,
80    stop: Arc<AtomicBool>,
81    tls: Arc<rustls::ClientConfig>,
82}
83
84/// The monitoring service.
85#[derive(Clone)]
86pub struct Monitors {
87    inner: Arc<Inner>,
88}
89
90fn fnv(s: &str) -> u64 {
91    s.bytes().fold(0xcbf2_9ce4_8422_2325, |h, b| {
92        (h ^ u64::from(b)).wrapping_mul(0x0100_0000_01b3)
93    })
94}
95
96/// A little randomness for the schedule, without a crate.
97fn jitter(max_ms: u64) -> i64 {
98    if max_ms == 0 {
99        return 0;
100    }
101    let n = std::time::SystemTime::now()
102        .duration_since(std::time::UNIX_EPOCH)
103        .map(|d| d.subsec_nanos())
104        .unwrap_or(0);
105    let r = fnv(&n.to_string()) % (2 * max_ms + 1);
106    r as i64 - max_ms as i64
107}
108
109impl Monitors {
110    pub fn new(
111        state: &Path,
112        apps: Apps,
113        secrets: Arc<Secrets>,
114        allow_private: AllowPrivate,
115        public_url: Option<String>,
116    ) -> Monitors {
117        Monitors {
118            inner: Arc::new(Inner {
119                state: state.to_path_buf(),
120                apps,
121                secrets,
122                allow_private,
123                public_url: public_url.map(|u| u.trim_end_matches('/').to_string()),
124                edit: Mutex::new(()),
125                dbs: Mutex::new(BTreeMap::new()),
126                slots: Mutex::new(BTreeMap::new()),
127                details: Mutex::new(VecDeque::new()),
128                stop: Arc::new(AtomicBool::new(false)),
129                tls: crate::net::default_tls(),
130            }),
131        }
132    }
133
134    /// Trust these TLS roots instead of the public set (tests).
135    pub fn with_tls(mut self, tls: Arc<rustls::ClientConfig>) -> Monitors {
136        if let Some(i) = Arc::get_mut(&mut self.inner) {
137            i.tls = tls;
138        }
139        self
140    }
141
142    fn path(&self, org: &OrgId, file: &str) -> PathBuf {
143        dir(&self.inner.state, org).join(file)
144    }
145
146    pub(crate) fn db(&self, org: &OrgId) -> Result<Arc<Mutex<Db>>> {
147        let mut dbs = self.inner.dbs.lock().unwrap();
148        if let Some(d) = dbs.get(org) {
149            return Ok(d.clone());
150        }
151        let d = Arc::new(Mutex::new(Db::open(&self.path(org, "monitors.db"))?));
152        dbs.insert(org.clone(), d.clone());
153        Ok(d)
154    }
155
156    /// The orgs with state on this daemon.
157    pub fn orgs(&self) -> Vec<OrgId> {
158        let mut out = vec![OrgId::default_org()];
159        if let Ok(rd) = std::fs::read_dir(self.inner.state.join("orgs")) {
160            for e in rd.flatten() {
161                if let Some(o) = e.file_name().to_str().and_then(|n| OrgId::new(n).ok()) {
162                    if !out.contains(&o) {
163                        out.push(o);
164                    }
165                }
166            }
167        }
168        out
169    }
170
171    // ---- definitions ----
172
173    pub fn list(&self, org: &OrgId) -> Result<Vec<Monitor>> {
174        read_json(&self.path(org, "monitors.json"))
175    }
176
177    fn save(&self, org: &OrgId, all: &[Monitor]) -> Result<()> {
178        write_json(&self.path(org, "monitors.json"), &all)
179    }
180
181    pub fn get(&self, org: &OrgId, name: &str) -> Result<Monitor> {
182        self.list(org)?
183            .into_iter()
184            .find(|m| m.name == name)
185            .ok_or_else(|| Error::NotFound(format!("monitor {name}")))
186    }
187
188    pub fn settings(&self, org: &OrgId) -> Result<Settings> {
189        read_json(&self.path(org, "settings.json"))
190    }
191
192    pub fn set_settings(&self, org: &OrgId, s: &Settings) -> Result<()> {
193        for a in &s.exclude_apps {
194            crate::app::validate_app_name(a)?;
195        }
196        let _g = self.inner.edit.lock().unwrap();
197        write_json(&self.path(org, "settings.json"), s)
198    }
199
200    /// What a monitor needs to exist: its app and its secrets.
201    fn check_refs(&self, org: &OrgId, m: &Monitor) -> Result<()> {
202        if m.kind == Kind::App {
203            let a = m.app.as_deref().unwrap_or_default();
204            self.inner.apps.get(org, a)?;
205        }
206        for s in m.secrets() {
207            self.inner.secrets.inspect(org, s).map_err(|e| {
208                if e.is_not_found() {
209                    Error::invalid(format!("no secret {s} in org {org}: create it first"))
210                } else {
211                    e
212                }
213            })?;
214        }
215        Ok(())
216    }
217
218    pub fn create(&self, org: &OrgId, mut m: Monitor) -> Result<Monitor> {
219        m.validate()?;
220        self.check_refs(org, &m)?;
221        let _g = self.inner.edit.lock().unwrap();
222        let mut all = self.list(org)?;
223        if all.iter().any(|x| x.name == m.name) {
224            return Err(Error::invalid(format!("monitor {} exists", m.name)));
225        }
226        if all.len() >= MAX_PER_ORG {
227            return Err(Error::invalid(format!(
228                "at most {MAX_PER_ORG} monitors per org"
229            )));
230        }
231        let now = now_ms() / 1000;
232        (m.created_at, m.updated_at) = (now, now);
233        all.push(m.clone());
234        self.save(org, &all)?;
235        self.due_soon(org, &m.name);
236        Ok(m)
237    }
238
239    /// Change fields of a monitor: `patch` holds the fields to set (null
240    /// puts a field back to its default). The name stays.
241    pub fn update(
242        &self,
243        org: &OrgId,
244        name: &str,
245        patch: serde_json::Map<String, Value>,
246    ) -> Result<Monitor> {
247        if patch.get("name").is_some_and(|n| n.as_str() != Some(name)) {
248            return Err(Error::invalid(
249                "a monitor's name cannot change; create another",
250            ));
251        }
252        let _g = self.inner.edit.lock().unwrap();
253        let mut all = self.list(org)?;
254        let cur = all
255            .iter_mut()
256            .find(|m| m.name == name)
257            .ok_or_else(|| Error::NotFound(format!("monitor {name}")))?;
258        let mut v = serde_json::to_value(&*cur)?;
259        let o = v.as_object_mut().expect("a monitor is an object");
260        for (k, x) in patch {
261            if matches!(k.as_str(), "created_at" | "updated_at" | "auto") {
262                continue;
263            }
264            if x.is_null() {
265                o.remove(&k);
266            } else {
267                o.insert(k, x);
268            }
269        }
270        let mut n: Monitor =
271            serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad monitor: {e}")))?;
272        n.validate()?;
273        self.check_refs(org, &n)?;
274        n.updated_at = now_ms() / 1000;
275        *cur = n.clone();
276        self.save(org, &all)?;
277        drop(_g);
278        self.edit_state(org, name, State::reset_counts)?;
279        self.due_soon(org, name);
280        Ok(n)
281    }
282
283    /// Remove a monitor and its history. An app's own monitor stays away:
284    /// the app joins the org's exclusions.
285    pub fn delete(&self, org: &OrgId, name: &str) -> Result<Monitor> {
286        let _g = self.inner.edit.lock().unwrap();
287        let mut all = self.list(org)?;
288        let i = all
289            .iter()
290            .position(|m| m.name == name)
291            .ok_or_else(|| Error::NotFound(format!("monitor {name}")))?;
292        let m = all.remove(i);
293        self.save(org, &all)?;
294        if m.auto {
295            let mut s = self.settings(org)?;
296            let app = m.app.clone().unwrap_or_default();
297            if !s.exclude_apps.contains(&app) {
298                s.exclude_apps.push(app);
299                write_json(&self.path(org, "settings.json"), &s)?;
300            }
301        }
302        drop(_g);
303        self.db(org)?.lock().unwrap().forget(name)?;
304        self.inner
305            .slots
306            .lock()
307            .unwrap()
308            .remove(&(org.clone(), name.to_string()));
309        Ok(m)
310    }
311
312    pub fn set_paused(&self, org: &OrgId, name: &str, paused: bool) -> Result<Monitor> {
313        let mut p = serde_json::Map::new();
314        p.insert("paused".into(), json!(paused));
315        self.update(org, name, p)
316    }
317
318    fn edit_state(&self, org: &OrgId, name: &str, f: impl FnOnce(&mut State)) -> Result<()> {
319        let db = self.db(org)?;
320        let db = db.lock().unwrap();
321        let mut s: Stored = db.load_state(name)?.unwrap_or_default();
322        f(&mut s.state);
323        db.save_state(name, &s)
324    }
325
326    pub(crate) fn stored(&self, org: &OrgId, name: &str) -> Result<Stored> {
327        Ok(self
328            .db(org)?
329            .lock()
330            .unwrap()
331            .load_state(name)?
332            .unwrap_or_default())
333    }
334
335    fn due_soon(&self, org: &OrgId, name: &str) {
336        let mut slots = self.inner.slots.lock().unwrap();
337        let s = slots.entry((org.clone(), name.to_string())).or_default();
338        s.next = now_ms() + 1000;
339    }
340
341    // ---- apps' own monitors ----
342
343    /// Give every app with a served domain its own monitor, and remove
344    /// those whose app is gone, has no domains, or opted out.
345    pub fn sync_auto(&self, org: &OrgId) -> Result<()> {
346        let settings = self.settings(org)?;
347        let apps = self.inner.apps.list(org)?;
348        let ctl = self.inner.apps.controller();
349        let served = |a: &crate::app::App| -> bool {
350            let Ok(stack) = a.spec.stack() else {
351                return false;
352            };
353            ctl.status(&crate::stack::qualified(org, &stack))
354                .ok()
355                .and_then(|s| s.services.into_iter().find(|x| x.service == a.spec.name))
356                .is_some_and(|s| s.domains.iter().any(|d| d.url.is_some()))
357        };
358        let wanted = |a: &crate::app::App| {
359            settings.auto_monitors
360                && !settings.exclude_apps.contains(&a.spec.name)
361                && !a.spec.domains.is_empty()
362        };
363        let _g = self.inner.edit.lock().unwrap();
364        let mut all = self.list(org)?;
365        let before = all.len();
366        let mut gone = Vec::new();
367        all.retain(|m| {
368            let keep = !m.auto
369                || apps
370                    .iter()
371                    .any(|a| Some(&a.spec.name) == m.app.as_ref() && wanted(a));
372            if !keep {
373                gone.push(m.name.clone());
374            }
375            keep
376        });
377        let mut added = Vec::new();
378        for a in apps.iter().filter(|a| wanted(a) && served(a)) {
379            let name = format!("{AUTO_PREFIX}{}", a.spec.name);
380            if all.len() >= MAX_PER_ORG || all.iter().any(|m| m.name == name) {
381                continue;
382            }
383            let mut m = Monitor::new(&name, Kind::App);
384            m.app = Some(a.spec.name.clone());
385            m.auto = true;
386            let now = now_ms() / 1000;
387            (m.created_at, m.updated_at) = (now, now);
388            added.push(name);
389            all.push(m);
390        }
391        if all.len() != before || !gone.is_empty() {
392            self.save(org, &all)?;
393        }
394        drop(_g);
395        for n in &gone {
396            self.db(org)?.lock().unwrap().forget(n)?;
397        }
398        for n in &added {
399            eprintln!("isb serve: monitor: {org}/{n}: watching the app's domain");
400            self.due_soon(org, n);
401        }
402        Ok(())
403    }
404
405    // ---- running ----
406
407    /// Start the scheduler and the workers.
408    pub fn start(&self) {
409        let (tx, rx) = std::sync::mpsc::sync_channel::<Job>(QUEUE);
410        let rx = Arc::new(Mutex::new(rx));
411        for i in 0..WORKERS {
412            let (me, rx) = (self.clone(), rx.clone());
413            let _ = std::thread::Builder::new()
414                .name(format!("isb-monitor-{i}"))
415                .spawn(move || me.worker(&rx));
416        }
417        let me = self.clone();
418        let _ = std::thread::Builder::new()
419            .name("isb-monitor".into())
420            .spawn(move || me.scheduler(&tx));
421    }
422
423    pub fn shutdown(&self) {
424        self.inner.stop.store(true, Ordering::SeqCst);
425    }
426
427    /// Set by [`Monitors::shutdown`]: what the heartbeat stops on too.
428    pub fn stopper(&self) -> Arc<AtomicBool> {
429        self.inner.stop.clone()
430    }
431
432    fn scheduler(&self, tx: &SyncSender<Job>) {
433        let (mut synced, mut rolled) = (0u64, now_ms());
434        while !self.inner.stop.load(Ordering::SeqCst) {
435            let now = now_ms();
436            if now.saturating_sub(synced) >= AUTO_SYNC_MS {
437                for o in self.orgs() {
438                    if let Err(e) = self.sync_auto(&o) {
439                        eprintln!("isb serve: monitor: {o}: apps' own monitors: {e}");
440                    }
441                }
442                synced = now;
443            }
444            if now.saturating_sub(rolled) >= ROLLUP_MS {
445                self.rollup_all(now);
446                rolled = now;
447            }
448            self.tick(now, tx);
449            std::thread::sleep(Duration::from_secs(1));
450        }
451    }
452
453    fn rollup_all(&self, now: u64) {
454        for o in self.orgs() {
455            if !self.path(&o, "monitors.db").exists() {
456                continue;
457            }
458            if let Err(e) = self.db(&o).and_then(|d| d.lock().unwrap().rollup(now)) {
459                eprintln!("isb serve: monitor: {o}: history rollup: {e}");
460            }
461        }
462    }
463
464    /// Queue every check that is due.
465    fn tick(&self, now: u64, tx: &SyncSender<Job>) {
466        let mut seen = Vec::new();
467        for org in self.orgs() {
468            for m in self.list(&org).unwrap_or_default() {
469                let key = (org.clone(), m.name.clone());
470                seen.push(key.clone());
471                if m.paused {
472                    continue;
473                }
474                let mut slots = self.inner.slots.lock().unwrap();
475                let s = slots.entry(key.clone()).or_insert_with(|| Slot {
476                    // The first check: a phase spread by name, within a minute.
477                    next: now + fnv(&format!("{}/{}", org, m.name)) % (m.interval.min(60) * 1000),
478                    running: false,
479                });
480                if s.running || now < s.next {
481                    continue;
482                }
483                match tx.try_send(Job {
484                    org: org.clone(),
485                    monitor: m,
486                }) {
487                    Ok(()) => s.running = true,
488                    Err(TrySendError::Full(_)) => return,
489                    Err(TrySendError::Disconnected(_)) => return,
490                }
491            }
492        }
493        self.inner
494            .slots
495            .lock()
496            .unwrap()
497            .retain(|k, _| seen.contains(k));
498    }
499
500    fn worker(&self, rx: &Mutex<Receiver<Job>>) {
501        loop {
502            let job = match rx.lock().unwrap().recv_timeout(Duration::from_secs(1)) {
503                Ok(j) => j,
504                Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
505                    if self.inner.stop.load(Ordering::SeqCst) {
506                        return;
507                    }
508                    continue;
509                }
510                Err(_) => return,
511            };
512            let started = now_ms();
513            let o = self.check(&job.org, &job.monitor);
514            if let Err(e) = self.record(&job.org, &job.monitor, o) {
515                eprintln!("isb serve: monitor: {}/{}: {e}", job.org, job.monitor.name);
516            }
517            let iv = job.monitor.interval * 1000;
518            let mut slots = self.inner.slots.lock().unwrap();
519            if let Some(s) = slots.get_mut(&(job.org, job.monitor.name)) {
520                s.running = false;
521                s.next = (started + iv).saturating_add_signed(jitter((iv / 20).min(2000)));
522            }
523        }
524    }
525
526    /// Run one check now. An app's monitor that has not been up yet does
527    /// not look at the app until it has a live deployment: the check fails
528    /// as "waiting", which the state machine keeps as pending.
529    pub fn check(&self, org: &OrgId, m: &Monitor) -> Outcome {
530        if m.kind == Kind::App
531            && self
532                .stored(org, &m.name)
533                .map(|s| s.state.status == Status::Pending)
534                .unwrap_or(true)
535            && !target::app_is_live(&self.inner.apps, org, m)
536        {
537            return Outcome {
538                at: now_ms(),
539                ok: false,
540                error: Some(WAITING_FOR_APP.into()),
541                ..Default::default()
542            };
543        }
544        let ctx = Ctx {
545            apps: &self.inner.apps,
546            ctl: self.inner.apps.controller(),
547            secrets: &self.inner.secrets,
548            allow_private: (self.inner.allow_private)(),
549            tls: self.inner.tls.clone(),
550        };
551        target::run(&ctx, org, m, now_ms())
552    }
553
554    /// Keep a check's result, step the state machine, raise events.
555    pub fn record(&self, org: &OrgId, m: &Monitor, o: Outcome) -> Result<()> {
556        // A monitor removed or paused while its check ran: drop the result.
557        match self.get(org, &m.name) {
558            Ok(cur) if !cur.paused => {}
559            _ => return Ok(()),
560        }
561        let db = self.db(org)?;
562        let db = db.lock().unwrap();
563        let mut s: Stored = db.load_state(&m.name)?.unwrap_or_default();
564        let t = Thresholds {
565            failures: m.failure_threshold,
566            recoveries: m.recovery_threshold,
567        };
568        let step = s.state.observe(o.ok, o.at, t);
569        db.insert(
570            &m.name,
571            &Check {
572                at: o.at,
573                ok: o.ok,
574                latency_ms: o.latency_ms,
575                status: o.status,
576                error: o.error.clone(),
577                pending: step.pending,
578            },
579        )?;
580        if let Notify::NeverUp { since } = step.notify {
581            let why = format!("never came up: {}", o.error.as_deref().unwrap_or("failed"));
582            db.open_incident(&m.name, since, Some(&why))?;
583        }
584        match step.changed {
585            Some(Status::Down) => {
586                let since = s.state.failing_since.unwrap_or(o.at);
587                db.open_incident(&m.name, since, o.error.as_deref())?;
588            }
589            Some(Status::Up) => db.close_incident(&m.name, o.at)?,
590            _ => {}
591        }
592        let cert = self.cert_due(m, &o, &mut s.state);
593        s.last = Some(o.clone());
594        db.save_state(&m.name, &s)?;
595        drop(db);
596        match step.notify {
597            Notify::Down { since, flapping } => {
598                self.emit_down(org, m, &o, since, flapping, s.state.fails)
599            }
600            Notify::NeverUp { since } => self.emit_never_up(org, m, &o, since, s.state.fails),
601            Notify::Up {
602                down_since,
603                downtime_ms,
604            } => self.emit_up(org, m, &o, down_since, downtime_ms),
605            Notify::None => {}
606        }
607        if let Some(days) = cert {
608            self.emit_cert(org, m, &o, days);
609        }
610        Ok(())
611    }
612
613    /// Days left on a certificate that is due a warning, once per
614    /// certificate.
615    fn cert_due(&self, m: &Monitor, o: &Outcome, s: &mut State) -> Option<i64> {
616        let exp = o.cert_expires?;
617        if m.cert_expiry_days == 0 || s.cert_warned == Some(exp) {
618            return None;
619        }
620        let days = (exp as i64 - (o.at / 1000) as i64).div_euclid(86_400);
621        (days <= i64::from(m.cert_expiry_days)).then(|| {
622            s.cert_warned = Some(exp);
623            days
624        })
625    }
626
627    // ---- events ----
628
629    /// Where a monitor's events go: an app's service (so channel rules on
630    /// apps and projects match), else `<org>/@monitors`, service = the
631    /// monitor.
632    fn subject(&self, org: &OrgId, m: &Monitor) -> (String, String) {
633        if let Some(a) = m
634            .app
635            .as_deref()
636            .and_then(|a| self.inner.apps.get(org, a).ok())
637        {
638            if let Ok(st) = a.spec.stack() {
639                return (crate::stack::qualified(org, &st), a.spec.name);
640            }
641        }
642        (crate::stack::qualified(org, "@monitors"), m.name.clone())
643    }
644
645    /// The web UI's page of a monitor, when isb knows its public URL.
646    pub fn link(&self, org: &OrgId, name: &str) -> Option<String> {
647        let base = self.inner.public_url.as_deref()?;
648        Some(format!("{base}/orgs/{org}/uptime/{name}"))
649    }
650
651    fn base_details(&self, org: &OrgId, m: &Monitor, o: &Outcome) -> Value {
652        json!({
653            "monitor": m.name,
654            "type": m.kind,
655            "app": m.app,
656            "url": o.url.clone().unwrap_or_else(|| m.target()),
657            "status": o.status,
658            "latency_ms": o.latency_ms,
659            "error": o.error,
660            "via": o.via,
661            "note": o.note,
662            "checked_at": o.at,
663            "link": self.link(org, &m.name),
664        })
665    }
666
667    fn emit(
668        &self,
669        org: &OrgId,
670        m: &Monitor,
671        kind: &str,
672        level: &str,
673        message: String,
674        details: Value,
675    ) {
676        let (stack, service) = self.subject(org, m);
677        {
678            let mut d = self.inner.details.lock().unwrap();
679            if d.len() >= DETAILS_KEPT {
680                d.pop_front();
681            }
682            d.push_back(((org.clone(), kind.to_string(), message.clone()), details));
683        }
684        self.inner
685            .apps
686            .controller()
687            .event(kind, level, &stack, &service, message);
688    }
689
690    fn emit_down(
691        &self,
692        org: &OrgId,
693        m: &Monitor,
694        o: &Outcome,
695        since: u64,
696        flapping: bool,
697        fails: u32,
698    ) {
699        let mut d = self.base_details(org, m, o);
700        d["down_since"] = json!(since);
701        d["failures"] = json!(fails);
702        d["flapping"] = json!(flapping);
703        self.emit(
704            org,
705            m,
706            "monitor.down",
707            "error",
708            down_message(m, o, fails, flapping),
709            d,
710        );
711    }
712
713    /// A new monitor that has not had one successful check in
714    /// [`NEVER_UP_MS`]: a `monitor.down` saying so, once.
715    fn emit_never_up(&self, org: &OrgId, m: &Monitor, o: &Outcome, since: u64, fails: u32) {
716        let mut d = self.base_details(org, m, o);
717        d["down_since"] = json!(since);
718        d["failures"] = json!(fails);
719        d["flapping"] = json!(false);
720        d["never_up"] = json!(true);
721        self.emit(
722            org,
723            m,
724            "monitor.down",
725            "error",
726            never_up_message(m, o, fails),
727            d,
728        );
729    }
730
731    fn emit_up(&self, org: &OrgId, m: &Monitor, o: &Outcome, down_since: u64, downtime_ms: u64) {
732        let mut d = self.base_details(org, m, o);
733        d["down_since"] = json!(down_since);
734        d["downtime_ms"] = json!(downtime_ms);
735        d["downtime"] = json!(human(downtime_ms));
736        self.emit(
737            org,
738            m,
739            "monitor.up",
740            "info",
741            up_message(m, o, downtime_ms),
742            d,
743        );
744    }
745
746    fn emit_cert(&self, org: &OrgId, m: &Monitor, o: &Outcome, days: i64) {
747        let mut d = self.base_details(org, m, o);
748        d["cert_expires_at"] = json!(o.cert_expires);
749        d["cert_days_left"] = json!(days);
750        self.emit(
751            org,
752            m,
753            "monitor.cert_expiring",
754            "warn",
755            cert_message(m, o, days),
756            d,
757        );
758    }
759
760    /// The details of an event this service raised, for the notifier.
761    pub fn details(&self, org: &OrgId, kind: &str, message: &str) -> Option<Value> {
762        let d = self.inner.details.lock().unwrap();
763        d.iter()
764            .rev()
765            .find(|((o, k, msg), _)| o == org && k == kind && msg == message)
766            .map(|(_, v)| v.clone())
767    }
768}
769
770fn what(m: &Monitor, o: &Outcome) -> String {
771    o.url.clone().unwrap_or_else(|| m.target())
772}
773
774pub fn down_message(m: &Monitor, o: &Outcome, fails: u32, flapping: bool) -> String {
775    let why = o.error.as_deref().unwrap_or("failed");
776    let mut s = format!(
777        "Monitor {} is DOWN: {}: {why} ({fails} failed check{} in a row)",
778        m.name,
779        what(m, o),
780        if fails == 1 { "" } else { "s" }
781    );
782    if flapping {
783        s.push_str("; it is flapping, so further changes are held until it is stable for 30 min");
784    }
785    s
786}
787
788pub fn never_up_message(m: &Monitor, o: &Outcome, fails: u32) -> String {
789    let why = o.error.as_deref().unwrap_or("failed");
790    format!(
791        "Monitor {} never came up: {}: {why} (no successful check in {} min, {fails} failed check{})",
792        m.name,
793        what(m, o),
794        NEVER_UP_MS / 60_000,
795        if fails == 1 { "" } else { "s" }
796    )
797}
798
799pub fn up_message(m: &Monitor, o: &Outcome, downtime_ms: u64) -> String {
800    let answer = match (o.status, o.latency_ms) {
801        (Some(st), Some(l)) => format!("answered HTTP {st} in {l} ms"),
802        (None, Some(l)) => format!("answered in {l} ms"),
803        _ => "answered".into(),
804    };
805    format!(
806        "Monitor {} is UP again after {}: {} {answer}",
807        m.name,
808        human(downtime_ms),
809        what(m, o)
810    )
811}
812
813pub fn cert_message(m: &Monitor, o: &Outcome, days: i64) -> String {
814    let on = o
815        .cert_expires
816        .map(|e| {
817            let d = (e / 86_400) as i64;
818            let (y, mo, da) = civil(d);
819            format!(" ({y:04}-{mo:02}-{da:02})")
820        })
821        .unwrap_or_default();
822    let when = if days < 0 {
823        "has expired".to_string()
824    } else {
825        format!("expires in {days} day{}", if days == 1 { "" } else { "s" })
826    };
827    format!(
828        "Monitor {}: the TLS certificate of {} {when}{on}",
829        m.name,
830        what(m, o)
831    )
832}
833
834/// The date of a day number since 1970-01-01.
835fn civil(z: i64) -> (i64, i64, i64) {
836    let z = z + 719_468;
837    let era = z.div_euclid(146_097);
838    let doe = z - era * 146_097;
839    let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
840    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
841    let mp = (5 * doy + 2) / 153;
842    let d = doy - (153 * mp + 2) / 5 + 1;
843    let m = if mp < 10 { mp + 3 } else { mp - 9 };
844    (yoe + era * 400 + i64::from(m <= 2), m, d)
845}
846
847#[cfg(test)]
848#[path = "service_tests.rs"]
849mod tests;