use common::event::{Event, EventHub, Origin};
use flume::Receiver;
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
pub type EventCallback = Box<dyn Fn(Event) + Send>;
type SubscriberList = Vec<(u64, EventCallback)>;
#[derive(Clone)]
pub struct EventHubClient {
subscribers: Arc<Mutex<HashMap<Origin, SubscriberList>>>,
receiver: Receiver<Event>,
next_subscriber_id: Arc<AtomicU64>,
}
impl EventHubClient {
pub fn new(event_hub: &EventHub) -> Self {
EventHubClient {
subscribers: Arc::new(Mutex::new(HashMap::new())),
receiver: event_hub.subscribe_receiver(),
next_subscriber_id: Arc::new(AtomicU64::new(0)),
}
}
pub fn subscribe<F>(&self, origin: Origin, callback: F) -> SubscriptionToken
where
F: Fn(Event) + Send + 'static,
{
let id = self.next_subscriber_id.fetch_add(1, Ordering::Relaxed);
{
let mut subs = self.subscribers.lock().unwrap();
subs.entry(origin.clone())
.or_default()
.push((id, Box::new(callback)));
}
SubscriptionToken {
subscribers: Arc::clone(&self.subscribers),
origin,
id,
}
}
pub fn start(&self, quit_signal: Arc<std::sync::atomic::AtomicBool>) {
let receiver = self.receiver.clone();
let subscribers = Arc::clone(&self.subscribers);
let quit_signal = Arc::clone(&quit_signal);
log::info!("EventHubClient starting event loop");
thread::spawn(move || {
log::info!("EventHubClient event loop started");
loop {
match receiver.recv_timeout(std::time::Duration::from_millis(200)) {
Ok(event) => {
log::debug!("EventHubClient received event: {:?}", event);
let subs = subscribers.lock().unwrap();
if let Some(callbacks) = subs.get(&event.origin) {
for (_id, callback) in callbacks {
callback(event.clone());
}
}
}
Err(flume::RecvTimeoutError::Timeout) => {
}
Err(flume::RecvTimeoutError::Disconnected) => {
log::info!("EventHubClient channel disconnected");
break;
}
}
if quit_signal.load(std::sync::atomic::Ordering::Relaxed) {
log::info!("EventHubClient quitting event loop");
break;
}
}
});
}
}
pub struct SubscriptionToken {
subscribers: Arc<Mutex<HashMap<Origin, SubscriberList>>>,
origin: Origin,
id: u64,
}
impl Drop for SubscriptionToken {
fn drop(&mut self) {
if let Ok(mut subs) = self.subscribers.lock()
&& let Some(list) = subs.get_mut(&self.origin)
{
list.retain(|(id, _)| *id != self.id);
if list.is_empty() {
subs.remove(&self.origin);
}
}
}
}