use tokio::sync::broadcast;
use crate::lfd::types::Event;
#[derive(Debug, Clone)]
pub struct EventHub {
sender: broadcast::Sender<Event>,
}
impl EventHub {
pub fn new(buffer: usize) -> Self {
let (sender, _) = broadcast::channel(buffer);
Self { sender }
}
pub fn send(&self, event: Event) {
let _ = self.sender.send(event);
}
pub fn subscribe(&self) -> broadcast::Receiver<Event> {
self.sender.subscribe()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::lfd::id::LfdId;
#[tokio::test]
async fn event_hub_delivers_events_in_order() {
let hub = EventHub::new(16);
let mut rx = hub.subscribe();
let wave_id = LfdId::from_raw("wave-1");
let run_id = LfdId::from_raw("run-1");
hub.send(Event::wave_started(wave_id.clone(), run_id.clone()));
hub.send(Event::wave_waiting(
wave_id.clone(),
run_id,
"review".to_string(),
None,
None,
));
hub.send(Event::wave_updated(wave_id));
let types: Vec<String> = (0..3)
.map(|_| {
let event = rx.try_recv().unwrap();
let json = serde_json::to_value(&event).unwrap();
json["type"].as_str().unwrap().to_string()
})
.collect();
assert_eq!(types, vec!["wave_started", "wave_waiting", "wave_updated"]);
}
}