1#![forbid(unsafe_code)]
2
3use core::fmt::Debug;
4
5use kithara_platform::tokio::sync::broadcast::error::{RecvError, TryRecvError};
6
7use crate::{Envelope, EventBus, EventMeta, TopicReceiver};
8
9pub trait Event: Clone + Debug + Send + Sync + 'static {}
14
15pub trait EventSet: Sized + Send + 'static {
26 type Receivers: Send;
27
28 fn publish(bus: &EventBus, meta: EventMeta, event: Self);
29
30 fn recv(
31 rx: &mut Self::Receivers,
32 ) -> impl Future<Output = Result<Envelope<Self>, RecvError>> + Send;
33
34 fn subscribe(bus: &EventBus) -> Self::Receivers;
35
36 fn try_recv(rx: &mut Self::Receivers) -> Result<Envelope<Self>, TryRecvError>;
39}
40
41impl<E: Event> EventSet for E {
42 type Receivers = TopicReceiver<E>;
43
44 fn publish(bus: &EventBus, meta: EventMeta, event: Self) {
45 bus.publish_stamped(meta, event);
46 }
47
48 fn recv(
49 rx: &mut Self::Receivers,
50 ) -> impl Future<Output = Result<Envelope<Self>, RecvError>> + Send {
51 rx.recv()
52 }
53
54 fn subscribe(bus: &EventBus) -> Self::Receivers {
55 TopicReceiver::new(bus)
56 }
57
58 fn try_recv(rx: &mut Self::Receivers) -> Result<Envelope<Self>, TryRecvError> {
59 rx.try_recv()
60 }
61}