use serde::{Deserialize, Serialize};
pub const FLAP_WINDOW_MS: u64 = 30 * 60 * 1000;
pub const FLAP_DOWNS: usize = 3;
pub const NEVER_UP_MS: u64 = 30 * 60 * 1000;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Status {
#[default]
Pending,
Up,
Down,
}
impl Status {
pub fn as_str(self) -> &'static str {
match self {
Status::Pending => "pending",
Status::Up => "up",
Status::Down => "down",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Thresholds {
pub failures: u32,
pub recoveries: u32,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(default)]
pub struct State {
pub status: Status,
pub since: u64,
pub fails: u32,
pub oks: u32,
#[serde(skip_serializing_if = "Option::is_none")]
pub failing_since: Option<u64>,
pub notified: Status,
#[serde(skip_serializing_if = "Option::is_none")]
pub down_since: Option<u64>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub downs: Vec<u64>,
pub flapping: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub cert_warned: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub first_check: Option<u64>,
#[serde(skip_serializing_if = "std::ops::Not::not")]
pub never_up: bool,
}
#[derive(Debug, Clone, PartialEq)]
pub struct Step {
pub changed: Option<Status>,
pub notify: Notify,
pub pending: bool,
}
#[derive(Debug, Clone, PartialEq)]
pub enum Notify {
None,
Down {
since: u64,
flapping: bool,
},
NeverUp {
since: u64,
},
Up {
down_since: u64,
downtime_ms: u64,
},
}
impl State {
pub fn observe(&mut self, ok: bool, now: u64, t: Thresholds) -> Step {
if self.status == Status::Pending && !ok {
return self.observe_pending(now);
}
let changed = self.count(ok, now, t);
if let Some(s) = changed {
self.status = s;
self.since = now;
if s == Status::Down {
self.downs.push(now);
}
}
self.downs
.retain(|d| now.saturating_sub(*d) < FLAP_WINDOW_MS);
let was = self.flapping;
if changed == Some(Status::Down) && self.downs.len() >= FLAP_DOWNS {
self.flapping = true;
} else if self.flapping && now.saturating_sub(self.since) >= FLAP_WINDOW_MS {
self.flapping = false;
}
let notify = self.decide(now, was);
Step {
changed,
notify,
pending: false,
}
}
fn observe_pending(&mut self, now: u64) -> Step {
let first = *self.first_check.get_or_insert(now);
self.fails = self.fails.saturating_add(1);
self.oks = 0;
self.failing_since.get_or_insert(now);
let mut notify = Notify::None;
if !self.never_up && now.saturating_sub(first) >= NEVER_UP_MS {
self.never_up = true;
self.notified = Status::Down;
self.down_since = Some(first);
notify = Notify::NeverUp { since: first };
}
Step {
changed: None,
notify,
pending: true,
}
}
fn count(&mut self, ok: bool, now: u64, t: Thresholds) -> Option<Status> {
if ok {
self.oks = self.oks.saturating_add(1);
self.fails = 0;
self.failing_since = None;
let enough = self.status == Status::Pending || self.oks >= t.recoveries.max(1);
if enough {
self.never_up = false;
self.first_check = None;
}
(self.status != Status::Up && enough).then_some(Status::Up)
} else {
self.fails = self.fails.saturating_add(1);
self.oks = 0;
self.failing_since.get_or_insert(now);
(self.status != Status::Down && self.fails >= t.failures.max(1)).then_some(Status::Down)
}
}
fn decide(&mut self, now: u64, was_flapping: bool) -> Notify {
if self.flapping && was_flapping {
return Notify::None;
}
match (self.status, self.notified) {
(Status::Down, n) if n != Status::Down => {
let since = self.failing_since.unwrap_or(now).min(now);
self.notified = Status::Down;
self.down_since = Some(since);
Notify::Down {
since,
flapping: self.flapping,
}
}
(Status::Up, Status::Down) => {
let down_since = self.down_since.take().unwrap_or(now);
self.notified = Status::Up;
Notify::Up {
down_since,
downtime_ms: now.saturating_sub(down_since),
}
}
(Status::Up, Status::Pending) => {
self.notified = Status::Up;
Notify::None
}
_ => Notify::None,
}
}
pub fn reset_counts(&mut self) {
self.fails = 0;
self.oks = 0;
self.failing_since = None;
if self.status == Status::Pending {
self.first_check = None;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
const T: Thresholds = Thresholds {
failures: 2,
recoveries: 2,
};
fn run(s: &mut State, checks: &[(bool, u64)]) -> Vec<Notify> {
checks
.iter()
.map(|(ok, at)| s.observe(*ok, *at, T).notify)
.filter(|n| *n != Notify::None)
.collect()
}
#[test]
fn thresholds_and_one_notification_per_incident() {
let mut s = State::default();
assert_eq!(run(&mut s, &[(true, 1000)]), vec![]);
assert_eq!(s.status, Status::Up);
assert_eq!(run(&mut s, &[(false, 2000), (true, 3000)]), vec![]);
assert_eq!(s.status, Status::Up);
let n = run(&mut s, &[(false, 4000), (false, 5000)]);
assert_eq!(
n,
vec![Notify::Down {
since: 4000,
flapping: false
}]
);
assert_eq!(s.status, Status::Down);
assert_eq!(run(&mut s, &[(false, 6000), (false, 7000)]), vec![]);
assert_eq!(
run(&mut s, &[(true, 8000), (false, 9000), (true, 10_000)]),
vec![]
);
assert_eq!(s.status, Status::Down);
let n = run(&mut s, &[(true, 11_000)]);
assert_eq!(
n,
vec![Notify::Up {
down_since: 4000,
downtime_ms: 7000
}]
);
assert_eq!(s.status, Status::Up);
assert_eq!(s.notified, Status::Up);
}
#[test]
fn a_new_monitor_waits_for_its_first_success() {
let mut s = State::default();
for at in [10, 20, 30, 40] {
let st = s.observe(false, at, T);
assert!(st.pending);
assert_eq!((st.changed, st.notify), (None, Notify::None));
}
assert_eq!(s.status, Status::Pending);
assert!(s.downs.is_empty());
let st = s.observe(true, 50, T);
assert!(!st.pending);
assert_eq!((st.changed, st.notify), (Some(Status::Up), Notify::None));
assert_eq!((s.first_check, s.never_up, s.fails), (None, false, 0));
assert_eq!(s.notified, Status::Up);
}
#[test]
fn pending_failures_do_not_count_toward_the_threshold() {
let mut s = State::default();
run(&mut s, &[(false, 10), (false, 20), (false, 30)]);
assert_eq!(run(&mut s, &[(true, 40), (false, 50)]), vec![]);
assert_eq!(s.status, Status::Up);
assert_eq!(s.fails, 1);
}
#[test]
fn real_downtime_after_up_still_pages() {
let mut s = State::default();
run(&mut s, &[(false, 10), (true, 20)]);
let n = run(&mut s, &[(false, 30), (false, 40)]);
assert_eq!(
n,
vec![Notify::Down {
since: 30,
flapping: false
}]
);
assert_eq!(s.status, Status::Down);
assert!(!s.observe(false, 50, T).pending);
}
#[test]
fn a_monitor_that_never_comes_up_says_so_once() {
let mut s = State::default();
assert_eq!(run(&mut s, &[(false, 1000)]), vec![]);
assert_eq!(run(&mut s, &[(false, 1000 + NEVER_UP_MS - 1)]), vec![]);
let n = run(
&mut s,
&[(false, 1000 + NEVER_UP_MS), (false, 2000 + NEVER_UP_MS)],
);
assert_eq!(n, vec![Notify::NeverUp { since: 1000 }]);
assert_eq!((s.status, s.never_up), (Status::Pending, true));
assert_eq!(run(&mut s, &[(false, 3000 + NEVER_UP_MS)]), vec![]);
let at = 4000 + NEVER_UP_MS;
let n = run(&mut s, &[(true, at)]);
assert_eq!(
n,
vec![Notify::Up {
down_since: 1000,
downtime_ms: at - 1000
}]
);
assert_eq!((s.status, s.never_up), (Status::Up, false));
}
#[test]
fn thresholds_of_one() {
let one = Thresholds {
failures: 1,
recoveries: 1,
};
let mut s = State::default();
assert_eq!(s.observe(true, 1, one).notify, Notify::None);
assert!(matches!(
s.observe(false, 2, one).notify,
Notify::Down { .. }
));
assert!(matches!(s.observe(true, 3, one).notify, Notify::Up { .. }));
let zero = Thresholds {
failures: 0,
recoveries: 0,
};
assert!(matches!(
s.observe(false, 4, zero).notify,
Notify::Down { .. }
));
}
#[test]
fn flapping_is_damped_and_settles() {
let one = Thresholds {
failures: 1,
recoveries: 1,
};
let mut s = State::default();
let mut sent = Vec::new();
let mut at = 1_000;
s.observe(true, at, one);
for _ in 0..10 {
for ok in [false, true] {
at += 60_000;
let n = s.observe(ok, at, one).notify;
if n != Notify::None {
sent.push(n);
}
}
}
assert_eq!(sent.len(), 5, "{sent:?}");
assert!(matches!(sent[4], Notify::Down { flapping: true, .. }));
assert!(s.flapping);
let mut later = Vec::new();
for _ in 0..40 {
at += 60_000;
let n = s.observe(true, at, one).notify;
if n != Notify::None {
later.push(n);
}
}
assert_eq!(later.len(), 1, "{later:?}");
assert!(matches!(later[0], Notify::Up { .. }));
assert!(!s.flapping);
assert_eq!(s.notified, Status::Up);
}
#[test]
fn flapping_that_settles_down_is_not_told_twice() {
let one = Thresholds {
failures: 1,
recoveries: 1,
};
let mut s = State::default();
let mut at = 0;
let mut sent = Vec::new();
for ok in [true, false, true, false, true, false, true, false] {
at += 1000;
let n = s.observe(ok, at, one).notify;
if n != Notify::None {
sent.push(n);
}
}
assert_eq!(s.notified, Status::Down);
let n_before = sent.len();
for _ in 0..40 {
at += 60_000;
let n = s.observe(false, at, one).notify;
if n != Notify::None {
sent.push(n);
}
}
assert_eq!(sent.len(), n_before, "{sent:?}");
assert!(!s.flapping);
}
#[test]
fn state_survives_a_restart() {
let mut s = State::default();
run(&mut s, &[(true, 1), (false, 2), (false, 3)]);
assert_eq!(s.status, Status::Down);
let saved = serde_json::to_string(&s).unwrap();
let mut back: State = serde_json::from_str(&saved).unwrap();
assert_eq!(back, s);
assert_eq!(run(&mut back, &[(false, 4), (false, 5)]), vec![]);
let n = run(&mut back, &[(true, 6), (true, 7)]);
assert_eq!(
n,
vec![Notify::Up {
down_since: 2,
downtime_ms: 5
}]
);
let old: State = serde_json::from_str(r#"{"status":"up"}"#).unwrap();
assert_eq!(old.status, Status::Up);
}
}