use std::io::{Read, Write};
use std::net::{TcpStream, ToSocketAddrs};
use std::sync::mpsc;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
const FAST_TICK: Duration = Duration::from_secs(30);
const SLOW_TICK_MULTIPLE: u32 = 60;
#[derive(Debug, Clone, Default)]
pub struct Tier3Snapshot {
pub public_ip: Option<String>,
pub pending_updates: Option<usize>,
pub net_rx_bps: Option<u64>,
pub net_tx_bps: Option<u64>,
pub ready: bool,
}
#[derive(Debug, Clone, Default)]
pub struct Tier3Handle(Arc<Mutex<Tier3Snapshot>>);
impl Tier3Handle {
pub fn snapshot(&self) -> Tier3Snapshot {
match self.0.lock() {
Ok(guard) => guard.clone(),
Err(poisoned) => poisoned.into_inner().clone(),
}
}
pub(crate) fn set(&self, snapshot: Tier3Snapshot) {
match self.0.lock() {
Ok(mut guard) => *guard = snapshot,
Err(poisoned) => *poisoned.into_inner() = snapshot,
}
}
}
pub fn spawn_refresh_thread(handle: Tier3Handle) {
std::thread::spawn(move || {
let mut networks = sysinfo::Networks::new_with_refreshed_list();
let mut last_tick = Instant::now();
let mut tick: u32 = 0;
loop {
let public_ip = fetch_public_ip();
networks.refresh(true);
let elapsed = last_tick.elapsed().as_secs_f64().max(1.0);
last_tick = Instant::now();
let (rx, tx) = networks
.list()
.values()
.fold((0u64, 0u64), |(rx, tx), n| (rx + n.received(), tx + n.transmitted()));
let pending_updates = if tick.is_multiple_of(SLOW_TICK_MULTIPLE) {
fetch_pending_updates()
} else {
handle.snapshot().pending_updates
};
handle.set(Tier3Snapshot {
public_ip,
pending_updates,
net_rx_bps: Some((rx as f64 / elapsed) as u64),
net_tx_bps: Some((tx as f64 / elapsed) as u64),
ready: true,
});
tick = tick.wrapping_add(1);
std::thread::sleep(FAST_TICK);
}
});
}
pub fn fetch_tier3_bounded(timeout: Duration) -> Tier3Snapshot {
let (tx, rx) = mpsc::channel();
std::thread::spawn(move || {
let snapshot = Tier3Snapshot {
public_ip: fetch_public_ip(),
pending_updates: fetch_pending_updates(),
net_rx_bps: None,
net_tx_bps: None,
ready: true,
};
let _ = tx.send(snapshot);
});
rx.recv_timeout(timeout).unwrap_or_default()
}
fn fetch_public_ip() -> Option<String> {
let addr = "api.ipify.org:80".to_socket_addrs().ok()?.next()?;
let mut stream = TcpStream::connect_timeout(&addr, Duration::from_secs(3)).ok()?;
stream.set_read_timeout(Some(Duration::from_secs(3))).ok()?;
stream.set_write_timeout(Some(Duration::from_secs(3))).ok()?;
stream.write_all(b"GET / HTTP/1.1\r\nHost: api.ipify.org\r\nConnection: close\r\n\r\n").ok()?;
let mut response = String::new();
stream.read_to_string(&mut response).ok()?;
let body = response.split("\r\n\r\n").nth(1)?.trim();
let looks_like_an_address = !body.is_empty()
&& body.len() <= 45 && body.chars().all(|c| c.is_ascii_hexdigit() || c == '.' || c == ':');
if looks_like_an_address { Some(body.to_string()) } else { None }
}
#[cfg(target_os = "macos")]
fn fetch_pending_updates() -> Option<usize> {
let out = std::process::Command::new("brew").arg("outdated").output().ok()?;
if !out.status.success() {
return None;
}
Some(String::from_utf8_lossy(&out.stdout).lines().filter(|l| !l.trim().is_empty()).count())
}
#[cfg(target_os = "linux")]
fn fetch_pending_updates() -> Option<usize> {
use std::path::Path;
if Path::new("/usr/bin/apt").exists() || Path::new("/usr/bin/apt-get").exists() {
let out = std::process::Command::new("apt").args(["list", "--upgradable"]).output().ok()?;
if !out.status.success() {
return None;
}
let text = String::from_utf8_lossy(&out.stdout);
return Some(text.lines().filter(|l| !l.trim().is_empty() && !l.starts_with("Listing")).count());
}
if Path::new("/usr/bin/dnf").exists() {
let out = std::process::Command::new("dnf").arg("check-update").output().ok()?;
if !matches!(out.status.code(), Some(0) | Some(100)) {
return None;
}
let text = String::from_utf8_lossy(&out.stdout);
return Some(text.lines().filter(|l| l.trim().split_whitespace().count() >= 3).count());
}
if Path::new("/usr/bin/pacman").exists() {
let out = std::process::Command::new("pacman").args(["-Qu"]).output().ok()?;
if !out.status.success() {
return None;
}
return Some(String::from_utf8_lossy(&out.stdout).lines().filter(|l| !l.trim().is_empty()).count());
}
None
}
#[cfg(target_os = "windows")]
fn fetch_pending_updates() -> Option<usize> {
let out = std::process::Command::new("winget").arg("upgrade").output().ok()?;
if !out.status.success() {
return None;
}
let text = String::from_utf8_lossy(&out.stdout);
let mut count = 0usize;
let mut past_header = false;
for line in text.lines() {
let trimmed = line.trim();
if !trimmed.is_empty() && trimmed.chars().all(|c| c == '-') {
past_header = true;
continue;
}
if past_header && !trimmed.is_empty() {
count += 1;
}
}
Some(count)
}
#[cfg(not(any(target_os = "macos", target_os = "linux", target_os = "windows")))]
fn fetch_pending_updates() -> Option<usize> {
None
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_handle_reports_not_ready() {
let handle = Tier3Handle::default();
assert!(!handle.snapshot().ready);
}
#[test]
fn bounded_fetch_never_waits_past_its_timeout() {
let start = Instant::now();
let snapshot = fetch_tier3_bounded(Duration::from_millis(1));
let elapsed = start.elapsed();
assert!(!snapshot.ready, "expected a timed-out snapshot, got: {snapshot:?}");
assert!(snapshot.public_ip.is_none());
assert!(snapshot.pending_updates.is_none());
assert!(elapsed < Duration::from_millis(500), "bounded fetch took too long: {elapsed:?}");
}
#[test]
fn set_then_snapshot_round_trips() {
let handle = Tier3Handle::default();
handle.set(Tier3Snapshot { public_ip: Some("203.0.113.7".to_string()), ready: true, ..Default::default() });
let snap = handle.snapshot();
assert!(snap.ready);
assert_eq!(snap.public_ip.as_deref(), Some("203.0.113.7"));
}
}