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