//! Delivery of round reports. Sinks only see `RoundReport`s; the monitor only publishes them. pub mod log; pub mod telegram; use std::{sync::Arc, time::Duration}; use tokio::{ sync::broadcast::{self, error::RecvError}, task::JoinSet, }; use tracing::warn; use crate::event::RoundReport; pub trait Sink: Send + 'static { fn handle(&mut self, report: &RoundReport) -> impl Future + Send; } /// Fans every report out to every sink. Each sink runs in its own task, so a slow or /// failing sink never delays the monitor or the other sinks. pub struct Hub { tx: broadcast::Sender>, sinks: JoinSet<()>, } impl Hub { pub fn new() -> Self { let (tx, _) = broadcast::channel(64); Self { tx, sinks: JoinSet::new() } } pub fn add(&mut self, mut sink: S) { let mut rx = self.tx.subscribe(); self.sinks.spawn(async move { loop { match rx.recv().await { Ok(report) => sink.handle(&report).await, Err(RecvError::Lagged(n)) => warn!("a notification sink fell behind and skipped {n} reports"), Err(RecvError::Closed) => break, } } }); } pub fn publish(&self, report: RoundReport) { // Errors only when no sink is subscribed, and then there is nobody to tell. let _ = self.tx.send(Arc::new(report)); } /// Stop accepting reports and give sinks up to `grace` to deliver what is queued. pub async fn close(self, grace: Duration) { let Self { tx, mut sinks } = self; drop(tx); let drained = tokio::time::timeout(grace, async { while sinks.join_next().await.is_some() {} }).await; if drained.is_err() { warn!("notification sinks did not finish within {grace:?}"); } } }