//! Métrologie du serveur UDP. //! //! Ce module expose : //! - [`UdpMetrics`] : compteurs atomiques lock-free (pas de contention dans la //! boucle de routage). //! - [`UdpMetricsSnapshot`] : lecture cohérente de tous les compteurs à un //! instant T, utilisable pour calculer des deltas. //! - [`UdpRates`] : taux moyens par seconde calculés entre deux snapshots. //! - [`spawn_reporter`] : tâche tokio de reporting périodique via `tracing`. //! //! # Métriques collectées //! //! | Compteur | Description | //! |--------------------|-----------------------------------------------| //! | `packets_received` | Datagrammes reçus | //! | `bytes_received` | Octets reçus (payload uniquement) | //! | `packets_sent` | Datagrammes retransmis vers des abonnés | //! | `bytes_sent` | Octets retransmis | //! | `packets_dropped` | Paquets ignorés (canal sans abonnés) | //! | `send_errors` | Échecs `send_to` | //! | `recv_errors` | Échecs `recv_from` (avant erreur fatale) | //! //! Chaque métrique est également disponible en taux moyen par seconde via //! [`UdpMetricsSnapshot::rates_since`]. use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; use std::time::{Duration, Instant}; use crate::metrics::{Metrics, MetricsSnapshot}; // ── Compteurs ──────────────────────────────────────────────────────────────── /// Compteurs atomiques du serveur UDP. /// /// Partagé via [`Arc`] entre la boucle de routage et le reporter périodique. /// Tous les accès utilisent [`Ordering::Relaxed`] : on accepte que les lectures /// voient des valeurs légèrement décalées entre compteurs (suffisant pour de la /// métrologie), ce qui évite tout overhead de synchronisation. #[derive(Debug, Default)] pub struct UdpMetrics { /// Nombre total de datagrammes reçus. pub packets_received: AtomicU64, /// Volume total d'octets reçus (payload des datagrammes). pub bytes_received: AtomicU64, /// Nombre total de datagrammes retransmis (somme sur tous les abonnés). pub packets_sent: AtomicU64, /// Volume total d'octets retransmis. pub bytes_sent: AtomicU64, /// Paquets ignorés car le canal ne possède aucun abonné. pub packets_dropped: AtomicU64, /// Nombre d'erreurs `send_to` (non fatales). pub send_errors: AtomicU64, /// Nombre d'erreurs `recv_from` enregistrées avant arrêt du serveur. pub recv_errors: AtomicU64, } impl UdpMetrics { /// Crée un jeu de métriques vide enroulé dans un [`Arc`]. pub fn new() -> Arc { Arc::new(Self::default()) } /// Enregistre la réception d'un datagramme de `bytes` octets. #[inline] pub fn inc_received(&self, bytes: u64) { self.packets_received.fetch_add(1, Ordering::Relaxed); self.bytes_received.fetch_add(bytes, Ordering::Relaxed); } /// Enregistre l'émission d'un datagramme de `bytes` octets vers un client. #[inline] pub fn inc_sent(&self, bytes: u64) { self.packets_sent.fetch_add(1, Ordering::Relaxed); self.bytes_sent.fetch_add(bytes, Ordering::Relaxed); } /// Enregistre un paquet ignoré (canal sans abonnés). #[inline] pub fn inc_dropped(&self) { self.packets_dropped.fetch_add(1, Ordering::Relaxed); } /// Enregistre un échec `send_to` non fatal. #[inline] pub fn inc_send_error(&self) { self.send_errors.fetch_add(1, Ordering::Relaxed); } /// Enregistre un échec `recv_from`. #[inline] pub fn inc_recv_error(&self) { self.recv_errors.fetch_add(1, Ordering::Relaxed); } /// Prend un instantané cohérent de tous les compteurs. pub fn snapshot(&self) -> UdpMetricsSnapshot { UdpMetricsSnapshot { taken_at: Instant::now(), packets_received: self.packets_received.load(Ordering::Relaxed), bytes_received: self.bytes_received.load(Ordering::Relaxed), packets_sent: self.packets_sent.load(Ordering::Relaxed), bytes_sent: self.bytes_sent.load(Ordering::Relaxed), packets_dropped: self.packets_dropped.load(Ordering::Relaxed), send_errors: self.send_errors.load(Ordering::Relaxed), recv_errors: self.recv_errors.load(Ordering::Relaxed), } } } impl Metrics for UdpMetrics { type Snapshot = UdpMetricsSnapshot; fn snapshot(&self) -> UdpMetricsSnapshot { self.snapshot() } } // ── Snapshot ───────────────────────────────────────────────────────────────── /// Lecture cohérente de l'ensemble des compteurs à un instant T. /// /// Permet de calculer des deltas et des taux entre deux points dans le temps /// sans bloquer la boucle de routage. #[derive(Debug, Clone, Copy)] pub struct UdpMetricsSnapshot { pub taken_at: Instant, pub packets_received: u64, pub bytes_received: u64, pub packets_sent: u64, pub bytes_sent: u64, pub packets_dropped: u64, pub send_errors: u64, pub recv_errors: u64, } impl UdpMetricsSnapshot { /// Calcule les taux moyens par seconde depuis un snapshot précédent. pub fn rates_since(&self, previous: &Self) -> UdpRates { let secs = self .taken_at .duration_since(previous.taken_at) .as_secs_f64() .max(f64::EPSILON); UdpRates { packets_received_per_sec: self .packets_received .saturating_sub(previous.packets_received) as f64 / secs, bytes_received_per_sec: self.bytes_received.saturating_sub(previous.bytes_received) as f64 / secs, packets_sent_per_sec: self.packets_sent.saturating_sub(previous.packets_sent) as f64 / secs, bytes_sent_per_sec: self.bytes_sent.saturating_sub(previous.bytes_sent) as f64 / secs, packets_dropped_per_sec: self .packets_dropped .saturating_sub(previous.packets_dropped) as f64 / secs, } } } impl MetricsSnapshot for UdpMetricsSnapshot { fn taken_at(&self) -> Instant { self.taken_at } } // ── Taux ───────────────────────────────────────────────────────────────────── /// Taux moyens par seconde calculés entre deux [`UdpMetricsSnapshot`]. #[derive(Debug, Clone, Copy)] pub struct UdpRates { /// Paquets reçus par seconde. pub packets_received_per_sec: f64, /// Octets reçus par seconde. pub bytes_received_per_sec: f64, /// Paquets envoyés par seconde. pub packets_sent_per_sec: f64, /// Octets envoyés par seconde. pub bytes_sent_per_sec: f64, /// Paquets ignorés par seconde. pub packets_dropped_per_sec: f64, } // ── Reporter périodique ─────────────────────────────────────────────────────── /// Lance une tâche Tokio qui logue les métriques toutes les `interval`. /// /// Chaque rapport inclut les compteurs cumulatifs **et** les taux moyens sur /// la fenêtre écoulée depuis le rapport précédent. /// /// # Exemple /// ```no_run /// use std::time::Duration; /// use std::sync::Arc; /// use oxspeak_server_lib::udp::metrics::{UdpMetrics, spawn_reporter}; /// /// #[tokio::main] /// async fn main() { /// let metrics = UdpMetrics::new(); /// spawn_reporter(Arc::clone(&metrics), Duration::from_secs(5)); /// } /// ``` pub fn spawn_reporter(metrics: Arc, interval: Duration) { tokio::spawn(async move { let mut ticker = tokio::time::interval(interval); // Le premier tick est immédiat ; on le consomme pour démarrer à t=0. ticker.tick().await; let mut prev_snapshot = metrics.snapshot(); loop { ticker.tick().await; let current = metrics.snapshot(); let rates = current.rates_since(&prev_snapshot); tracing::info!( // ── Cumulatifs ── pkts_rx = current.packets_received, bytes_rx = current.bytes_received, pkts_tx = current.packets_sent, bytes_tx = current.bytes_sent, pkts_dropped = current.packets_dropped, send_errors = current.send_errors, recv_errors = current.recv_errors, // ── Taux / s ── pkts_rx_s = format!("{:.1}", rates.packets_received_per_sec), bytes_rx_s = format!("{:.0}", rates.bytes_received_per_sec), pkts_tx_s = format!("{:.1}", rates.packets_sent_per_sec), bytes_tx_s = format!("{:.0}", rates.bytes_sent_per_sec), pkts_dropped_s = format!("{:.1}", rates.packets_dropped_per_sec), "UDP metrics" ); prev_snapshot = current; } }); }