Skip to main content

isb_apps/monitor/
state.rs

1//! A monitor's state machine: failure and recovery thresholds (hysteresis),
2//! one notification per change (de-duplication) and flap damping. No I/O,
3//! so every rule is tested on its own.
4//!
5//! - `failure_threshold` failed checks in a row take a monitor down;
6//!   `recovery_threshold` successful ones bring it back up. A new monitor's
7//!   first success makes it up at once.
8//! - A new monitor is `pending` until its first success. Failures before
9//!   that are *pending* checks, not downtime: no threshold, no incident, no
10//!   `down`. If it has only failed for [`NEVER_UP_MS`], channels hear once
11//!   that it *never came up* (a `down`), and the monitor stays pending
12//!   (flagged `never_up`) until a check succeeds, which then sends an `up`.
13//! - Channels hear `down` once per incident and `up` once after it, always
14//!   alternating: an up is only sent after a down was.
15//! - A monitor that went down [`FLAP_DOWNS`] times within [`FLAP_WINDOW_MS`]
16//!   is *flapping*: the change that makes it so is still sent (saying so),
17//!   then nothing more until it has held one state for [`FLAP_WINDOW_MS`];
18//!   then the state it settled in is sent if channels last heard otherwise.
19//! - An app or stack service scaled to 0 (stopped) is not checked: its
20//!   monitor is *stopped*, quietly. Stopping closes any open incident
21//!   without telling channels, and starting again begins as a new monitor
22//!   does, pending until its first success.
23
24use serde::{Deserialize, Serialize};
25
26/// How long a monitor must hold one state to stop flapping, and the window
27/// downs are counted in.
28pub const FLAP_WINDOW_MS: u64 = 30 * 60 * 1000;
29/// Downs within the window that make a monitor flapping.
30pub const FLAP_DOWNS: usize = 3;
31/// How long a new monitor may only fail before it is called "never came up".
32pub const NEVER_UP_MS: u64 = 30 * 60 * 1000;
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
35#[serde(rename_all = "lowercase")]
36pub enum Status {
37    /// Not checked yet (or channels were never told anything).
38    #[default]
39    Pending,
40    Up,
41    Down,
42}
43
44impl Status {
45    pub fn as_str(self) -> &'static str {
46        match self {
47            Status::Pending => "pending",
48            Status::Up => "up",
49            Status::Down => "down",
50        }
51    }
52}
53
54/// How many checks in a row change the status.
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub struct Thresholds {
57    pub failures: u32,
58    pub recoveries: u32,
59}
60
61/// A monitor's state, kept across daemon restarts so a restart never sends
62/// a down (or an up) twice.
63#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
64#[serde(default)]
65pub struct State {
66    pub status: Status,
67    /// Unix milliseconds the status began.
68    pub since: u64,
69    /// Failed checks in a row.
70    pub fails: u32,
71    /// Successful checks in a row.
72    pub oks: u32,
73    /// When the current run of failures began.
74    #[serde(skip_serializing_if = "Option::is_none")]
75    pub failing_since: Option<u64>,
76    /// What channels were last told.
77    pub notified: Status,
78    /// When the incident channels were told about began.
79    #[serde(skip_serializing_if = "Option::is_none")]
80    pub down_since: Option<u64>,
81    /// Unix milliseconds of recent downs, for flap damping.
82    #[serde(skip_serializing_if = "Vec::is_empty")]
83    pub downs: Vec<u64>,
84    pub flapping: bool,
85    /// The certificate (by its expiry, unix seconds) already warned about.
86    #[serde(skip_serializing_if = "Option::is_none")]
87    pub cert_warned: Option<u64>,
88    /// When the first check of a pending monitor ran.
89    #[serde(skip_serializing_if = "Option::is_none")]
90    pub first_check: Option<u64>,
91    /// Pending for [`NEVER_UP_MS`] with only failures.
92    #[serde(skip_serializing_if = "std::ops::Not::not")]
93    pub never_up: bool,
94    /// Its app or stack service is scaled to 0, since [`State::since`].
95    #[serde(skip_serializing_if = "std::ops::Not::not")]
96    pub stopped: bool,
97}
98
99/// What a check changed, and what to tell channels.
100#[derive(Debug, Clone, PartialEq)]
101pub struct Step {
102    /// The new status, when it changed.
103    pub changed: Option<Status>,
104    pub notify: Notify,
105    /// The check failed while the monitor was waiting for its first
106    /// success: kept as pending, not as a failure.
107    pub pending: bool,
108}
109
110#[derive(Debug, Clone, PartialEq)]
111pub enum Notify {
112    None,
113    /// Down since `since` (the first failed check of the incident).
114    Down {
115        since: u64,
116        flapping: bool,
117    },
118    /// Still no success `NEVER_UP_MS` after the first check (`since`).
119    NeverUp {
120        since: u64,
121    },
122    /// Up again; it was down from `down_since`.
123    Up {
124        down_since: u64,
125        downtime_ms: u64,
126    },
127}
128
129impl State {
130    /// Take one check's result at `now` (unix milliseconds).
131    pub fn observe(&mut self, ok: bool, now: u64, t: Thresholds) -> Step {
132        // Checked again: its target runs.
133        self.stopped = false;
134        if self.status == Status::Pending && !ok {
135            return self.observe_pending(now);
136        }
137        let changed = self.count(ok, now, t);
138        if let Some(s) = changed {
139            self.status = s;
140            self.since = now;
141            if s == Status::Down {
142                self.downs.push(now);
143            }
144        }
145        self.downs
146            .retain(|d| now.saturating_sub(*d) < FLAP_WINDOW_MS);
147        let was = self.flapping;
148        if changed == Some(Status::Down) && self.downs.len() >= FLAP_DOWNS {
149            self.flapping = true;
150        } else if self.flapping && now.saturating_sub(self.since) >= FLAP_WINDOW_MS {
151            self.flapping = false;
152        }
153        let notify = self.decide(now, was);
154        Step {
155            changed,
156            notify,
157            pending: false,
158        }
159    }
160
161    /// A failed check before the first success: counted, never a down.
162    fn observe_pending(&mut self, now: u64) -> Step {
163        let first = *self.first_check.get_or_insert(now);
164        self.fails = self.fails.saturating_add(1);
165        self.oks = 0;
166        self.failing_since.get_or_insert(now);
167        let mut notify = Notify::None;
168        if !self.never_up && now.saturating_sub(first) >= NEVER_UP_MS {
169            self.never_up = true;
170            self.notified = Status::Down;
171            self.down_since = Some(first);
172            notify = Notify::NeverUp { since: first };
173        }
174        Step {
175            changed: None,
176            notify,
177            pending: true,
178        }
179    }
180
181    /// Count the check; the status it now calls for, when that is a change.
182    fn count(&mut self, ok: bool, now: u64, t: Thresholds) -> Option<Status> {
183        if ok {
184            self.oks = self.oks.saturating_add(1);
185            self.fails = 0;
186            self.failing_since = None;
187            let enough = self.status == Status::Pending || self.oks >= t.recoveries.max(1);
188            if enough {
189                self.never_up = false;
190                self.first_check = None;
191            }
192            (self.status != Status::Up && enough).then_some(Status::Up)
193        } else {
194            self.fails = self.fails.saturating_add(1);
195            self.oks = 0;
196            self.failing_since.get_or_insert(now);
197            (self.status != Status::Down && self.fails >= t.failures.max(1)).then_some(Status::Down)
198        }
199    }
200
201    fn decide(&mut self, now: u64, was_flapping: bool) -> Notify {
202        // Held while flapping; the change that started it still goes out.
203        if self.flapping && was_flapping {
204            return Notify::None;
205        }
206        match (self.status, self.notified) {
207            (Status::Down, n) if n != Status::Down => {
208                let since = self.failing_since.unwrap_or(now).min(now);
209                self.notified = Status::Down;
210                self.down_since = Some(since);
211                Notify::Down {
212                    since,
213                    flapping: self.flapping,
214                }
215            }
216            (Status::Up, Status::Down) => {
217                let down_since = self.down_since.take().unwrap_or(now);
218                self.notified = Status::Up;
219                Notify::Up {
220                    down_since,
221                    downtime_ms: now.saturating_sub(down_since),
222                }
223            }
224            (Status::Up, Status::Pending) => {
225                // The first up is not news.
226                self.notified = Status::Up;
227                Notify::None
228            }
229            _ => Notify::None,
230        }
231    }
232
233    /// Its target was scaled to 0 at `now`: forget the run it was on, so
234    /// starting again waits for the first success like a new monitor.
235    /// Returns whether it was not stopped already.
236    pub fn stop(&mut self, now: u64) -> bool {
237        if self.stopped {
238            return false;
239        }
240        *self = State {
241            since: now,
242            stopped: true,
243            ..State::default()
244        };
245        true
246    }
247
248    /// Start counting afresh (a resumed or edited monitor), keeping what
249    /// channels were told.
250    pub fn reset_counts(&mut self) {
251        self.fails = 0;
252        self.oks = 0;
253        self.failing_since = None;
254        if self.status == Status::Pending {
255            self.first_check = None;
256        }
257    }
258}
259
260#[cfg(test)]
261mod tests {
262    use super::*;
263
264    const T: Thresholds = Thresholds {
265        failures: 2,
266        recoveries: 2,
267    };
268
269    fn run(s: &mut State, checks: &[(bool, u64)]) -> Vec<Notify> {
270        checks
271            .iter()
272            .map(|(ok, at)| s.observe(*ok, *at, T).notify)
273            .filter(|n| *n != Notify::None)
274            .collect()
275    }
276
277    #[test]
278    fn thresholds_and_one_notification_per_incident() {
279        let mut s = State::default();
280        // The first success makes a new monitor up, quietly.
281        assert_eq!(run(&mut s, &[(true, 1000)]), vec![]);
282        assert_eq!(s.status, Status::Up);
283        // One failure is not enough.
284        assert_eq!(run(&mut s, &[(false, 2000), (true, 3000)]), vec![]);
285        assert_eq!(s.status, Status::Up);
286        // Two are; the incident starts at the first.
287        let n = run(&mut s, &[(false, 4000), (false, 5000)]);
288        assert_eq!(
289            n,
290            vec![Notify::Down {
291                since: 4000,
292                flapping: false
293            }]
294        );
295        assert_eq!(s.status, Status::Down);
296        // More failures say nothing more.
297        assert_eq!(run(&mut s, &[(false, 6000), (false, 7000)]), vec![]);
298        // One success does not recover it (hysteresis); a failure resets it.
299        assert_eq!(
300            run(&mut s, &[(true, 8000), (false, 9000), (true, 10_000)]),
301            vec![]
302        );
303        assert_eq!(s.status, Status::Down);
304        // Two in a row do, with the downtime from the first failure.
305        let n = run(&mut s, &[(true, 11_000)]);
306        assert_eq!(
307            n,
308            vec![Notify::Up {
309                down_since: 4000,
310                downtime_ms: 7000
311            }]
312        );
313        assert_eq!(s.status, Status::Up);
314        assert_eq!(s.notified, Status::Up);
315    }
316
317    #[test]
318    fn a_new_monitor_waits_for_its_first_success() {
319        let mut s = State::default();
320        // Failures before the first success are pending, never down.
321        for at in [10, 20, 30, 40] {
322            let st = s.observe(false, at, T);
323            assert!(st.pending);
324            assert_eq!((st.changed, st.notify), (None, Notify::None));
325        }
326        assert_eq!(s.status, Status::Pending);
327        assert!(s.downs.is_empty());
328        // The first success is up, quietly, and clears the waiting.
329        let st = s.observe(true, 50, T);
330        assert!(!st.pending);
331        assert_eq!((st.changed, st.notify), (Some(Status::Up), Notify::None));
332        assert_eq!((s.first_check, s.never_up, s.fails), (None, false, 0));
333        assert_eq!(s.notified, Status::Up);
334    }
335
336    #[test]
337    fn pending_failures_do_not_count_toward_the_threshold() {
338        let mut s = State::default();
339        run(&mut s, &[(false, 10), (false, 20), (false, 30)]);
340        // Up on the first success; one failure then is still below 2.
341        assert_eq!(run(&mut s, &[(true, 40), (false, 50)]), vec![]);
342        assert_eq!(s.status, Status::Up);
343        assert_eq!(s.fails, 1);
344    }
345
346    #[test]
347    fn real_downtime_after_up_still_pages() {
348        let mut s = State::default();
349        run(&mut s, &[(false, 10), (true, 20)]);
350        let n = run(&mut s, &[(false, 30), (false, 40)]);
351        assert_eq!(
352            n,
353            vec![Notify::Down {
354                since: 30,
355                flapping: false
356            }]
357        );
358        assert_eq!(s.status, Status::Down);
359        assert!(!s.observe(false, 50, T).pending);
360    }
361
362    #[test]
363    fn a_monitor_that_never_comes_up_says_so_once() {
364        let mut s = State::default();
365        assert_eq!(run(&mut s, &[(false, 1000)]), vec![]);
366        // Just short of the window: nothing.
367        assert_eq!(run(&mut s, &[(false, 1000 + NEVER_UP_MS - 1)]), vec![]);
368        let n = run(
369            &mut s,
370            &[(false, 1000 + NEVER_UP_MS), (false, 2000 + NEVER_UP_MS)],
371        );
372        assert_eq!(n, vec![Notify::NeverUp { since: 1000 }]);
373        // Still pending, flagged; no second message.
374        assert_eq!((s.status, s.never_up), (Status::Pending, true));
375        assert_eq!(run(&mut s, &[(false, 3000 + NEVER_UP_MS)]), vec![]);
376        // When it finally answers, channels hear it is up.
377        let at = 4000 + NEVER_UP_MS;
378        let n = run(&mut s, &[(true, at)]);
379        assert_eq!(
380            n,
381            vec![Notify::Up {
382                down_since: 1000,
383                downtime_ms: at - 1000
384            }]
385        );
386        assert_eq!((s.status, s.never_up), (Status::Up, false));
387    }
388
389    #[test]
390    fn thresholds_of_one() {
391        let one = Thresholds {
392            failures: 1,
393            recoveries: 1,
394        };
395        let mut s = State::default();
396        assert_eq!(s.observe(true, 1, one).notify, Notify::None);
397        assert!(matches!(
398            s.observe(false, 2, one).notify,
399            Notify::Down { .. }
400        ));
401        assert!(matches!(s.observe(true, 3, one).notify, Notify::Up { .. }));
402        // Zero is treated as one.
403        let zero = Thresholds {
404            failures: 0,
405            recoveries: 0,
406        };
407        assert!(matches!(
408            s.observe(false, 4, zero).notify,
409            Notify::Down { .. }
410        ));
411    }
412
413    #[test]
414    fn flapping_is_damped_and_settles() {
415        let one = Thresholds {
416            failures: 1,
417            recoveries: 1,
418        };
419        let mut s = State::default();
420        let mut sent = Vec::new();
421        let mut at = 1_000;
422        s.observe(true, at, one);
423        // Down and up every minute, ten times.
424        for _ in 0..10 {
425            for ok in [false, true] {
426                at += 60_000;
427                let n = s.observe(ok, at, one).notify;
428                if n != Notify::None {
429                    sent.push(n);
430                }
431            }
432        }
433        // down, up, down, up, then the third down (flapping), then quiet.
434        assert_eq!(sent.len(), 5, "{sent:?}");
435        assert!(matches!(sent[4], Notify::Down { flapping: true, .. }));
436        assert!(s.flapping);
437        // It stays up: nothing until it has held for the window, then the
438        // up channels have not heard yet.
439        let mut later = Vec::new();
440        for _ in 0..40 {
441            at += 60_000;
442            let n = s.observe(true, at, one).notify;
443            if n != Notify::None {
444                later.push(n);
445            }
446        }
447        assert_eq!(later.len(), 1, "{later:?}");
448        assert!(matches!(later[0], Notify::Up { .. }));
449        assert!(!s.flapping);
450        assert_eq!(s.notified, Status::Up);
451    }
452
453    #[test]
454    fn flapping_that_settles_down_is_not_told_twice() {
455        let one = Thresholds {
456            failures: 1,
457            recoveries: 1,
458        };
459        let mut s = State::default();
460        let mut at = 0;
461        let mut sent = Vec::new();
462        for ok in [true, false, true, false, true, false, true, false] {
463            at += 1000;
464            let n = s.observe(ok, at, one).notify;
465            if n != Notify::None {
466                sent.push(n);
467            }
468        }
469        // The last flap left it down, and channels heard the flapping down.
470        assert_eq!(s.notified, Status::Down);
471        let n_before = sent.len();
472        for _ in 0..40 {
473            at += 60_000;
474            let n = s.observe(false, at, one).notify;
475            if n != Notify::None {
476                sent.push(n);
477            }
478        }
479        assert_eq!(sent.len(), n_before, "{sent:?}");
480        assert!(!s.flapping);
481    }
482
483    #[test]
484    fn stopping_is_quiet_and_starting_begins_anew() {
485        let mut s = State::default();
486        run(&mut s, &[(true, 1), (false, 2), (false, 3)]);
487        assert_eq!(s.notified, Status::Down);
488        assert!(s.stop(4));
489        assert!(!s.stop(5), "already stopped");
490        assert_eq!((s.status, s.since, s.stopped), (Status::Pending, 4, true));
491        // Failures while it comes back up are pending, never a down.
492        assert_eq!(
493            run(&mut s, &[(false, 10), (false, 11), (false, 12)]),
494            vec![]
495        );
496        assert!(!s.stopped);
497        assert_eq!(s.status, Status::Pending);
498        // Its first success makes it up, and no up goes out: the down
499        // channels heard belonged to the run before the stop.
500        assert_eq!(run(&mut s, &[(true, 13)]), vec![]);
501        assert_eq!(s.status, Status::Up);
502    }
503
504    #[test]
505    fn state_survives_a_restart() {
506        let mut s = State::default();
507        run(&mut s, &[(true, 1), (false, 2), (false, 3)]);
508        assert_eq!(s.status, Status::Down);
509        let saved = serde_json::to_string(&s).unwrap();
510        let mut back: State = serde_json::from_str(&saved).unwrap();
511        assert_eq!(back, s);
512        // Still down after the restart: no second down.
513        assert_eq!(run(&mut back, &[(false, 4), (false, 5)]), vec![]);
514        let n = run(&mut back, &[(true, 6), (true, 7)]);
515        assert_eq!(
516            n,
517            vec![Notify::Up {
518                down_since: 2,
519                downtime_ms: 5
520            }]
521        );
522        // An old state file without the newer fields still loads.
523        let old: State = serde_json::from_str(r#"{"status":"up"}"#).unwrap();
524        assert_eq!(old.status, Status::Up);
525    }
526}