Skip to main content

clawless_core/event/
channel.rs

1use std::fmt;
2
3use tokio::sync::mpsc;
4
5use super::Event;
6use super::receiver::EventReceiver;
7use super::sender::EventSender;
8
9/// Error returned when sending an event fails
10///
11/// A send fails when the [`EventReceiver`] has been dropped, meaning the consumer is no longer
12/// listening. The error carries the unsent [`Event`] so callers can log or inspect what was lost.
13// r[impl event.transport.error]
14#[derive(Debug)]
15pub struct SendError(pub Event);
16
17impl fmt::Display for SendError {
18    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
19        f.write_str("event channel closed")
20    }
21}
22
23impl std::error::Error for SendError {}
24
25/// Creates a bounded event channel
26///
27/// Returns a paired sender and receiver. The sender is clonable; the receiver is not. When all
28/// senders are dropped, the receiver's [`recv`] method returns `None`.
29///
30/// The channel is bounded with a capacity of 256 events. If the consumer falls behind, producers
31/// will await until space is available, providing natural back-pressure.
32///
33/// [`recv`]: EventReceiver::recv
34// r[impl event.transport.async]
35// r[impl event.transport.multi-producer]
36// r[impl event.transport.single-consumer]
37// r[impl event.transport.ordered]
38// r[impl event.transport.backpressure]
39// r[impl event.transport.completion]
40// r[impl event.transport.drain]
41pub fn event_channel() -> (EventSender, EventReceiver) {
42    let (tx, rx) = mpsc::channel(256);
43    (EventSender::new(tx), EventReceiver::new(rx))
44}
45
46#[cfg(test)]
47mod tests {
48    // An assertion in a test panics by design. A `# Panics` section on every test
49    // would repeat that and give the reader no information.
50    #![allow(clippy::missing_panics_doc)]
51
52    use super::*;
53
54    // r[verify event.transport.drain]
55    // r[verify event.transport.completion]
56    #[tokio::test]
57    async fn event_channel_recv_drains_buffered_events_after_sender_dropped() {
58        let (sender, mut receiver) = event_channel();
59
60        sender
61            .send(Event::Message("first".to_string()))
62            .await
63            .expect("should send");
64        sender
65            .send(Event::Message("second".to_string()))
66            .await
67            .expect("should send");
68        drop(sender);
69
70        let first = receiver.recv().await.expect("should receive first");
71        let second = receiver.recv().await.expect("should receive second");
72        let done = receiver.recv().await;
73
74        assert!(matches!(first, Event::Message(ref s) if s == "first"));
75        assert!(matches!(second, Event::Message(ref s) if s == "second"));
76        assert!(done.is_none());
77    }
78
79    // r[verify event.transport.completion]
80    #[tokio::test]
81    async fn event_channel_recv_returns_none_when_all_senders_dropped() {
82        let (sender, mut receiver) = event_channel();
83
84        drop(sender);
85
86        let result = receiver.recv().await;
87
88        assert!(result.is_none());
89    }
90
91    // r[verify event.transport.error]
92    #[tokio::test]
93    async fn event_channel_send_after_receiver_dropped_returns_error() {
94        let (sender, receiver) = event_channel();
95
96        drop(receiver);
97
98        let error = sender
99            .send(Event::Message("hello".to_string()))
100            .await
101            .expect_err("should fail");
102
103        assert!(matches!(error.0, Event::Message(ref s) if s == "hello"));
104    }
105
106    // r[verify event.transport.async]
107    // r[verify event.transport.ordered]
108    #[tokio::test]
109    async fn event_channel_send_and_recv_delivers_event() {
110        let (sender, mut receiver) = event_channel();
111
112        sender
113            .send(Event::Message("hello".to_string()))
114            .await
115            .expect("should send");
116
117        let event = receiver.recv().await.expect("should receive");
118
119        assert!(matches!(event, Event::Message(ref s) if s == "hello"));
120    }
121
122    #[test]
123    fn send_error_display_shows_closed_message() {
124        let error = SendError(Event::Message("lost".to_string()));
125
126        let message = error.to_string();
127
128        assert_eq!(message, "event channel closed");
129    }
130
131    #[test]
132    fn send_error_is_std_error() {
133        fn assert_error<T: std::error::Error>() {}
134        assert_error::<SendError>();
135    }
136
137    #[test]
138    fn trait_send_error_send() {
139        fn assert_send<T: Send>() {}
140        assert_send::<SendError>();
141    }
142
143    #[test]
144    fn trait_send_error_sync() {
145        fn assert_sync<T: Sync>() {}
146        assert_sync::<SendError>();
147    }
148
149    #[test]
150    fn trait_send_error_unpin() {
151        fn assert_unpin<T: Unpin>() {}
152        assert_unpin::<SendError>();
153    }
154}