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::{Kind, MAX_PER_ORG, Monitor, Settings, auto, 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 and stacks 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/// Why a stack service's monitor has not checked yet.
45pub const WAITING_FOR_SERVICE: &str = "waiting for the service's first replica in rotation";
46
47/// Is the address policy relaxed (the platform's private-targets setting)?
48pub type AllowPrivate = Arc<dyn Fn() -> bool + Send + Sync>;
49
50/// A monitor's state as kept: the state machine and its last check.
51#[derive(Debug, Clone, Default, Serialize, Deserialize)]
52#[serde(default)]
53pub struct Stored {
54    pub state: State,
55    pub last: Option<Outcome>,
56}
57
58#[derive(Debug, Clone, Copy, Default)]
59struct Slot {
60    next: u64,
61    running: bool,
62}
63
64struct Job {
65    org: OrgId,
66    monitor: Monitor,
67}
68
69type DetailKey = (OrgId, String, String);
70
71struct Inner {
72    state: PathBuf,
73    apps: Apps,
74    secrets: Arc<Secrets>,
75    allow_private: AllowPrivate,
76    public_url: Option<String>,
77    /// Held across read-modify-write of an org's files.
78    edit: Mutex<()>,
79    dbs: Mutex<BTreeMap<OrgId, Arc<Mutex<Db>>>>,
80    slots: Mutex<BTreeMap<(OrgId, String), Slot>>,
81    details: Mutex<VecDeque<(DetailKey, Value)>>,
82    stop: Arc<AtomicBool>,
83    tls: Arc<rustls::ClientConfig>,
84}
85
86/// The monitoring service.
87#[derive(Clone)]
88pub struct Monitors {
89    inner: Arc<Inner>,
90}
91
92pub(super) fn fnv(s: &str) -> u64 {
93    s.bytes().fold(0xcbf2_9ce4_8422_2325, |h, b| {
94        (h ^ u64::from(b)).wrapping_mul(0x0100_0000_01b3)
95    })
96}
97
98/// A little randomness for the schedule, without a crate.
99fn jitter(max_ms: u64) -> i64 {
100    if max_ms == 0 {
101        return 0;
102    }
103    let n = std::time::SystemTime::now()
104        .duration_since(std::time::UNIX_EPOCH)
105        .map(|d| d.subsec_nanos())
106        .unwrap_or(0);
107    let r = fnv(&n.to_string()) % (2 * max_ms + 1);
108    r as i64 - max_ms as i64
109}
110
111impl Monitors {
112    pub fn new(
113        state: &Path,
114        apps: Apps,
115        secrets: Arc<Secrets>,
116        allow_private: AllowPrivate,
117        public_url: Option<String>,
118    ) -> Monitors {
119        Monitors {
120            inner: Arc::new(Inner {
121                state: state.to_path_buf(),
122                apps,
123                secrets,
124                allow_private,
125                public_url: public_url.map(|u| u.trim_end_matches('/').to_string()),
126                edit: Mutex::new(()),
127                dbs: Mutex::new(BTreeMap::new()),
128                slots: Mutex::new(BTreeMap::new()),
129                details: Mutex::new(VecDeque::new()),
130                stop: Arc::new(AtomicBool::new(false)),
131                tls: crate::net::default_tls(),
132            }),
133        }
134    }
135
136    /// Trust these TLS roots instead of the public set (tests).
137    pub fn with_tls(mut self, tls: Arc<rustls::ClientConfig>) -> Monitors {
138        if let Some(i) = Arc::get_mut(&mut self.inner) {
139            i.tls = tls;
140        }
141        self
142    }
143
144    fn path(&self, org: &OrgId, file: &str) -> PathBuf {
145        dir(&self.inner.state, org).join(file)
146    }
147
148    pub(crate) fn db(&self, org: &OrgId) -> Result<Arc<Mutex<Db>>> {
149        let mut dbs = self.inner.dbs.lock().unwrap();
150        if let Some(d) = dbs.get(org) {
151            return Ok(d.clone());
152        }
153        let d = Arc::new(Mutex::new(Db::open(&self.path(org, "monitors.db"))?));
154        dbs.insert(org.clone(), d.clone());
155        Ok(d)
156    }
157
158    /// The orgs with state on this daemon.
159    pub fn orgs(&self) -> Vec<OrgId> {
160        let mut out = vec![OrgId::default_org()];
161        if let Ok(rd) = std::fs::read_dir(self.inner.state.join("orgs")) {
162            for e in rd.flatten() {
163                if let Some(o) = e.file_name().to_str().and_then(|n| OrgId::new(n).ok()) {
164                    if !out.contains(&o) {
165                        out.push(o);
166                    }
167                }
168            }
169        }
170        out
171    }
172
173    // ---- definitions ----
174
175    pub fn list(&self, org: &OrgId) -> Result<Vec<Monitor>> {
176        read_json(&self.path(org, "monitors.json"))
177    }
178
179    fn save(&self, org: &OrgId, all: &[Monitor]) -> Result<()> {
180        write_json(&self.path(org, "monitors.json"), &all)
181    }
182
183    pub fn get(&self, org: &OrgId, name: &str) -> Result<Monitor> {
184        self.list(org)?
185            .into_iter()
186            .find(|m| m.name == name)
187            .ok_or_else(|| Error::NotFound(format!("monitor {name}")))
188    }
189
190    pub fn settings(&self, org: &OrgId) -> Result<Settings> {
191        read_json(&self.path(org, "settings.json"))
192    }
193
194    pub fn set_settings(&self, org: &OrgId, s: &Settings) -> Result<()> {
195        for a in &s.exclude_apps {
196            crate::app::validate_app_name(a)?;
197        }
198        for x in &s.exclude_services {
199            let ok = x.split_once('/').is_some_and(|(st, sv)| {
200                crate::stack::validate_stack_name(st).is_ok()
201                    && super::validate_service_name(sv).is_ok()
202            });
203            if !ok {
204                return Err(Error::invalid(format!(
205                    "exclude_services: {x:?} is not <stack>/<service>"
206                )));
207            }
208        }
209        let _g = self.inner.edit.lock().unwrap();
210        write_json(&self.path(org, "settings.json"), s)
211    }
212
213    /// What a monitor needs to exist: its app or stack service, and its
214    /// secrets.
215    fn check_refs(&self, org: &OrgId, m: &Monitor) -> Result<()> {
216        if m.kind == Kind::App {
217            let a = m.app.as_deref().unwrap_or_default();
218            self.inner.apps.get(org, a)?;
219        }
220        if m.kind == Kind::Service {
221            let (st, sv) = (
222                m.stack.as_deref().unwrap_or_default(),
223                m.service.as_deref().unwrap_or_default(),
224            );
225            let d = self
226                .inner
227                .apps
228                .controller()
229                .definition(&crate::stack::qualified(org, st))
230                .map_err(|_| Error::NotFound(format!("stack {st} in org {org}")))?;
231            if !d.file.services.contains_key(sv) {
232                return Err(Error::NotFound(format!("service {sv} in stack {st}")));
233            }
234        }
235        for s in m.secrets() {
236            self.inner.secrets.inspect(org, s).map_err(|e| {
237                if e.is_not_found() {
238                    Error::invalid(format!("no secret {s} in org {org}: create it first"))
239                } else {
240                    e
241                }
242            })?;
243        }
244        Ok(())
245    }
246
247    pub fn create(&self, org: &OrgId, mut m: Monitor) -> Result<Monitor> {
248        m.validate()?;
249        self.check_refs(org, &m)?;
250        let _g = self.inner.edit.lock().unwrap();
251        let mut all = self.list(org)?;
252        if all.iter().any(|x| x.name == m.name) {
253            return Err(Error::invalid(format!("monitor {} exists", m.name)));
254        }
255        if all.len() >= MAX_PER_ORG {
256            return Err(Error::invalid(format!(
257                "at most {MAX_PER_ORG} monitors per org"
258            )));
259        }
260        let now = now_ms() / 1000;
261        (m.created_at, m.updated_at) = (now, now);
262        all.push(m.clone());
263        self.save(org, &all)?;
264        self.due_soon(org, &m.name);
265        Ok(m)
266    }
267
268    /// Change fields of a monitor: `patch` holds the fields to set (null
269    /// puts a field back to its default). The name stays.
270    pub fn update(
271        &self,
272        org: &OrgId,
273        name: &str,
274        patch: serde_json::Map<String, Value>,
275    ) -> Result<Monitor> {
276        if patch.get("name").is_some_and(|n| n.as_str() != Some(name)) {
277            return Err(Error::invalid(
278                "a monitor's name cannot change; create another",
279            ));
280        }
281        let _g = self.inner.edit.lock().unwrap();
282        let mut all = self.list(org)?;
283        let cur = all
284            .iter_mut()
285            .find(|m| m.name == name)
286            .ok_or_else(|| Error::NotFound(format!("monitor {name}")))?;
287        let mut v = serde_json::to_value(&*cur)?;
288        let o = v.as_object_mut().expect("a monitor is an object");
289        for (k, x) in patch {
290            if matches!(k.as_str(), "created_at" | "updated_at" | "auto") {
291                continue;
292            }
293            if x.is_null() {
294                o.remove(&k);
295            } else {
296                o.insert(k, x);
297            }
298        }
299        let mut n: Monitor =
300            serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad monitor: {e}")))?;
301        n.validate()?;
302        self.check_refs(org, &n)?;
303        n.updated_at = now_ms() / 1000;
304        *cur = n.clone();
305        self.save(org, &all)?;
306        drop(_g);
307        self.edit_state(org, name, State::reset_counts)?;
308        self.due_soon(org, name);
309        Ok(n)
310    }
311
312    /// Remove a monitor and its history. An app's or a stack service's own
313    /// monitor stays away: it joins the org's exclusions.
314    pub fn delete(&self, org: &OrgId, name: &str) -> Result<Monitor> {
315        let _g = self.inner.edit.lock().unwrap();
316        let mut all = self.list(org)?;
317        let i = all
318            .iter()
319            .position(|m| m.name == name)
320            .ok_or_else(|| Error::NotFound(format!("monitor {name}")))?;
321        let m = all.remove(i);
322        self.save(org, &all)?;
323        if m.auto {
324            let mut s = self.settings(org)?;
325            let (list, x) = match m.kind {
326                Kind::Service => (
327                    &mut s.exclude_services,
328                    auto::exclusion(
329                        m.stack.as_deref().unwrap_or_default(),
330                        m.service.as_deref().unwrap_or_default(),
331                    ),
332                ),
333                _ => (&mut s.exclude_apps, m.app.clone().unwrap_or_default()),
334            };
335            if !list.contains(&x) {
336                list.push(x);
337                write_json(&self.path(org, "settings.json"), &s)?;
338            }
339        }
340        drop(_g);
341        self.db(org)?.lock().unwrap().forget(name)?;
342        self.inner
343            .slots
344            .lock()
345            .unwrap()
346            .remove(&(org.clone(), name.to_string()));
347        Ok(m)
348    }
349
350    pub fn set_paused(&self, org: &OrgId, name: &str, paused: bool) -> Result<Monitor> {
351        let mut p = serde_json::Map::new();
352        p.insert("paused".into(), json!(paused));
353        self.update(org, name, p)
354    }
355
356    fn edit_state(&self, org: &OrgId, name: &str, f: impl FnOnce(&mut State)) -> Result<()> {
357        let db = self.db(org)?;
358        let db = db.lock().unwrap();
359        let mut s: Stored = db.load_state(name)?.unwrap_or_default();
360        f(&mut s.state);
361        db.save_state(name, &s)
362    }
363
364    pub(crate) fn stored(&self, org: &OrgId, name: &str) -> Result<Stored> {
365        Ok(self
366            .db(org)?
367            .lock()
368            .unwrap()
369            .load_state(name)?
370            .unwrap_or_default())
371    }
372
373    fn due_soon(&self, org: &OrgId, name: &str) {
374        let mut slots = self.inner.slots.lock().unwrap();
375        let s = slots.entry((org.clone(), name.to_string())).or_default();
376        s.next = now_ms() + 1000;
377    }
378
379    // ---- apps' and stack services' own monitors ----
380
381    /// Give every app and compose stack service with a served domain its
382    /// own monitor, and remove those whose target is gone, has no domains,
383    /// or opted out ([`auto`]).
384    pub fn sync_auto(&self, org: &OrgId) -> Result<()> {
385        let settings = self.settings(org)?;
386        let mut cands = auto::apps(&self.inner.apps, org)?;
387        cands.extend(auto::stack_services(&self.inner.apps, org));
388        let _g = self.inner.edit.lock().unwrap();
389        let mut all = self.list(org)?;
390        let (gone, added) = auto::reconcile(&mut all, &settings, &cands, now_ms() / 1000);
391        if !gone.is_empty() || !added.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 its 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}: 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 or stack service's monitor that has not
527    /// been up yet does not look at it until it has a live deployment: the
528    /// check fails as "waiting", which the state machine keeps as pending.
529    pub fn check(&self, org: &OrgId, m: &Monitor) -> Outcome {
530        if m.follows()
531            && self
532                .stored(org, &m.name)
533                .map(|s| s.state.status == Status::Pending)
534                .unwrap_or(true)
535            && !target::is_live(&self.inner.apps, org, m)
536        {
537            let why = if m.kind == Kind::App {
538                WAITING_FOR_APP
539            } else {
540                WAITING_FOR_SERVICE
541            };
542            return Outcome {
543                at: now_ms(),
544                ok: false,
545                error: Some(why.into()),
546                ..Default::default()
547            };
548        }
549        let ctx = Ctx {
550            apps: &self.inner.apps,
551            ctl: self.inner.apps.controller(),
552            secrets: &self.inner.secrets,
553            allow_private: (self.inner.allow_private)(),
554            tls: self.inner.tls.clone(),
555        };
556        target::run(&ctx, org, m, now_ms())
557    }
558
559    /// Keep a check's result, step the state machine, raise events.
560    pub fn record(&self, org: &OrgId, m: &Monitor, o: Outcome) -> Result<()> {
561        // A monitor removed or paused while its check ran: drop the result.
562        match self.get(org, &m.name) {
563            Ok(cur) if !cur.paused => {}
564            _ => return Ok(()),
565        }
566        let db = self.db(org)?;
567        let db = db.lock().unwrap();
568        let mut s: Stored = db.load_state(&m.name)?.unwrap_or_default();
569        let t = Thresholds {
570            failures: m.failure_threshold,
571            recoveries: m.recovery_threshold,
572        };
573        let step = s.state.observe(o.ok, o.at, t);
574        db.insert(
575            &m.name,
576            &Check {
577                at: o.at,
578                ok: o.ok,
579                latency_ms: o.latency_ms,
580                status: o.status,
581                error: o.error.clone(),
582                pending: step.pending,
583            },
584        )?;
585        if let Notify::NeverUp { since } = step.notify {
586            let why = format!("never came up: {}", o.error.as_deref().unwrap_or("failed"));
587            db.open_incident(&m.name, since, Some(&why))?;
588        }
589        match step.changed {
590            Some(Status::Down) => {
591                let since = s.state.failing_since.unwrap_or(o.at);
592                db.open_incident(&m.name, since, o.error.as_deref())?;
593            }
594            Some(Status::Up) => db.close_incident(&m.name, o.at)?,
595            _ => {}
596        }
597        let cert = self.cert_due(m, &o, &mut s.state);
598        s.last = Some(o.clone());
599        db.save_state(&m.name, &s)?;
600        drop(db);
601        match step.notify {
602            Notify::Down { since, flapping } => {
603                self.emit_down(org, m, &o, since, flapping, s.state.fails)
604            }
605            Notify::NeverUp { since } => self.emit_never_up(org, m, &o, since, s.state.fails),
606            Notify::Up {
607                down_since,
608                downtime_ms,
609            } => self.emit_up(org, m, &o, down_since, downtime_ms),
610            Notify::None => {}
611        }
612        if let Some(days) = cert {
613            self.emit_cert(org, m, &o, days);
614        }
615        Ok(())
616    }
617
618    /// Days left on a certificate that is due a warning, once per
619    /// certificate.
620    fn cert_due(&self, m: &Monitor, o: &Outcome, s: &mut State) -> Option<i64> {
621        let exp = o.cert_expires?;
622        if m.cert_expiry_days == 0 || s.cert_warned == Some(exp) {
623            return None;
624        }
625        let days = (exp as i64 - (o.at / 1000) as i64).div_euclid(86_400);
626        (days <= i64::from(m.cert_expiry_days)).then(|| {
627            s.cert_warned = Some(exp);
628            days
629        })
630    }
631
632    // ---- events ----
633
634    /// Where a monitor's events go: an app's service (so channel rules on
635    /// apps and projects match), a stack service's own, else
636    /// `<org>/@monitors`, service = the monitor.
637    fn subject(&self, org: &OrgId, m: &Monitor) -> (String, String) {
638        if let (Kind::Service, Some(st), Some(sv)) = (m.kind, &m.stack, &m.service) {
639            return (crate::stack::qualified(org, st), sv.clone());
640        }
641        if let Some(a) = m
642            .app
643            .as_deref()
644            .and_then(|a| self.inner.apps.get(org, a).ok())
645        {
646            if let Ok(st) = a.spec.stack() {
647                return (crate::stack::qualified(org, &st), a.spec.name);
648            }
649        }
650        (crate::stack::qualified(org, "@monitors"), m.name.clone())
651    }
652
653    /// The web UI's page of a monitor, when isb knows its public URL.
654    pub fn link(&self, org: &OrgId, name: &str) -> Option<String> {
655        let base = self.inner.public_url.as_deref()?;
656        Some(format!("{base}/orgs/{org}/uptime/{name}"))
657    }
658
659    fn base_details(&self, org: &OrgId, m: &Monitor, o: &Outcome) -> Value {
660        json!({
661            "monitor": m.name,
662            "type": m.kind,
663            "app": m.app,
664            "stack": m.stack,
665            "service": m.service,
666            "url": o.url.clone().unwrap_or_else(|| m.target()),
667            "status": o.status,
668            "latency_ms": o.latency_ms,
669            "error": o.error,
670            "via": o.via,
671            "note": o.note,
672            "checked_at": o.at,
673            "link": self.link(org, &m.name),
674        })
675    }
676
677    fn emit(
678        &self,
679        org: &OrgId,
680        m: &Monitor,
681        kind: &str,
682        level: &str,
683        message: String,
684        details: Value,
685    ) {
686        let (stack, service) = self.subject(org, m);
687        {
688            let mut d = self.inner.details.lock().unwrap();
689            if d.len() >= DETAILS_KEPT {
690                d.pop_front();
691            }
692            d.push_back(((org.clone(), kind.to_string(), message.clone()), details));
693        }
694        self.inner
695            .apps
696            .controller()
697            .event(kind, level, &stack, &service, message);
698    }
699
700    fn emit_down(
701        &self,
702        org: &OrgId,
703        m: &Monitor,
704        o: &Outcome,
705        since: u64,
706        flapping: bool,
707        fails: u32,
708    ) {
709        let mut d = self.base_details(org, m, o);
710        d["down_since"] = json!(since);
711        d["failures"] = json!(fails);
712        d["flapping"] = json!(flapping);
713        self.emit(
714            org,
715            m,
716            "monitor.down",
717            "error",
718            down_message(m, o, fails, flapping),
719            d,
720        );
721    }
722
723    /// A new monitor that has not had one successful check in
724    /// [`NEVER_UP_MS`]: a `monitor.down` saying so, once.
725    fn emit_never_up(&self, org: &OrgId, m: &Monitor, o: &Outcome, since: u64, fails: u32) {
726        let mut d = self.base_details(org, m, o);
727        d["down_since"] = json!(since);
728        d["failures"] = json!(fails);
729        d["flapping"] = json!(false);
730        d["never_up"] = json!(true);
731        self.emit(
732            org,
733            m,
734            "monitor.down",
735            "error",
736            never_up_message(m, o, fails),
737            d,
738        );
739    }
740
741    fn emit_up(&self, org: &OrgId, m: &Monitor, o: &Outcome, down_since: u64, downtime_ms: u64) {
742        let mut d = self.base_details(org, m, o);
743        d["down_since"] = json!(down_since);
744        d["downtime_ms"] = json!(downtime_ms);
745        d["downtime"] = json!(human(downtime_ms));
746        self.emit(
747            org,
748            m,
749            "monitor.up",
750            "info",
751            up_message(m, o, downtime_ms),
752            d,
753        );
754    }
755
756    fn emit_cert(&self, org: &OrgId, m: &Monitor, o: &Outcome, days: i64) {
757        let mut d = self.base_details(org, m, o);
758        d["cert_expires_at"] = json!(o.cert_expires);
759        d["cert_days_left"] = json!(days);
760        self.emit(
761            org,
762            m,
763            "monitor.cert_expiring",
764            "warn",
765            cert_message(m, o, days),
766            d,
767        );
768    }
769
770    /// The details of an event this service raised, for the notifier.
771    pub fn details(&self, org: &OrgId, kind: &str, message: &str) -> Option<Value> {
772        let d = self.inner.details.lock().unwrap();
773        d.iter()
774            .rev()
775            .find(|((o, k, msg), _)| o == org && k == kind && msg == message)
776            .map(|(_, v)| v.clone())
777    }
778}
779
780fn what(m: &Monitor, o: &Outcome) -> String {
781    o.url.clone().unwrap_or_else(|| m.target())
782}
783
784pub fn down_message(m: &Monitor, o: &Outcome, fails: u32, flapping: bool) -> String {
785    let why = o.error.as_deref().unwrap_or("failed");
786    let mut s = format!(
787        "Monitor {} is DOWN: {}: {why} ({fails} failed check{} in a row)",
788        m.name,
789        what(m, o),
790        if fails == 1 { "" } else { "s" }
791    );
792    if flapping {
793        s.push_str("; it is flapping, so further changes are held until it is stable for 30 min");
794    }
795    s
796}
797
798pub fn never_up_message(m: &Monitor, o: &Outcome, fails: u32) -> String {
799    let why = o.error.as_deref().unwrap_or("failed");
800    format!(
801        "Monitor {} never came up: {}: {why} (no successful check in {} min, {fails} failed check{})",
802        m.name,
803        what(m, o),
804        NEVER_UP_MS / 60_000,
805        if fails == 1 { "" } else { "s" }
806    )
807}
808
809pub fn up_message(m: &Monitor, o: &Outcome, downtime_ms: u64) -> String {
810    let answer = match (o.status, o.latency_ms) {
811        (Some(st), Some(l)) => format!("answered HTTP {st} in {l} ms"),
812        (None, Some(l)) => format!("answered in {l} ms"),
813        _ => "answered".into(),
814    };
815    format!(
816        "Monitor {} is UP again after {}: {} {answer}",
817        m.name,
818        human(downtime_ms),
819        what(m, o)
820    )
821}
822
823pub fn cert_message(m: &Monitor, o: &Outcome, days: i64) -> String {
824    let on = o
825        .cert_expires
826        .map(|e| {
827            let d = (e / 86_400) as i64;
828            let (y, mo, da) = civil(d);
829            format!(" ({y:04}-{mo:02}-{da:02})")
830        })
831        .unwrap_or_default();
832    let when = if days < 0 {
833        "has expired".to_string()
834    } else {
835        format!("expires in {days} day{}", if days == 1 { "" } else { "s" })
836    };
837    format!(
838        "Monitor {}: the TLS certificate of {} {when}{on}",
839        m.name,
840        what(m, o)
841    )
842}
843
844/// The date of a day number since 1970-01-01.
845fn civil(z: i64) -> (i64, i64, i64) {
846    let z = z + 719_468;
847    let era = z.div_euclid(146_097);
848    let doe = z - era * 146_097;
849    let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
850    let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
851    let mp = (5 * doy + 2) / 153;
852    let d = doy - (153 * mp + 2) / 5 + 1;
853    let m = if mp < 10 { mp + 3 } else { mp - 9 };
854    (yoe + era * 400 + i64::from(m <= 2), m, d)
855}
856
857#[cfg(test)]
858#[path = "service_tests.rs"]
859mod tests;