ironflow_engine/notify/
audit_log.rs1use 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
14pub struct AuditLogSubscriber {
34 store: Arc<dyn AuditLogStore>,
35}
36
37impl AuditLogSubscriber {
38 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}