use std::sync::Arc;
use std::sync::OnceLock;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::time::{Duration, Instant};
pub(crate) struct HealthState {
started: OnceLock<Instant>,
sockets_open: AtomicBool,
first_event: AtomicBool,
last_event_ms: AtomicU64,
active_flows: AtomicUsize,
packets: AtomicU64,
drops: AtomicU64,
handler_errors: AtomicU64,
backend_errors: AtomicU64,
}
impl HealthState {
pub(crate) fn new() -> Arc<Self> {
Arc::new(Self {
started: OnceLock::new(),
sockets_open: AtomicBool::new(false),
first_event: AtomicBool::new(false),
last_event_ms: AtomicU64::new(0),
active_flows: AtomicUsize::new(0),
packets: AtomicU64::new(0),
drops: AtomicU64::new(0),
handler_errors: AtomicU64::new(0),
backend_errors: AtomicU64::new(0),
})
}
pub(crate) fn record_handler_error(&self) {
self.handler_errors.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn record_backend_error(&self) {
self.backend_errors.fetch_add(1, Ordering::Relaxed);
}
pub(crate) fn mark_started(&self) {
let _ = self.started.set(Instant::now());
}
pub(crate) fn mark_sockets_open(&self) {
self.sockets_open.store(true, Ordering::Relaxed);
}
pub(crate) fn record_event(&self, active_flows: usize) {
let elapsed = self
.started
.get()
.map(|s| s.elapsed().as_millis() as u64)
.unwrap_or(0);
self.last_event_ms.store(elapsed, Ordering::Relaxed);
self.active_flows.store(active_flows, Ordering::Relaxed);
self.first_event.store(true, Ordering::Relaxed);
}
pub(crate) fn record_totals(&self, packets: u64, drops: u64) {
self.packets.store(packets, Ordering::Relaxed);
self.drops.store(drops, Ordering::Relaxed);
}
fn uptime(&self) -> Duration {
self.started
.get()
.map(Instant::elapsed)
.unwrap_or(Duration::ZERO)
}
fn last_event_age(&self) -> Option<Duration> {
if !self.first_event.load(Ordering::Relaxed) {
return None;
}
let last = Duration::from_millis(self.last_event_ms.load(Ordering::Relaxed));
Some(self.uptime().saturating_sub(last))
}
}
#[derive(Clone)]
pub struct MonitorHealth {
inner: Arc<HealthState>,
}
impl MonitorHealth {
pub(crate) fn new(inner: Arc<HealthState>) -> Self {
Self { inner }
}
pub fn uptime(&self) -> Duration {
self.inner.uptime()
}
pub fn last_event_age(&self) -> Option<Duration> {
self.inner.last_event_age()
}
pub fn active_flows(&self) -> usize {
self.inner.active_flows.load(Ordering::Relaxed)
}
pub fn packets(&self) -> u64 {
self.inner.packets.load(Ordering::Relaxed)
}
pub fn drops(&self) -> u64 {
self.inner.drops.load(Ordering::Relaxed)
}
pub fn handler_errors(&self) -> u64 {
self.inner.handler_errors.load(Ordering::Relaxed)
}
pub fn backend_errors(&self) -> u64 {
self.inner.backend_errors.load(Ordering::Relaxed)
}
pub fn has_seen_traffic(&self) -> bool {
self.inner.first_event.load(Ordering::Relaxed)
}
pub fn is_ready(&self) -> bool {
self.inner.sockets_open.load(Ordering::Relaxed)
}
pub fn is_live(&self, window: Duration) -> bool {
match self.last_event_age() {
Some(age) => age < window,
None => self.uptime() < window,
}
}
#[cfg(feature = "metrics")]
pub fn record_metrics(&self) {
metrics::gauge!(crate::metrics::GAUGE_HANDLER_ERRORS).set(self.handler_errors() as f64);
metrics::gauge!(crate::metrics::GAUGE_BACKEND_ERRORS).set(self.backend_errors() as f64);
metrics::gauge!(crate::metrics::GAUGE_ACTIVE_FLOWS).set(self.active_flows() as f64);
}
pub fn snapshot(&self) -> MonitorHealthSnapshot {
MonitorHealthSnapshot {
uptime: self.uptime(),
last_event_age: self.last_event_age(),
active_flows: self.active_flows(),
packets: self.packets(),
drops: self.drops(),
handler_errors: self.handler_errors(),
backend_errors: self.backend_errors(),
ready: self.is_ready(),
seen_traffic: self.has_seen_traffic(),
}
}
}
impl std::fmt::Debug for MonitorHealth {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("MonitorHealth")
.field("uptime", &self.uptime())
.field("ready", &self.is_ready())
.field("active_flows", &self.active_flows())
.finish_non_exhaustive()
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct MonitorHealthSnapshot {
pub uptime: Duration,
pub last_event_age: Option<Duration>,
pub active_flows: usize,
pub packets: u64,
pub drops: u64,
pub handler_errors: u64,
pub backend_errors: u64,
pub ready: bool,
pub seen_traffic: bool,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn unstarted_monitor_reads_zero_uptime_and_no_event() {
let h = MonitorHealth::new(HealthState::new());
assert_eq!(h.uptime(), Duration::ZERO);
assert_eq!(h.last_event_age(), None);
assert!(!h.is_ready());
assert!(!h.has_seen_traffic());
}
#[test]
fn readiness_flips_on_sockets_open_not_on_traffic() {
let state = HealthState::new();
let h = MonitorHealth::new(state.clone());
state.mark_started();
assert!(!h.is_ready(), "not ready before sockets open");
state.mark_sockets_open();
assert!(
h.is_ready(),
"ready once sockets open, even with no traffic"
);
assert!(!h.has_seen_traffic());
}
#[test]
fn liveness_uses_startup_grace_then_last_event_age() {
let state = HealthState::new();
let h = MonitorHealth::new(state.clone());
state.mark_started();
assert!(h.is_live(Duration::from_secs(60)));
assert!(!h.is_live(Duration::ZERO));
state.record_event(3);
assert!(h.has_seen_traffic());
assert_eq!(h.active_flows(), 3);
assert!(h.is_live(Duration::from_secs(60)));
assert!(!h.is_live(Duration::ZERO));
}
#[test]
fn totals_track_latest_telemetry_sample() {
let state = HealthState::new();
let h = MonitorHealth::new(state.clone());
assert_eq!(h.packets(), 0);
state.record_totals(1000, 5);
assert_eq!(h.packets(), 1000);
assert_eq!(h.drops(), 5);
state.record_totals(2500, 12);
assert_eq!(h.packets(), 2500);
assert_eq!(h.drops(), 12);
}
#[test]
fn resilience_counters_accumulate() {
let state = HealthState::new();
let h = MonitorHealth::new(state.clone());
assert_eq!(h.handler_errors(), 0);
assert_eq!(h.backend_errors(), 0);
state.record_handler_error();
state.record_handler_error();
state.record_backend_error();
assert_eq!(h.handler_errors(), 2);
assert_eq!(h.backend_errors(), 1);
}
#[test]
fn snapshot_is_internally_consistent() {
let state = HealthState::new();
let h = MonitorHealth::new(state.clone());
state.mark_started();
state.mark_sockets_open();
state.record_event(7);
state.record_totals(42, 1);
state.record_handler_error();
let snap = h.snapshot();
assert!(snap.ready);
assert!(snap.seen_traffic);
assert_eq!(snap.active_flows, 7);
assert_eq!(snap.packets, 42);
assert_eq!(snap.drops, 1);
assert_eq!(snap.handler_errors, 1);
assert_eq!(snap.backend_errors, 0);
assert!(snap.last_event_age.is_some());
}
#[cfg(feature = "metrics")]
#[test]
fn record_metrics_is_a_noop_without_a_recorder() {
let state = HealthState::new();
let h = MonitorHealth::new(state.clone());
state.record_handler_error();
h.record_metrics(); }
}