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