//! The core: keeps one group's A records pointed at its healthy members. It talks to the //! outside world only through the `DnsProvider` and `Probe` traits, and describes what //! happened as a `RoundReport` without knowing who reads it. mod health; mod members; mod reconcile; use std::{ collections::{BTreeSet, HashSet}, net::IpAddr, sync::Arc, }; use anyhow::{Context, Result, bail}; use futures::future::join_all; use crate::{ config::{Group, ProbeConfig}, dns::DnsProvider, event::{Event, IpStatus, RoundReport}, probe::Probe, }; use health::{Change, Status, Tracker}; use members::Members; pub struct GroupMonitor { cfg: Group, dns: D, probe: Arc

, dry_run: bool, tracker: Tracker, /// Last round's member records; `None` until the first round has run. prev_members: Option, all_down: bool, failing: bool, } impl GroupMonitor { pub fn new(cfg: Group, probe_cfg: &ProbeConfig, dns: D, probe: Arc

, dry_run: bool) -> Self { Self { cfg, dns, probe, dry_run, tracker: Tracker::new(probe_cfg.fail_threshold, probe_cfg.rise_threshold), prev_members: None, all_down: false, failing: false, } } pub async fn run_round(&mut self) -> RoundReport { let mut report = RoundReport { group: self.cfg.name.clone(), record_type: self.cfg.record_type.as_str(), dry_run: self.dry_run, ips: Vec::new(), events: Vec::new(), dns_now: None, error: None, }; match self.check(&mut report).await { Ok(()) => { if self.failing { self.failing = false; report.events.push(Event::CheckRecovered); } } Err(e) => { let error = format!("{e:#}"); if !self.failing { self.failing = true; report.events.push(Event::CheckFailed { error: error.clone() }); } report.error = Some(error); } } report } /// Any read error aborts the round before DNS is touched. async fn check(&mut self, report: &mut RoundReport) -> Result<()> { let members = self.resolve().await?; match &self.prev_members { None => report.events.push(Event::Started { missing: members.iter().filter(|(_, ips)| ips.is_empty()).map(|(m, _)| m.clone()).collect(), }), Some(prev) => report.events.extend( members::diff(prev, &members) .into_iter() .map(|(member, old, new)| Event::MemberChanged { member, old, new }), ), } let first_round = self.prev_members.is_none(); let by_ip = members::by_ip(&members); self.prev_members = Some(members); let outcomes = join_all(by_ip.keys().map(|ip| self.probe.probe(*ip, self.cfg.tcp_port))).await; self.tracker.retain(&by_ip.keys().copied().collect::>()); let mut desired = BTreeSet::new(); for ((ip, names), outcome) in by_ip.into_iter().zip(outcomes) { let (status, change) = self.tracker.observe(ip, outcome.is_ok()); let up = status == Status::Up; // A new IP that is up needs no event of its own: MemberChanged and DnsUpdated cover it. if change == Change::Flipped || (change == Change::New && !up && !first_round) { report.events.push(Event::HealthChanged { ip, members: names.clone(), up, probe: outcome.clone(), }); } if up { desired.insert(ip); } report.ips.push(IpStatus { ip, members: names, up, probe: outcome }); } let name = &self.cfg.name; let kind = self.cfg.record_type.as_str(); let records = self.dns.records(name).await?; if records.iter().any(|r| r.kind == "CNAME") { bail!("{name} is a CNAME; delete it by hand before dns-monitor can manage {kind} records"); } let current = records .iter() .filter(|r| r.kind == kind) .map(|r| Ok((r.id.clone(), r.content.parse().with_context(|| format!("{kind} record of {name}: {}", r.content))?))) .collect::>>()?; let current_ips = || current.iter().map(|(_, ip)| *ip).collect::>().into_iter().collect(); let Some(plan) = reconcile::plan(&desired, ¤t) else { if !self.all_down { self.all_down = true; report.events.push(Event::NoHealthyMembers); } report.dns_now = Some(current_ips()); return Ok(()); }; self.all_down = false; let (mut added, mut removed) = (Vec::new(), Vec::new()); let applied = self.apply(&plan, &mut added, &mut removed).await; // Report whatever was written, even if a later write failed. if !added.is_empty() || !removed.is_empty() { report.events.push(Event::DnsUpdated { added, removed }); } applied?; report.dns_now = Some(if self.dry_run { current_ips() } else { desired.into_iter().collect() }); Ok(()) } /// Member name -> IPs. Literal IP members stand for themselves. async fn resolve(&self) -> Result { let kind = self.cfg.record_type.as_str(); let mut members = Members::new(); for m in &self.cfg.members { let ips = match m.parse::() { Ok(ip) => BTreeSet::from([ip]), Err(_) => self .dns .records(m) .await? .into_iter() .filter(|r| r.kind == kind) .map(|r| r.content.parse().with_context(|| format!("{kind} record of {m}: {}", r.content))) .collect::>()?, }; members.insert(m.clone(), ips); } Ok(members) } /// Add before remove, so the name never resolves to nothing mid-update. async fn apply( &self, plan: &reconcile::Plan, added: &mut Vec, removed: &mut Vec, ) -> Result<()> { for ip in &plan.add { if !self.dry_run { self.dns.create(&self.cfg.name, *ip, self.cfg.ttl).await?; } added.push(*ip); } for (id, ip) in &plan.remove { if !self.dry_run { self.dns.delete(id).await?; } removed.push(*ip); } Ok(()) } } #[cfg(test)] mod tests { use std::{ net::{Ipv4Addr, Ipv6Addr}, sync::{ Mutex, atomic::{AtomicBool, AtomicUsize, Ordering}, }, time::Duration, }; use super::*; use crate::{config::RecordType, dns::DnsRecord, event::ProbeOutcome}; const NAME: &str = "g.example.com"; const A: &str = "a.example.com"; const B: &str = "b.example.com"; const PLACEHOLDER: IpAddr = IpAddr::V4(Ipv4Addr::new(1, 1, 1, 1)); fn ip(last: u8) -> IpAddr { IpAddr::V4(Ipv4Addr::new(192, 0, 2, last)) } #[derive(Default)] struct FakeDns { records: Mutex>, next_id: AtomicUsize, writes: AtomicUsize, failing: AtomicBool, } impl FakeDns { fn set(&self, name: &str, ips: &[IpAddr]) { let mut records = self.records.lock().unwrap(); records.retain(|(n, _)| n != name); for ip in ips { records.push((name.to_string(), self.record(*ip))); } } fn ips(&self, name: &str) -> BTreeSet { let records = self.records.lock().unwrap(); records.iter().filter(|(n, _)| n == name).map(|(_, r)| r.content.parse().unwrap()).collect() } fn record(&self, ip: IpAddr) -> DnsRecord { DnsRecord { id: self.next_id.fetch_add(1, Ordering::Relaxed).to_string(), kind: if ip.is_ipv4() { "A" } else { "AAAA" }.to_string(), content: ip.to_string(), } } } impl DnsProvider for Arc { async fn records(&self, name: &str) -> Result> { if self.failing.load(Ordering::Relaxed) { bail!("API down"); } let records = self.records.lock().unwrap(); Ok(records.iter().filter(|(n, _)| n == name).map(|(_, r)| r.clone()).collect()) } async fn create(&self, name: &str, ip: IpAddr, _ttl: u32) -> Result<()> { self.writes.fetch_add(1, Ordering::Relaxed); let record = self.record(ip); self.records.lock().unwrap().push((name.to_string(), record)); Ok(()) } async fn delete(&self, record_id: &str) -> Result<()> { self.writes.fetch_add(1, Ordering::Relaxed); self.records.lock().unwrap().retain(|(_, r)| r.id != record_id); Ok(()) } } #[derive(Default)] struct FakeProbe { down: Mutex>, } impl FakeProbe { fn set_down(&self, ip: IpAddr, down: bool) { let mut set = self.down.lock().unwrap(); if down { set.insert(ip); } else { set.remove(&ip); } } } impl Probe for FakeProbe { async fn probe(&self, ip: IpAddr, _tcp_port: Option) -> ProbeOutcome { if self.down.lock().unwrap().contains(&ip) { Err("timeout".to_string()) } else { Ok(Duration::from_millis(5)) } } } /// Group g = {a -> .1, b -> .2}; g itself starts out as a placeholder record. fn setup(dry_run: bool) -> (GroupMonitor, FakeProbe>, Arc, Arc) { let dns = Arc::new(FakeDns::default()); dns.set(A, &[ip(1)]); dns.set(B, &[ip(2)]); dns.set(NAME, &[PLACEHOLDER]); let probe = Arc::new(FakeProbe::default()); let cfg = Group { zone: "example.com".to_string(), name: NAME.to_string(), record_type: RecordType::A, members: vec![A.to_string(), B.to_string()], ttl: 60, tcp_port: None, }; let monitor = GroupMonitor::new(cfg, &ProbeConfig::default(), dns.clone(), probe.clone(), dry_run); (monitor, dns, probe) } #[tokio::test] async fn first_round_starts_and_points_dns_at_members() { let (mut m, dns, _) = setup(false); let r = m.run_round().await; assert_eq!( r.events, vec![ Event::Started { missing: vec![] }, Event::DnsUpdated { added: vec![ip(1), ip(2)], removed: vec![PLACEHOLDER] }, ] ); assert_eq!(r.ips.len(), 2); assert_eq!(r.dns_now, Some(vec![ip(1), ip(2)])); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(1), ip(2)])); } #[tokio::test] async fn steady_state_is_silent_and_writes_nothing() { let (mut m, dns, _) = setup(false); m.run_round().await; let writes = dns.writes.load(Ordering::Relaxed); let r = m.run_round().await; assert_eq!(r.events, vec![]); assert_eq!(dns.writes.load(Ordering::Relaxed), writes); } #[tokio::test] async fn dead_node_is_dropped_after_threshold_and_readded_after_recovery() { let (mut m, dns, probe) = setup(false); m.run_round().await; probe.set_down(ip(2), true); assert_eq!(m.run_round().await.events, vec![]); assert_eq!(m.run_round().await.events, vec![]); let r = m.run_round().await; assert_eq!( r.events, vec![ Event::HealthChanged { ip: ip(2), members: vec![B.to_string()], up: false, probe: Err("timeout".to_string()), }, Event::DnsUpdated { added: vec![], removed: vec![ip(2)] }, ] ); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(1)])); probe.set_down(ip(2), false); assert_eq!(m.run_round().await.events, vec![]); let r = m.run_round().await; assert!(matches!( &r.events[..], [Event::HealthChanged { up: true, .. }, Event::DnsUpdated { added, .. }] if *added == vec![ip(2)] )); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(1), ip(2)])); } #[tokio::test] async fn member_pointed_at_another_node_is_reported() { let (mut m, dns, _) = setup(false); m.run_round().await; // What the cross-node cron jobs do: a's name temporarily gets b's IP. dns.set(A, &[ip(2)]); let r = m.run_round().await; assert_eq!( r.events, vec![ Event::MemberChanged { member: A.to_string(), old: vec![ip(1)], new: vec![ip(2)] }, Event::DnsUpdated { added: vec![], removed: vec![ip(1)] }, ] ); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(2)])); } #[tokio::test] async fn new_ip_that_is_down_is_reported_and_kept_out() { let (mut m, dns, probe) = setup(false); m.run_round().await; probe.set_down(ip(9), true); dns.set(A, &[ip(9)]); let r = m.run_round().await; assert_eq!(r.events[0], Event::MemberChanged { member: A.to_string(), old: vec![ip(1)], new: vec![ip(9)] }); assert!(matches!(r.events[1], Event::HealthChanged { up: false, .. })); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(2)])); } #[tokio::test] async fn all_down_leaves_dns_alone_and_alerts_once() { let (mut m, dns, probe) = setup(false); m.run_round().await; probe.set_down(ip(1), true); probe.set_down(ip(2), true); m.run_round().await; m.run_round().await; let r = m.run_round().await; assert_eq!(r.events.len(), 3, "{:?}", r.events); assert_eq!(r.events.last(), Some(&Event::NoHealthyMembers)); assert_eq!(m.run_round().await.events, vec![]); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(1), ip(2)])); } #[tokio::test] async fn api_errors_alert_once_then_recover() { let (mut m, dns, _) = setup(false); m.run_round().await; dns.failing.store(true, Ordering::Relaxed); let r = m.run_round().await; assert!(matches!(&r.events[..], [Event::CheckFailed { .. }])); assert!(r.error.is_some()); let r = m.run_round().await; assert_eq!(r.events, vec![]); assert!(r.error.is_some()); dns.failing.store(false, Ordering::Relaxed); assert_eq!(m.run_round().await.events, vec![Event::CheckRecovered]); } #[tokio::test] async fn dry_run_reports_changes_without_writing() { let (mut m, dns, _) = setup(true); let r = m.run_round().await; assert!(r.events.contains(&Event::DnsUpdated { added: vec![ip(1), ip(2)], removed: vec![PLACEHOLDER] })); assert_eq!(r.dns_now, Some(vec![PLACEHOLDER])); assert_eq!(dns.ips(NAME), BTreeSet::from([PLACEHOLDER])); assert_eq!(dns.writes.load(Ordering::Relaxed), 0); } #[tokio::test] async fn aaaa_group_manages_only_aaaa_records() { let v6 = |last: u16| IpAddr::V6(Ipv6Addr::new(0x2001, 0xdb8, 0, 0, 0, 0, 0, last)); let dns = Arc::new(FakeDns::default()); // Members also have A records, and the group name has an A record of its own: // neither may be probed or touched. dns.set(A, &[ip(1), v6(1)]); dns.set(B, &[v6(2)]); dns.set(NAME, &[ip(9), v6(9)]); let probe = Arc::new(FakeProbe::default()); probe.set_down(v6(2), true); let cfg = Group { zone: "example.com".to_string(), name: NAME.to_string(), record_type: RecordType::AAAA, members: vec![A.to_string(), B.to_string()], ttl: 60, tcp_port: None, }; let mut m = GroupMonitor::new(cfg, &ProbeConfig::default(), dns.clone(), probe, false); let r = m.run_round().await; assert_eq!(r.record_type, "AAAA"); assert_eq!(r.ips.iter().map(|s| s.ip).collect::>(), vec![v6(1), v6(2)]); assert!(r.events.contains(&Event::DnsUpdated { added: vec![v6(1)], removed: vec![v6(9)] })); assert_eq!(dns.ips(NAME), BTreeSet::from([ip(9), v6(1)])); } }