use pamoja_core::{EventBus, Result};
use tokio::sync::broadcast;
pub struct BroadcastBus<E> {
sender: broadcast::Sender<E>,
receiver: broadcast::Receiver<E>,
}
impl<E: Clone> BroadcastBus<E> {
pub fn new(capacity: usize) -> Self {
let (sender, receiver) = broadcast::channel(capacity.max(1));
Self { sender, receiver }
}
pub fn subscribe(&self) -> Self {
Self {
sender: self.sender.clone(),
receiver: self.sender.subscribe(),
}
}
}
impl<E: Clone> EventBus for BroadcastBus<E> {
type Event = E;
async fn publish(&self, event: Self::Event) -> Result<()> {
let _ = self.sender.send(event);
Ok(())
}
async fn next_event(&mut self) -> Result<Option<Self::Event>> {
loop {
match self.receiver.recv().await {
Ok(event) => return Ok(Some(event)),
Err(broadcast::error::RecvError::Closed) => return Ok(None),
Err(broadcast::error::RecvError::Lagged(_)) => continue,
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn delivers_a_published_event() {
let mut bus = BroadcastBus::new(8);
bus.publish(1).await.expect("publish");
assert_eq!(bus.next_event().await.expect("next"), Some(1));
}
#[tokio::test]
async fn fans_out_to_every_subscriber() {
let bus = BroadcastBus::new(8);
let mut first = bus.subscribe();
let mut second = bus.subscribe();
bus.publish("event").await.expect("publish");
assert_eq!(first.next_event().await.expect("next"), Some("event"));
assert_eq!(second.next_event().await.expect("next"), Some("event"));
}
#[tokio::test]
async fn a_lagging_subscriber_skips_dropped_events_and_resumes() {
let bus = BroadcastBus::new(2);
let mut subscriber = bus.subscribe();
for value in 0..5 {
bus.publish(value).await.expect("publish");
}
let mut seen = Vec::new();
seen.push(subscriber.next_event().await.expect("next").expect("event"));
seen.push(subscriber.next_event().await.expect("next").expect("event"));
assert_eq!(seen, vec![3, 4]);
}
}