use crate::scope::Scope;
use futures_core::future::BoxFuture;
use serde_json::Value;
use std::sync::Arc;
use tracing::{error, warn};
pub type AnyError = Box<dyn std::error::Error + Send + Sync + 'static>;
pub type EventName = String;
pub type EventHandler =
Arc<dyn Fn(&Scope, Value) -> BoxFuture<'static, Result<(), AnyError>> + Send + Sync>;
#[derive(Clone, Default)]
pub struct EventBus {
handlers: Arc<Vec<(EventName, EventHandler)>>,
}
impl EventBus {
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn on(self, name: impl Into<EventName>, handler: EventHandler) -> Self {
let mut handlers: Vec<(EventName, EventHandler)> = match Arc::try_unwrap(self.handlers) {
Ok(handlers) => handlers,
Err(shared) => (*shared).clone(),
};
handlers.push((name.into(), handler));
Self {
handlers: Arc::new(handlers),
}
}
pub fn handlers(&self) -> &[(EventName, EventHandler)] {
&self.handlers
}
#[allow(clippy::needless_pass_by_value)]
pub fn emit_in(&self, scope: &Scope, name: &str, payload: Value) {
let matched: Vec<&(EventName, EventHandler)> = self
.handlers
.iter()
.filter(|(handler_name, _)| handler_name == name)
.collect();
for (_, handler) in &matched {
let fut = handler(scope, payload.clone());
let event = name.to_string();
scope.defer.wait_until(Box::pin(async move {
if let Err(err) = fut.await {
error!(event = %event, error = %err, "event handler failed");
}
}));
}
if matched.is_empty() && !self.handlers.is_empty() {
warn!(event = %name, "emitted event has no registered handler");
}
}
}