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)]
439#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
440pub struct LogLineEvent {
441 #[serde(default)]
450 pub id: Uuid,
451 pub run_id: Uuid,
453 pub step_id: Uuid,
455 pub step_name: String,
457 pub stream: LogStream,
459 pub line: String,
461 pub at: DateTime<Utc>,
463}
464
465#[derive(Debug, Clone, Serialize, Deserialize)]
482#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
483pub struct UserSignedInEvent {
484 pub user_id: Uuid,
486 pub username: String,
488 pub at: DateTime<Utc>,
490}
491
492#[derive(Debug, Clone, Serialize, Deserialize)]
509#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
510pub struct UserSignedUpEvent {
511 pub user_id: Uuid,
513 pub username: String,
515 pub at: DateTime<Utc>,
517}
518
519#[derive(Debug, Clone, Serialize, Deserialize)]
536#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
537pub struct UserSignedOutEvent {
538 pub user_id: Uuid,
540 pub at: DateTime<Utc>,
542}
543
544#[derive(Debug, Clone, Serialize, Deserialize)]
576#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
577#[serde(tag = "type", rename_all = "snake_case")]
578pub enum Event {
579 RunCreated(RunCreatedEvent),
582
583 RunStatusChanged(RunStatusChangedEvent),
585
586 RunFailed(RunFailedEvent),
592
593 RunBudgetExceeded(RunBudgetExceededEvent),
601
602 RetryForced(RetryForcedEvent),
609
610 StepCompleted(StepCompletedEvent),
613
614 StepFailed(StepFailedEvent),
616
617 ApprovalRequested(ApprovalRequestedEvent),
620
621 ApprovalGranted(ApprovalGrantedEvent),
623
624 ApprovalRejected(ApprovalRejectedEvent),
626
627 ApprovalEscalated(ApprovalEscalatedEvent),
629
630 LogLine(LogLineEvent),
636
637 UserSignedIn(UserSignedInEvent),
640
641 UserSignedUp(UserSignedUpEvent),
643
644 UserSignedOut(UserSignedOutEvent),
646}
647
648impl Event {
649 pub const RUN_CREATED: &'static str = "run_created";
651 pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
653 pub const RUN_FAILED: &'static str = "run_failed";
655 pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
657 pub const RETRY_FORCED: &'static str = "retry_forced";
659 pub const STEP_COMPLETED: &'static str = "step_completed";
661 pub const STEP_FAILED: &'static str = "step_failed";
663 pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
665 pub const APPROVAL_GRANTED: &'static str = "approval_granted";
667 pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
669 pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
671 pub const LOG_LINE: &'static str = "log_line";
673 pub const USER_SIGNED_IN: &'static str = "user_signed_in";
675 pub const USER_SIGNED_UP: &'static str = "user_signed_up";
677 pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
679
680 pub const ALL: &'static [&'static str] = &[
696 Self::RUN_CREATED,
697 Self::RUN_STATUS_CHANGED,
698 Self::RUN_FAILED,
699 Self::RUN_BUDGET_EXCEEDED,
700 Self::STEP_COMPLETED,
701 Self::STEP_FAILED,
702 Self::APPROVAL_REQUESTED,
703 Self::APPROVAL_GRANTED,
704 Self::APPROVAL_REJECTED,
705 Self::APPROVAL_ESCALATED,
706 Self::LOG_LINE,
707 Self::USER_SIGNED_IN,
708 Self::USER_SIGNED_UP,
709 Self::USER_SIGNED_OUT,
710 Self::RETRY_FORCED,
711 ];
712
713 #[deny(unreachable_patterns)]
732 pub fn event_type(&self) -> &'static str {
733 match self {
734 Event::RunCreated(_) => Self::RUN_CREATED,
735 Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
736 Event::RunFailed(_) => Self::RUN_FAILED,
737 Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
738 Event::RetryForced(_) => Self::RETRY_FORCED,
739 Event::StepCompleted(_) => Self::STEP_COMPLETED,
740 Event::StepFailed(_) => Self::STEP_FAILED,
741 Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
742 Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
743 Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
744 Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
745 Event::LogLine(_) => Self::LOG_LINE,
746 Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
747 Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
748 Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
749 }
750 }
751
752 #[deny(unreachable_patterns)]
775 pub fn run_id(&self) -> Option<Uuid> {
776 match self {
777 Event::RunCreated(e) => Some(e.run_id),
778 Event::RunStatusChanged(e) => Some(e.run_id),
779 Event::RunFailed(e) => Some(e.run_id),
780 Event::RunBudgetExceeded(e) => Some(e.run_id),
781 Event::RetryForced(e) => Some(e.run_id),
782 Event::StepCompleted(e) => Some(e.run_id),
783 Event::StepFailed(e) => Some(e.run_id),
784 Event::ApprovalRequested(e) => Some(e.run_id),
785 Event::ApprovalGranted(e) => Some(e.run_id),
786 Event::ApprovalRejected(e) => Some(e.run_id),
787 Event::ApprovalEscalated(e) => Some(e.run_id),
788 Event::LogLine(e) => Some(e.run_id),
789 Event::UserSignedIn(_) | Event::UserSignedUp(_) | Event::UserSignedOut(_) => None,
790 }
791 }
792
793 #[deny(unreachable_patterns)]
821 pub fn step_id(&self) -> Option<Uuid> {
822 match self {
823 Event::StepCompleted(e) => Some(e.step_id),
824 Event::StepFailed(e) => Some(e.step_id),
825 Event::ApprovalRequested(e) => Some(e.step_id),
826 Event::ApprovalEscalated(e) => Some(e.step_id),
827 Event::RunCreated(_)
828 | Event::RunStatusChanged(_)
829 | Event::RunFailed(_)
830 | Event::RunBudgetExceeded(_)
831 | Event::RetryForced(_)
832 | Event::ApprovalGranted(_)
833 | Event::ApprovalRejected(_)
834 | Event::LogLine(_)
835 | Event::UserSignedIn(_)
836 | Event::UserSignedUp(_)
837 | Event::UserSignedOut(_) => None,
838 }
839 }
840
841 #[deny(unreachable_patterns)]
864 pub fn user_id(&self) -> Option<Uuid> {
865 match self {
866 Event::UserSignedIn(e) => Some(e.user_id),
867 Event::UserSignedUp(e) => Some(e.user_id),
868 Event::UserSignedOut(e) => Some(e.user_id),
869 Event::RunCreated(_)
870 | Event::RunStatusChanged(_)
871 | Event::RunFailed(_)
872 | Event::RunBudgetExceeded(_)
873 | Event::RetryForced(_)
874 | Event::StepCompleted(_)
875 | Event::StepFailed(_)
876 | Event::ApprovalRequested(_)
877 | Event::ApprovalGranted(_)
878 | Event::ApprovalRejected(_)
879 | Event::ApprovalEscalated(_)
880 | Event::LogLine(_) => None,
881 }
882 }
883}
884
885#[cfg(test)]
886mod tests {
887 use super::*;
888
889 #[test]
890 fn run_status_changed_serde_roundtrip() {
891 let event = Event::RunStatusChanged(RunStatusChangedEvent {
892 run_id: Uuid::now_v7(),
893 workflow_name: "deploy".to_string(),
894 from: RunStatus::Running,
895 to: RunStatus::Completed,
896 error: None,
897 cost_usd: Decimal::new(42, 2),
898 duration_ms: 5000,
899 labels: HashMap::new(),
900 at: Utc::now(),
901 });
902
903 let json = serde_json::to_string(&event).expect("serialize");
904 let back: Event = serde_json::from_str(&json).expect("deserialize");
905
906 assert_eq!(back.event_type(), "run_status_changed");
907 assert!(json.contains("\"type\":\"run_status_changed\""));
908 }
909
910 #[test]
911 fn run_failed_serde_roundtrip() {
912 let event = Event::RunFailed(RunFailedEvent {
913 run_id: Uuid::now_v7(),
914 workflow_name: "deploy".to_string(),
915 error: Some("step crashed".to_string()),
916 cost_usd: Decimal::new(10, 2),
917 duration_ms: 3000,
918 labels: HashMap::new(),
919 at: Utc::now(),
920 });
921
922 let json = serde_json::to_string(&event).expect("serialize");
923 let back: Event = serde_json::from_str(&json).expect("deserialize");
924
925 assert_eq!(back.event_type(), "run_failed");
926 assert!(json.contains("\"type\":\"run_failed\""));
927 assert!(json.contains("step crashed"));
928 }
929
930 #[test]
931 fn run_budget_exceeded_serde_roundtrip() {
932 let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
933 run_id: Uuid::now_v7(),
934 workflow_name: "deploy".to_string(),
935 limit_usd: Decimal::new(200, 2),
936 spent_usd: Decimal::new(180, 2),
937 step_budget_usd: Decimal::new(50, 2),
938 at: Utc::now(),
939 });
940
941 let json = serde_json::to_string(&event).expect("serialize");
942 let back: Event = serde_json::from_str(&json).expect("deserialize");
943
944 assert_eq!(back.event_type(), "run_budget_exceeded");
945 assert!(json.contains("\"type\":\"run_budget_exceeded\""));
946 assert!(json.contains("limit_usd"));
947 assert!(json.contains("step_budget_usd"));
948 }
949
950 #[test]
951 fn all_contains_run_budget_exceeded() {
952 assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
953 }
954
955 #[test]
956 fn user_signed_in_serde_roundtrip() {
957 let event = Event::UserSignedIn(UserSignedInEvent {
958 user_id: Uuid::now_v7(),
959 username: "alice".to_string(),
960 at: Utc::now(),
961 });
962
963 let json = serde_json::to_string(&event).expect("serialize");
964 let back: Event = serde_json::from_str(&json).expect("deserialize");
965
966 assert_eq!(back.event_type(), "user_signed_in");
967 assert!(json.contains("alice"));
968 }
969
970 #[test]
971 fn step_failed_serde_roundtrip() {
972 let event = Event::StepFailed(StepFailedEvent {
973 run_id: Uuid::now_v7(),
974 step_id: Uuid::now_v7(),
975 step_name: "build".to_string(),
976 kind: StepKind::Shell,
977 error: "exit code 1".to_string(),
978 at: Utc::now(),
979 });
980
981 let json = serde_json::to_string(&event).expect("serialize");
982 let back: Event = serde_json::from_str(&json).expect("deserialize");
983
984 assert_eq!(back.event_type(), "step_failed");
985 }
986
987 #[test]
988 fn approval_requested_serde_roundtrip() {
989 let event = Event::ApprovalRequested(ApprovalRequestedEvent {
990 run_id: Uuid::now_v7(),
991 step_id: Uuid::now_v7(),
992 message: "Deploy to prod?".to_string(),
993 at: Utc::now(),
994 });
995
996 let json = serde_json::to_string(&event).expect("serialize");
997 assert!(json.contains("approval_requested"));
998 }
999
1000 #[test]
1001 fn log_line_serde_roundtrip() {
1002 let event = Event::LogLine(LogLineEvent {
1003 id: Uuid::now_v7(),
1004 run_id: Uuid::now_v7(),
1005 step_id: Uuid::now_v7(),
1006 step_name: "build".to_string(),
1007 stream: LogStream::Stdout,
1008 line: "Compiling ironflow v0.1.0".to_string(),
1009 at: Utc::now(),
1010 });
1011
1012 let json = serde_json::to_string(&event).expect("serialize");
1013 let back: Event = serde_json::from_str(&json).expect("deserialize");
1014
1015 assert_eq!(back.event_type(), "log_line");
1016 assert!(json.contains("\"type\":\"log_line\""));
1017 assert!(json.contains("Compiling ironflow"));
1018 }
1019
1020 #[test]
1025 fn legacy_flat_json_deserializes_into_typed_payload() {
1026 let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1027 .parse()
1028 .expect("valid uuid");
1029
1030 let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1031 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1032 match event {
1033 Event::RunCreated(e) => {
1034 assert_eq!(e.run_id, run_id);
1035 assert_eq!(e.workflow_name, "deploy");
1036 }
1037 other => panic!("expected RunCreated, got {other:?}"),
1038 }
1039
1040 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"}"#;
1041 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1042 match event {
1043 Event::RunStatusChanged(e) => {
1044 assert_eq!(e.from, RunStatus::Running);
1045 assert_eq!(e.to, RunStatus::Completed);
1046 assert_eq!(e.cost_usd, Decimal::new(5, 1));
1047 assert_eq!(e.duration_ms, 5000);
1048 assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1049 }
1050 other => panic!("expected RunStatusChanged, got {other:?}"),
1051 }
1052
1053 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"}"#;
1055 let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1056 match event {
1057 Event::RunStatusChanged(e) => {
1058 assert!(e.labels.is_empty());
1059 assert_eq!(e.error.as_deref(), Some("boom"));
1060 }
1061 other => panic!("expected RunStatusChanged, got {other:?}"),
1062 }
1063
1064 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"}"#;
1065 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1066 match event {
1067 Event::RunFailed(e) => {
1068 assert_eq!(e.error.as_deref(), Some("boom"));
1069 assert!(e.labels.is_empty());
1070 }
1071 other => panic!("expected RunFailed, got {other:?}"),
1072 }
1073
1074 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"}"#;
1075 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1076 match event {
1077 Event::StepFailed(e) => {
1078 assert_eq!(e.kind, StepKind::Shell);
1079 assert_eq!(e.error, "exit code 1");
1080 }
1081 other => panic!("expected StepFailed, got {other:?}"),
1082 }
1083
1084 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"}"#;
1085 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1086 match event {
1087 Event::LogLine(e) => {
1088 assert_eq!(e.stream, LogStream::Stdout);
1089 assert_eq!(e.line, "hello");
1090 assert_eq!(e.id, Uuid::nil());
1093 }
1094 other => panic!("expected LogLine, got {other:?}"),
1095 }
1096
1097 let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1098 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1099 match event {
1100 Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1101 other => panic!("expected UserSignedIn, got {other:?}"),
1102 }
1103 }
1104
1105 #[test]
1108 fn serialized_event_is_flat_with_type_tag() {
1109 let run_id = Uuid::now_v7();
1110 let event = Event::RunCreated(RunCreatedEvent {
1111 run_id,
1112 workflow_name: "deploy".to_string(),
1113 at: Utc::now(),
1114 });
1115
1116 let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1117 let object = value.as_object().expect("event serializes to an object");
1118
1119 assert_eq!(
1120 object.get("type").and_then(|v| v.as_str()),
1121 Some("run_created")
1122 );
1123 assert_eq!(
1124 object.get("workflow_name").and_then(|v| v.as_str()),
1125 Some("deploy")
1126 );
1127 assert_eq!(
1128 object.get("run_id").and_then(|v| v.as_str()),
1129 Some(run_id.to_string().as_str())
1130 );
1131 assert!(object.contains_key("at"));
1132 assert_eq!(object.len(), 4, "no nesting: {object:?}");
1133 assert!(!object.contains_key("RunCreated"));
1134 }
1135
1136 #[test]
1137 fn run_id_returns_some_for_run_events() {
1138 let run_id = Uuid::now_v7();
1139 let now = Utc::now();
1140
1141 let events = vec![
1142 Event::RunCreated(RunCreatedEvent {
1143 run_id,
1144 workflow_name: "w".to_string(),
1145 at: now,
1146 }),
1147 Event::RunStatusChanged(RunStatusChangedEvent {
1148 run_id,
1149 workflow_name: "w".to_string(),
1150 from: RunStatus::Pending,
1151 to: RunStatus::Running,
1152 error: None,
1153 cost_usd: Decimal::ZERO,
1154 duration_ms: 0,
1155 labels: HashMap::new(),
1156 at: now,
1157 }),
1158 Event::RunFailed(RunFailedEvent {
1159 run_id,
1160 workflow_name: "w".to_string(),
1161 error: None,
1162 cost_usd: Decimal::ZERO,
1163 duration_ms: 0,
1164 labels: HashMap::new(),
1165 at: now,
1166 }),
1167 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1168 run_id,
1169 workflow_name: "w".to_string(),
1170 limit_usd: Decimal::ZERO,
1171 spent_usd: Decimal::ZERO,
1172 step_budget_usd: Decimal::ZERO,
1173 at: now,
1174 }),
1175 Event::RetryForced(RetryForcedEvent {
1176 run_id,
1177 workflow_name: "w".to_string(),
1178 original_version: "1".to_string(),
1179 current_version: "2".to_string(),
1180 at: now,
1181 }),
1182 Event::StepCompleted(StepCompletedEvent {
1183 run_id,
1184 step_id: Uuid::now_v7(),
1185 step_name: "s".to_string(),
1186 kind: StepKind::Shell,
1187 duration_ms: 0,
1188 cost_usd: Decimal::ZERO,
1189 at: now,
1190 }),
1191 Event::StepFailed(StepFailedEvent {
1192 run_id,
1193 step_id: Uuid::now_v7(),
1194 step_name: "s".to_string(),
1195 kind: StepKind::Shell,
1196 error: "e".to_string(),
1197 at: now,
1198 }),
1199 Event::ApprovalRequested(ApprovalRequestedEvent {
1200 run_id,
1201 step_id: Uuid::now_v7(),
1202 message: "ok?".to_string(),
1203 at: now,
1204 }),
1205 Event::ApprovalGranted(ApprovalGrantedEvent {
1206 run_id,
1207 approved_by: "alice".to_string(),
1208 at: now,
1209 }),
1210 Event::ApprovalRejected(ApprovalRejectedEvent {
1211 run_id,
1212 rejected_by: "bob".to_string(),
1213 at: now,
1214 }),
1215 Event::LogLine(LogLineEvent {
1216 id: Uuid::now_v7(),
1217 run_id,
1218 step_id: Uuid::now_v7(),
1219 step_name: "s".to_string(),
1220 stream: LogStream::Stdout,
1221 line: "l".to_string(),
1222 at: now,
1223 }),
1224 ];
1225
1226 for event in &events {
1227 assert_eq!(
1228 event.run_id(),
1229 Some(run_id),
1230 "{} should carry a run_id",
1231 event.event_type()
1232 );
1233 }
1234 }
1235
1236 #[test]
1237 fn run_id_returns_none_for_auth_events() {
1238 let user_id = Uuid::now_v7();
1239 let now = Utc::now();
1240
1241 let events = vec![
1242 Event::UserSignedIn(UserSignedInEvent {
1243 user_id,
1244 username: "alice".to_string(),
1245 at: now,
1246 }),
1247 Event::UserSignedUp(UserSignedUpEvent {
1248 user_id,
1249 username: "alice".to_string(),
1250 at: now,
1251 }),
1252 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1253 ];
1254
1255 for event in &events {
1256 assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1257 }
1258 }
1259
1260 #[test]
1261 fn step_id_returns_some_only_for_step_events() {
1262 let step_id = Uuid::now_v7();
1263 let run_id = Uuid::now_v7();
1264 let now = Utc::now();
1265
1266 let with_step = vec![
1267 Event::StepCompleted(StepCompletedEvent {
1268 run_id,
1269 step_id,
1270 step_name: "s".to_string(),
1271 kind: StepKind::Shell,
1272 duration_ms: 0,
1273 cost_usd: Decimal::ZERO,
1274 at: now,
1275 }),
1276 Event::StepFailed(StepFailedEvent {
1277 run_id,
1278 step_id,
1279 step_name: "s".to_string(),
1280 kind: StepKind::Shell,
1281 error: "e".to_string(),
1282 at: now,
1283 }),
1284 Event::ApprovalRequested(ApprovalRequestedEvent {
1285 run_id,
1286 step_id,
1287 message: "ok?".to_string(),
1288 at: now,
1289 }),
1290 ];
1291
1292 for event in &with_step {
1293 assert_eq!(
1294 event.step_id(),
1295 Some(step_id),
1296 "{} should carry a step_id",
1297 event.event_type()
1298 );
1299 }
1300
1301 let without_step = vec![
1302 Event::RunCreated(RunCreatedEvent {
1303 run_id,
1304 workflow_name: "w".to_string(),
1305 at: now,
1306 }),
1307 Event::ApprovalGranted(ApprovalGrantedEvent {
1308 run_id,
1309 approved_by: "alice".to_string(),
1310 at: now,
1311 }),
1312 Event::LogLine(LogLineEvent {
1315 id: Uuid::now_v7(),
1316 run_id,
1317 step_id,
1318 step_name: "s".to_string(),
1319 stream: LogStream::Stdout,
1320 line: "l".to_string(),
1321 at: now,
1322 }),
1323 Event::UserSignedOut(UserSignedOutEvent {
1324 user_id: Uuid::now_v7(),
1325 at: now,
1326 }),
1327 ];
1328
1329 for event in &without_step {
1330 assert_eq!(
1331 event.step_id(),
1332 None,
1333 "{} should not carry a step_id",
1334 event.event_type()
1335 );
1336 }
1337 }
1338
1339 #[test]
1340 fn user_id_returns_some_only_for_auth_events() {
1341 let user_id = Uuid::now_v7();
1342 let run_id = Uuid::now_v7();
1343 let now = Utc::now();
1344
1345 let auth = vec![
1346 Event::UserSignedIn(UserSignedInEvent {
1347 user_id,
1348 username: "alice".to_string(),
1349 at: now,
1350 }),
1351 Event::UserSignedUp(UserSignedUpEvent {
1352 user_id,
1353 username: "alice".to_string(),
1354 at: now,
1355 }),
1356 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1357 ];
1358
1359 for event in &auth {
1360 assert_eq!(
1361 event.user_id(),
1362 Some(user_id),
1363 "{} should carry a user_id",
1364 event.event_type()
1365 );
1366 }
1367
1368 let non_auth = vec![
1369 Event::RunCreated(RunCreatedEvent {
1370 run_id,
1371 workflow_name: "w".to_string(),
1372 at: now,
1373 }),
1374 Event::StepFailed(StepFailedEvent {
1375 run_id,
1376 step_id: Uuid::now_v7(),
1377 step_name: "s".to_string(),
1378 kind: StepKind::Shell,
1379 error: "e".to_string(),
1380 at: now,
1381 }),
1382 ];
1383
1384 for event in &non_auth {
1385 assert_eq!(
1386 event.user_id(),
1387 None,
1388 "{} should not carry a user_id",
1389 event.event_type()
1390 );
1391 }
1392 }
1393
1394 #[test]
1395 fn event_type_all_variants() {
1396 let id = Uuid::now_v7();
1397 let now = Utc::now();
1398
1399 let cases: Vec<(Event, &str)> = vec![
1400 (
1401 Event::RunCreated(RunCreatedEvent {
1402 run_id: id,
1403 workflow_name: "w".to_string(),
1404 at: now,
1405 }),
1406 "run_created",
1407 ),
1408 (
1409 Event::RunStatusChanged(RunStatusChangedEvent {
1410 run_id: id,
1411 workflow_name: "w".to_string(),
1412 from: RunStatus::Pending,
1413 to: RunStatus::Running,
1414 error: None,
1415 cost_usd: Decimal::ZERO,
1416 duration_ms: 0,
1417 labels: HashMap::new(),
1418 at: now,
1419 }),
1420 "run_status_changed",
1421 ),
1422 (
1423 Event::RunFailed(RunFailedEvent {
1424 run_id: id,
1425 workflow_name: "w".to_string(),
1426 error: Some("boom".to_string()),
1427 cost_usd: Decimal::ZERO,
1428 duration_ms: 0,
1429 labels: HashMap::new(),
1430 at: now,
1431 }),
1432 "run_failed",
1433 ),
1434 (
1435 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1436 run_id: id,
1437 workflow_name: "w".to_string(),
1438 limit_usd: Decimal::new(200, 2),
1439 spent_usd: Decimal::new(180, 2),
1440 step_budget_usd: Decimal::new(50, 2),
1441 at: now,
1442 }),
1443 "run_budget_exceeded",
1444 ),
1445 (
1446 Event::RetryForced(RetryForcedEvent {
1447 run_id: id,
1448 workflow_name: "w".to_string(),
1449 original_version: "1".to_string(),
1450 current_version: "2".to_string(),
1451 at: now,
1452 }),
1453 "retry_forced",
1454 ),
1455 (
1456 Event::StepCompleted(StepCompletedEvent {
1457 run_id: id,
1458 step_id: id,
1459 step_name: "s".to_string(),
1460 kind: StepKind::Shell,
1461 duration_ms: 0,
1462 cost_usd: Decimal::ZERO,
1463 at: now,
1464 }),
1465 "step_completed",
1466 ),
1467 (
1468 Event::StepFailed(StepFailedEvent {
1469 run_id: id,
1470 step_id: id,
1471 step_name: "s".to_string(),
1472 kind: StepKind::Shell,
1473 error: "err".to_string(),
1474 at: now,
1475 }),
1476 "step_failed",
1477 ),
1478 (
1479 Event::ApprovalRequested(ApprovalRequestedEvent {
1480 run_id: id,
1481 step_id: id,
1482 message: "ok?".to_string(),
1483 at: now,
1484 }),
1485 "approval_requested",
1486 ),
1487 (
1488 Event::ApprovalGranted(ApprovalGrantedEvent {
1489 run_id: id,
1490 approved_by: "alice".to_string(),
1491 at: now,
1492 }),
1493 "approval_granted",
1494 ),
1495 (
1496 Event::ApprovalRejected(ApprovalRejectedEvent {
1497 run_id: id,
1498 rejected_by: "bob".to_string(),
1499 at: now,
1500 }),
1501 "approval_rejected",
1502 ),
1503 (
1504 Event::ApprovalEscalated(ApprovalEscalatedEvent {
1505 run_id: id,
1506 step_id: id,
1507 step_name: "prod-gate".to_string(),
1508 stage: 0,
1509 policy: "auto_reject".to_string(),
1510 action: "rejected".to_string(),
1511 reason: "approval deadline of 3600s expired".to_string(),
1512 assignee: None,
1513 at: now,
1514 }),
1515 "approval_escalated",
1516 ),
1517 (
1518 Event::LogLine(LogLineEvent {
1519 id,
1520 run_id: id,
1521 step_id: id,
1522 step_name: "build".to_string(),
1523 stream: LogStream::Stdout,
1524 line: "Compiling ironflow v0.1.0".to_string(),
1525 at: now,
1526 }),
1527 "log_line",
1528 ),
1529 (
1530 Event::UserSignedIn(UserSignedInEvent {
1531 user_id: id,
1532 username: "u".to_string(),
1533 at: now,
1534 }),
1535 "user_signed_in",
1536 ),
1537 (
1538 Event::UserSignedUp(UserSignedUpEvent {
1539 user_id: id,
1540 username: "u".to_string(),
1541 at: now,
1542 }),
1543 "user_signed_up",
1544 ),
1545 (
1546 Event::UserSignedOut(UserSignedOutEvent {
1547 user_id: id,
1548 at: now,
1549 }),
1550 "user_signed_out",
1551 ),
1552 ];
1553
1554 assert_eq!(
1555 cases.len(),
1556 Event::ALL.len(),
1557 "every variant must be covered"
1558 );
1559
1560 for (event, expected_type) in cases {
1561 assert_eq!(event.event_type(), expected_type);
1562 }
1563 }
1564
1565 #[test]
1566 fn approval_escalated_serde_roundtrip() {
1567 let run_id = Uuid::now_v7();
1568 let step_id = Uuid::now_v7();
1569 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1570 run_id,
1571 step_id,
1572 step_name: "prod-gate".to_string(),
1573 stage: 1,
1574 policy: "escalate".to_string(),
1575 action: "reassigned to sre-oncall".to_string(),
1576 reason: "approval deadline of 3600s expired".to_string(),
1577 assignee: Some(Assignee::group("sre-oncall")),
1578 at: Utc::now(),
1579 });
1580
1581 let json = serde_json::to_string(&event).expect("serialize");
1582 assert!(
1583 json.contains("\"type\":\"approval_escalated\""),
1584 "got {json}"
1585 );
1586
1587 let back: Event = serde_json::from_str(&json).expect("deserialize");
1588 let Event::ApprovalEscalated(payload) = back else {
1589 panic!("expected an approval_escalated event");
1590 };
1591 assert_eq!(payload.run_id, run_id);
1592 assert_eq!(payload.stage, 1);
1593 assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1594 }
1595
1596 #[test]
1597 fn approval_escalated_carries_run_and_step_ids() {
1598 let run_id = Uuid::now_v7();
1599 let step_id = Uuid::now_v7();
1600 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1601 run_id,
1602 step_id,
1603 step_name: "prod-gate".to_string(),
1604 stage: 0,
1605 policy: "notify".to_string(),
1606 action: "notified 1 target".to_string(),
1607 reason: "approval deadline of 60s expired".to_string(),
1608 assignee: None,
1609 at: Utc::now(),
1610 });
1611
1612 assert_eq!(event.run_id(), Some(run_id));
1613 assert_eq!(event.step_id(), Some(step_id));
1614 assert_eq!(event.user_id(), None);
1615 }
1616}