Skip to main content

ironflow_runtime/trigger/
event.rs

1//! Domain-event trigger for workflow chaining.
2//!
3//! [`EventTrigger`] subscribes to the [`EventPublisher`](ironflow_engine::notify::EventPublisher)
4//! and creates a new run whenever a matching event fires. This allows
5//! declarative workflow chaining without external infrastructure.
6//!
7//! # Anti-loop protection
8//!
9//! Each triggered run carries a `_chain_depth` label. When the depth
10//! reaches `max_chain_depth`, the trigger ignores the event and logs a
11//! warning. This prevents two workflows from triggering each other
12//! indefinitely.
13//!
14//! # Examples
15//!
16//! ```no_run
17//! use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
18//! use ironflow_store::entities::EventKind;
19//!
20//! let trigger = EventTrigger::new(vec![
21//!     EventTriggerRule {
22//!         on_event: EventKind::RunFailed,
23//!         source_workflow: "deploy".to_string(),
24//!         target_workflow: "rollback".to_string(),
25//!         max_chain_depth: 3,
26//!     },
27//! ]);
28//! ```
29
30use 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
41/// Label key used to track chaining depth on triggered runs.
42pub const CHAIN_DEPTH_LABEL: &str = "_chain_depth";
43
44/// A rule that maps a domain event to a workflow to trigger.
45///
46/// # Examples
47///
48/// ```
49/// use ironflow_runtime::trigger::event::EventTriggerRule;
50/// use ironflow_store::entities::EventKind;
51///
52/// let rule = EventTriggerRule {
53///     on_event: EventKind::RunFailed,
54///     source_workflow: "deploy".to_string(),
55///     target_workflow: "rollback".to_string(),
56///     max_chain_depth: 3,
57///  };
58/// assert_eq!(rule.target_workflow, "rollback");
59/// ```
60#[derive(Debug, Clone)]
61pub struct EventTriggerRule {
62    /// The event kind to react to.
63    pub on_event: EventKind,
64    /// Only react to events from this workflow.
65    pub source_workflow: String,
66    /// The workflow to trigger.
67    pub target_workflow: String,
68    /// Maximum chaining depth (default 3). Beyond this, the event is
69    /// ignored and logged.
70    pub max_chain_depth: u8,
71}
72
73/// A trigger that reacts to internal domain events.
74///
75/// Register this trigger with the runtime via
76/// [`Runtime::trigger`](crate::runtime::Runtime::trigger). It must also
77/// be registered as an [`EventSubscriber`] on the
78/// [`EventPublisher`](ironflow_engine::notify::EventPublisher) so it
79/// receives events.
80///
81/// # Examples
82///
83/// ```no_run
84/// use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
85/// use ironflow_store::entities::EventKind;
86///
87/// let trigger = EventTrigger::new(vec![
88///     EventTriggerRule {
89///         on_event: EventKind::RunFailed,
90///         source_workflow: "deploy".to_string(),
91///         target_workflow: "rollback".to_string(),
92///         max_chain_depth: 3,
93///     },
94/// ]);
95/// ```
96pub struct EventTrigger {
97    rules: Vec<EventTriggerRule>,
98    /// Internal channel from the EventSubscriber side to the Trigger side.
99    event_tx: mpsc::Sender<InternalEvent>,
100    event_rx: tokio::sync::Mutex<mpsc::Receiver<InternalEvent>>,
101}
102
103/// Internal representation of a domain event relevant to the trigger.
104#[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    /// Create a new event trigger with the given rules.
114    ///
115    /// # Examples
116    ///
117    /// ```
118    /// use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
119    /// use ironflow_store::entities::EventKind;
120    ///
121    /// let trigger = EventTrigger::new(vec![
122    ///     EventTriggerRule {
123    ///         on_event: EventKind::RunFailed,
124    ///         source_workflow: "deploy".to_string(),
125    ///         target_workflow: "rollback".to_string(),
126    ///         max_chain_depth: 3,
127    ///     },
128    /// ]);
129    /// ```
130    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    /// The event kinds this trigger listens for (for EventPublisher subscription).
140    ///
141    /// # Examples
142    ///
143    /// ```
144    /// use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
145    /// use ironflow_store::entities::EventKind;
146    ///
147    /// let trigger = EventTrigger::new(vec![
148    ///     EventTriggerRule {
149    ///         on_event: EventKind::RunFailed,
150    ///         source_workflow: "deploy".to_string(),
151    ///         target_workflow: "rollback".to_string(),
152    ///         max_chain_depth: 3,
153    ///     },
154    /// ]);
155    /// let kinds = trigger.subscribed_event_types();
156    /// assert!(kinds.contains(&"run_failed"));
157    /// ```
158    pub fn subscribed_event_types(&self) -> Vec<&'static str> {
159        self.rules.iter().map(|r| r.on_event.as_str()).collect()
160    }
161
162    /// Find matching rules for a given event.
163    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    /// Build the payload for a triggered run.
171    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    /// Extract the chain depth from an event, defaulting to 0.
184    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        // Send an internal event through the subscriber path
346        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        // Send an event from a different workflow
392        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        // Give it time to process
406        tokio::time::sleep(Duration::from_millis(100)).await;
407        token.cancel();
408        let _ = handle.await;
409
410        // No event should have been emitted
411        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        // RunStatusChanged instead of RunFailed
422        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        // Call the EventSubscriber::handle method
511        EventSubscriber::handle(&trigger, &event).await;
512
513        // The internal channel should have the event
514        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        // Trigger should be running
546        tokio::time::sleep(Duration::from_millis(50)).await;
547        assert!(!handle.is_finished());
548
549        // Cancel and verify clean shutdown
550        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}