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, .. } => 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}