Skip to main content

ironflow_engine/notify/
subscriber.rs

1//! [`EventSubscriber`] trait -- react to domain events.
2
3use std::future::Future;
4use std::pin::Pin;
5
6use super::Event;
7
8/// Boxed future returned by [`EventSubscriber::handle`].
9pub type SubscriberFuture<'a> = Pin<Box<dyn Future<Output = ()> + Send + 'a>>;
10
11/// A subscriber that reacts to domain events.
12///
13/// Implement this trait to create custom notification channels (Slack,
14/// Discord, PagerDuty, etc.). The engine broadcasts events to all
15/// registered subscribers via [`EventPublisher`](super::EventPublisher).
16///
17/// # Contract
18///
19/// - [`handle`](EventSubscriber::handle) is called only for events that
20///   match the filter configured at subscription time.
21/// - Implementations must not block -- heavy work should be spawned.
22/// - Errors are logged internally; they must not propagate.
23///
24/// # Examples
25///
26/// ```no_run
27/// use ironflow_engine::notify::{EventSubscriber, Event, SubscriberFuture};
28///
29/// struct LogSubscriber;
30///
31/// impl EventSubscriber for LogSubscriber {
32///     fn name(&self) -> &str { "log" }
33///
34///     fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
35///         Box::pin(async move {
36///             println!("[{}] {:?}", event.event_type(), event);
37///         })
38///     }
39/// }
40/// ```
41pub trait EventSubscriber: Send + Sync {
42    /// A short identifier for this subscriber (used in logs).
43    fn name(&self) -> &str;
44
45    /// Handle a domain event.
46    ///
47    /// Only called for events matching the filter set at subscription
48    /// time. The subscriber does not need to filter.
49    fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a>;
50}
51
52#[cfg(test)]
53mod tests {
54    use super::*;
55    use crate::notify::RunCreatedEvent;
56
57    struct TestSubscriber {
58        name: String,
59    }
60
61    impl EventSubscriber for TestSubscriber {
62        fn name(&self) -> &str {
63            &self.name
64        }
65
66        fn handle<'a>(&'a self, _event: &'a Event) -> SubscriberFuture<'a> {
67            Box::pin(async move {
68                // No-op for testing
69            })
70        }
71    }
72
73    struct CountingSubscriber {
74        name: String,
75    }
76
77    impl EventSubscriber for CountingSubscriber {
78        fn name(&self) -> &str {
79            &self.name
80        }
81
82        fn handle<'a>(&'a self, _event: &'a Event) -> SubscriberFuture<'a> {
83            Box::pin(async move {
84                // Simulates async work
85                tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
86            })
87        }
88    }
89
90    #[test]
91    fn subscriber_has_identifier_name() {
92        let sub = TestSubscriber {
93            name: "test_sub".to_string(),
94        };
95        assert_eq!(sub.name(), "test_sub");
96    }
97
98    #[test]
99    fn subscriber_name_is_consistent() {
100        let sub = TestSubscriber {
101            name: "my_subscriber".to_string(),
102        };
103        assert_eq!(sub.name(), "my_subscriber");
104        assert_eq!(sub.name(), "my_subscriber");
105    }
106
107    #[test]
108    fn different_subscribers_have_different_names() {
109        let sub1 = TestSubscriber {
110            name: "sub1".to_string(),
111        };
112        let sub2 = TestSubscriber {
113            name: "sub2".to_string(),
114        };
115
116        assert_ne!(sub1.name(), sub2.name());
117    }
118
119    #[tokio::test]
120    async fn subscriber_handle_completes_successfully() {
121        use chrono::Utc;
122        let sub = TestSubscriber {
123            name: "test".to_string(),
124        };
125
126        // Create a dummy event for testing
127        let event = Event::RunCreated(RunCreatedEvent {
128            run_id: uuid::Uuid::now_v7(),
129            workflow_name: "test-wf".to_string(),
130            at: Utc::now(),
131        });
132
133        // Should complete without error
134        sub.handle(&event).await;
135    }
136
137    #[tokio::test]
138    async fn subscriber_handle_is_async() {
139        use chrono::Utc;
140        let sub = CountingSubscriber {
141            name: "async_test".to_string(),
142        };
143
144        let event = Event::RunCreated(RunCreatedEvent {
145            run_id: uuid::Uuid::now_v7(),
146            workflow_name: "test".to_string(),
147            at: Utc::now(),
148        });
149
150        let start = std::time::Instant::now();
151        sub.handle(&event).await;
152        let elapsed = start.elapsed();
153
154        // Should have taken at least 1ms due to the sleep
155        assert!(elapsed.as_millis() >= 1);
156    }
157
158    #[tokio::test]
159    async fn multiple_subscribers_can_handle_same_event() {
160        use chrono::Utc;
161        let sub1 = TestSubscriber {
162            name: "sub1".to_string(),
163        };
164        let sub2 = TestSubscriber {
165            name: "sub2".to_string(),
166        };
167
168        let event = Event::RunCreated(RunCreatedEvent {
169            run_id: uuid::Uuid::now_v7(),
170            workflow_name: "test".to_string(),
171            at: Utc::now(),
172        });
173
174        // Both should handle without issue
175        sub1.handle(&event).await;
176        sub2.handle(&event).await;
177    }
178
179    #[test]
180    fn subscriber_implements_send_sync() {
181        fn assert_send_sync<T: Send + Sync>() {}
182        assert_send_sync::<TestSubscriber>();
183        assert_send_sync::<CountingSubscriber>();
184    }
185
186    #[tokio::test]
187    async fn subscriber_future_is_boxed() {
188        use chrono::Utc;
189        let sub = TestSubscriber {
190            name: "boxed_test".to_string(),
191        };
192
193        let event = Event::RunCreated(RunCreatedEvent {
194            run_id: uuid::Uuid::now_v7(),
195            workflow_name: "test".to_string(),
196            at: Utc::now(),
197        });
198
199        let future = sub.handle(&event);
200        // The future should be a Pin<Box<_>> and awaitable
201        let _ = future.await;
202    }
203
204    #[test]
205    fn subscriber_name_borrowed_lifetime() {
206        let sub = TestSubscriber {
207            name: "lifetime_test".to_string(),
208        };
209
210        let name1 = sub.name();
211        let name2 = sub.name();
212
213        // Both should be valid references
214        assert_eq!(name1, name2);
215        assert_eq!(name1, "lifetime_test");
216    }
217}