use std::sync::atomic::{AtomicU64, Ordering};
use tokio::sync::mpsc;
use super::delivery::EntityEvent;
pub struct RealtimeBroadcastObserver {
event_tx: mpsc::Sender<EntityEvent>,
events_dropped: AtomicU64,
}
impl RealtimeBroadcastObserver {
#[must_use]
pub fn new(capacity: usize) -> (Self, mpsc::Receiver<EntityEvent>) {
let (tx, rx) = mpsc::channel(capacity);
(
Self {
event_tx: tx,
events_dropped: AtomicU64::new(0),
},
rx,
)
}
pub fn on_mutation_complete(&self, event: EntityEvent) {
if self.event_tx.try_send(event).is_err() {
self.events_dropped.fetch_add(1, Ordering::Relaxed);
}
}
#[must_use]
pub fn events_dropped_total(&self) -> u64 {
self.events_dropped.load(Ordering::Relaxed)
}
}