Skip to main content

kithara_events/
event.rs

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
9/// A value that can travel on its own channel.
10///
11/// Implemented only through `#[derive(Event)]` outside this crate; the
12/// `derivable_event` idiom check denies a hand-written impl anywhere else.
13pub trait Event: Clone + Debug + Send + Sync + 'static {}
14
15/// One or more [`Event`] types a consumer wants on a single receiver.
16///
17/// Every `Event` is a one-member set through the blanket impl below; a
18/// multi-member set is a consumer-local enum with `#[derive(EventSet)]`.
19///
20/// A multi-member receiver polls its members in declaration order and yields
21/// the first that has an event queued, so it reports no order between
22/// members: an earlier-declared member preempts one that published first.
23/// Publication order holds inside a member, which is where a consumer reads
24/// it from — a one-member receiver.
25pub 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    /// # Errors
37    /// Returns `Empty`, lag information, or `Closed` when all members close.
38    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}