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