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::{ApprovalRequirement, Assignee, RunStatus, StepKind};
12
13fn default_approval_count() -> u32 {
16 1
17}
18
19#[derive(Debug, Clone, Serialize, Deserialize)]
36#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
37pub struct RunCreatedEvent {
38 pub run_id: Uuid,
40 pub workflow_name: String,
42 pub at: DateTime<Utc>,
44}
45
46#[derive(Debug, Clone, Serialize, Deserialize)]
73#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
74pub struct RunStatusChangedEvent {
75 pub run_id: Uuid,
77 pub workflow_name: String,
79 pub from: RunStatus,
81 pub to: RunStatus,
83 pub error: Option<String>,
85 pub cost_usd: Decimal,
87 pub duration_ms: u64,
89 #[serde(default)]
91 pub labels: HashMap<String, String>,
92 pub at: DateTime<Utc>,
94}
95
96#[derive(Debug, Clone, Serialize, Deserialize)]
120#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
121pub struct RunFailedEvent {
122 pub run_id: Uuid,
124 pub workflow_name: String,
126 pub error: Option<String>,
128 pub cost_usd: Decimal,
130 pub duration_ms: u64,
132 #[serde(default)]
134 pub labels: HashMap<String, String>,
135 pub at: DateTime<Utc>,
137}
138
139#[derive(Debug, Clone, Serialize, Deserialize)]
160#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
161pub struct RunBudgetExceededEvent {
162 pub run_id: Uuid,
164 pub workflow_name: String,
166 pub limit_usd: Decimal,
168 pub spent_usd: Decimal,
170 pub step_budget_usd: Decimal,
172 pub at: DateTime<Utc>,
174}
175
176#[derive(Debug, Clone, Serialize, Deserialize)]
195#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
196pub struct RetryForcedEvent {
197 pub run_id: Uuid,
199 pub workflow_name: String,
201 pub original_version: String,
203 pub current_version: String,
205 pub at: DateTime<Utc>,
207}
208
209#[derive(Debug, Clone, Serialize, Deserialize)]
232#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
233pub struct StepCompletedEvent {
234 pub run_id: Uuid,
236 pub step_id: Uuid,
238 pub step_name: String,
240 #[cfg_attr(feature = "openapi", schema(value_type = String))]
242 pub kind: StepKind,
243 pub duration_ms: u64,
245 pub cost_usd: Decimal,
247 pub at: DateTime<Utc>,
249}
250
251#[derive(Debug, Clone, Serialize, Deserialize)]
272#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
273pub struct StepFailedEvent {
274 pub run_id: Uuid,
276 pub step_id: Uuid,
278 pub step_name: String,
280 #[cfg_attr(feature = "openapi", schema(value_type = String))]
282 pub kind: StepKind,
283 pub error: String,
285 pub at: DateTime<Utc>,
287}
288
289#[derive(Debug, Clone, Serialize, Deserialize)]
311#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
312pub struct ApprovalRequestedEvent {
313 pub run_id: Uuid,
315 pub step_id: Uuid,
317 pub message: String,
319 #[serde(default)]
322 pub requirement: Option<ApprovalRequirement>,
323 pub at: DateTime<Utc>,
325}
326
327#[derive(Debug, Clone, Serialize, Deserialize)]
353#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
354pub struct ApprovalGrantedEvent {
355 pub run_id: Uuid,
357 #[serde(default)]
359 pub step_id: Option<Uuid>,
360 pub approved_by: String,
362 #[serde(default = "default_approval_count")]
364 pub approvals_received: u32,
365 #[serde(default = "default_approval_count")]
367 pub approvals_required: u32,
368 #[serde(default)]
370 pub requirement: Option<ApprovalRequirement>,
371 pub at: DateTime<Utc>,
373}
374
375#[derive(Debug, Clone, Serialize, Deserialize)]
394#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
395pub struct ApprovalRejectedEvent {
396 pub run_id: Uuid,
398 #[serde(default)]
400 pub step_id: Option<Uuid>,
401 pub rejected_by: String,
403 #[serde(default)]
405 pub requirement: Option<ApprovalRequirement>,
406 pub at: DateTime<Utc>,
408}
409
410#[derive(Debug, Clone, Serialize, Deserialize)]
438#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
439pub struct ApprovalEscalatedEvent {
440 pub run_id: Uuid,
442 pub step_id: Uuid,
444 pub step_name: String,
446 pub stage: u32,
448 pub policy: String,
450 pub action: String,
452 pub reason: String,
454 #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
456 pub assignee: Option<Assignee>,
457 pub at: DateTime<Utc>,
459}
460
461#[derive(Debug, Clone, Serialize, Deserialize)]
482#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
483pub struct LogLineEvent {
484 #[serde(default)]
493 pub id: Uuid,
494 pub run_id: Uuid,
496 pub step_id: Uuid,
498 pub step_name: String,
500 pub stream: LogStream,
502 pub line: String,
504 pub at: DateTime<Utc>,
506}
507
508#[derive(Debug, Clone, Serialize, Deserialize)]
525#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
526pub struct UserSignedInEvent {
527 pub user_id: Uuid,
529 pub username: String,
531 pub at: DateTime<Utc>,
533}
534
535#[derive(Debug, Clone, Serialize, Deserialize)]
552#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
553pub struct UserSignedUpEvent {
554 pub user_id: Uuid,
556 pub username: String,
558 pub at: DateTime<Utc>,
560}
561
562#[derive(Debug, Clone, Serialize, Deserialize)]
579#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
580pub struct UserSignedOutEvent {
581 pub user_id: Uuid,
583 pub at: DateTime<Utc>,
585}
586
587#[derive(Debug, Clone, Serialize, Deserialize)]
619#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
620#[serde(tag = "type", rename_all = "snake_case")]
621pub enum Event {
622 RunCreated(RunCreatedEvent),
625
626 RunStatusChanged(RunStatusChangedEvent),
628
629 RunFailed(RunFailedEvent),
635
636 RunBudgetExceeded(RunBudgetExceededEvent),
644
645 RetryForced(RetryForcedEvent),
652
653 StepCompleted(StepCompletedEvent),
656
657 StepFailed(StepFailedEvent),
659
660 ApprovalRequested(ApprovalRequestedEvent),
663
664 ApprovalGranted(ApprovalGrantedEvent),
666
667 ApprovalRejected(ApprovalRejectedEvent),
669
670 ApprovalEscalated(ApprovalEscalatedEvent),
672
673 LogLine(LogLineEvent),
679
680 UserSignedIn(UserSignedInEvent),
683
684 UserSignedUp(UserSignedUpEvent),
686
687 UserSignedOut(UserSignedOutEvent),
689}
690
691impl Event {
692 pub const RUN_CREATED: &'static str = "run_created";
694 pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
696 pub const RUN_FAILED: &'static str = "run_failed";
698 pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
700 pub const RETRY_FORCED: &'static str = "retry_forced";
702 pub const STEP_COMPLETED: &'static str = "step_completed";
704 pub const STEP_FAILED: &'static str = "step_failed";
706 pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
708 pub const APPROVAL_GRANTED: &'static str = "approval_granted";
710 pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
712 pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
714 pub const LOG_LINE: &'static str = "log_line";
716 pub const USER_SIGNED_IN: &'static str = "user_signed_in";
718 pub const USER_SIGNED_UP: &'static str = "user_signed_up";
720 pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
722
723 pub const ALL: &'static [&'static str] = &[
739 Self::RUN_CREATED,
740 Self::RUN_STATUS_CHANGED,
741 Self::RUN_FAILED,
742 Self::RUN_BUDGET_EXCEEDED,
743 Self::STEP_COMPLETED,
744 Self::STEP_FAILED,
745 Self::APPROVAL_REQUESTED,
746 Self::APPROVAL_GRANTED,
747 Self::APPROVAL_REJECTED,
748 Self::APPROVAL_ESCALATED,
749 Self::LOG_LINE,
750 Self::USER_SIGNED_IN,
751 Self::USER_SIGNED_UP,
752 Self::USER_SIGNED_OUT,
753 Self::RETRY_FORCED,
754 ];
755
756 #[deny(unreachable_patterns)]
775 pub fn event_type(&self) -> &'static str {
776 match self {
777 Event::RunCreated(_) => Self::RUN_CREATED,
778 Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
779 Event::RunFailed(_) => Self::RUN_FAILED,
780 Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
781 Event::RetryForced(_) => Self::RETRY_FORCED,
782 Event::StepCompleted(_) => Self::STEP_COMPLETED,
783 Event::StepFailed(_) => Self::STEP_FAILED,
784 Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
785 Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
786 Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
787 Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
788 Event::LogLine(_) => Self::LOG_LINE,
789 Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
790 Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
791 Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
792 }
793 }
794
795 #[deny(unreachable_patterns)]
818 pub fn run_id(&self) -> Option<Uuid> {
819 match self {
820 Event::RunCreated(e) => Some(e.run_id),
821 Event::RunStatusChanged(e) => Some(e.run_id),
822 Event::RunFailed(e) => Some(e.run_id),
823 Event::RunBudgetExceeded(e) => Some(e.run_id),
824 Event::RetryForced(e) => Some(e.run_id),
825 Event::StepCompleted(e) => Some(e.run_id),
826 Event::StepFailed(e) => Some(e.run_id),
827 Event::ApprovalRequested(e) => Some(e.run_id),
828 Event::ApprovalGranted(e) => Some(e.run_id),
829 Event::ApprovalRejected(e) => Some(e.run_id),
830 Event::ApprovalEscalated(e) => Some(e.run_id),
831 Event::LogLine(e) => Some(e.run_id),
832 Event::UserSignedIn(_) | Event::UserSignedUp(_) | Event::UserSignedOut(_) => None,
833 }
834 }
835
836 #[deny(unreachable_patterns)]
866 pub fn step_id(&self) -> Option<Uuid> {
867 match self {
868 Event::StepCompleted(e) => Some(e.step_id),
869 Event::StepFailed(e) => Some(e.step_id),
870 Event::ApprovalRequested(e) => Some(e.step_id),
871 Event::ApprovalEscalated(e) => Some(e.step_id),
872 Event::ApprovalGranted(e) => e.step_id,
873 Event::ApprovalRejected(e) => e.step_id,
874 Event::RunCreated(_)
875 | Event::RunStatusChanged(_)
876 | Event::RunFailed(_)
877 | Event::RunBudgetExceeded(_)
878 | Event::RetryForced(_)
879 | Event::LogLine(_)
880 | Event::UserSignedIn(_)
881 | Event::UserSignedUp(_)
882 | Event::UserSignedOut(_) => None,
883 }
884 }
885
886 #[deny(unreachable_patterns)]
909 pub fn user_id(&self) -> Option<Uuid> {
910 match self {
911 Event::UserSignedIn(e) => Some(e.user_id),
912 Event::UserSignedUp(e) => Some(e.user_id),
913 Event::UserSignedOut(e) => Some(e.user_id),
914 Event::RunCreated(_)
915 | Event::RunStatusChanged(_)
916 | Event::RunFailed(_)
917 | Event::RunBudgetExceeded(_)
918 | Event::RetryForced(_)
919 | Event::StepCompleted(_)
920 | Event::StepFailed(_)
921 | Event::ApprovalRequested(_)
922 | Event::ApprovalGranted(_)
923 | Event::ApprovalRejected(_)
924 | Event::ApprovalEscalated(_)
925 | Event::LogLine(_) => None,
926 }
927 }
928}
929
930#[cfg(test)]
931mod tests {
932 use super::*;
933
934 #[test]
935 fn run_status_changed_serde_roundtrip() {
936 let event = Event::RunStatusChanged(RunStatusChangedEvent {
937 run_id: Uuid::now_v7(),
938 workflow_name: "deploy".to_string(),
939 from: RunStatus::Running,
940 to: RunStatus::Completed,
941 error: None,
942 cost_usd: Decimal::new(42, 2),
943 duration_ms: 5000,
944 labels: HashMap::new(),
945 at: Utc::now(),
946 });
947
948 let json = serde_json::to_string(&event).expect("serialize");
949 let back: Event = serde_json::from_str(&json).expect("deserialize");
950
951 assert_eq!(back.event_type(), "run_status_changed");
952 assert!(json.contains("\"type\":\"run_status_changed\""));
953 }
954
955 #[test]
956 fn run_failed_serde_roundtrip() {
957 let event = Event::RunFailed(RunFailedEvent {
958 run_id: Uuid::now_v7(),
959 workflow_name: "deploy".to_string(),
960 error: Some("step crashed".to_string()),
961 cost_usd: Decimal::new(10, 2),
962 duration_ms: 3000,
963 labels: HashMap::new(),
964 at: Utc::now(),
965 });
966
967 let json = serde_json::to_string(&event).expect("serialize");
968 let back: Event = serde_json::from_str(&json).expect("deserialize");
969
970 assert_eq!(back.event_type(), "run_failed");
971 assert!(json.contains("\"type\":\"run_failed\""));
972 assert!(json.contains("step crashed"));
973 }
974
975 #[test]
976 fn run_budget_exceeded_serde_roundtrip() {
977 let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
978 run_id: Uuid::now_v7(),
979 workflow_name: "deploy".to_string(),
980 limit_usd: Decimal::new(200, 2),
981 spent_usd: Decimal::new(180, 2),
982 step_budget_usd: Decimal::new(50, 2),
983 at: Utc::now(),
984 });
985
986 let json = serde_json::to_string(&event).expect("serialize");
987 let back: Event = serde_json::from_str(&json).expect("deserialize");
988
989 assert_eq!(back.event_type(), "run_budget_exceeded");
990 assert!(json.contains("\"type\":\"run_budget_exceeded\""));
991 assert!(json.contains("limit_usd"));
992 assert!(json.contains("step_budget_usd"));
993 }
994
995 #[test]
996 fn all_contains_run_budget_exceeded() {
997 assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
998 }
999
1000 #[test]
1001 fn user_signed_in_serde_roundtrip() {
1002 let event = Event::UserSignedIn(UserSignedInEvent {
1003 user_id: Uuid::now_v7(),
1004 username: "alice".to_string(),
1005 at: Utc::now(),
1006 });
1007
1008 let json = serde_json::to_string(&event).expect("serialize");
1009 let back: Event = serde_json::from_str(&json).expect("deserialize");
1010
1011 assert_eq!(back.event_type(), "user_signed_in");
1012 assert!(json.contains("alice"));
1013 }
1014
1015 #[test]
1016 fn step_failed_serde_roundtrip() {
1017 let event = Event::StepFailed(StepFailedEvent {
1018 run_id: Uuid::now_v7(),
1019 step_id: Uuid::now_v7(),
1020 step_name: "build".to_string(),
1021 kind: StepKind::Shell,
1022 error: "exit code 1".to_string(),
1023 at: Utc::now(),
1024 });
1025
1026 let json = serde_json::to_string(&event).expect("serialize");
1027 let back: Event = serde_json::from_str(&json).expect("deserialize");
1028
1029 assert_eq!(back.event_type(), "step_failed");
1030 }
1031
1032 #[test]
1033 fn legacy_approval_granted_defaults_to_a_single_vote() {
1034 let raw = r#"{"type":"approval_granted","run_id":"01890000-0000-7000-8000-000000000000","approved_by":"alice","at":"2026-01-01T00:00:00Z"}"#;
1035 let event: Event = serde_json::from_str(raw).expect("deserialize");
1036 let Event::ApprovalGranted(event) = event else {
1037 panic!("expected approval_granted");
1038 };
1039
1040 assert_eq!(event.step_id, None);
1041 assert_eq!(event.approvals_received, 1);
1042 assert_eq!(event.approvals_required, 1);
1043 assert!(event.requirement.is_none());
1044 }
1045
1046 #[test]
1047 fn legacy_approval_requested_and_rejected_have_no_requirement() {
1048 let raw = r#"{"type":"approval_requested","run_id":"01890000-0000-7000-8000-000000000000","step_id":"01890000-0000-7000-8000-000000000001","message":"ok?","at":"2026-01-01T00:00:00Z"}"#;
1049 let requested: Event = serde_json::from_str(raw).expect("deserialize");
1050 let Event::ApprovalRequested(requested) = requested else {
1051 panic!("expected approval_requested");
1052 };
1053 assert!(requested.requirement.is_none());
1054
1055 let raw = r#"{"type":"approval_rejected","run_id":"01890000-0000-7000-8000-000000000000","rejected_by":"bob","at":"2026-01-01T00:00:00Z"}"#;
1056 let rejected: Event = serde_json::from_str(raw).expect("deserialize");
1057 let Event::ApprovalRejected(rejected) = rejected else {
1058 panic!("expected approval_rejected");
1059 };
1060 assert_eq!(rejected.step_id, None);
1061 assert!(rejected.requirement.is_none());
1062 }
1063
1064 #[test]
1065 fn approval_granted_roundtrips_the_vote_counts() {
1066 let requirement = ApprovalRequirement {
1067 reason: Some("amount > 10k".to_string()),
1068 required_approvers: 2,
1069 approver_groups: vec!["finance".to_string()],
1070 };
1071 let event = Event::ApprovalGranted(ApprovalGrantedEvent {
1072 run_id: Uuid::now_v7(),
1073 step_id: Some(Uuid::now_v7()),
1074 approved_by: "alice".to_string(),
1075 approvals_received: 1,
1076 approvals_required: 2,
1077 requirement: Some(requirement.clone()),
1078 at: Utc::now(),
1079 });
1080
1081 let json = serde_json::to_string(&event).expect("serialize");
1082 let back: Event = serde_json::from_str(&json).expect("deserialize");
1083 let Event::ApprovalGranted(back) = back else {
1084 panic!("expected approval_granted");
1085 };
1086 assert_eq!(back.approvals_received, 1);
1087 assert_eq!(back.approvals_required, 2);
1088 assert_eq!(back.requirement, Some(requirement));
1089 }
1090
1091 #[test]
1092 fn approval_requested_serde_roundtrip() {
1093 let event = Event::ApprovalRequested(ApprovalRequestedEvent {
1094 run_id: Uuid::now_v7(),
1095 step_id: Uuid::now_v7(),
1096 message: "Deploy to prod?".to_string(),
1097 requirement: None,
1098 at: Utc::now(),
1099 });
1100
1101 let json = serde_json::to_string(&event).expect("serialize");
1102 assert!(json.contains("approval_requested"));
1103 }
1104
1105 #[test]
1106 fn log_line_serde_roundtrip() {
1107 let event = Event::LogLine(LogLineEvent {
1108 id: Uuid::now_v7(),
1109 run_id: Uuid::now_v7(),
1110 step_id: Uuid::now_v7(),
1111 step_name: "build".to_string(),
1112 stream: LogStream::Stdout,
1113 line: "Compiling ironflow v0.1.0".to_string(),
1114 at: Utc::now(),
1115 });
1116
1117 let json = serde_json::to_string(&event).expect("serialize");
1118 let back: Event = serde_json::from_str(&json).expect("deserialize");
1119
1120 assert_eq!(back.event_type(), "log_line");
1121 assert!(json.contains("\"type\":\"log_line\""));
1122 assert!(json.contains("Compiling ironflow"));
1123 }
1124
1125 #[test]
1130 fn legacy_flat_json_deserializes_into_typed_payload() {
1131 let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1132 .parse()
1133 .expect("valid uuid");
1134
1135 let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1136 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1137 match event {
1138 Event::RunCreated(e) => {
1139 assert_eq!(e.run_id, run_id);
1140 assert_eq!(e.workflow_name, "deploy");
1141 }
1142 other => panic!("expected RunCreated, got {other:?}"),
1143 }
1144
1145 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"}"#;
1146 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1147 match event {
1148 Event::RunStatusChanged(e) => {
1149 assert_eq!(e.from, RunStatus::Running);
1150 assert_eq!(e.to, RunStatus::Completed);
1151 assert_eq!(e.cost_usd, Decimal::new(5, 1));
1152 assert_eq!(e.duration_ms, 5000);
1153 assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1154 }
1155 other => panic!("expected RunStatusChanged, got {other:?}"),
1156 }
1157
1158 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"}"#;
1160 let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1161 match event {
1162 Event::RunStatusChanged(e) => {
1163 assert!(e.labels.is_empty());
1164 assert_eq!(e.error.as_deref(), Some("boom"));
1165 }
1166 other => panic!("expected RunStatusChanged, got {other:?}"),
1167 }
1168
1169 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"}"#;
1170 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1171 match event {
1172 Event::RunFailed(e) => {
1173 assert_eq!(e.error.as_deref(), Some("boom"));
1174 assert!(e.labels.is_empty());
1175 }
1176 other => panic!("expected RunFailed, got {other:?}"),
1177 }
1178
1179 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"}"#;
1180 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1181 match event {
1182 Event::StepFailed(e) => {
1183 assert_eq!(e.kind, StepKind::Shell);
1184 assert_eq!(e.error, "exit code 1");
1185 }
1186 other => panic!("expected StepFailed, got {other:?}"),
1187 }
1188
1189 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"}"#;
1190 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1191 match event {
1192 Event::LogLine(e) => {
1193 assert_eq!(e.stream, LogStream::Stdout);
1194 assert_eq!(e.line, "hello");
1195 assert_eq!(e.id, Uuid::nil());
1198 }
1199 other => panic!("expected LogLine, got {other:?}"),
1200 }
1201
1202 let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1203 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1204 match event {
1205 Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1206 other => panic!("expected UserSignedIn, got {other:?}"),
1207 }
1208 }
1209
1210 #[test]
1213 fn serialized_event_is_flat_with_type_tag() {
1214 let run_id = Uuid::now_v7();
1215 let event = Event::RunCreated(RunCreatedEvent {
1216 run_id,
1217 workflow_name: "deploy".to_string(),
1218 at: Utc::now(),
1219 });
1220
1221 let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1222 let object = value.as_object().expect("event serializes to an object");
1223
1224 assert_eq!(
1225 object.get("type").and_then(|v| v.as_str()),
1226 Some("run_created")
1227 );
1228 assert_eq!(
1229 object.get("workflow_name").and_then(|v| v.as_str()),
1230 Some("deploy")
1231 );
1232 assert_eq!(
1233 object.get("run_id").and_then(|v| v.as_str()),
1234 Some(run_id.to_string().as_str())
1235 );
1236 assert!(object.contains_key("at"));
1237 assert_eq!(object.len(), 4, "no nesting: {object:?}");
1238 assert!(!object.contains_key("RunCreated"));
1239 }
1240
1241 #[test]
1242 fn run_id_returns_some_for_run_events() {
1243 let run_id = Uuid::now_v7();
1244 let now = Utc::now();
1245
1246 let events = vec![
1247 Event::RunCreated(RunCreatedEvent {
1248 run_id,
1249 workflow_name: "w".to_string(),
1250 at: now,
1251 }),
1252 Event::RunStatusChanged(RunStatusChangedEvent {
1253 run_id,
1254 workflow_name: "w".to_string(),
1255 from: RunStatus::Pending,
1256 to: RunStatus::Running,
1257 error: None,
1258 cost_usd: Decimal::ZERO,
1259 duration_ms: 0,
1260 labels: HashMap::new(),
1261 at: now,
1262 }),
1263 Event::RunFailed(RunFailedEvent {
1264 run_id,
1265 workflow_name: "w".to_string(),
1266 error: None,
1267 cost_usd: Decimal::ZERO,
1268 duration_ms: 0,
1269 labels: HashMap::new(),
1270 at: now,
1271 }),
1272 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1273 run_id,
1274 workflow_name: "w".to_string(),
1275 limit_usd: Decimal::ZERO,
1276 spent_usd: Decimal::ZERO,
1277 step_budget_usd: Decimal::ZERO,
1278 at: now,
1279 }),
1280 Event::RetryForced(RetryForcedEvent {
1281 run_id,
1282 workflow_name: "w".to_string(),
1283 original_version: "1".to_string(),
1284 current_version: "2".to_string(),
1285 at: now,
1286 }),
1287 Event::StepCompleted(StepCompletedEvent {
1288 run_id,
1289 step_id: Uuid::now_v7(),
1290 step_name: "s".to_string(),
1291 kind: StepKind::Shell,
1292 duration_ms: 0,
1293 cost_usd: Decimal::ZERO,
1294 at: now,
1295 }),
1296 Event::StepFailed(StepFailedEvent {
1297 run_id,
1298 step_id: Uuid::now_v7(),
1299 step_name: "s".to_string(),
1300 kind: StepKind::Shell,
1301 error: "e".to_string(),
1302 at: now,
1303 }),
1304 Event::ApprovalRequested(ApprovalRequestedEvent {
1305 run_id,
1306 step_id: Uuid::now_v7(),
1307 message: "ok?".to_string(),
1308 requirement: None,
1309 at: now,
1310 }),
1311 Event::ApprovalGranted(ApprovalGrantedEvent {
1312 run_id,
1313 step_id: None,
1314 approved_by: "alice".to_string(),
1315 approvals_received: 1,
1316 approvals_required: 1,
1317 requirement: None,
1318 at: now,
1319 }),
1320 Event::ApprovalRejected(ApprovalRejectedEvent {
1321 run_id,
1322 step_id: None,
1323 rejected_by: "bob".to_string(),
1324 requirement: None,
1325 at: now,
1326 }),
1327 Event::LogLine(LogLineEvent {
1328 id: Uuid::now_v7(),
1329 run_id,
1330 step_id: Uuid::now_v7(),
1331 step_name: "s".to_string(),
1332 stream: LogStream::Stdout,
1333 line: "l".to_string(),
1334 at: now,
1335 }),
1336 ];
1337
1338 for event in &events {
1339 assert_eq!(
1340 event.run_id(),
1341 Some(run_id),
1342 "{} should carry a run_id",
1343 event.event_type()
1344 );
1345 }
1346 }
1347
1348 #[test]
1349 fn run_id_returns_none_for_auth_events() {
1350 let user_id = Uuid::now_v7();
1351 let now = Utc::now();
1352
1353 let events = vec![
1354 Event::UserSignedIn(UserSignedInEvent {
1355 user_id,
1356 username: "alice".to_string(),
1357 at: now,
1358 }),
1359 Event::UserSignedUp(UserSignedUpEvent {
1360 user_id,
1361 username: "alice".to_string(),
1362 at: now,
1363 }),
1364 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1365 ];
1366
1367 for event in &events {
1368 assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1369 }
1370 }
1371
1372 #[test]
1373 fn step_id_returns_some_only_for_step_events() {
1374 let step_id = Uuid::now_v7();
1375 let run_id = Uuid::now_v7();
1376 let now = Utc::now();
1377
1378 let with_step = vec![
1379 Event::StepCompleted(StepCompletedEvent {
1380 run_id,
1381 step_id,
1382 step_name: "s".to_string(),
1383 kind: StepKind::Shell,
1384 duration_ms: 0,
1385 cost_usd: Decimal::ZERO,
1386 at: now,
1387 }),
1388 Event::StepFailed(StepFailedEvent {
1389 run_id,
1390 step_id,
1391 step_name: "s".to_string(),
1392 kind: StepKind::Shell,
1393 error: "e".to_string(),
1394 at: now,
1395 }),
1396 Event::ApprovalRequested(ApprovalRequestedEvent {
1397 run_id,
1398 step_id,
1399 message: "ok?".to_string(),
1400 requirement: None,
1401 at: now,
1402 }),
1403 Event::ApprovalGranted(ApprovalGrantedEvent {
1404 run_id,
1405 step_id: Some(step_id),
1406 approved_by: "alice".to_string(),
1407 approvals_received: 1,
1408 approvals_required: 2,
1409 requirement: None,
1410 at: now,
1411 }),
1412 Event::ApprovalRejected(ApprovalRejectedEvent {
1413 run_id,
1414 step_id: Some(step_id),
1415 rejected_by: "bob".to_string(),
1416 requirement: None,
1417 at: now,
1418 }),
1419 ];
1420
1421 for event in &with_step {
1422 assert_eq!(
1423 event.step_id(),
1424 Some(step_id),
1425 "{} should carry a step_id",
1426 event.event_type()
1427 );
1428 }
1429
1430 let without_step = vec![
1431 Event::RunCreated(RunCreatedEvent {
1432 run_id,
1433 workflow_name: "w".to_string(),
1434 at: now,
1435 }),
1436 Event::ApprovalGranted(ApprovalGrantedEvent {
1438 run_id,
1439 step_id: None,
1440 approved_by: "alice".to_string(),
1441 approvals_received: 1,
1442 approvals_required: 1,
1443 requirement: None,
1444 at: now,
1445 }),
1446 Event::LogLine(LogLineEvent {
1449 id: Uuid::now_v7(),
1450 run_id,
1451 step_id,
1452 step_name: "s".to_string(),
1453 stream: LogStream::Stdout,
1454 line: "l".to_string(),
1455 at: now,
1456 }),
1457 Event::UserSignedOut(UserSignedOutEvent {
1458 user_id: Uuid::now_v7(),
1459 at: now,
1460 }),
1461 ];
1462
1463 for event in &without_step {
1464 assert_eq!(
1465 event.step_id(),
1466 None,
1467 "{} should not carry a step_id",
1468 event.event_type()
1469 );
1470 }
1471 }
1472
1473 #[test]
1474 fn user_id_returns_some_only_for_auth_events() {
1475 let user_id = Uuid::now_v7();
1476 let run_id = Uuid::now_v7();
1477 let now = Utc::now();
1478
1479 let auth = vec![
1480 Event::UserSignedIn(UserSignedInEvent {
1481 user_id,
1482 username: "alice".to_string(),
1483 at: now,
1484 }),
1485 Event::UserSignedUp(UserSignedUpEvent {
1486 user_id,
1487 username: "alice".to_string(),
1488 at: now,
1489 }),
1490 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1491 ];
1492
1493 for event in &auth {
1494 assert_eq!(
1495 event.user_id(),
1496 Some(user_id),
1497 "{} should carry a user_id",
1498 event.event_type()
1499 );
1500 }
1501
1502 let non_auth = vec![
1503 Event::RunCreated(RunCreatedEvent {
1504 run_id,
1505 workflow_name: "w".to_string(),
1506 at: now,
1507 }),
1508 Event::StepFailed(StepFailedEvent {
1509 run_id,
1510 step_id: Uuid::now_v7(),
1511 step_name: "s".to_string(),
1512 kind: StepKind::Shell,
1513 error: "e".to_string(),
1514 at: now,
1515 }),
1516 ];
1517
1518 for event in &non_auth {
1519 assert_eq!(
1520 event.user_id(),
1521 None,
1522 "{} should not carry a user_id",
1523 event.event_type()
1524 );
1525 }
1526 }
1527
1528 #[test]
1529 fn event_type_all_variants() {
1530 let id = Uuid::now_v7();
1531 let now = Utc::now();
1532
1533 let cases: Vec<(Event, &str)> = vec![
1534 (
1535 Event::RunCreated(RunCreatedEvent {
1536 run_id: id,
1537 workflow_name: "w".to_string(),
1538 at: now,
1539 }),
1540 "run_created",
1541 ),
1542 (
1543 Event::RunStatusChanged(RunStatusChangedEvent {
1544 run_id: id,
1545 workflow_name: "w".to_string(),
1546 from: RunStatus::Pending,
1547 to: RunStatus::Running,
1548 error: None,
1549 cost_usd: Decimal::ZERO,
1550 duration_ms: 0,
1551 labels: HashMap::new(),
1552 at: now,
1553 }),
1554 "run_status_changed",
1555 ),
1556 (
1557 Event::RunFailed(RunFailedEvent {
1558 run_id: id,
1559 workflow_name: "w".to_string(),
1560 error: Some("boom".to_string()),
1561 cost_usd: Decimal::ZERO,
1562 duration_ms: 0,
1563 labels: HashMap::new(),
1564 at: now,
1565 }),
1566 "run_failed",
1567 ),
1568 (
1569 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1570 run_id: id,
1571 workflow_name: "w".to_string(),
1572 limit_usd: Decimal::new(200, 2),
1573 spent_usd: Decimal::new(180, 2),
1574 step_budget_usd: Decimal::new(50, 2),
1575 at: now,
1576 }),
1577 "run_budget_exceeded",
1578 ),
1579 (
1580 Event::RetryForced(RetryForcedEvent {
1581 run_id: id,
1582 workflow_name: "w".to_string(),
1583 original_version: "1".to_string(),
1584 current_version: "2".to_string(),
1585 at: now,
1586 }),
1587 "retry_forced",
1588 ),
1589 (
1590 Event::StepCompleted(StepCompletedEvent {
1591 run_id: id,
1592 step_id: id,
1593 step_name: "s".to_string(),
1594 kind: StepKind::Shell,
1595 duration_ms: 0,
1596 cost_usd: Decimal::ZERO,
1597 at: now,
1598 }),
1599 "step_completed",
1600 ),
1601 (
1602 Event::StepFailed(StepFailedEvent {
1603 run_id: id,
1604 step_id: id,
1605 step_name: "s".to_string(),
1606 kind: StepKind::Shell,
1607 error: "err".to_string(),
1608 at: now,
1609 }),
1610 "step_failed",
1611 ),
1612 (
1613 Event::ApprovalRequested(ApprovalRequestedEvent {
1614 run_id: id,
1615 step_id: id,
1616 message: "ok?".to_string(),
1617 requirement: None,
1618 at: now,
1619 }),
1620 "approval_requested",
1621 ),
1622 (
1623 Event::ApprovalGranted(ApprovalGrantedEvent {
1624 run_id: id,
1625 step_id: Some(id),
1626 approved_by: "alice".to_string(),
1627 approvals_received: 1,
1628 approvals_required: 1,
1629 requirement: None,
1630 at: now,
1631 }),
1632 "approval_granted",
1633 ),
1634 (
1635 Event::ApprovalRejected(ApprovalRejectedEvent {
1636 run_id: id,
1637 step_id: Some(id),
1638 rejected_by: "bob".to_string(),
1639 requirement: None,
1640 at: now,
1641 }),
1642 "approval_rejected",
1643 ),
1644 (
1645 Event::ApprovalEscalated(ApprovalEscalatedEvent {
1646 run_id: id,
1647 step_id: id,
1648 step_name: "prod-gate".to_string(),
1649 stage: 0,
1650 policy: "auto_reject".to_string(),
1651 action: "rejected".to_string(),
1652 reason: "approval deadline of 3600s expired".to_string(),
1653 assignee: None,
1654 at: now,
1655 }),
1656 "approval_escalated",
1657 ),
1658 (
1659 Event::LogLine(LogLineEvent {
1660 id,
1661 run_id: id,
1662 step_id: id,
1663 step_name: "build".to_string(),
1664 stream: LogStream::Stdout,
1665 line: "Compiling ironflow v0.1.0".to_string(),
1666 at: now,
1667 }),
1668 "log_line",
1669 ),
1670 (
1671 Event::UserSignedIn(UserSignedInEvent {
1672 user_id: id,
1673 username: "u".to_string(),
1674 at: now,
1675 }),
1676 "user_signed_in",
1677 ),
1678 (
1679 Event::UserSignedUp(UserSignedUpEvent {
1680 user_id: id,
1681 username: "u".to_string(),
1682 at: now,
1683 }),
1684 "user_signed_up",
1685 ),
1686 (
1687 Event::UserSignedOut(UserSignedOutEvent {
1688 user_id: id,
1689 at: now,
1690 }),
1691 "user_signed_out",
1692 ),
1693 ];
1694
1695 assert_eq!(
1696 cases.len(),
1697 Event::ALL.len(),
1698 "every variant must be covered"
1699 );
1700
1701 for (event, expected_type) in cases {
1702 assert_eq!(event.event_type(), expected_type);
1703 }
1704 }
1705
1706 #[test]
1707 fn approval_escalated_serde_roundtrip() {
1708 let run_id = Uuid::now_v7();
1709 let step_id = Uuid::now_v7();
1710 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1711 run_id,
1712 step_id,
1713 step_name: "prod-gate".to_string(),
1714 stage: 1,
1715 policy: "escalate".to_string(),
1716 action: "reassigned to sre-oncall".to_string(),
1717 reason: "approval deadline of 3600s expired".to_string(),
1718 assignee: Some(Assignee::group("sre-oncall")),
1719 at: Utc::now(),
1720 });
1721
1722 let json = serde_json::to_string(&event).expect("serialize");
1723 assert!(
1724 json.contains("\"type\":\"approval_escalated\""),
1725 "got {json}"
1726 );
1727
1728 let back: Event = serde_json::from_str(&json).expect("deserialize");
1729 let Event::ApprovalEscalated(payload) = back else {
1730 panic!("expected an approval_escalated event");
1731 };
1732 assert_eq!(payload.run_id, run_id);
1733 assert_eq!(payload.stage, 1);
1734 assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1735 }
1736
1737 #[test]
1738 fn approval_escalated_carries_run_and_step_ids() {
1739 let run_id = Uuid::now_v7();
1740 let step_id = Uuid::now_v7();
1741 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1742 run_id,
1743 step_id,
1744 step_name: "prod-gate".to_string(),
1745 stage: 0,
1746 policy: "notify".to_string(),
1747 action: "notified 1 target".to_string(),
1748 reason: "approval deadline of 60s expired".to_string(),
1749 assignee: None,
1750 at: Utc::now(),
1751 });
1752
1753 assert_eq!(event.run_id(), Some(run_id));
1754 assert_eq!(event.step_id(), Some(step_id));
1755 assert_eq!(event.user_id(), None);
1756 }
1757}