use crate::{
closure_handler_wrapper::ClosureHandlerWrapper,
dispatched_event::DispatchedEvent,
event_dispatcher::event_dispatcher,
event_listener::{merge_subscribers, unsubscribe, SubscriberList, LOG_TITLE},
};
use async_trait::async_trait;
use futures::future::BoxFuture;
#[async_trait]
pub trait Dispatchable:
serde::Serialize + serde::de::DeserializeOwned + Clone + Send + Sync
{
fn event() -> String
where
Self: Sized,
{
std::any::type_name::<Self>().to_string()
}
fn dispatch_event(self)
where
Self: Sized + 'static,
{
event_dispatcher().dispatch(self);
}
fn dispatch_event_as(self, name: &str) {
event_dispatcher().dispatch_str(name, self);
}
fn supports_cluster(&self) -> bool {
true
}
fn serialize_event(&self) -> String {
let event = DispatchedEvent::new(
serde_json::to_string(self).expect("could not serialize event"),
Self::event(),
);
serde_json::to_string(&event).unwrap()
}
async fn subscribe<H: EventHandler + Default>()
where
Self: Sized,
{
crate::setup().await;
let event: String = Self::event();
let the_handler = H::default().to_handler();
let mut subscriber = SubscriberList::new();
log::trace!(
target: LOG_TITLE,
"registered handler: {:?}, for event: {:?}",
&event,
&the_handler.handler_id()
);
subscriber.insert(event, vec![the_handler]);
merge_subscribers(subscriber).await;
}
async fn subscribe_with(handler: impl EventHandler) {
crate::setup().await;
let event: String = Self::event();
let the_handler = handler.to_handler();
let mut subscriber = SubscriberList::new();
log::trace!(
target: LOG_TITLE,
"registered handler: {:?}, for event: {:?}",
&event,
&the_handler.handler_id()
);
subscriber.insert(event, vec![the_handler]);
merge_subscribers(subscriber).await;
}
async fn subscribe_fn(
handler: impl Fn(DispatchedEvent) -> BoxFuture<'static, ()> + Send + Sync + 'static,
) {
let wrapper = ClosureHandlerWrapper(handler);
Self::subscribe_with(wrapper).await;
}
async fn unsubscribe<H: EventHandler + Default>() {
crate::setup().await;
let the_handler = H::default().to_handler();
unsubscribe(Self::event(), the_handler.handler_id()).await;
}
}
#[async_trait]
pub trait EventHandler: Send + Sync + 'static {
async fn handle(&self, event: DispatchedEvent);
fn to_handler(self) -> Box<Self>
where
Self: Sized,
{
Box::new(self)
}
fn handler_id(&self) -> String {
std::any::type_name::<Self>().to_string()
}
fn execute_once(&self) -> bool {
false
}
fn propagate(&self) -> bool {
true
}
}
mod test {
use super::*;
#[derive(Clone, serde::Serialize, serde::Deserialize)]
struct UserCreated {
id: u32,
}
impl Dispatchable for UserCreated {}
#[tokio::test]
async fn test_event_dispatching() {
UserCreated::subscribe::<HandleUserCreated>().await;
event_dispatcher()
.dispatch_sync(UserCreated { id: 200 })
.await;
}
#[tokio::test]
async fn test_event_dispatching_with() {
UserCreated::subscribe_with(HandleUserCreated).await;
event_dispatcher()
.dispatch_sync(UserCreated { id: 200 })
.await;
}
#[derive(Default)]
#[allow(dead_code)]
struct HandleUserCreated;
#[async_trait]
impl EventHandler for HandleUserCreated {
async fn handle(&self, dispatched: DispatchedEvent) {
let the_event = dispatched.the_event();
assert_eq!(the_event.is_none(), false);
let event: UserCreated = the_event.unwrap();
assert_eq!(event.id, 200);
}
}
}