Skip to main content

stateset_embedded/events/
bus.rs

1//! Event bus for in-process pub/sub using tokio broadcast channels
2
3use stateset_core::CommerceEvent;
4use std::sync::atomic::{AtomicU64, Ordering};
5use std::task::{Context, Poll};
6use tokio::sync::broadcast;
7
8/// Event bus for broadcasting events to multiple subscribers
9pub struct EventBus {
10    sender: broadcast::Sender<CommerceEvent>,
11    events_published: AtomicU64,
12    events_publish_failures: AtomicU64,
13}
14
15impl std::fmt::Debug for EventBus {
16    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
17        f.debug_struct("EventBus")
18            .field("events_published", &self.events_published.load(Ordering::Relaxed))
19            .field("receiver_count", &self.sender.receiver_count())
20            .finish_non_exhaustive()
21    }
22}
23
24impl EventBus {
25    /// Create a new event bus with the specified channel capacity
26    #[must_use]
27    pub fn new(capacity: usize) -> Self {
28        let (sender, _) = broadcast::channel(capacity);
29        Self {
30            sender,
31            events_published: AtomicU64::new(0),
32            events_publish_failures: AtomicU64::new(0),
33        }
34    }
35
36    /// Publish an event to all subscribers
37    #[inline]
38    pub fn publish(&self, event: CommerceEvent) -> usize {
39        self.events_published.fetch_add(1, Ordering::Relaxed);
40        match self.sender.send(event) {
41            Ok(receivers) => receivers,
42            Err(error) => {
43                self.events_publish_failures.fetch_add(1, Ordering::Relaxed);
44                // Only allocate event_type string in the error path
45                let event_type = error.0.event_type().to_string();
46                let receiver_count = self.sender.receiver_count();
47                if receiver_count == 0 {
48                    tracing::debug!(
49                        event_type,
50                        error = %error,
51                        receiver_count,
52                        "Dropped event publish: no active subscribers"
53                    );
54                } else {
55                    tracing::warn!(
56                        event_type,
57                        error = %error,
58                        receiver_count,
59                        "Failed to publish event to in-process subscribers"
60                    );
61                }
62                0
63            }
64        }
65    }
66
67    /// Subscribe to events from this bus
68    pub fn subscribe(&self) -> EventSubscription {
69        EventSubscription { receiver: EventReceiver::new(self.sender.subscribe()) }
70    }
71
72    /// Get the number of active receivers
73    pub fn receiver_count(&self) -> usize {
74        self.sender.receiver_count()
75    }
76
77    /// Get total number of events published
78    pub fn events_published(&self) -> u64 {
79        self.events_published.load(Ordering::Relaxed)
80    }
81
82    /// Get the number of events that failed to publish to the in-process bus
83    pub fn events_publish_failures(&self) -> u64 {
84        self.events_publish_failures.load(Ordering::Relaxed)
85    }
86}
87
88/// Wrapper around broadcast receiver with convenience methods
89pub struct EventReceiver {
90    inner: broadcast::Receiver<CommerceEvent>,
91}
92
93impl std::fmt::Debug for EventReceiver {
94    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
95        f.debug_struct("EventReceiver").finish_non_exhaustive()
96    }
97}
98
99impl EventReceiver {
100    const fn new(receiver: broadcast::Receiver<CommerceEvent>) -> Self {
101        Self { inner: receiver }
102    }
103
104    /// Receive the next event, waiting if necessary
105    pub async fn recv(&mut self) -> Option<CommerceEvent> {
106        loop {
107            match self.inner.recv().await {
108                Ok(event) => return Some(event),
109                Err(broadcast::error::RecvError::Lagged(skipped)) => {
110                    // Log that we skipped some events due to slow consumer
111                    tracing::warn!(skipped, "Event receiver lagged, skipped events");
112                    continue;
113                }
114                Err(broadcast::error::RecvError::Closed) => return None,
115            }
116        }
117    }
118
119    /// Try to receive an event without waiting
120    pub fn try_recv(&mut self) -> Option<CommerceEvent> {
121        self.inner.try_recv().ok()
122    }
123
124    fn poll_recv(&mut self, cx: &mut Context<'_>) -> Poll<Option<CommerceEvent>> {
125        loop {
126            let mut recv = std::pin::pin!(self.inner.recv());
127            match std::future::Future::poll(recv.as_mut(), cx) {
128                Poll::Ready(Ok(event)) => return Poll::Ready(Some(event)),
129                Poll::Ready(Err(broadcast::error::RecvError::Lagged(skipped))) => {
130                    tracing::warn!(skipped, "Event receiver lagged, skipped events");
131                    continue;
132                }
133                Poll::Ready(Err(broadcast::error::RecvError::Closed)) => return Poll::Ready(None),
134                Poll::Pending => return Poll::Pending,
135            }
136        }
137    }
138}
139
140/// An event subscription that can be used to receive events
141pub struct EventSubscription {
142    receiver: EventReceiver,
143}
144
145impl std::fmt::Debug for EventSubscription {
146    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
147        f.debug_struct("EventSubscription").finish_non_exhaustive()
148    }
149}
150
151impl EventSubscription {
152    /// Receive the next event
153    pub async fn recv(&mut self) -> Option<CommerceEvent> {
154        self.receiver.recv().await
155    }
156
157    /// Try to receive without waiting
158    pub fn try_recv(&mut self) -> Option<CommerceEvent> {
159        self.receiver.try_recv()
160    }
161}
162
163// Implement Stream trait for use with StreamExt
164impl futures::Stream for EventSubscription {
165    type Item = CommerceEvent;
166
167    fn poll_next(
168        mut self: std::pin::Pin<&mut Self>,
169        cx: &mut Context<'_>,
170    ) -> std::task::Poll<Option<Self::Item>> {
171        self.receiver.poll_recv(cx)
172    }
173}
174
175#[cfg(test)]
176mod tests {
177    use super::*;
178    use chrono::Utc;
179    use rust_decimal_macros::dec;
180
181    #[tokio::test]
182    async fn test_event_bus_publish_subscribe() {
183        let bus = EventBus::new(16);
184
185        let mut sub1 = bus.subscribe();
186        let mut sub2 = bus.subscribe();
187
188        let event = CommerceEvent::OrderCreated {
189            order_id: stateset_core::OrderId::new(),
190            customer_id: stateset_core::CustomerId::new(),
191            total_amount: dec!(100.00),
192            item_count: 2,
193            timestamp: Utc::now(),
194        };
195
196        // Publish should reach both subscribers
197        let receivers = bus.publish(event);
198        assert_eq!(receivers, 2);
199
200        // Both subscribers should receive the event
201        let received1 = sub1.try_recv();
202        let received2 = sub2.try_recv();
203
204        assert!(received1.is_some());
205        assert!(received2.is_some());
206    }
207
208    #[tokio::test]
209    async fn test_event_bus_no_subscribers() {
210        let bus = EventBus::new(16);
211
212        let event = CommerceEvent::CustomerCreated {
213            customer_id: stateset_core::CustomerId::new(),
214            email: "test@example.com".to_string(),
215            timestamp: Utc::now(),
216        };
217
218        // Should not panic even with no subscribers
219        let receivers = bus.publish(event);
220        assert_eq!(receivers, 0);
221    }
222
223    #[test]
224    fn test_receiver_count() {
225        let bus = EventBus::new(16);
226        assert_eq!(bus.receiver_count(), 0);
227
228        let _sub1 = bus.subscribe();
229        assert_eq!(bus.receiver_count(), 1);
230
231        let _sub2 = bus.subscribe();
232        assert_eq!(bus.receiver_count(), 2);
233    }
234
235    #[test]
236    fn test_event_bus_publish_failure_tracking() {
237        let bus = EventBus::new(16);
238
239        let event = CommerceEvent::CustomerCreated {
240            customer_id: stateset_core::CustomerId::new(),
241            email: "test@example.com".to_string(),
242            timestamp: Utc::now(),
243        };
244
245        let receivers = bus.publish(event);
246        assert_eq!(receivers, 0);
247        assert_eq!(bus.events_published(), 1);
248        assert_eq!(bus.events_publish_failures(), 1);
249    }
250
251    #[tokio::test]
252    async fn test_event_subscription_stream_next_wakes_on_publish() {
253        use futures::StreamExt;
254
255        let bus = EventBus::new(16);
256        let mut sub = bus.subscribe();
257        let event = CommerceEvent::OrderCreated {
258            order_id: stateset_core::OrderId::new(),
259            customer_id: stateset_core::CustomerId::new(),
260            total_amount: dec!(100.00),
261            item_count: 2,
262            timestamp: Utc::now(),
263        };
264
265        let bus_for_publish = bus;
266        let publish_task = tokio::spawn(async move {
267            tokio::time::sleep(std::time::Duration::from_millis(25)).await;
268            bus_for_publish.publish(event.clone());
269        });
270
271        let got = tokio::time::timeout(std::time::Duration::from_millis(500), sub.next())
272            .await
273            .expect("timed out while waiting for published event");
274
275        publish_task.await.expect("publisher task failed");
276        assert!(got.is_some());
277    }
278}