1use serde::{Deserialize, Serialize};
25
26pub const FLAP_WINDOW_MS: u64 = 30 * 60 * 1000;
29pub const FLAP_DOWNS: usize = 3;
31pub 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 #[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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub struct Thresholds {
57 pub failures: u32,
58 pub recoveries: u32,
59}
60
61#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
64#[serde(default)]
65pub struct State {
66 pub status: Status,
67 pub since: u64,
69 pub fails: u32,
71 pub oks: u32,
73 #[serde(skip_serializing_if = "Option::is_none")]
75 pub failing_since: Option<u64>,
76 pub notified: Status,
78 #[serde(skip_serializing_if = "Option::is_none")]
80 pub down_since: Option<u64>,
81 #[serde(skip_serializing_if = "Vec::is_empty")]
83 pub downs: Vec<u64>,
84 pub flapping: bool,
85 #[serde(skip_serializing_if = "Option::is_none")]
87 pub cert_warned: Option<u64>,
88 #[serde(skip_serializing_if = "Option::is_none")]
90 pub first_check: Option<u64>,
91 #[serde(skip_serializing_if = "std::ops::Not::not")]
93 pub never_up: bool,
94 #[serde(skip_serializing_if = "std::ops::Not::not")]
96 pub stopped: bool,
97}
98
99#[derive(Debug, Clone, PartialEq)]
101pub struct Step {
102 pub changed: Option<Status>,
104 pub notify: Notify,
105 pub pending: bool,
108}
109
110#[derive(Debug, Clone, PartialEq)]
111pub enum Notify {
112 None,
113 Down {
115 since: u64,
116 flapping: bool,
117 },
118 NeverUp {
120 since: u64,
121 },
122 Up {
124 down_since: u64,
125 downtime_ms: u64,
126 },
127}
128
129impl State {
130 pub fn observe(&mut self, ok: bool, now: u64, t: Thresholds) -> Step {
132 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 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 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 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 self.notified = Status::Up;
227 Notify::None
228 }
229 _ => Notify::None,
230 }
231 }
232
233 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 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 assert_eq!(run(&mut s, &[(true, 1000)]), vec![]);
282 assert_eq!(s.status, Status::Up);
283 assert_eq!(run(&mut s, &[(false, 2000), (true, 3000)]), vec![]);
285 assert_eq!(s.status, Status::Up);
286 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 assert_eq!(run(&mut s, &[(false, 6000), (false, 7000)]), vec![]);
298 assert_eq!(
300 run(&mut s, &[(true, 8000), (false, 9000), (true, 10_000)]),
301 vec![]
302 );
303 assert_eq!(s.status, Status::Down);
304 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 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 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 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 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 assert_eq!((s.status, s.never_up), (Status::Pending, true));
375 assert_eq!(run(&mut s, &[(false, 3000 + NEVER_UP_MS)]), vec![]);
376 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 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 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 assert_eq!(sent.len(), 5, "{sent:?}");
435 assert!(matches!(sent[4], Notify::Down { flapping: true, .. }));
436 assert!(s.flapping);
437 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 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 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 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 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 let old: State = serde_json::from_str(r#"{"status":"up"}"#).unwrap();
524 assert_eq!(old.status, Status::Up);
525 }
526}