ironflow_engine/notify/
publisher.rs1use std::sync::Arc;
4
5use tokio::spawn;
6
7use super::{Event, EventSubscriber};
8
9struct Subscription {
11 subscriber: Arc<dyn EventSubscriber>,
12 event_types: Vec<&'static str>,
13}
14
15impl Subscription {
16 fn accepts(&self, event: &Event) -> bool {
18 self.event_types.contains(&event.event_type())
19 }
20}
21
22pub struct EventPublisher {
40 subscriptions: Vec<Subscription>,
41}
42
43impl EventPublisher {
44 pub fn new() -> Self {
55 Self {
56 subscriptions: Vec::new(),
57 }
58 }
59
60 pub fn subscribe(
88 &mut self,
89 subscriber: impl EventSubscriber + 'static,
90 event_types: &[&'static str],
91 ) {
92 self.subscriptions.push(Subscription {
93 subscriber: Arc::new(subscriber),
94 event_types: event_types.to_vec(),
95 });
96 }
97
98 pub fn subscriber_count(&self) -> usize {
100 self.subscriptions.len()
101 }
102
103 pub fn publish(&self, event: Event) {
108 for subscription in &self.subscriptions {
109 if !subscription.accepts(&event) {
110 continue;
111 }
112 let subscriber = subscription.subscriber.clone();
113 let event = event.clone();
114 spawn(async move {
115 subscriber.handle(&event).await;
116 });
117 }
118 }
119}
120
121impl Default for EventPublisher {
122 fn default() -> Self {
123 Self::new()
124 }
125}
126
127#[cfg(test)]
128mod tests {
129 use std::collections::HashMap;
130 use std::sync::atomic::{AtomicU32, Ordering};
131 use std::time::Duration;
132
133 use super::*;
134 use crate::notify::{
135 RunStatusChangedEvent, SubscriberFuture, UserSignedInEvent, WebhookSubscriber,
136 };
137 use rust_decimal::Decimal;
138 use tokio::time::sleep;
139
140 use chrono::Utc;
141 use ironflow_store::models::RunStatus;
142 use uuid::Uuid;
143
144 fn sample_run_status_changed() -> Event {
145 Event::RunStatusChanged(RunStatusChangedEvent {
146 run_id: Uuid::now_v7(),
147 workflow_name: "deploy".to_string(),
148 from: RunStatus::Running,
149 to: RunStatus::Completed,
150 error: None,
151 cost_usd: Decimal::new(42, 2),
152 duration_ms: 5000,
153 labels: HashMap::new(),
154 at: Utc::now(),
155 })
156 }
157
158 fn sample_user_signed_in() -> Event {
159 Event::UserSignedIn(UserSignedInEvent {
160 user_id: Uuid::now_v7(),
161 username: "alice".to_string(),
162 at: Utc::now(),
163 })
164 }
165
166 #[test]
167 fn starts_empty() {
168 let publisher = EventPublisher::new();
169 assert_eq!(publisher.subscriber_count(), 0);
170 }
171
172 #[test]
173 fn subscribe_increments_count() {
174 let mut publisher = EventPublisher::new();
175 publisher.subscribe(
176 WebhookSubscriber::new("https://example.com"),
177 &[Event::RUN_STATUS_CHANGED],
178 );
179 assert_eq!(publisher.subscriber_count(), 1);
180 }
181
182 #[test]
183 fn publish_with_no_subscribers_is_noop() {
184 let publisher = EventPublisher::new();
185 publisher.publish(sample_run_status_changed());
186 }
187
188 #[test]
189 fn default_is_empty() {
190 let publisher = EventPublisher::default();
191 assert_eq!(publisher.subscriber_count(), 0);
192 }
193
194 struct CountingSubscriber {
195 count: AtomicU32,
196 }
197
198 impl CountingSubscriber {
199 fn new() -> Self {
200 Self {
201 count: AtomicU32::new(0),
202 }
203 }
204
205 fn count(&self) -> u32 {
206 self.count.load(Ordering::SeqCst)
207 }
208 }
209
210 impl EventSubscriber for CountingSubscriber {
211 fn name(&self) -> &str {
212 "counting"
213 }
214
215 fn handle<'a>(&'a self, _event: &'a Event) -> SubscriberFuture<'a> {
216 Box::pin(async move {
217 self.count.fetch_add(1, Ordering::SeqCst);
218 })
219 }
220 }
221
222 #[tokio::test]
223 async fn subscriber_receives_matching_events() {
224 let subscriber = Arc::new(CountingSubscriber::new());
225 let mut publisher = EventPublisher::new();
226
227 struct ArcSub(Arc<CountingSubscriber>);
228 impl EventSubscriber for ArcSub {
229 fn name(&self) -> &str {
230 self.0.name()
231 }
232 fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
233 self.0.handle(event)
234 }
235 }
236
237 publisher.subscribe(ArcSub(subscriber.clone()), &[Event::RUN_STATUS_CHANGED]);
238
239 publisher.publish(sample_run_status_changed()); publisher.publish(sample_user_signed_in()); sleep(Duration::from_millis(50)).await;
243
244 assert_eq!(subscriber.count(), 1);
245 }
246
247 #[tokio::test]
248 async fn all_filter_matches_everything() {
249 let subscriber = Arc::new(CountingSubscriber::new());
250 let mut publisher = EventPublisher::new();
251
252 struct ArcSub(Arc<CountingSubscriber>);
253 impl EventSubscriber for ArcSub {
254 fn name(&self) -> &str {
255 self.0.name()
256 }
257 fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
258 self.0.handle(event)
259 }
260 }
261
262 publisher.subscribe(ArcSub(subscriber.clone()), Event::ALL);
263
264 publisher.publish(sample_run_status_changed());
265 publisher.publish(sample_user_signed_in());
266
267 sleep(Duration::from_millis(50)).await;
268
269 assert_eq!(subscriber.count(), 2);
270 }
271
272 #[tokio::test]
273 async fn empty_filter_matches_nothing() {
274 let subscriber = Arc::new(CountingSubscriber::new());
275 let mut publisher = EventPublisher::new();
276
277 struct ArcSub(Arc<CountingSubscriber>);
278 impl EventSubscriber for ArcSub {
279 fn name(&self) -> &str {
280 self.0.name()
281 }
282 fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
283 self.0.handle(event)
284 }
285 }
286
287 publisher.subscribe(ArcSub(subscriber.clone()), &[]);
288
289 publisher.publish(sample_run_status_changed());
290 publisher.publish(sample_user_signed_in());
291
292 sleep(Duration::from_millis(50)).await;
293
294 assert_eq!(subscriber.count(), 0);
295 }
296}