Skip to main content

ironflow_engine/notify/
audit_log.rs

1//! [`AuditLogSubscriber`] -- persists every event to an [`AuditLogStore`].
2
3use std::str::FromStr;
4use std::sync::Arc;
5
6use serde_json::to_value;
7use tracing::error;
8use uuid::Uuid;
9
10use ironflow_store::audit_log_store::AuditLogStore;
11use ironflow_store::entities::{EventKind, NewAuditLogEntry};
12
13use super::{Event, EventSubscriber, SubscriberFuture};
14
15/// Subscriber that persists every received event as an audit log entry.
16///
17/// Extracts contextual IDs (run, step, user) from the event payload
18/// for efficient filtering, and stores the full event as JSON.
19///
20/// # Examples
21///
22/// ```no_run
23/// use std::sync::Arc;
24/// use ironflow_engine::notify::{AuditLogSubscriber, Event, EventPublisher};
25/// use ironflow_store::memory::InMemoryStore;
26///
27/// let store = Arc::new(InMemoryStore::new());
28/// let mut publisher = EventPublisher::new();
29/// publisher.subscribe(
30///     AuditLogSubscriber::new(store),
31///     Event::ALL,
32/// );
33/// ```
34pub struct AuditLogSubscriber {
35    store: Arc<dyn AuditLogStore>,
36}
37
38impl AuditLogSubscriber {
39    /// Create a new subscriber backed by the given store.
40    ///
41    /// # Examples
42    ///
43    /// ```
44    /// use std::sync::Arc;
45    /// use ironflow_engine::notify::AuditLogSubscriber;
46    /// use ironflow_store::memory::InMemoryStore;
47    ///
48    /// let store = Arc::new(InMemoryStore::new());
49    /// let subscriber = AuditLogSubscriber::new(store);
50    /// ```
51    pub fn new(store: Arc<dyn AuditLogStore>) -> Self {
52        Self { store }
53    }
54}
55
56fn extract_run_id(event: &Event) -> Option<Uuid> {
57    match event {
58        Event::RunCreated { run_id, .. }
59        | Event::RunStatusChanged { run_id, .. }
60        | Event::RunFailed { run_id, .. }
61        | Event::RunBudgetExceeded { run_id, .. }
62        | Event::StepCompleted { run_id, .. }
63        | Event::StepFailed { run_id, .. }
64        | Event::ApprovalRequested { run_id, .. }
65        | Event::ApprovalGranted { run_id, .. }
66        | Event::ApprovalRejected { run_id, .. }
67        | Event::LogLine { run_id, .. } => Some(*run_id),
68        Event::UserSignedIn { .. } | Event::UserSignedUp { .. } | Event::UserSignedOut { .. } => {
69            None
70        }
71    }
72}
73
74fn extract_step_id(event: &Event) -> Option<Uuid> {
75    match event {
76        Event::StepCompleted { step_id, .. }
77        | Event::StepFailed { step_id, .. }
78        | Event::ApprovalRequested { step_id, .. } => Some(*step_id),
79        _ => None,
80    }
81}
82
83fn extract_user_id(event: &Event) -> Option<Uuid> {
84    match event {
85        Event::UserSignedIn { user_id, .. }
86        | Event::UserSignedUp { user_id, .. }
87        | Event::UserSignedOut { user_id, .. } => Some(*user_id),
88        _ => None,
89    }
90}
91
92impl EventSubscriber for AuditLogSubscriber {
93    fn name(&self) -> &str {
94        "audit_log"
95    }
96
97    fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
98        Box::pin(async move {
99            let event_kind = match EventKind::from_str(event.event_type()) {
100                Ok(k) => k,
101                Err(e) => {
102                    error!(error = %e, event_type = event.event_type(), "unknown event kind for audit log");
103                    return;
104                }
105            };
106
107            let payload = match to_value(event) {
108                Ok(v) => v,
109                Err(e) => {
110                    error!(error = %e, event_type = event.event_type(), "failed to serialize event for audit log");
111                    return;
112                }
113            };
114
115            let entry = NewAuditLogEntry {
116                event_type: event_kind,
117                payload,
118                run_id: extract_run_id(event),
119                step_id: extract_step_id(event),
120                user_id: extract_user_id(event),
121            };
122
123            if let Err(e) = self.store.append_audit_log(entry).await {
124                error!(error = %e, event_type = event.event_type(), "failed to persist audit log entry");
125            }
126        })
127    }
128}
129
130#[cfg(test)]
131mod tests {
132    use std::sync::Arc;
133    use std::time::Duration;
134
135    use chrono::Utc;
136    use rust_decimal::Decimal;
137    use uuid::Uuid;
138
139    use ironflow_store::audit_log_store::AuditLogStore;
140    use ironflow_store::entities::{AuditLogFilter, EventKind};
141    use ironflow_store::memory::InMemoryStore;
142    use ironflow_store::models::RunStatus;
143
144    use super::*;
145    use crate::notify::{EventPublisher, EventSubscriber};
146
147    fn sample_run_status_changed() -> Event {
148        Event::RunStatusChanged {
149            run_id: Uuid::now_v7(),
150            workflow_name: "deploy".to_string(),
151            from: RunStatus::Running,
152            to: RunStatus::Completed,
153            error: None,
154            cost_usd: Decimal::new(42, 2),
155            duration_ms: 5000,
156            at: Utc::now(),
157        }
158    }
159
160    fn sample_user_signed_in() -> Event {
161        Event::UserSignedIn {
162            user_id: Uuid::now_v7(),
163            username: "alice".to_string(),
164            at: Utc::now(),
165        }
166    }
167
168    fn sample_step_failed() -> Event {
169        Event::StepFailed {
170            run_id: Uuid::now_v7(),
171            step_id: Uuid::now_v7(),
172            step_name: "build".to_string(),
173            kind: ironflow_store::models::StepKind::Shell,
174            error: "exit code 1".to_string(),
175            at: Utc::now(),
176        }
177    }
178
179    #[test]
180    fn name_is_audit_log() {
181        let store = Arc::new(InMemoryStore::new());
182        let subscriber = AuditLogSubscriber::new(store);
183        assert_eq!(subscriber.name(), "audit_log");
184    }
185
186    #[test]
187    fn extract_run_id_from_run_event() {
188        let event = sample_run_status_changed();
189        assert!(extract_run_id(&event).is_some());
190    }
191
192    #[test]
193    fn extract_run_id_from_user_event_is_none() {
194        let event = sample_user_signed_in();
195        assert!(extract_run_id(&event).is_none());
196    }
197
198    #[test]
199    fn extract_step_id_from_step_event() {
200        let event = sample_step_failed();
201        assert!(extract_step_id(&event).is_some());
202    }
203
204    #[test]
205    fn extract_step_id_from_run_event_is_none() {
206        let event = sample_run_status_changed();
207        assert!(extract_step_id(&event).is_none());
208    }
209
210    #[test]
211    fn extract_user_id_from_user_event() {
212        let event = sample_user_signed_in();
213        assert!(extract_user_id(&event).is_some());
214    }
215
216    #[test]
217    fn extract_user_id_from_run_event_is_none() {
218        let event = sample_run_status_changed();
219        assert!(extract_user_id(&event).is_none());
220    }
221
222    #[tokio::test]
223    async fn handle_persists_event() {
224        let store = Arc::new(InMemoryStore::new());
225        let subscriber = AuditLogSubscriber::new(store.clone());
226
227        let event = sample_run_status_changed();
228        subscriber.handle(&event).await;
229
230        let page = store
231            .list_audit_logs(AuditLogFilter::default(), 1, 20)
232            .await
233            .unwrap();
234
235        assert_eq!(page.items.len(), 1);
236        assert_eq!(page.items[0].event_type, EventKind::RunStatusChanged);
237        assert!(page.items[0].run_id.is_some());
238        assert!(page.items[0].step_id.is_none());
239        assert!(page.items[0].user_id.is_none());
240    }
241
242    #[tokio::test]
243    async fn handle_persists_step_event_with_ids() {
244        let store = Arc::new(InMemoryStore::new());
245        let subscriber = AuditLogSubscriber::new(store.clone());
246
247        let event = sample_step_failed();
248        subscriber.handle(&event).await;
249
250        let page = store
251            .list_audit_logs(AuditLogFilter::default(), 1, 20)
252            .await
253            .unwrap();
254
255        assert_eq!(page.items.len(), 1);
256        assert_eq!(page.items[0].event_type, EventKind::StepFailed);
257        assert!(page.items[0].run_id.is_some());
258        assert!(page.items[0].step_id.is_some());
259    }
260
261    #[tokio::test]
262    async fn handle_persists_user_event_with_user_id() {
263        let store = Arc::new(InMemoryStore::new());
264        let subscriber = AuditLogSubscriber::new(store.clone());
265
266        let event = sample_user_signed_in();
267        subscriber.handle(&event).await;
268
269        let page = store
270            .list_audit_logs(AuditLogFilter::default(), 1, 20)
271            .await
272            .unwrap();
273
274        assert_eq!(page.items.len(), 1);
275        assert_eq!(page.items[0].event_type, EventKind::UserSignedIn);
276        assert!(page.items[0].user_id.is_some());
277        assert!(page.items[0].run_id.is_none());
278    }
279
280    #[tokio::test]
281    async fn publisher_dispatches_to_audit_log_subscriber() {
282        let store = Arc::new(InMemoryStore::new());
283        let mut publisher = EventPublisher::new();
284
285        publisher.subscribe(AuditLogSubscriber::new(store.clone()), Event::ALL);
286
287        publisher.publish(sample_run_status_changed());
288        publisher.publish(sample_user_signed_in());
289        publisher.publish(sample_step_failed());
290
291        tokio::time::sleep(Duration::from_millis(100)).await;
292
293        let page = store
294            .list_audit_logs(AuditLogFilter::default(), 1, 20)
295            .await
296            .unwrap();
297
298        assert_eq!(page.items.len(), 3);
299    }
300
301    #[tokio::test]
302    async fn full_event_payload_is_preserved() {
303        let store = Arc::new(InMemoryStore::new());
304        let subscriber = AuditLogSubscriber::new(store.clone());
305
306        let run_id = Uuid::now_v7();
307        let event = Event::RunFailed {
308            run_id,
309            workflow_name: "deploy".to_string(),
310            error: Some("step crashed".to_string()),
311            cost_usd: Decimal::new(10, 2),
312            duration_ms: 3000,
313            at: Utc::now(),
314        };
315        subscriber.handle(&event).await;
316
317        let page = store
318            .list_audit_logs(AuditLogFilter::default(), 1, 20)
319            .await
320            .unwrap();
321
322        let payload = &page.items[0].payload;
323        assert_eq!(payload["type"], "run_failed");
324        assert_eq!(payload["workflow_name"], "deploy");
325        assert_eq!(payload["error"], "step crashed");
326    }
327}