ironflow_runtime/trigger/
event.rs1use serde_json::json;
31use tokio::sync::mpsc;
32use tokio_util::sync::CancellationToken;
33use tracing::{info, warn};
34use uuid::Uuid;
35
36use ironflow_engine::notify::{Event, EventSubscriber, SubscriberFuture};
37use ironflow_store::entities::{EventKind, TriggerKind};
38
39use super::{Trigger, TriggerEvent, TriggerFuture, TriggerSink};
40
41pub const CHAIN_DEPTH_LABEL: &str = "_chain_depth";
43
44#[derive(Debug, Clone)]
61pub struct EventTriggerRule {
62 pub on_event: EventKind,
64 pub source_workflow: String,
66 pub target_workflow: String,
68 pub max_chain_depth: u8,
71}
72
73pub struct EventTrigger {
97 rules: Vec<EventTriggerRule>,
98 event_tx: mpsc::Sender<InternalEvent>,
100 event_rx: tokio::sync::Mutex<mpsc::Receiver<InternalEvent>>,
101}
102
103#[derive(Debug)]
105struct InternalEvent {
106 run_id: Uuid,
107 workflow_name: String,
108 event_kind: EventKind,
109 error: Option<String>,
110}
111
112impl EventTrigger {
113 pub fn new(rules: Vec<EventTriggerRule>) -> Self {
131 let (event_tx, event_rx) = mpsc::channel(256);
132 Self {
133 rules,
134 event_tx,
135 event_rx: tokio::sync::Mutex::new(event_rx),
136 }
137 }
138
139 pub fn subscribed_event_types(&self) -> Vec<&'static str> {
159 self.rules.iter().map(|r| r.on_event.as_str()).collect()
160 }
161
162 fn matching_rules(&self, event_kind: EventKind, workflow_name: &str) -> Vec<&EventTriggerRule> {
164 self.rules
165 .iter()
166 .filter(|r| r.on_event == event_kind && r.source_workflow == workflow_name)
167 .collect()
168 }
169
170 fn build_payload(
172 source_run_id: Uuid,
173 source_workflow: &str,
174 error: &Option<String>,
175 ) -> serde_json::Value {
176 json!({
177 "source_run_id": source_run_id,
178 "source_workflow": source_workflow,
179 "error": error,
180 })
181 }
182
183 fn chain_depth_from_event(_event_kind: &EventKind) -> u8 {
185 0
186 }
187}
188
189impl Trigger for EventTrigger {
190 fn name(&self) -> &str {
191 "event-trigger"
192 }
193
194 fn start<'a>(&'a self, sink: TriggerSink, token: &'a CancellationToken) -> TriggerFuture<'a> {
195 Box::pin(async move {
196 let mut rx = self.event_rx.lock().await;
197 loop {
198 tokio::select! {
199 _ = token.cancelled() => {
200 info!("event trigger shutting down");
201 return Ok(());
202 }
203 event = rx.recv() => {
204 let Some(event) = event else {
205 return Ok(());
206 };
207 let rules = self.matching_rules(event.event_kind, &event.workflow_name);
208 for rule in rules {
209 let depth = Self::chain_depth_from_event(&rule.on_event);
210 if depth >= rule.max_chain_depth {
211 warn!(
212 source_workflow = %event.workflow_name,
213 target_workflow = %rule.target_workflow,
214 chain_depth = depth,
215 max_chain_depth = rule.max_chain_depth,
216 "chain depth exceeded, ignoring event"
217 );
218 continue;
219 }
220
221 let payload = Self::build_payload(
222 event.run_id,
223 &event.workflow_name,
224 &event.error,
225 );
226
227 let trigger_event = TriggerEvent {
228 workflow_name: rule.target_workflow.clone(),
229 payload,
230 trigger_kind: TriggerKind::RunEvent {
231 source_run_id: event.run_id,
232 event_kind: rule.on_event.as_str().to_string(),
233 },
234 };
235
236 if let Err(e) = sink.send(trigger_event).await {
237 warn!(error = %e, "failed to emit trigger event");
238 } else {
239 info!(
240 source_workflow = %event.workflow_name,
241 target_workflow = %rule.target_workflow,
242 source_run_id = %event.run_id,
243 "event trigger fired"
244 );
245 }
246 }
247 }
248 }
249 }
250 })
251 }
252}
253
254impl EventSubscriber for EventTrigger {
255 fn name(&self) -> &str {
256 "event-trigger"
257 }
258
259 fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
260 Box::pin(async move {
261 let internal = match event {
262 Event::RunFailed {
263 run_id,
264 workflow_name,
265 error,
266 ..
267 } => InternalEvent {
268 run_id: *run_id,
269 workflow_name: workflow_name.clone(),
270 event_kind: EventKind::RunFailed,
271 error: error.clone(),
272 },
273 Event::RunStatusChanged {
274 run_id,
275 workflow_name,
276 error,
277 ..
278 } => InternalEvent {
279 run_id: *run_id,
280 workflow_name: workflow_name.clone(),
281 event_kind: EventKind::RunStatusChanged,
282 error: error.clone(),
283 },
284 Event::StepFailed {
285 run_id,
286 step_name,
287 error,
288 ..
289 } => InternalEvent {
290 run_id: *run_id,
291 workflow_name: step_name.clone(),
292 event_kind: EventKind::StepFailed,
293 error: Some(error.clone()),
294 },
295 Event::ApprovalRejected {
296 run_id,
297 rejected_by,
298 ..
299 } => InternalEvent {
300 run_id: *run_id,
301 workflow_name: String::new(),
302 event_kind: EventKind::ApprovalRejected,
303 error: Some(format!("rejected by {rejected_by}")),
304 },
305 _ => return,
306 };
307
308 if self.event_tx.send(internal).await.is_err() {
309 warn!("event trigger receiver dropped, event lost");
310 }
311 })
312 }
313}
314
315#[cfg(test)]
316mod tests {
317 use std::time::Duration;
318
319 use chrono::Utc;
320 use rust_decimal::Decimal;
321 use tokio::time::timeout;
322
323 use super::*;
324
325 fn make_trigger(rules: Vec<EventTriggerRule>) -> EventTrigger {
326 EventTrigger::new(rules)
327 }
328
329 fn deploy_to_rollback_rule() -> EventTriggerRule {
330 EventTriggerRule {
331 on_event: EventKind::RunFailed,
332 source_workflow: "deploy".to_string(),
333 target_workflow: "rollback".to_string(),
334 max_chain_depth: 3,
335 }
336 }
337
338 #[tokio::test]
339 async fn event_trigger_fires_on_matching_run_failed() {
340 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
341 let (sink, mut rx) = TriggerSink::channel(16);
342 let token = CancellationToken::new();
343 let token_clone = token.clone();
344
345 let run_id = Uuid::now_v7();
347 trigger
348 .event_tx
349 .send(InternalEvent {
350 run_id,
351 workflow_name: "deploy".to_string(),
352 event_kind: EventKind::RunFailed,
353 error: Some("step crashed".to_string()),
354 })
355 .await
356 .unwrap();
357
358 let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
359
360 let event = timeout(Duration::from_secs(2), rx.recv())
361 .await
362 .expect("timed out")
363 .expect("channel closed");
364
365 assert_eq!(event.workflow_name, "rollback");
366 assert!(matches!(event.trigger_kind, TriggerKind::RunEvent { .. }));
367 if let TriggerKind::RunEvent {
368 source_run_id,
369 event_kind,
370 } = &event.trigger_kind
371 {
372 assert_eq!(*source_run_id, run_id);
373 assert_eq!(event_kind, "run_failed");
374 }
375
376 let payload = &event.payload;
377 assert_eq!(payload["source_workflow"], "deploy");
378 assert_eq!(payload["error"], "step crashed");
379
380 token.cancel();
381 let _ = handle.await;
382 }
383
384 #[tokio::test]
385 async fn event_trigger_ignores_non_matching_workflow() {
386 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
387 let (sink, mut rx) = TriggerSink::channel(16);
388 let token = CancellationToken::new();
389 let token_clone = token.clone();
390
391 trigger
393 .event_tx
394 .send(InternalEvent {
395 run_id: Uuid::now_v7(),
396 workflow_name: "build".to_string(),
397 event_kind: EventKind::RunFailed,
398 error: None,
399 })
400 .await
401 .unwrap();
402
403 let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
404
405 tokio::time::sleep(Duration::from_millis(100)).await;
407 token.cancel();
408 let _ = handle.await;
409
410 assert!(rx.try_recv().is_err());
412 }
413
414 #[tokio::test]
415 async fn event_trigger_ignores_non_matching_event_kind() {
416 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
417 let (sink, mut rx) = TriggerSink::channel(16);
418 let token = CancellationToken::new();
419 let token_clone = token.clone();
420
421 trigger
423 .event_tx
424 .send(InternalEvent {
425 run_id: Uuid::now_v7(),
426 workflow_name: "deploy".to_string(),
427 event_kind: EventKind::RunStatusChanged,
428 error: None,
429 })
430 .await
431 .unwrap();
432
433 let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
434
435 tokio::time::sleep(Duration::from_millis(100)).await;
436 token.cancel();
437 let _ = handle.await;
438
439 assert!(rx.try_recv().is_err());
440 }
441
442 #[tokio::test]
443 async fn event_trigger_payload_contains_source_info() {
444 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
445 let (sink, mut rx) = TriggerSink::channel(16);
446 let token = CancellationToken::new();
447 let token_clone = token.clone();
448
449 let run_id = Uuid::now_v7();
450 trigger
451 .event_tx
452 .send(InternalEvent {
453 run_id,
454 workflow_name: "deploy".to_string(),
455 event_kind: EventKind::RunFailed,
456 error: Some("timeout".to_string()),
457 })
458 .await
459 .unwrap();
460
461 let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
462
463 let event = timeout(Duration::from_secs(2), rx.recv())
464 .await
465 .expect("timed out")
466 .expect("channel closed");
467
468 assert_eq!(event.payload["source_run_id"], run_id.to_string());
469 assert_eq!(event.payload["source_workflow"], "deploy");
470 assert_eq!(event.payload["error"], "timeout");
471
472 token.cancel();
473 let _ = handle.await;
474 }
475
476 #[test]
477 fn subscribed_event_types_reflects_rules() {
478 let trigger = make_trigger(vec![
479 EventTriggerRule {
480 on_event: EventKind::RunFailed,
481 source_workflow: "a".to_string(),
482 target_workflow: "b".to_string(),
483 max_chain_depth: 3,
484 },
485 EventTriggerRule {
486 on_event: EventKind::StepFailed,
487 source_workflow: "c".to_string(),
488 target_workflow: "d".to_string(),
489 max_chain_depth: 3,
490 },
491 ]);
492 let types = trigger.subscribed_event_types();
493 assert!(types.contains(&"run_failed"));
494 assert!(types.contains(&"step_failed"));
495 }
496
497 #[tokio::test]
498 async fn event_subscriber_forwards_run_failed() {
499 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
500
501 let event = Event::RunFailed {
502 run_id: Uuid::now_v7(),
503 workflow_name: "deploy".to_string(),
504 error: Some("crash".to_string()),
505 cost_usd: Decimal::ZERO,
506 duration_ms: 0,
507 at: Utc::now(),
508 };
509
510 EventSubscriber::handle(&trigger, &event).await;
512
513 let mut rx = trigger.event_rx.lock().await;
515 let internal = rx.try_recv().unwrap();
516 assert_eq!(internal.workflow_name, "deploy");
517 assert_eq!(internal.event_kind, EventKind::RunFailed);
518 }
519
520 #[tokio::test]
521 async fn event_subscriber_ignores_irrelevant_events() {
522 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
523
524 let event = Event::RunCreated {
525 run_id: Uuid::now_v7(),
526 workflow_name: "deploy".to_string(),
527 at: Utc::now(),
528 };
529
530 EventSubscriber::handle(&trigger, &event).await;
531
532 let mut rx = trigger.event_rx.lock().await;
533 assert!(rx.try_recv().is_err());
534 }
535
536 #[tokio::test]
537 async fn graceful_shutdown() {
538 let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
539 let (sink, _rx) = TriggerSink::channel(16);
540 let token = CancellationToken::new();
541 let token_clone = token.clone();
542
543 let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
544
545 tokio::time::sleep(Duration::from_millis(50)).await;
547 assert!(!handle.is_finished());
548
549 token.cancel();
551 let result = timeout(Duration::from_secs(2), handle)
552 .await
553 .expect("timed out")
554 .expect("task panicked");
555 assert!(result.is_ok());
556 }
557}