use std::{ collections::{HashMap, HashSet}, fmt, net::IpAddr, }; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Status { Up, Down, } impl fmt::Display for Status { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.write_str(match self { Status::Up => "UP", Status::Down => "DOWN", }) } } #[derive(Debug, PartialEq, Eq)] pub enum Change { None, /// First time this IP was seen (fresh start, or a member's IP just changed). New, Flipped, } #[derive(Debug)] struct Entry { status: Status, ok_streak: u32, fail_streak: u32, } /// Debounces probe results per IP so one lost check doesn't pull a node out of DNS. pub struct Tracker { fail_threshold: u32, rise_threshold: u32, entries: HashMap, } impl Tracker { pub fn new(fail_threshold: u32, rise_threshold: u32) -> Self { Self { fail_threshold, rise_threshold, entries: HashMap::new(), } } pub fn observe(&mut self, ip: IpAddr, ok: bool) -> (Status, Change) { let Some(e) = self.entries.get_mut(&ip) else { // A new IP is judged on its first probe alone, so a node that just changed IP is // picked up within one interval instead of waiting out rise_threshold. let status = if ok { Status::Up } else { Status::Down }; self.entries.insert( ip, Entry { status, ok_streak: ok as u32, fail_streak: !ok as u32, }, ); return (status, Change::New); }; let before = e.status; if ok { e.ok_streak += 1; e.fail_streak = 0; if e.status == Status::Down && e.ok_streak >= self.rise_threshold { e.status = Status::Up; } } else { e.fail_streak += 1; e.ok_streak = 0; if e.status == Status::Up && e.fail_streak >= self.fail_threshold { e.status = Status::Down; } } let change = if e.status == before { Change::None } else { Change::Flipped }; (e.status, change) } /// Forget IPs that no longer belong to any member. pub fn retain(&mut self, keep: &HashSet) { self.entries.retain(|ip, _| keep.contains(ip)); } } #[cfg(test)] mod tests { use super::*; use std::net::Ipv4Addr; const IP: IpAddr = IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)); #[test] fn new_ip_is_judged_on_first_probe() { let mut t = Tracker::new(3, 2); assert_eq!(t.observe(IP, true), (Status::Up, Change::New)); let other = IpAddr::V4(Ipv4Addr::new(192, 0, 2, 2)); assert_eq!(t.observe(other, false), (Status::Down, Change::New)); } #[test] fn goes_down_only_after_fail_threshold_consecutive_failures() { let mut t = Tracker::new(3, 2); t.observe(IP, true); assert_eq!(t.observe(IP, false), (Status::Up, Change::None)); assert_eq!(t.observe(IP, false), (Status::Up, Change::None)); // A success resets the streak. assert_eq!(t.observe(IP, true), (Status::Up, Change::None)); assert_eq!(t.observe(IP, false), (Status::Up, Change::None)); assert_eq!(t.observe(IP, false), (Status::Up, Change::None)); assert_eq!(t.observe(IP, false), (Status::Down, Change::Flipped)); assert_eq!(t.observe(IP, false), (Status::Down, Change::None)); } #[test] fn comes_back_only_after_rise_threshold_consecutive_successes() { let mut t = Tracker::new(1, 2); t.observe(IP, true); assert_eq!(t.observe(IP, false), (Status::Down, Change::Flipped)); assert_eq!(t.observe(IP, true), (Status::Down, Change::None)); assert_eq!(t.observe(IP, false), (Status::Down, Change::None)); assert_eq!(t.observe(IP, true), (Status::Down, Change::None)); assert_eq!(t.observe(IP, true), (Status::Up, Change::Flipped)); } #[test] fn retain_forgets_old_ips() { let mut t = Tracker::new(3, 2); t.observe(IP, false); t.retain(&HashSet::new()); // Seen again after being forgotten: treated as new, so one success brings it up. assert_eq!(t.observe(IP, true), (Status::Up, Change::New)); } }