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//!         conditions: vec![],
27//!     },
28//! ]);
29//! ```
30
31use std::collections::HashMap;
32use std::fmt;
33use std::panic::{AssertUnwindSafe, catch_unwind};
34use std::sync::Arc;
35
36use rust_decimal::Decimal;
37use serde_json::json;
38use tokio::sync::mpsc;
39use tokio_util::sync::CancellationToken;
40use tracing::{info, warn};
41use uuid::Uuid;
42
43use ironflow_engine::notify::{Event, EventSubscriber, SubscriberFuture};
44use ironflow_store::entities::{EventKind, TriggerKind};
45
46use super::{Trigger, TriggerEvent, TriggerFuture, TriggerSink};
47
48/// Label key used to track chaining depth on triggered runs.
49pub const CHAIN_DEPTH_LABEL: &str = "_chain_depth";
50
51/// Context available to [`TriggerCondition`]s when evaluating whether a
52/// rule should fire.
53///
54/// Exposes metadata from the source run: labels, error message,
55/// aggregated cost and duration.
56///
57/// # Examples
58///
59/// ```
60/// use std::collections::HashMap;
61/// use ironflow_runtime::trigger::event::TriggerContext;
62/// use rust_decimal::Decimal;
63///
64/// let ctx = TriggerContext {
65///     labels: HashMap::from([("env".to_string(), "prod".to_string())]),
66///     error: Some("timeout".to_string()),
67///     cost_usd: Decimal::new(42, 2),
68///     duration_ms: 5000,
69/// };
70/// assert_eq!(ctx.labels.get("env").unwrap(), "prod");
71/// ```
72#[derive(Debug, Clone)]
73pub struct TriggerContext {
74    /// Labels of the source run.
75    pub labels: HashMap<String, String>,
76    /// Error message of the source run (if any).
77    pub error: Option<String>,
78    /// Aggregated cost in USD of the source run.
79    pub cost_usd: Decimal,
80    /// Aggregated duration in milliseconds of the source run.
81    pub duration_ms: u64,
82}
83
84/// A condition that must be satisfied for an [`EventTriggerRule`] to fire.
85///
86/// All conditions on a rule are evaluated with AND semantics: the rule
87/// fires only when every condition returns `true`. For OR semantics,
88/// create separate rules.
89///
90/// # Examples
91///
92/// ```
93/// use ironflow_runtime::trigger::event::TriggerCondition;
94///
95/// let cond = TriggerCondition::Label {
96///     key: "env".to_string(),
97///     value: "prod".to_string(),
98/// };
99/// assert!(matches!(cond, TriggerCondition::Label { .. }));
100/// ```
101pub enum TriggerCondition {
102    /// The source run must carry a label with this exact key-value pair.
103    Label {
104        /// Label key to check.
105        key: String,
106        /// Expected label value.
107        value: String,
108    },
109    /// An arbitrary predicate evaluated against the [`TriggerContext`].
110    ///
111    /// Uses [`Arc`] so the condition is `Clone + Send + Sync`.
112    Expression(Arc<dyn Fn(&TriggerContext) -> bool + Send + Sync>),
113}
114
115impl TriggerCondition {
116    /// Evaluate this condition against the given context.
117    ///
118    /// # Examples
119    ///
120    /// ```
121    /// use std::collections::HashMap;
122    /// use ironflow_runtime::trigger::event::{TriggerCondition, TriggerContext};
123    /// use rust_decimal::Decimal;
124    ///
125    /// let ctx = TriggerContext {
126    ///     labels: HashMap::from([("env".to_string(), "prod".to_string())]),
127    ///     error: None,
128    ///     cost_usd: Decimal::ZERO,
129    ///     duration_ms: 0,
130    /// };
131    /// let cond = TriggerCondition::Label {
132    ///     key: "env".to_string(),
133    ///     value: "prod".to_string(),
134    /// };
135    /// assert!(cond.evaluate(&ctx));
136    /// ```
137    ///
138    /// # Panics
139    ///
140    /// Does not panic. If an [`Expression`](TriggerCondition::Expression)
141    /// closure panics, the panic is caught and the condition evaluates to
142    /// `false`.
143    pub fn evaluate(&self, ctx: &TriggerContext) -> bool {
144        match self {
145            TriggerCondition::Label { key, value } => {
146                ctx.labels.get(key).is_some_and(|v| v == value)
147            }
148            TriggerCondition::Expression(f) => match catch_unwind(AssertUnwindSafe(|| f(ctx))) {
149                Ok(result) => result,
150                Err(_) => {
151                    warn!("expression condition panicked, treating as false");
152                    false
153                }
154            },
155        }
156    }
157}
158
159impl Clone for TriggerCondition {
160    fn clone(&self) -> Self {
161        match self {
162            TriggerCondition::Label { key, value } => TriggerCondition::Label {
163                key: key.clone(),
164                value: value.clone(),
165            },
166            TriggerCondition::Expression(f) => TriggerCondition::Expression(Arc::clone(f)),
167        }
168    }
169}
170
171impl fmt::Debug for TriggerCondition {
172    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
173        match self {
174            TriggerCondition::Label { key, value } => f
175                .debug_struct("Label")
176                .field("key", key)
177                .field("value", value)
178                .finish(),
179            TriggerCondition::Expression(_) => f.write_str("Expression(<closure>)"),
180        }
181    }
182}
183
184/// A rule that maps a domain event to a workflow to trigger.
185///
186/// # Examples
187///
188/// ```
189/// use ironflow_runtime::trigger::event::EventTriggerRule;
190/// use ironflow_store::entities::EventKind;
191///
192/// let rule = EventTriggerRule {
193///     on_event: EventKind::RunFailed,
194///     source_workflow: "deploy".to_string(),
195///     target_workflow: "rollback".to_string(),
196///     max_chain_depth: 3,
197///     conditions: vec![],
198///  };
199/// assert_eq!(rule.target_workflow, "rollback");
200/// ```
201#[derive(Debug, Clone)]
202pub struct EventTriggerRule {
203    /// The event kind to react to.
204    pub on_event: EventKind,
205    /// Only react to events from this workflow.
206    pub source_workflow: String,
207    /// The workflow to trigger.
208    pub target_workflow: String,
209    /// Maximum chaining depth (default 3). Beyond this, the event is
210    /// ignored and logged.
211    pub max_chain_depth: u8,
212    /// Optional conditions that must all match for the rule to fire.
213    ///
214    /// Evaluated with AND semantics. An empty list means the rule fires
215    /// unconditionally (backward compatible).
216    pub conditions: Vec<TriggerCondition>,
217}
218
219/// A trigger that reacts to internal domain events.
220///
221/// Register this trigger with the runtime via
222/// [`Runtime::trigger`](crate::runtime::Runtime::trigger). It must also
223/// be registered as an [`EventSubscriber`] on the
224/// [`EventPublisher`](ironflow_engine::notify::EventPublisher) so it
225/// receives events.
226///
227/// # Examples
228///
229/// ```no_run
230/// use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
231/// use ironflow_store::entities::EventKind;
232///
233/// let trigger = EventTrigger::new(vec![
234///     EventTriggerRule {
235///         on_event: EventKind::RunFailed,
236///         source_workflow: "deploy".to_string(),
237///         target_workflow: "rollback".to_string(),
238///         max_chain_depth: 3,
239///         conditions: vec![],
240///     },
241/// ]);
242/// ```
243pub struct EventTrigger {
244    rules: Vec<EventTriggerRule>,
245    /// Internal channel from the EventSubscriber side to the Trigger side.
246    event_tx: mpsc::Sender<InternalEvent>,
247    event_rx: tokio::sync::Mutex<mpsc::Receiver<InternalEvent>>,
248}
249
250/// Internal representation of a domain event relevant to the trigger.
251#[derive(Debug)]
252struct InternalEvent {
253    run_id: Uuid,
254    workflow_name: String,
255    event_kind: EventKind,
256    error: Option<String>,
257    labels: HashMap<String, String>,
258    cost_usd: Decimal,
259    duration_ms: u64,
260}
261
262impl EventTrigger {
263    /// Create a new event trigger with the given rules.
264    ///
265    /// # Examples
266    ///
267    /// ```
268    /// use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
269    /// use ironflow_store::entities::EventKind;
270    ///
271    /// let trigger = EventTrigger::new(vec![
272    ///     EventTriggerRule {
273    ///         on_event: EventKind::RunFailed,
274    ///         source_workflow: "deploy".to_string(),
275    ///         target_workflow: "rollback".to_string(),
276    ///         max_chain_depth: 3,
277    ///         conditions: vec![],
278    ///     },
279    /// ]);
280    /// ```
281    pub fn new(rules: Vec<EventTriggerRule>) -> Self {
282        let (event_tx, event_rx) = mpsc::channel(256);
283        Self {
284            rules,
285            event_tx,
286            event_rx: tokio::sync::Mutex::new(event_rx),
287        }
288    }
289
290    /// The event kinds this trigger listens for (for EventPublisher subscription).
291    ///
292    /// # Examples
293    ///
294    /// ```
295    /// use ironflow_runtime::trigger::event::{EventTrigger, EventTriggerRule};
296    /// use ironflow_store::entities::EventKind;
297    ///
298    /// let trigger = EventTrigger::new(vec![
299    ///     EventTriggerRule {
300    ///         on_event: EventKind::RunFailed,
301    ///         source_workflow: "deploy".to_string(),
302    ///         target_workflow: "rollback".to_string(),
303    ///         max_chain_depth: 3,
304    ///         conditions: vec![],
305    ///     },
306    /// ]);
307    /// let kinds = trigger.subscribed_event_types();
308    /// assert!(kinds.contains(&"run_failed"));
309    /// ```
310    pub fn subscribed_event_types(&self) -> Vec<&'static str> {
311        self.rules.iter().map(|r| r.on_event.as_str()).collect()
312    }
313
314    /// Find matching rules for a given event.
315    fn matching_rules(&self, event_kind: EventKind, workflow_name: &str) -> Vec<&EventTriggerRule> {
316        self.rules
317            .iter()
318            .filter(|r| r.on_event == event_kind && r.source_workflow == workflow_name)
319            .collect()
320    }
321
322    /// Build the payload for a triggered run.
323    fn build_payload(
324        source_run_id: Uuid,
325        source_workflow: &str,
326        error: &Option<String>,
327    ) -> serde_json::Value {
328        json!({
329            "source_run_id": source_run_id,
330            "source_workflow": source_workflow,
331            "error": error,
332        })
333    }
334
335    /// Extract the chain depth from an event, defaulting to 0.
336    fn chain_depth_from_event(_event_kind: &EventKind) -> u8 {
337        0
338    }
339}
340
341impl Trigger for EventTrigger {
342    fn name(&self) -> &str {
343        "event-trigger"
344    }
345
346    fn start<'a>(&'a self, sink: TriggerSink, token: &'a CancellationToken) -> TriggerFuture<'a> {
347        Box::pin(async move {
348            let mut rx = self.event_rx.lock().await;
349            loop {
350                tokio::select! {
351                    _ = token.cancelled() => {
352                        info!("event trigger shutting down");
353                        return Ok(());
354                    }
355                    event = rx.recv() => {
356                        let Some(event) = event else {
357                            return Ok(());
358                        };
359                        let rules = self.matching_rules(event.event_kind, &event.workflow_name);
360                        if rules.is_empty() {
361                            continue;
362                        }
363
364                        let trigger_ctx = TriggerContext {
365                            labels: event.labels.clone(),
366                            error: event.error.clone(),
367                            cost_usd: event.cost_usd,
368                            duration_ms: event.duration_ms,
369                        };
370
371                        for rule in rules {
372                            let depth = Self::chain_depth_from_event(&rule.on_event);
373                            if depth >= rule.max_chain_depth {
374                                warn!(
375                                    source_workflow = %event.workflow_name,
376                                    target_workflow = %rule.target_workflow,
377                                    chain_depth = depth,
378                                    max_chain_depth = rule.max_chain_depth,
379                                    "chain depth exceeded, ignoring event"
380                                );
381                                continue;
382                            }
383
384                            if !rule.conditions.iter().all(|c| c.evaluate(&trigger_ctx)) {
385                                info!(
386                                    source_workflow = %event.workflow_name,
387                                    target_workflow = %rule.target_workflow,
388                                    "conditions not met, skipping rule"
389                                );
390                                continue;
391                            }
392
393                            let payload = Self::build_payload(
394                                event.run_id,
395                                &event.workflow_name,
396                                &event.error,
397                            );
398
399                            let trigger_event = TriggerEvent {
400                                workflow_name: rule.target_workflow.clone(),
401                                payload,
402                                trigger_kind: TriggerKind::RunEvent {
403                                    source_run_id: event.run_id,
404                                    event_kind: rule.on_event.as_str().to_string(),
405                                },
406                            };
407
408                            if let Err(e) = sink.send(trigger_event).await {
409                                warn!(error = %e, "failed to emit trigger event");
410                            } else {
411                                info!(
412                                    source_workflow = %event.workflow_name,
413                                    target_workflow = %rule.target_workflow,
414                                    source_run_id = %event.run_id,
415                                    "event trigger fired"
416                                );
417                            }
418                        }
419                    }
420                }
421            }
422        })
423    }
424}
425
426impl EventSubscriber for EventTrigger {
427    fn name(&self) -> &str {
428        "event-trigger"
429    }
430
431    #[deny(unreachable_patterns)]
432    fn handle<'a>(&'a self, event: &'a Event) -> SubscriberFuture<'a> {
433        Box::pin(async move {
434            let internal = match event {
435                Event::RunFailed(e) => InternalEvent {
436                    run_id: e.run_id,
437                    workflow_name: e.workflow_name.clone(),
438                    event_kind: EventKind::RunFailed,
439                    error: e.error.clone(),
440                    labels: e.labels.clone(),
441                    cost_usd: e.cost_usd,
442                    duration_ms: e.duration_ms,
443                },
444                Event::RunStatusChanged(e) => InternalEvent {
445                    run_id: e.run_id,
446                    workflow_name: e.workflow_name.clone(),
447                    event_kind: EventKind::RunStatusChanged,
448                    error: e.error.clone(),
449                    labels: e.labels.clone(),
450                    cost_usd: e.cost_usd,
451                    duration_ms: e.duration_ms,
452                },
453                // `step_name` lands in `workflow_name` here, matching the
454                // pre-refactor behaviour of this subscriber.
455                Event::StepFailed(e) => InternalEvent {
456                    run_id: e.run_id,
457                    workflow_name: e.step_name.clone(),
458                    event_kind: EventKind::StepFailed,
459                    error: Some(e.error.clone()),
460                    labels: HashMap::new(),
461                    cost_usd: Decimal::ZERO,
462                    duration_ms: 0,
463                },
464                Event::ApprovalRejected(e) => InternalEvent {
465                    run_id: e.run_id,
466                    workflow_name: String::new(),
467                    event_kind: EventKind::ApprovalRejected,
468                    error: Some(format!("rejected by {}", e.rejected_by)),
469                    labels: HashMap::new(),
470                    cost_usd: Decimal::ZERO,
471                    duration_ms: 0,
472                },
473                _ => return,
474            };
475
476            if self.event_tx.send(internal).await.is_err() {
477                warn!("event trigger receiver dropped, event lost");
478            }
479        })
480    }
481}
482
483#[cfg(test)]
484mod tests {
485    use std::time::Duration;
486
487    use chrono::Utc;
488    use ironflow_engine::notify::{RunCreatedEvent, RunFailedEvent};
489    use rust_decimal::Decimal;
490    use tokio::time::timeout;
491
492    use super::*;
493
494    fn make_trigger(rules: Vec<EventTriggerRule>) -> EventTrigger {
495        EventTrigger::new(rules)
496    }
497
498    fn deploy_to_rollback_rule() -> EventTriggerRule {
499        EventTriggerRule {
500            on_event: EventKind::RunFailed,
501            source_workflow: "deploy".to_string(),
502            target_workflow: "rollback".to_string(),
503            max_chain_depth: 3,
504            conditions: vec![],
505        }
506    }
507
508    fn internal_event(
509        run_id: Uuid,
510        workflow_name: &str,
511        event_kind: EventKind,
512        error: Option<String>,
513    ) -> InternalEvent {
514        InternalEvent {
515            run_id,
516            workflow_name: workflow_name.to_string(),
517            event_kind,
518            error,
519            labels: HashMap::new(),
520            cost_usd: Decimal::ZERO,
521            duration_ms: 0,
522        }
523    }
524
525    #[tokio::test]
526    async fn event_trigger_fires_on_matching_run_failed() {
527        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
528        let (sink, mut rx) = TriggerSink::channel(16);
529        let token = CancellationToken::new();
530        let token_clone = token.clone();
531
532        let run_id = Uuid::now_v7();
533        trigger
534            .event_tx
535            .send(internal_event(
536                run_id,
537                "deploy",
538                EventKind::RunFailed,
539                Some("step crashed".to_string()),
540            ))
541            .await
542            .unwrap();
543
544        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
545
546        let event = timeout(Duration::from_secs(2), rx.recv())
547            .await
548            .expect("timed out")
549            .expect("channel closed");
550
551        assert_eq!(event.workflow_name, "rollback");
552        assert!(matches!(event.trigger_kind, TriggerKind::RunEvent { .. }));
553        if let TriggerKind::RunEvent {
554            source_run_id,
555            event_kind,
556        } = &event.trigger_kind
557        {
558            assert_eq!(*source_run_id, run_id);
559            assert_eq!(event_kind, "run_failed");
560        }
561
562        let payload = &event.payload;
563        assert_eq!(payload["source_workflow"], "deploy");
564        assert_eq!(payload["error"], "step crashed");
565
566        token.cancel();
567        let _ = handle.await;
568    }
569
570    #[tokio::test]
571    async fn event_trigger_ignores_non_matching_workflow() {
572        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
573        let (sink, mut rx) = TriggerSink::channel(16);
574        let token = CancellationToken::new();
575        let token_clone = token.clone();
576
577        trigger
578            .event_tx
579            .send(internal_event(
580                Uuid::now_v7(),
581                "build",
582                EventKind::RunFailed,
583                None,
584            ))
585            .await
586            .unwrap();
587
588        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
589
590        // Give it time to process
591        tokio::time::sleep(Duration::from_millis(100)).await;
592        token.cancel();
593        let _ = handle.await;
594
595        // No event should have been emitted
596        assert!(rx.try_recv().is_err());
597    }
598
599    #[tokio::test]
600    async fn event_trigger_ignores_non_matching_event_kind() {
601        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
602        let (sink, mut rx) = TriggerSink::channel(16);
603        let token = CancellationToken::new();
604        let token_clone = token.clone();
605
606        trigger
607            .event_tx
608            .send(internal_event(
609                Uuid::now_v7(),
610                "deploy",
611                EventKind::RunStatusChanged,
612                None,
613            ))
614            .await
615            .unwrap();
616
617        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
618
619        tokio::time::sleep(Duration::from_millis(100)).await;
620        token.cancel();
621        let _ = handle.await;
622
623        assert!(rx.try_recv().is_err());
624    }
625
626    #[tokio::test]
627    async fn event_trigger_payload_contains_source_info() {
628        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
629        let (sink, mut rx) = TriggerSink::channel(16);
630        let token = CancellationToken::new();
631        let token_clone = token.clone();
632
633        let run_id = Uuid::now_v7();
634        trigger
635            .event_tx
636            .send(internal_event(
637                run_id,
638                "deploy",
639                EventKind::RunFailed,
640                Some("timeout".to_string()),
641            ))
642            .await
643            .unwrap();
644
645        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
646
647        let event = timeout(Duration::from_secs(2), rx.recv())
648            .await
649            .expect("timed out")
650            .expect("channel closed");
651
652        assert_eq!(event.payload["source_run_id"], run_id.to_string());
653        assert_eq!(event.payload["source_workflow"], "deploy");
654        assert_eq!(event.payload["error"], "timeout");
655
656        token.cancel();
657        let _ = handle.await;
658    }
659
660    #[test]
661    fn subscribed_event_types_reflects_rules() {
662        let trigger = make_trigger(vec![
663            EventTriggerRule {
664                on_event: EventKind::RunFailed,
665                source_workflow: "a".to_string(),
666                target_workflow: "b".to_string(),
667                max_chain_depth: 3,
668                conditions: vec![],
669            },
670            EventTriggerRule {
671                on_event: EventKind::StepFailed,
672                source_workflow: "c".to_string(),
673                target_workflow: "d".to_string(),
674                max_chain_depth: 3,
675                conditions: vec![],
676            },
677        ]);
678        let types = trigger.subscribed_event_types();
679        assert!(types.contains(&"run_failed"));
680        assert!(types.contains(&"step_failed"));
681    }
682
683    #[tokio::test]
684    async fn event_subscriber_forwards_run_failed() {
685        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
686
687        let event = Event::RunFailed(RunFailedEvent {
688            run_id: Uuid::now_v7(),
689            workflow_name: "deploy".to_string(),
690            error: Some("crash".to_string()),
691            cost_usd: Decimal::ZERO,
692            duration_ms: 0,
693            labels: HashMap::new(),
694            at: Utc::now(),
695        });
696
697        // Call the EventSubscriber::handle method
698        EventSubscriber::handle(&trigger, &event).await;
699
700        // The internal channel should have the event
701        let mut rx = trigger.event_rx.lock().await;
702        let internal = rx.try_recv().unwrap();
703        assert_eq!(internal.workflow_name, "deploy");
704        assert_eq!(internal.event_kind, EventKind::RunFailed);
705    }
706
707    #[tokio::test]
708    async fn event_subscriber_ignores_irrelevant_events() {
709        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
710
711        let event = Event::RunCreated(RunCreatedEvent {
712            run_id: Uuid::now_v7(),
713            workflow_name: "deploy".to_string(),
714            at: Utc::now(),
715        });
716
717        EventSubscriber::handle(&trigger, &event).await;
718
719        let mut rx = trigger.event_rx.lock().await;
720        assert!(rx.try_recv().is_err());
721    }
722
723    #[tokio::test]
724    async fn graceful_shutdown() {
725        let trigger = make_trigger(vec![deploy_to_rollback_rule()]);
726        let (sink, _rx) = TriggerSink::channel(16);
727        let token = CancellationToken::new();
728        let token_clone = token.clone();
729
730        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
731
732        // Trigger should be running
733        tokio::time::sleep(Duration::from_millis(50)).await;
734        assert!(!handle.is_finished());
735
736        // Cancel and verify clean shutdown
737        token.cancel();
738        let result = timeout(Duration::from_secs(2), handle)
739            .await
740            .expect("timed out")
741            .expect("task panicked");
742        assert!(result.is_ok());
743    }
744
745    fn internal_event_with_labels(
746        run_id: Uuid,
747        workflow_name: &str,
748        event_kind: EventKind,
749        error: Option<String>,
750        labels: HashMap<String, String>,
751    ) -> InternalEvent {
752        InternalEvent {
753            run_id,
754            workflow_name: workflow_name.to_string(),
755            event_kind,
756            error,
757            labels,
758            cost_usd: Decimal::new(42, 2),
759            duration_ms: 5000,
760        }
761    }
762
763    #[tokio::test]
764    async fn condition_label_matches() {
765        let rule = EventTriggerRule {
766            on_event: EventKind::RunFailed,
767            source_workflow: "deploy".to_string(),
768            target_workflow: "rollback".to_string(),
769            max_chain_depth: 3,
770            conditions: vec![TriggerCondition::Label {
771                key: "env".to_string(),
772                value: "prod".to_string(),
773            }],
774        };
775        let trigger = make_trigger(vec![rule]);
776        let (sink, mut rx) = TriggerSink::channel(16);
777        let token = CancellationToken::new();
778        let token_clone = token.clone();
779
780        let run_id = Uuid::now_v7();
781        trigger
782            .event_tx
783            .send(internal_event_with_labels(
784                run_id,
785                "deploy",
786                EventKind::RunFailed,
787                Some("crash".to_string()),
788                HashMap::from([("env".to_string(), "prod".to_string())]),
789            ))
790            .await
791            .unwrap();
792
793        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
794
795        let event = timeout(Duration::from_secs(2), rx.recv())
796            .await
797            .expect("timed out")
798            .expect("channel closed");
799
800        assert_eq!(event.workflow_name, "rollback");
801        token.cancel();
802        let _ = handle.await;
803    }
804
805    #[tokio::test]
806    async fn condition_label_absent_no_fire() {
807        let rule = EventTriggerRule {
808            on_event: EventKind::RunFailed,
809            source_workflow: "deploy".to_string(),
810            target_workflow: "rollback".to_string(),
811            max_chain_depth: 3,
812            conditions: vec![TriggerCondition::Label {
813                key: "env".to_string(),
814                value: "prod".to_string(),
815            }],
816        };
817        let trigger = make_trigger(vec![rule]);
818        let (sink, mut rx) = TriggerSink::channel(16);
819        let token = CancellationToken::new();
820        let token_clone = token.clone();
821
822        trigger
823            .event_tx
824            .send(internal_event_with_labels(
825                Uuid::now_v7(),
826                "deploy",
827                EventKind::RunFailed,
828                None,
829                HashMap::new(),
830            ))
831            .await
832            .unwrap();
833
834        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
835
836        tokio::time::sleep(Duration::from_millis(100)).await;
837        token.cancel();
838        let _ = handle.await;
839
840        assert!(rx.try_recv().is_err());
841    }
842
843    #[tokio::test]
844    async fn condition_label_wrong_value_no_fire() {
845        let rule = EventTriggerRule {
846            on_event: EventKind::RunFailed,
847            source_workflow: "deploy".to_string(),
848            target_workflow: "rollback".to_string(),
849            max_chain_depth: 3,
850            conditions: vec![TriggerCondition::Label {
851                key: "env".to_string(),
852                value: "prod".to_string(),
853            }],
854        };
855        let trigger = make_trigger(vec![rule]);
856        let (sink, mut rx) = TriggerSink::channel(16);
857        let token = CancellationToken::new();
858        let token_clone = token.clone();
859
860        trigger
861            .event_tx
862            .send(internal_event_with_labels(
863                Uuid::now_v7(),
864                "deploy",
865                EventKind::RunFailed,
866                None,
867                HashMap::from([("env".to_string(), "staging".to_string())]),
868            ))
869            .await
870            .unwrap();
871
872        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
873
874        tokio::time::sleep(Duration::from_millis(100)).await;
875        token.cancel();
876        let _ = handle.await;
877
878        assert!(rx.try_recv().is_err());
879    }
880
881    #[tokio::test]
882    async fn multiple_conditions_all_match() {
883        let rule = EventTriggerRule {
884            on_event: EventKind::RunFailed,
885            source_workflow: "deploy".to_string(),
886            target_workflow: "rollback".to_string(),
887            max_chain_depth: 3,
888            conditions: vec![
889                TriggerCondition::Label {
890                    key: "env".to_string(),
891                    value: "prod".to_string(),
892                },
893                TriggerCondition::Label {
894                    key: "region".to_string(),
895                    value: "eu-west-1".to_string(),
896                },
897            ],
898        };
899        let trigger = make_trigger(vec![rule]);
900        let (sink, mut rx) = TriggerSink::channel(16);
901        let token = CancellationToken::new();
902        let token_clone = token.clone();
903
904        trigger
905            .event_tx
906            .send(internal_event_with_labels(
907                Uuid::now_v7(),
908                "deploy",
909                EventKind::RunFailed,
910                None,
911                HashMap::from([
912                    ("env".to_string(), "prod".to_string()),
913                    ("region".to_string(), "eu-west-1".to_string()),
914                ]),
915            ))
916            .await
917            .unwrap();
918
919        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
920
921        let event = timeout(Duration::from_secs(2), rx.recv())
922            .await
923            .expect("timed out")
924            .expect("channel closed");
925
926        assert_eq!(event.workflow_name, "rollback");
927        token.cancel();
928        let _ = handle.await;
929    }
930
931    #[tokio::test]
932    async fn multiple_conditions_one_fails() {
933        let rule = EventTriggerRule {
934            on_event: EventKind::RunFailed,
935            source_workflow: "deploy".to_string(),
936            target_workflow: "rollback".to_string(),
937            max_chain_depth: 3,
938            conditions: vec![
939                TriggerCondition::Label {
940                    key: "env".to_string(),
941                    value: "prod".to_string(),
942                },
943                TriggerCondition::Label {
944                    key: "region".to_string(),
945                    value: "eu-west-1".to_string(),
946                },
947            ],
948        };
949        let trigger = make_trigger(vec![rule]);
950        let (sink, mut rx) = TriggerSink::channel(16);
951        let token = CancellationToken::new();
952        let token_clone = token.clone();
953
954        trigger
955            .event_tx
956            .send(internal_event_with_labels(
957                Uuid::now_v7(),
958                "deploy",
959                EventKind::RunFailed,
960                None,
961                HashMap::from([("env".to_string(), "prod".to_string())]),
962            ))
963            .await
964            .unwrap();
965
966        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
967
968        tokio::time::sleep(Duration::from_millis(100)).await;
969        token.cancel();
970        let _ = handle.await;
971
972        assert!(rx.try_recv().is_err());
973    }
974
975    #[tokio::test]
976    async fn empty_conditions_backward_compat() {
977        let rule = EventTriggerRule {
978            on_event: EventKind::RunFailed,
979            source_workflow: "deploy".to_string(),
980            target_workflow: "rollback".to_string(),
981            max_chain_depth: 3,
982            conditions: vec![],
983        };
984        let trigger = make_trigger(vec![rule]);
985        let (sink, mut rx) = TriggerSink::channel(16);
986        let token = CancellationToken::new();
987        let token_clone = token.clone();
988
989        trigger
990            .event_tx
991            .send(internal_event(
992                Uuid::now_v7(),
993                "deploy",
994                EventKind::RunFailed,
995                Some("boom".to_string()),
996            ))
997            .await
998            .unwrap();
999
1000        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
1001
1002        let event = timeout(Duration::from_secs(2), rx.recv())
1003            .await
1004            .expect("timed out")
1005            .expect("channel closed");
1006
1007        assert_eq!(event.workflow_name, "rollback");
1008        token.cancel();
1009        let _ = handle.await;
1010    }
1011
1012    #[tokio::test]
1013    async fn expression_condition_with_context() {
1014        let rule = EventTriggerRule {
1015            on_event: EventKind::RunFailed,
1016            source_workflow: "deploy".to_string(),
1017            target_workflow: "rollback".to_string(),
1018            max_chain_depth: 3,
1019            conditions: vec![TriggerCondition::Expression(Arc::new(|ctx| {
1020                ctx.cost_usd > Decimal::new(10, 2) && ctx.duration_ms > 1000
1021            }))],
1022        };
1023        let trigger = make_trigger(vec![rule]);
1024        let (sink, mut rx) = TriggerSink::channel(16);
1025        let token = CancellationToken::new();
1026        let token_clone = token.clone();
1027
1028        trigger
1029            .event_tx
1030            .send(internal_event_with_labels(
1031                Uuid::now_v7(),
1032                "deploy",
1033                EventKind::RunFailed,
1034                None,
1035                HashMap::new(),
1036            ))
1037            .await
1038            .unwrap();
1039
1040        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
1041
1042        let event = timeout(Duration::from_secs(2), rx.recv())
1043            .await
1044            .expect("timed out")
1045            .expect("channel closed");
1046
1047        assert_eq!(event.workflow_name, "rollback");
1048        token.cancel();
1049        let _ = handle.await;
1050    }
1051
1052    #[tokio::test]
1053    async fn expression_returns_false_no_fire() {
1054        let rule = EventTriggerRule {
1055            on_event: EventKind::RunFailed,
1056            source_workflow: "deploy".to_string(),
1057            target_workflow: "rollback".to_string(),
1058            max_chain_depth: 3,
1059            conditions: vec![TriggerCondition::Expression(Arc::new(|ctx| {
1060                ctx.cost_usd > Decimal::new(100, 0)
1061            }))],
1062        };
1063        let trigger = make_trigger(vec![rule]);
1064        let (sink, mut rx) = TriggerSink::channel(16);
1065        let token = CancellationToken::new();
1066        let token_clone = token.clone();
1067
1068        trigger
1069            .event_tx
1070            .send(internal_event_with_labels(
1071                Uuid::now_v7(),
1072                "deploy",
1073                EventKind::RunFailed,
1074                None,
1075                HashMap::new(),
1076            ))
1077            .await
1078            .unwrap();
1079
1080        let handle = tokio::spawn(async move { trigger.start(sink, &token_clone).await });
1081
1082        tokio::time::sleep(Duration::from_millis(100)).await;
1083        token.cancel();
1084        let _ = handle.await;
1085
1086        assert!(rx.try_recv().is_err());
1087    }
1088}