1use std::collections::HashMap;
4
5use chrono::{DateTime, Utc};
6use rust_decimal::Decimal;
7use serde::{Deserialize, Serialize};
8use uuid::Uuid;
9
10pub use ironflow_store::entities::LogStream;
11use ironflow_store::models::{Assignee, RunStatus, StepKind};
12
13#[derive(Debug, Clone, Serialize, Deserialize)]
30#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
31pub struct RunCreatedEvent {
32 pub run_id: Uuid,
34 pub workflow_name: String,
36 pub at: DateTime<Utc>,
38}
39
40#[derive(Debug, Clone, Serialize, Deserialize)]
67#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
68pub struct RunStatusChangedEvent {
69 pub run_id: Uuid,
71 pub workflow_name: String,
73 pub from: RunStatus,
75 pub to: RunStatus,
77 pub error: Option<String>,
79 pub cost_usd: Decimal,
81 pub duration_ms: u64,
83 #[serde(default)]
85 pub labels: HashMap<String, String>,
86 pub at: DateTime<Utc>,
88}
89
90#[derive(Debug, Clone, Serialize, Deserialize)]
114#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
115pub struct RunFailedEvent {
116 pub run_id: Uuid,
118 pub workflow_name: String,
120 pub error: Option<String>,
122 pub cost_usd: Decimal,
124 pub duration_ms: u64,
126 #[serde(default)]
128 pub labels: HashMap<String, String>,
129 pub at: DateTime<Utc>,
131}
132
133#[derive(Debug, Clone, Serialize, Deserialize)]
154#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
155pub struct RunBudgetExceededEvent {
156 pub run_id: Uuid,
158 pub workflow_name: String,
160 pub limit_usd: Decimal,
162 pub spent_usd: Decimal,
164 pub step_budget_usd: Decimal,
166 pub at: DateTime<Utc>,
168}
169
170#[derive(Debug, Clone, Serialize, Deserialize)]
189#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
190pub struct RetryForcedEvent {
191 pub run_id: Uuid,
193 pub workflow_name: String,
195 pub original_version: String,
197 pub current_version: String,
199 pub at: DateTime<Utc>,
201}
202
203#[derive(Debug, Clone, Serialize, Deserialize)]
226#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
227pub struct StepCompletedEvent {
228 pub run_id: Uuid,
230 pub step_id: Uuid,
232 pub step_name: String,
234 #[cfg_attr(feature = "openapi", schema(value_type = String))]
236 pub kind: StepKind,
237 pub duration_ms: u64,
239 pub cost_usd: Decimal,
241 pub at: DateTime<Utc>,
243}
244
245#[derive(Debug, Clone, Serialize, Deserialize)]
266#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
267pub struct StepFailedEvent {
268 pub run_id: Uuid,
270 pub step_id: Uuid,
272 pub step_name: String,
274 #[cfg_attr(feature = "openapi", schema(value_type = String))]
276 pub kind: StepKind,
277 pub error: String,
279 pub at: DateTime<Utc>,
281}
282
283#[derive(Debug, Clone, Serialize, Deserialize)]
301#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
302pub struct ApprovalRequestedEvent {
303 pub run_id: Uuid,
305 pub step_id: Uuid,
307 pub message: String,
309 pub at: DateTime<Utc>,
311}
312
313#[derive(Debug, Clone, Serialize, Deserialize)]
330#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
331pub struct ApprovalGrantedEvent {
332 pub run_id: Uuid,
334 pub approved_by: String,
336 pub at: DateTime<Utc>,
338}
339
340#[derive(Debug, Clone, Serialize, Deserialize)]
357#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
358pub struct ApprovalRejectedEvent {
359 pub run_id: Uuid,
361 pub rejected_by: String,
363 pub at: DateTime<Utc>,
365}
366
367#[derive(Debug, Clone, Serialize, Deserialize)]
395#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
396pub struct ApprovalEscalatedEvent {
397 pub run_id: Uuid,
399 pub step_id: Uuid,
401 pub step_name: String,
403 pub stage: u32,
405 pub policy: String,
407 pub action: String,
409 pub reason: String,
411 #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
413 pub assignee: Option<Assignee>,
414 pub at: DateTime<Utc>,
416}
417
418#[derive(Debug, Clone, Serialize, Deserialize)]
438#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
439pub struct LogLineEvent {
440 pub run_id: Uuid,
442 pub step_id: Uuid,
444 pub step_name: String,
446 pub stream: LogStream,
448 pub line: String,
450 pub at: DateTime<Utc>,
452}
453
454#[derive(Debug, Clone, Serialize, Deserialize)]
471#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
472pub struct UserSignedInEvent {
473 pub user_id: Uuid,
475 pub username: String,
477 pub at: DateTime<Utc>,
479}
480
481#[derive(Debug, Clone, Serialize, Deserialize)]
498#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
499pub struct UserSignedUpEvent {
500 pub user_id: Uuid,
502 pub username: String,
504 pub at: DateTime<Utc>,
506}
507
508#[derive(Debug, Clone, Serialize, Deserialize)]
525#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
526pub struct UserSignedOutEvent {
527 pub user_id: Uuid,
529 pub at: DateTime<Utc>,
531}
532
533#[derive(Debug, Clone, Serialize, Deserialize)]
565#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
566#[serde(tag = "type", rename_all = "snake_case")]
567pub enum Event {
568 RunCreated(RunCreatedEvent),
571
572 RunStatusChanged(RunStatusChangedEvent),
574
575 RunFailed(RunFailedEvent),
581
582 RunBudgetExceeded(RunBudgetExceededEvent),
590
591 RetryForced(RetryForcedEvent),
598
599 StepCompleted(StepCompletedEvent),
602
603 StepFailed(StepFailedEvent),
605
606 ApprovalRequested(ApprovalRequestedEvent),
609
610 ApprovalGranted(ApprovalGrantedEvent),
612
613 ApprovalRejected(ApprovalRejectedEvent),
615
616 ApprovalEscalated(ApprovalEscalatedEvent),
618
619 LogLine(LogLineEvent),
625
626 UserSignedIn(UserSignedInEvent),
629
630 UserSignedUp(UserSignedUpEvent),
632
633 UserSignedOut(UserSignedOutEvent),
635}
636
637impl Event {
638 pub const RUN_CREATED: &'static str = "run_created";
640 pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
642 pub const RUN_FAILED: &'static str = "run_failed";
644 pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
646 pub const RETRY_FORCED: &'static str = "retry_forced";
648 pub const STEP_COMPLETED: &'static str = "step_completed";
650 pub const STEP_FAILED: &'static str = "step_failed";
652 pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
654 pub const APPROVAL_GRANTED: &'static str = "approval_granted";
656 pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
658 pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
660 pub const LOG_LINE: &'static str = "log_line";
662 pub const USER_SIGNED_IN: &'static str = "user_signed_in";
664 pub const USER_SIGNED_UP: &'static str = "user_signed_up";
666 pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
668
669 pub const ALL: &'static [&'static str] = &[
685 Self::RUN_CREATED,
686 Self::RUN_STATUS_CHANGED,
687 Self::RUN_FAILED,
688 Self::RUN_BUDGET_EXCEEDED,
689 Self::STEP_COMPLETED,
690 Self::STEP_FAILED,
691 Self::APPROVAL_REQUESTED,
692 Self::APPROVAL_GRANTED,
693 Self::APPROVAL_REJECTED,
694 Self::APPROVAL_ESCALATED,
695 Self::LOG_LINE,
696 Self::USER_SIGNED_IN,
697 Self::USER_SIGNED_UP,
698 Self::USER_SIGNED_OUT,
699 Self::RETRY_FORCED,
700 ];
701
702 #[deny(unreachable_patterns)]
721 pub fn event_type(&self) -> &'static str {
722 match self {
723 Event::RunCreated(_) => Self::RUN_CREATED,
724 Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
725 Event::RunFailed(_) => Self::RUN_FAILED,
726 Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
727 Event::RetryForced(_) => Self::RETRY_FORCED,
728 Event::StepCompleted(_) => Self::STEP_COMPLETED,
729 Event::StepFailed(_) => Self::STEP_FAILED,
730 Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
731 Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
732 Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
733 Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
734 Event::LogLine(_) => Self::LOG_LINE,
735 Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
736 Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
737 Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
738 }
739 }
740
741 #[deny(unreachable_patterns)]
764 pub fn run_id(&self) -> Option<Uuid> {
765 match self {
766 Event::RunCreated(e) => Some(e.run_id),
767 Event::RunStatusChanged(e) => Some(e.run_id),
768 Event::RunFailed(e) => Some(e.run_id),
769 Event::RunBudgetExceeded(e) => Some(e.run_id),
770 Event::RetryForced(e) => Some(e.run_id),
771 Event::StepCompleted(e) => Some(e.run_id),
772 Event::StepFailed(e) => Some(e.run_id),
773 Event::ApprovalRequested(e) => Some(e.run_id),
774 Event::ApprovalGranted(e) => Some(e.run_id),
775 Event::ApprovalRejected(e) => Some(e.run_id),
776 Event::ApprovalEscalated(e) => Some(e.run_id),
777 Event::LogLine(e) => Some(e.run_id),
778 Event::UserSignedIn(_) | Event::UserSignedUp(_) | Event::UserSignedOut(_) => None,
779 }
780 }
781
782 #[deny(unreachable_patterns)]
810 pub fn step_id(&self) -> Option<Uuid> {
811 match self {
812 Event::StepCompleted(e) => Some(e.step_id),
813 Event::StepFailed(e) => Some(e.step_id),
814 Event::ApprovalRequested(e) => Some(e.step_id),
815 Event::ApprovalEscalated(e) => Some(e.step_id),
816 Event::RunCreated(_)
817 | Event::RunStatusChanged(_)
818 | Event::RunFailed(_)
819 | Event::RunBudgetExceeded(_)
820 | Event::RetryForced(_)
821 | Event::ApprovalGranted(_)
822 | Event::ApprovalRejected(_)
823 | Event::LogLine(_)
824 | Event::UserSignedIn(_)
825 | Event::UserSignedUp(_)
826 | Event::UserSignedOut(_) => None,
827 }
828 }
829
830 #[deny(unreachable_patterns)]
853 pub fn user_id(&self) -> Option<Uuid> {
854 match self {
855 Event::UserSignedIn(e) => Some(e.user_id),
856 Event::UserSignedUp(e) => Some(e.user_id),
857 Event::UserSignedOut(e) => Some(e.user_id),
858 Event::RunCreated(_)
859 | Event::RunStatusChanged(_)
860 | Event::RunFailed(_)
861 | Event::RunBudgetExceeded(_)
862 | Event::RetryForced(_)
863 | Event::StepCompleted(_)
864 | Event::StepFailed(_)
865 | Event::ApprovalRequested(_)
866 | Event::ApprovalGranted(_)
867 | Event::ApprovalRejected(_)
868 | Event::ApprovalEscalated(_)
869 | Event::LogLine(_) => None,
870 }
871 }
872}
873
874#[cfg(test)]
875mod tests {
876 use super::*;
877
878 #[test]
879 fn run_status_changed_serde_roundtrip() {
880 let event = Event::RunStatusChanged(RunStatusChangedEvent {
881 run_id: Uuid::now_v7(),
882 workflow_name: "deploy".to_string(),
883 from: RunStatus::Running,
884 to: RunStatus::Completed,
885 error: None,
886 cost_usd: Decimal::new(42, 2),
887 duration_ms: 5000,
888 labels: HashMap::new(),
889 at: Utc::now(),
890 });
891
892 let json = serde_json::to_string(&event).expect("serialize");
893 let back: Event = serde_json::from_str(&json).expect("deserialize");
894
895 assert_eq!(back.event_type(), "run_status_changed");
896 assert!(json.contains("\"type\":\"run_status_changed\""));
897 }
898
899 #[test]
900 fn run_failed_serde_roundtrip() {
901 let event = Event::RunFailed(RunFailedEvent {
902 run_id: Uuid::now_v7(),
903 workflow_name: "deploy".to_string(),
904 error: Some("step crashed".to_string()),
905 cost_usd: Decimal::new(10, 2),
906 duration_ms: 3000,
907 labels: HashMap::new(),
908 at: Utc::now(),
909 });
910
911 let json = serde_json::to_string(&event).expect("serialize");
912 let back: Event = serde_json::from_str(&json).expect("deserialize");
913
914 assert_eq!(back.event_type(), "run_failed");
915 assert!(json.contains("\"type\":\"run_failed\""));
916 assert!(json.contains("step crashed"));
917 }
918
919 #[test]
920 fn run_budget_exceeded_serde_roundtrip() {
921 let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
922 run_id: Uuid::now_v7(),
923 workflow_name: "deploy".to_string(),
924 limit_usd: Decimal::new(200, 2),
925 spent_usd: Decimal::new(180, 2),
926 step_budget_usd: Decimal::new(50, 2),
927 at: Utc::now(),
928 });
929
930 let json = serde_json::to_string(&event).expect("serialize");
931 let back: Event = serde_json::from_str(&json).expect("deserialize");
932
933 assert_eq!(back.event_type(), "run_budget_exceeded");
934 assert!(json.contains("\"type\":\"run_budget_exceeded\""));
935 assert!(json.contains("limit_usd"));
936 assert!(json.contains("step_budget_usd"));
937 }
938
939 #[test]
940 fn all_contains_run_budget_exceeded() {
941 assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
942 }
943
944 #[test]
945 fn user_signed_in_serde_roundtrip() {
946 let event = Event::UserSignedIn(UserSignedInEvent {
947 user_id: Uuid::now_v7(),
948 username: "alice".to_string(),
949 at: Utc::now(),
950 });
951
952 let json = serde_json::to_string(&event).expect("serialize");
953 let back: Event = serde_json::from_str(&json).expect("deserialize");
954
955 assert_eq!(back.event_type(), "user_signed_in");
956 assert!(json.contains("alice"));
957 }
958
959 #[test]
960 fn step_failed_serde_roundtrip() {
961 let event = Event::StepFailed(StepFailedEvent {
962 run_id: Uuid::now_v7(),
963 step_id: Uuid::now_v7(),
964 step_name: "build".to_string(),
965 kind: StepKind::Shell,
966 error: "exit code 1".to_string(),
967 at: Utc::now(),
968 });
969
970 let json = serde_json::to_string(&event).expect("serialize");
971 let back: Event = serde_json::from_str(&json).expect("deserialize");
972
973 assert_eq!(back.event_type(), "step_failed");
974 }
975
976 #[test]
977 fn approval_requested_serde_roundtrip() {
978 let event = Event::ApprovalRequested(ApprovalRequestedEvent {
979 run_id: Uuid::now_v7(),
980 step_id: Uuid::now_v7(),
981 message: "Deploy to prod?".to_string(),
982 at: Utc::now(),
983 });
984
985 let json = serde_json::to_string(&event).expect("serialize");
986 assert!(json.contains("approval_requested"));
987 }
988
989 #[test]
990 fn log_line_serde_roundtrip() {
991 let event = Event::LogLine(LogLineEvent {
992 run_id: Uuid::now_v7(),
993 step_id: Uuid::now_v7(),
994 step_name: "build".to_string(),
995 stream: LogStream::Stdout,
996 line: "Compiling ironflow v0.1.0".to_string(),
997 at: Utc::now(),
998 });
999
1000 let json = serde_json::to_string(&event).expect("serialize");
1001 let back: Event = serde_json::from_str(&json).expect("deserialize");
1002
1003 assert_eq!(back.event_type(), "log_line");
1004 assert!(json.contains("\"type\":\"log_line\""));
1005 assert!(json.contains("Compiling ironflow"));
1006 }
1007
1008 #[test]
1013 fn legacy_flat_json_deserializes_into_typed_payload() {
1014 let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1015 .parse()
1016 .expect("valid uuid");
1017
1018 let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1019 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1020 match event {
1021 Event::RunCreated(e) => {
1022 assert_eq!(e.run_id, run_id);
1023 assert_eq!(e.workflow_name, "deploy");
1024 }
1025 other => panic!("expected RunCreated, got {other:?}"),
1026 }
1027
1028 let raw = r#"{"type":"run_status_changed","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","from":"running","to":"completed","error":null,"cost_usd":0.5,"duration_ms":5000,"labels":{"env":"prod"},"at":"2026-01-01T00:00:00Z"}"#;
1029 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1030 match event {
1031 Event::RunStatusChanged(e) => {
1032 assert_eq!(e.from, RunStatus::Running);
1033 assert_eq!(e.to, RunStatus::Completed);
1034 assert_eq!(e.cost_usd, Decimal::new(5, 1));
1035 assert_eq!(e.duration_ms, 5000);
1036 assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1037 }
1038 other => panic!("expected RunStatusChanged, got {other:?}"),
1039 }
1040
1041 let raw = r#"{"type":"run_status_changed","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","from":"running","to":"failed","error":"boom","cost_usd":0,"duration_ms":0,"at":"2026-01-01T00:00:00Z"}"#;
1043 let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1044 match event {
1045 Event::RunStatusChanged(e) => {
1046 assert!(e.labels.is_empty());
1047 assert_eq!(e.error.as_deref(), Some("boom"));
1048 }
1049 other => panic!("expected RunStatusChanged, got {other:?}"),
1050 }
1051
1052 let raw = r#"{"type":"run_failed","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","error":"boom","cost_usd":0.25,"duration_ms":3000,"at":"2026-01-01T00:00:00Z"}"#;
1053 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1054 match event {
1055 Event::RunFailed(e) => {
1056 assert_eq!(e.error.as_deref(), Some("boom"));
1057 assert!(e.labels.is_empty());
1058 }
1059 other => panic!("expected RunFailed, got {other:?}"),
1060 }
1061
1062 let raw = r#"{"type":"step_failed","run_id":"01890000-0000-7000-8000-000000000000","step_id":"01890000-0000-7000-8000-000000000001","step_name":"build","kind":"shell","error":"exit code 1","at":"2026-01-01T00:00:00Z"}"#;
1063 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1064 match event {
1065 Event::StepFailed(e) => {
1066 assert_eq!(e.kind, StepKind::Shell);
1067 assert_eq!(e.error, "exit code 1");
1068 }
1069 other => panic!("expected StepFailed, got {other:?}"),
1070 }
1071
1072 let raw = r#"{"type":"log_line","run_id":"01890000-0000-7000-8000-000000000000","step_id":"01890000-0000-7000-8000-000000000001","step_name":"build","stream":"stdout","line":"hello","at":"2026-01-01T00:00:00Z"}"#;
1073 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1074 match event {
1075 Event::LogLine(e) => {
1076 assert_eq!(e.stream, LogStream::Stdout);
1077 assert_eq!(e.line, "hello");
1078 }
1079 other => panic!("expected LogLine, got {other:?}"),
1080 }
1081
1082 let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1083 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1084 match event {
1085 Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1086 other => panic!("expected UserSignedIn, got {other:?}"),
1087 }
1088 }
1089
1090 #[test]
1093 fn serialized_event_is_flat_with_type_tag() {
1094 let run_id = Uuid::now_v7();
1095 let event = Event::RunCreated(RunCreatedEvent {
1096 run_id,
1097 workflow_name: "deploy".to_string(),
1098 at: Utc::now(),
1099 });
1100
1101 let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1102 let object = value.as_object().expect("event serializes to an object");
1103
1104 assert_eq!(
1105 object.get("type").and_then(|v| v.as_str()),
1106 Some("run_created")
1107 );
1108 assert_eq!(
1109 object.get("workflow_name").and_then(|v| v.as_str()),
1110 Some("deploy")
1111 );
1112 assert_eq!(
1113 object.get("run_id").and_then(|v| v.as_str()),
1114 Some(run_id.to_string().as_str())
1115 );
1116 assert!(object.contains_key("at"));
1117 assert_eq!(object.len(), 4, "no nesting: {object:?}");
1118 assert!(!object.contains_key("RunCreated"));
1119 }
1120
1121 #[test]
1122 fn run_id_returns_some_for_run_events() {
1123 let run_id = Uuid::now_v7();
1124 let now = Utc::now();
1125
1126 let events = vec![
1127 Event::RunCreated(RunCreatedEvent {
1128 run_id,
1129 workflow_name: "w".to_string(),
1130 at: now,
1131 }),
1132 Event::RunStatusChanged(RunStatusChangedEvent {
1133 run_id,
1134 workflow_name: "w".to_string(),
1135 from: RunStatus::Pending,
1136 to: RunStatus::Running,
1137 error: None,
1138 cost_usd: Decimal::ZERO,
1139 duration_ms: 0,
1140 labels: HashMap::new(),
1141 at: now,
1142 }),
1143 Event::RunFailed(RunFailedEvent {
1144 run_id,
1145 workflow_name: "w".to_string(),
1146 error: None,
1147 cost_usd: Decimal::ZERO,
1148 duration_ms: 0,
1149 labels: HashMap::new(),
1150 at: now,
1151 }),
1152 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1153 run_id,
1154 workflow_name: "w".to_string(),
1155 limit_usd: Decimal::ZERO,
1156 spent_usd: Decimal::ZERO,
1157 step_budget_usd: Decimal::ZERO,
1158 at: now,
1159 }),
1160 Event::RetryForced(RetryForcedEvent {
1161 run_id,
1162 workflow_name: "w".to_string(),
1163 original_version: "1".to_string(),
1164 current_version: "2".to_string(),
1165 at: now,
1166 }),
1167 Event::StepCompleted(StepCompletedEvent {
1168 run_id,
1169 step_id: Uuid::now_v7(),
1170 step_name: "s".to_string(),
1171 kind: StepKind::Shell,
1172 duration_ms: 0,
1173 cost_usd: Decimal::ZERO,
1174 at: now,
1175 }),
1176 Event::StepFailed(StepFailedEvent {
1177 run_id,
1178 step_id: Uuid::now_v7(),
1179 step_name: "s".to_string(),
1180 kind: StepKind::Shell,
1181 error: "e".to_string(),
1182 at: now,
1183 }),
1184 Event::ApprovalRequested(ApprovalRequestedEvent {
1185 run_id,
1186 step_id: Uuid::now_v7(),
1187 message: "ok?".to_string(),
1188 at: now,
1189 }),
1190 Event::ApprovalGranted(ApprovalGrantedEvent {
1191 run_id,
1192 approved_by: "alice".to_string(),
1193 at: now,
1194 }),
1195 Event::ApprovalRejected(ApprovalRejectedEvent {
1196 run_id,
1197 rejected_by: "bob".to_string(),
1198 at: now,
1199 }),
1200 Event::LogLine(LogLineEvent {
1201 run_id,
1202 step_id: Uuid::now_v7(),
1203 step_name: "s".to_string(),
1204 stream: LogStream::Stdout,
1205 line: "l".to_string(),
1206 at: now,
1207 }),
1208 ];
1209
1210 for event in &events {
1211 assert_eq!(
1212 event.run_id(),
1213 Some(run_id),
1214 "{} should carry a run_id",
1215 event.event_type()
1216 );
1217 }
1218 }
1219
1220 #[test]
1221 fn run_id_returns_none_for_auth_events() {
1222 let user_id = Uuid::now_v7();
1223 let now = Utc::now();
1224
1225 let events = vec![
1226 Event::UserSignedIn(UserSignedInEvent {
1227 user_id,
1228 username: "alice".to_string(),
1229 at: now,
1230 }),
1231 Event::UserSignedUp(UserSignedUpEvent {
1232 user_id,
1233 username: "alice".to_string(),
1234 at: now,
1235 }),
1236 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1237 ];
1238
1239 for event in &events {
1240 assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1241 }
1242 }
1243
1244 #[test]
1245 fn step_id_returns_some_only_for_step_events() {
1246 let step_id = Uuid::now_v7();
1247 let run_id = Uuid::now_v7();
1248 let now = Utc::now();
1249
1250 let with_step = vec![
1251 Event::StepCompleted(StepCompletedEvent {
1252 run_id,
1253 step_id,
1254 step_name: "s".to_string(),
1255 kind: StepKind::Shell,
1256 duration_ms: 0,
1257 cost_usd: Decimal::ZERO,
1258 at: now,
1259 }),
1260 Event::StepFailed(StepFailedEvent {
1261 run_id,
1262 step_id,
1263 step_name: "s".to_string(),
1264 kind: StepKind::Shell,
1265 error: "e".to_string(),
1266 at: now,
1267 }),
1268 Event::ApprovalRequested(ApprovalRequestedEvent {
1269 run_id,
1270 step_id,
1271 message: "ok?".to_string(),
1272 at: now,
1273 }),
1274 ];
1275
1276 for event in &with_step {
1277 assert_eq!(
1278 event.step_id(),
1279 Some(step_id),
1280 "{} should carry a step_id",
1281 event.event_type()
1282 );
1283 }
1284
1285 let without_step = vec![
1286 Event::RunCreated(RunCreatedEvent {
1287 run_id,
1288 workflow_name: "w".to_string(),
1289 at: now,
1290 }),
1291 Event::ApprovalGranted(ApprovalGrantedEvent {
1292 run_id,
1293 approved_by: "alice".to_string(),
1294 at: now,
1295 }),
1296 Event::LogLine(LogLineEvent {
1299 run_id,
1300 step_id,
1301 step_name: "s".to_string(),
1302 stream: LogStream::Stdout,
1303 line: "l".to_string(),
1304 at: now,
1305 }),
1306 Event::UserSignedOut(UserSignedOutEvent {
1307 user_id: Uuid::now_v7(),
1308 at: now,
1309 }),
1310 ];
1311
1312 for event in &without_step {
1313 assert_eq!(
1314 event.step_id(),
1315 None,
1316 "{} should not carry a step_id",
1317 event.event_type()
1318 );
1319 }
1320 }
1321
1322 #[test]
1323 fn user_id_returns_some_only_for_auth_events() {
1324 let user_id = Uuid::now_v7();
1325 let run_id = Uuid::now_v7();
1326 let now = Utc::now();
1327
1328 let auth = vec![
1329 Event::UserSignedIn(UserSignedInEvent {
1330 user_id,
1331 username: "alice".to_string(),
1332 at: now,
1333 }),
1334 Event::UserSignedUp(UserSignedUpEvent {
1335 user_id,
1336 username: "alice".to_string(),
1337 at: now,
1338 }),
1339 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1340 ];
1341
1342 for event in &auth {
1343 assert_eq!(
1344 event.user_id(),
1345 Some(user_id),
1346 "{} should carry a user_id",
1347 event.event_type()
1348 );
1349 }
1350
1351 let non_auth = vec![
1352 Event::RunCreated(RunCreatedEvent {
1353 run_id,
1354 workflow_name: "w".to_string(),
1355 at: now,
1356 }),
1357 Event::StepFailed(StepFailedEvent {
1358 run_id,
1359 step_id: Uuid::now_v7(),
1360 step_name: "s".to_string(),
1361 kind: StepKind::Shell,
1362 error: "e".to_string(),
1363 at: now,
1364 }),
1365 ];
1366
1367 for event in &non_auth {
1368 assert_eq!(
1369 event.user_id(),
1370 None,
1371 "{} should not carry a user_id",
1372 event.event_type()
1373 );
1374 }
1375 }
1376
1377 #[test]
1378 fn event_type_all_variants() {
1379 let id = Uuid::now_v7();
1380 let now = Utc::now();
1381
1382 let cases: Vec<(Event, &str)> = vec![
1383 (
1384 Event::RunCreated(RunCreatedEvent {
1385 run_id: id,
1386 workflow_name: "w".to_string(),
1387 at: now,
1388 }),
1389 "run_created",
1390 ),
1391 (
1392 Event::RunStatusChanged(RunStatusChangedEvent {
1393 run_id: id,
1394 workflow_name: "w".to_string(),
1395 from: RunStatus::Pending,
1396 to: RunStatus::Running,
1397 error: None,
1398 cost_usd: Decimal::ZERO,
1399 duration_ms: 0,
1400 labels: HashMap::new(),
1401 at: now,
1402 }),
1403 "run_status_changed",
1404 ),
1405 (
1406 Event::RunFailed(RunFailedEvent {
1407 run_id: id,
1408 workflow_name: "w".to_string(),
1409 error: Some("boom".to_string()),
1410 cost_usd: Decimal::ZERO,
1411 duration_ms: 0,
1412 labels: HashMap::new(),
1413 at: now,
1414 }),
1415 "run_failed",
1416 ),
1417 (
1418 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1419 run_id: id,
1420 workflow_name: "w".to_string(),
1421 limit_usd: Decimal::new(200, 2),
1422 spent_usd: Decimal::new(180, 2),
1423 step_budget_usd: Decimal::new(50, 2),
1424 at: now,
1425 }),
1426 "run_budget_exceeded",
1427 ),
1428 (
1429 Event::RetryForced(RetryForcedEvent {
1430 run_id: id,
1431 workflow_name: "w".to_string(),
1432 original_version: "1".to_string(),
1433 current_version: "2".to_string(),
1434 at: now,
1435 }),
1436 "retry_forced",
1437 ),
1438 (
1439 Event::StepCompleted(StepCompletedEvent {
1440 run_id: id,
1441 step_id: id,
1442 step_name: "s".to_string(),
1443 kind: StepKind::Shell,
1444 duration_ms: 0,
1445 cost_usd: Decimal::ZERO,
1446 at: now,
1447 }),
1448 "step_completed",
1449 ),
1450 (
1451 Event::StepFailed(StepFailedEvent {
1452 run_id: id,
1453 step_id: id,
1454 step_name: "s".to_string(),
1455 kind: StepKind::Shell,
1456 error: "err".to_string(),
1457 at: now,
1458 }),
1459 "step_failed",
1460 ),
1461 (
1462 Event::ApprovalRequested(ApprovalRequestedEvent {
1463 run_id: id,
1464 step_id: id,
1465 message: "ok?".to_string(),
1466 at: now,
1467 }),
1468 "approval_requested",
1469 ),
1470 (
1471 Event::ApprovalGranted(ApprovalGrantedEvent {
1472 run_id: id,
1473 approved_by: "alice".to_string(),
1474 at: now,
1475 }),
1476 "approval_granted",
1477 ),
1478 (
1479 Event::ApprovalRejected(ApprovalRejectedEvent {
1480 run_id: id,
1481 rejected_by: "bob".to_string(),
1482 at: now,
1483 }),
1484 "approval_rejected",
1485 ),
1486 (
1487 Event::ApprovalEscalated(ApprovalEscalatedEvent {
1488 run_id: id,
1489 step_id: id,
1490 step_name: "prod-gate".to_string(),
1491 stage: 0,
1492 policy: "auto_reject".to_string(),
1493 action: "rejected".to_string(),
1494 reason: "approval deadline of 3600s expired".to_string(),
1495 assignee: None,
1496 at: now,
1497 }),
1498 "approval_escalated",
1499 ),
1500 (
1501 Event::LogLine(LogLineEvent {
1502 run_id: id,
1503 step_id: id,
1504 step_name: "build".to_string(),
1505 stream: LogStream::Stdout,
1506 line: "Compiling ironflow v0.1.0".to_string(),
1507 at: now,
1508 }),
1509 "log_line",
1510 ),
1511 (
1512 Event::UserSignedIn(UserSignedInEvent {
1513 user_id: id,
1514 username: "u".to_string(),
1515 at: now,
1516 }),
1517 "user_signed_in",
1518 ),
1519 (
1520 Event::UserSignedUp(UserSignedUpEvent {
1521 user_id: id,
1522 username: "u".to_string(),
1523 at: now,
1524 }),
1525 "user_signed_up",
1526 ),
1527 (
1528 Event::UserSignedOut(UserSignedOutEvent {
1529 user_id: id,
1530 at: now,
1531 }),
1532 "user_signed_out",
1533 ),
1534 ];
1535
1536 assert_eq!(
1537 cases.len(),
1538 Event::ALL.len(),
1539 "every variant must be covered"
1540 );
1541
1542 for (event, expected_type) in cases {
1543 assert_eq!(event.event_type(), expected_type);
1544 }
1545 }
1546
1547 #[test]
1548 fn approval_escalated_serde_roundtrip() {
1549 let run_id = Uuid::now_v7();
1550 let step_id = Uuid::now_v7();
1551 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1552 run_id,
1553 step_id,
1554 step_name: "prod-gate".to_string(),
1555 stage: 1,
1556 policy: "escalate".to_string(),
1557 action: "reassigned to sre-oncall".to_string(),
1558 reason: "approval deadline of 3600s expired".to_string(),
1559 assignee: Some(Assignee::group("sre-oncall")),
1560 at: Utc::now(),
1561 });
1562
1563 let json = serde_json::to_string(&event).expect("serialize");
1564 assert!(
1565 json.contains("\"type\":\"approval_escalated\""),
1566 "got {json}"
1567 );
1568
1569 let back: Event = serde_json::from_str(&json).expect("deserialize");
1570 let Event::ApprovalEscalated(payload) = back else {
1571 panic!("expected an approval_escalated event");
1572 };
1573 assert_eq!(payload.run_id, run_id);
1574 assert_eq!(payload.stage, 1);
1575 assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1576 }
1577
1578 #[test]
1579 fn approval_escalated_carries_run_and_step_ids() {
1580 let run_id = Uuid::now_v7();
1581 let step_id = Uuid::now_v7();
1582 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1583 run_id,
1584 step_id,
1585 step_name: "prod-gate".to_string(),
1586 stage: 0,
1587 policy: "notify".to_string(),
1588 action: "notified 1 target".to_string(),
1589 reason: "approval deadline of 60s expired".to_string(),
1590 assignee: None,
1591 at: Utc::now(),
1592 });
1593
1594 assert_eq!(event.run_id(), Some(run_id));
1595 assert_eq!(event.step_id(), Some(step_id));
1596 assert_eq!(event.user_id(), None);
1597 }
1598}