1use serde::{Deserialize, Serialize};
21
22pub const FLAP_WINDOW_MS: u64 = 30 * 60 * 1000;
25pub const FLAP_DOWNS: usize = 3;
27pub 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 #[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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
52pub struct Thresholds {
53 pub failures: u32,
54 pub recoveries: u32,
55}
56
57#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
60#[serde(default)]
61pub struct State {
62 pub status: Status,
63 pub since: u64,
65 pub fails: u32,
67 pub oks: u32,
69 #[serde(skip_serializing_if = "Option::is_none")]
71 pub failing_since: Option<u64>,
72 pub notified: Status,
74 #[serde(skip_serializing_if = "Option::is_none")]
76 pub down_since: Option<u64>,
77 #[serde(skip_serializing_if = "Vec::is_empty")]
79 pub downs: Vec<u64>,
80 pub flapping: bool,
81 #[serde(skip_serializing_if = "Option::is_none")]
83 pub cert_warned: Option<u64>,
84 #[serde(skip_serializing_if = "Option::is_none")]
86 pub first_check: Option<u64>,
87 #[serde(skip_serializing_if = "std::ops::Not::not")]
89 pub never_up: bool,
90}
91
92#[derive(Debug, Clone, PartialEq)]
94pub struct Step {
95 pub changed: Option<Status>,
97 pub notify: Notify,
98 pub pending: bool,
101}
102
103#[derive(Debug, Clone, PartialEq)]
104pub enum Notify {
105 None,
106 Down {
108 since: u64,
109 flapping: bool,
110 },
111 NeverUp {
113 since: u64,
114 },
115 Up {
117 down_since: u64,
118 downtime_ms: u64,
119 },
120}
121
122impl State {
123 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 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 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 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 self.notified = Status::Up;
218 Notify::None
219 }
220 _ => Notify::None,
221 }
222 }
223
224 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 assert_eq!(run(&mut s, &[(true, 1000)]), vec![]);
258 assert_eq!(s.status, Status::Up);
259 assert_eq!(run(&mut s, &[(false, 2000), (true, 3000)]), vec![]);
261 assert_eq!(s.status, Status::Up);
262 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 assert_eq!(run(&mut s, &[(false, 6000), (false, 7000)]), vec![]);
274 assert_eq!(
276 run(&mut s, &[(true, 8000), (false, 9000), (true, 10_000)]),
277 vec![]
278 );
279 assert_eq!(s.status, Status::Down);
280 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 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 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 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 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 assert_eq!((s.status, s.never_up), (Status::Pending, true));
351 assert_eq!(run(&mut s, &[(false, 3000 + NEVER_UP_MS)]), vec![]);
352 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 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 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 assert_eq!(sent.len(), 5, "{sent:?}");
411 assert!(matches!(sent[4], Notify::Down { flapping: true, .. }));
412 assert!(s.flapping);
413 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 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 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 let old: State = serde_json::from_str(r#"{"status":"up"}"#).unwrap();
479 assert_eq!(old.status, Status::Up);
480 }
481}