stateset_embedded/events/
bus.rs1use stateset_core::CommerceEvent;
4use std::sync::atomic::{AtomicU64, Ordering};
5use std::task::{Context, Poll};
6use tokio::sync::broadcast;
7
8pub 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 #[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 #[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 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 pub fn subscribe(&self) -> EventSubscription {
69 EventSubscription { receiver: EventReceiver::new(self.sender.subscribe()) }
70 }
71
72 pub fn receiver_count(&self) -> usize {
74 self.sender.receiver_count()
75 }
76
77 pub fn events_published(&self) -> u64 {
79 self.events_published.load(Ordering::Relaxed)
80 }
81
82 pub fn events_publish_failures(&self) -> u64 {
84 self.events_publish_failures.load(Ordering::Relaxed)
85 }
86}
87
88pub 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 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 tracing::warn!(skipped, "Event receiver lagged, skipped events");
112 continue;
113 }
114 Err(broadcast::error::RecvError::Closed) => return None,
115 }
116 }
117 }
118
119 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
140pub 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 pub async fn recv(&mut self) -> Option<CommerceEvent> {
154 self.receiver.recv().await
155 }
156
157 pub fn try_recv(&mut self) -> Option<CommerceEvent> {
159 self.receiver.try_recv()
160 }
161}
162
163impl 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 let receivers = bus.publish(event);
198 assert_eq!(receivers, 2);
199
200 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 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}