1use 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
48pub const CHAIN_DEPTH_LABEL: &str = "_chain_depth";
50
51#[derive(Debug, Clone)]
73pub struct TriggerContext {
74 pub labels: HashMap<String, String>,
76 pub error: Option<String>,
78 pub cost_usd: Decimal,
80 pub duration_ms: u64,
82}
83
84pub enum TriggerCondition {
102 Label {
104 key: String,
106 value: String,
108 },
109 Expression(Arc<dyn Fn(&TriggerContext) -> bool + Send + Sync>),
113}
114
115impl TriggerCondition {
116 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#[derive(Debug, Clone)]
202pub struct EventTriggerRule {
203 pub on_event: EventKind,
205 pub source_workflow: String,
207 pub target_workflow: String,
209 pub max_chain_depth: u8,
212 pub conditions: Vec<TriggerCondition>,
217}
218
219pub struct EventTrigger {
244 rules: Vec<EventTriggerRule>,
245 event_tx: mpsc::Sender<InternalEvent>,
247 event_rx: tokio::sync::Mutex<mpsc::Receiver<InternalEvent>>,
248}
249
250#[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 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 pub fn subscribed_event_types(&self) -> Vec<&'static str> {
311 self.rules.iter().map(|r| r.on_event.as_str()).collect()
312 }
313
314 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 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 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 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 tokio::time::sleep(Duration::from_millis(100)).await;
592 token.cancel();
593 let _ = handle.await;
594
595 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 EventSubscriber::handle(&trigger, &event).await;
699
700 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 tokio::time::sleep(Duration::from_millis(50)).await;
734 assert!(!handle.is_finished());
735
736 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}