use std::{ net::IpAddr, sync::atomic::{AtomicU16, Ordering}, time::Duration, }; use anyhow::{Context, Result}; use surge_ping::{ICMP, PingIdentifier, PingSequence}; use tokio::{net::TcpStream, time::timeout}; use crate::{config::ProbeConfig, event::ProbeOutcome}; pub trait Probe { /// `Ok(rtt)` if `ip` is alive (and accepts `tcp_port`, when given), `Err(reason)` otherwise. fn probe(&self, ip: IpAddr, tcp_port: Option) -> impl Future; } /// ICMP / ICMPv6 echo, plus an optional TCP connect. pub struct IcmpProber { v4: surge_ping::Client, v6: surge_ping::Client, count: u16, timeout: Duration, next_ident: AtomicU16, } impl IcmpProber { pub fn new(cfg: &ProbeConfig) -> Result { let open = |kind| { surge_ping::Client::new(&surge_ping::Config::builder().kind(kind).build()) .with_context(|| format!("opening {kind:?} socket (needs root or CAP_NET_RAW)")) }; Ok(Self { v4: open(ICMP::V4)?, v6: open(ICMP::V6)?, count: cfg.ping_count, timeout: Duration::from_millis(cfg.ping_timeout_ms), next_ident: AtomicU16::new(std::process::id() as u16), }) } async fn ping(&self, ip: IpAddr) -> ProbeOutcome { let client = if ip.is_ipv4() { &self.v4 } else { &self.v6 }; // Concurrent pingers on one socket are told apart by identifier. let ident = PingIdentifier(self.next_ident.fetch_add(1, Ordering::Relaxed)); let mut pinger = client.pinger(ip, ident).await; pinger.timeout(self.timeout); let mut last_err = String::new(); for seq in 0..self.count { match pinger.ping(PingSequence(seq), &[0; 56]).await { Ok((_, rtt)) => return Ok(rtt), Err(e) => last_err = e.to_string(), } } Err(format!("ping: {} lost ({last_err})", self.count)) } } impl Probe for IcmpProber { async fn probe(&self, ip: IpAddr, tcp_port: Option) -> ProbeOutcome { let rtt = self.ping(ip).await?; if let Some(port) = tcp_port { match timeout(self.timeout, TcpStream::connect((ip, port))).await { Ok(Ok(_)) => {} Ok(Err(e)) => return Err(format!("tcp/{port}: {e}")), Err(_) => return Err(format!("tcp/{port}: timed out")), } } Ok(rtt) } }