Skip to main content

ironflow_engine/notify/
publisher.rs

1//! [`EventPublisher`] -- broadcasts events to filtered subscribers.
2
3use std::sync::Arc;
4
5use tokio::spawn;
6
7use super::{Event, EventSubscriber};
8
9/// A subscriber paired with its event type filter.
10struct Subscription {
11    subscriber: Arc<dyn EventSubscriber>,
12    event_types: Vec<&'static str>,
13}
14
15impl Subscription {
16    /// Returns `true` if this subscription accepts the given event.
17    fn accepts(&self, event: &Event) -> bool {
18        self.event_types.contains(&event.event_type())
19    }
20}
21
22/// Broadcasts [`Event`]s to registered [`EventSubscriber`]s.
23///
24/// Each subscriber is paired with an event type filter at subscription
25/// time. Only matching events are dispatched. Each call runs in a
26/// spawned task so that slow subscribers do not block the engine.
27///
28/// # Examples
29///
30/// ```no_run
31/// use ironflow_engine::notify::{EventPublisher, WebhookSubscriber, Event};
32///
33/// let mut publisher = EventPublisher::new();
34/// publisher.subscribe(
35///     WebhookSubscriber::new("https://hooks.example.com/events"),
36///     &[Event::RUN_STATUS_CHANGED, Event::STEP_FAILED],
37/// );
38/// ```
39pub struct EventPublisher {
40    subscriptions: Vec<Subscription>,
41}
42
43impl EventPublisher {
44    /// Create an empty publisher with no subscribers.
45    ///
46    /// # Examples
47    ///
48    /// ```
49    /// use ironflow_engine::notify::EventPublisher;
50    ///
51    /// let publisher = EventPublisher::new();
52    /// assert_eq!(publisher.subscriber_count(), 0);
53    /// ```
54    pub fn new() -> Self {
55        Self {
56            subscriptions: Vec::new(),
57        }
58    }
59
60    /// Register a subscriber with an event type filter.
61    ///
62    /// The subscriber is called only for events whose
63    /// [`event_type()`](Event::event_type) is in `event_types`.
64    /// Pass [`Event::ALL`] to receive every event.
65    ///
66    /// Use the `Event::*` constants for the filter values.
67    ///
68    /// # Examples
69    ///
70    /// ```no_run
71    /// use ironflow_engine::notify::{EventPublisher, WebhookSubscriber, Event};
72    ///
73    /// let mut publisher = EventPublisher::new();
74    ///
75    /// // Only on specific event types:
76    /// publisher.subscribe(
77    ///     WebhookSubscriber::new("https://example.com/hook"),
78    ///     &[Event::RUN_STATUS_CHANGED, Event::STEP_FAILED],
79    /// );
80    ///
81    /// // On all events:
82    /// publisher.subscribe(
83    ///     WebhookSubscriber::new("https://example.com/all"),
84    ///     Event::ALL,
85    /// );
86    /// ```
87    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    /// Number of registered subscribers.
99    pub fn subscriber_count(&self) -> usize {
100        self.subscriptions.len()
101    }
102
103    /// Broadcast an event to all matching subscribers.
104    ///
105    /// Each matching subscriber runs in its own spawned task. This
106    /// method returns immediately and never blocks.
107    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()); // matches
240        publisher.publish(sample_user_signed_in()); // filtered out
241
242        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}