1use 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
15pub struct AuditLogSubscriber {
35 store: Arc<dyn AuditLogStore>,
36}
37
38impl AuditLogSubscriber {
39 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}