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 rule_index: Some(0),
1068 condition: Some("payload.amount > 10000".to_string()),
1069 required_approvers: 2,
1070 approver_groups: vec!["finance".to_string()],
1071 evaluated: Vec::new(),
1072 };
1073 let event = Event::ApprovalGranted(ApprovalGrantedEvent {
1074 run_id: Uuid::now_v7(),
1075 step_id: Some(Uuid::now_v7()),
1076 approved_by: "alice".to_string(),
1077 approvals_received: 1,
1078 approvals_required: 2,
1079 requirement: Some(requirement.clone()),
1080 at: Utc::now(),
1081 });
1082
1083 let json = serde_json::to_string(&event).expect("serialize");
1084 let back: Event = serde_json::from_str(&json).expect("deserialize");
1085 let Event::ApprovalGranted(back) = back else {
1086 panic!("expected approval_granted");
1087 };
1088 assert_eq!(back.approvals_received, 1);
1089 assert_eq!(back.approvals_required, 2);
1090 assert_eq!(back.requirement, Some(requirement));
1091 }
1092
1093 #[test]
1094 fn approval_requested_serde_roundtrip() {
1095 let event = Event::ApprovalRequested(ApprovalRequestedEvent {
1096 run_id: Uuid::now_v7(),
1097 step_id: Uuid::now_v7(),
1098 message: "Deploy to prod?".to_string(),
1099 requirement: None,
1100 at: Utc::now(),
1101 });
1102
1103 let json = serde_json::to_string(&event).expect("serialize");
1104 assert!(json.contains("approval_requested"));
1105 }
1106
1107 #[test]
1108 fn log_line_serde_roundtrip() {
1109 let event = Event::LogLine(LogLineEvent {
1110 id: Uuid::now_v7(),
1111 run_id: Uuid::now_v7(),
1112 step_id: Uuid::now_v7(),
1113 step_name: "build".to_string(),
1114 stream: LogStream::Stdout,
1115 line: "Compiling ironflow v0.1.0".to_string(),
1116 at: Utc::now(),
1117 });
1118
1119 let json = serde_json::to_string(&event).expect("serialize");
1120 let back: Event = serde_json::from_str(&json).expect("deserialize");
1121
1122 assert_eq!(back.event_type(), "log_line");
1123 assert!(json.contains("\"type\":\"log_line\""));
1124 assert!(json.contains("Compiling ironflow"));
1125 }
1126
1127 #[test]
1132 fn legacy_flat_json_deserializes_into_typed_payload() {
1133 let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1134 .parse()
1135 .expect("valid uuid");
1136
1137 let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1138 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1139 match event {
1140 Event::RunCreated(e) => {
1141 assert_eq!(e.run_id, run_id);
1142 assert_eq!(e.workflow_name, "deploy");
1143 }
1144 other => panic!("expected RunCreated, got {other:?}"),
1145 }
1146
1147 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"}"#;
1148 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1149 match event {
1150 Event::RunStatusChanged(e) => {
1151 assert_eq!(e.from, RunStatus::Running);
1152 assert_eq!(e.to, RunStatus::Completed);
1153 assert_eq!(e.cost_usd, Decimal::new(5, 1));
1154 assert_eq!(e.duration_ms, 5000);
1155 assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1156 }
1157 other => panic!("expected RunStatusChanged, got {other:?}"),
1158 }
1159
1160 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"}"#;
1162 let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1163 match event {
1164 Event::RunStatusChanged(e) => {
1165 assert!(e.labels.is_empty());
1166 assert_eq!(e.error.as_deref(), Some("boom"));
1167 }
1168 other => panic!("expected RunStatusChanged, got {other:?}"),
1169 }
1170
1171 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"}"#;
1172 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1173 match event {
1174 Event::RunFailed(e) => {
1175 assert_eq!(e.error.as_deref(), Some("boom"));
1176 assert!(e.labels.is_empty());
1177 }
1178 other => panic!("expected RunFailed, got {other:?}"),
1179 }
1180
1181 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"}"#;
1182 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1183 match event {
1184 Event::StepFailed(e) => {
1185 assert_eq!(e.kind, StepKind::Shell);
1186 assert_eq!(e.error, "exit code 1");
1187 }
1188 other => panic!("expected StepFailed, got {other:?}"),
1189 }
1190
1191 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"}"#;
1192 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1193 match event {
1194 Event::LogLine(e) => {
1195 assert_eq!(e.stream, LogStream::Stdout);
1196 assert_eq!(e.line, "hello");
1197 assert_eq!(e.id, Uuid::nil());
1200 }
1201 other => panic!("expected LogLine, got {other:?}"),
1202 }
1203
1204 let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1205 let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1206 match event {
1207 Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1208 other => panic!("expected UserSignedIn, got {other:?}"),
1209 }
1210 }
1211
1212 #[test]
1215 fn serialized_event_is_flat_with_type_tag() {
1216 let run_id = Uuid::now_v7();
1217 let event = Event::RunCreated(RunCreatedEvent {
1218 run_id,
1219 workflow_name: "deploy".to_string(),
1220 at: Utc::now(),
1221 });
1222
1223 let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1224 let object = value.as_object().expect("event serializes to an object");
1225
1226 assert_eq!(
1227 object.get("type").and_then(|v| v.as_str()),
1228 Some("run_created")
1229 );
1230 assert_eq!(
1231 object.get("workflow_name").and_then(|v| v.as_str()),
1232 Some("deploy")
1233 );
1234 assert_eq!(
1235 object.get("run_id").and_then(|v| v.as_str()),
1236 Some(run_id.to_string().as_str())
1237 );
1238 assert!(object.contains_key("at"));
1239 assert_eq!(object.len(), 4, "no nesting: {object:?}");
1240 assert!(!object.contains_key("RunCreated"));
1241 }
1242
1243 #[test]
1244 fn run_id_returns_some_for_run_events() {
1245 let run_id = Uuid::now_v7();
1246 let now = Utc::now();
1247
1248 let events = vec![
1249 Event::RunCreated(RunCreatedEvent {
1250 run_id,
1251 workflow_name: "w".to_string(),
1252 at: now,
1253 }),
1254 Event::RunStatusChanged(RunStatusChangedEvent {
1255 run_id,
1256 workflow_name: "w".to_string(),
1257 from: RunStatus::Pending,
1258 to: RunStatus::Running,
1259 error: None,
1260 cost_usd: Decimal::ZERO,
1261 duration_ms: 0,
1262 labels: HashMap::new(),
1263 at: now,
1264 }),
1265 Event::RunFailed(RunFailedEvent {
1266 run_id,
1267 workflow_name: "w".to_string(),
1268 error: None,
1269 cost_usd: Decimal::ZERO,
1270 duration_ms: 0,
1271 labels: HashMap::new(),
1272 at: now,
1273 }),
1274 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1275 run_id,
1276 workflow_name: "w".to_string(),
1277 limit_usd: Decimal::ZERO,
1278 spent_usd: Decimal::ZERO,
1279 step_budget_usd: Decimal::ZERO,
1280 at: now,
1281 }),
1282 Event::RetryForced(RetryForcedEvent {
1283 run_id,
1284 workflow_name: "w".to_string(),
1285 original_version: "1".to_string(),
1286 current_version: "2".to_string(),
1287 at: now,
1288 }),
1289 Event::StepCompleted(StepCompletedEvent {
1290 run_id,
1291 step_id: Uuid::now_v7(),
1292 step_name: "s".to_string(),
1293 kind: StepKind::Shell,
1294 duration_ms: 0,
1295 cost_usd: Decimal::ZERO,
1296 at: now,
1297 }),
1298 Event::StepFailed(StepFailedEvent {
1299 run_id,
1300 step_id: Uuid::now_v7(),
1301 step_name: "s".to_string(),
1302 kind: StepKind::Shell,
1303 error: "e".to_string(),
1304 at: now,
1305 }),
1306 Event::ApprovalRequested(ApprovalRequestedEvent {
1307 run_id,
1308 step_id: Uuid::now_v7(),
1309 message: "ok?".to_string(),
1310 requirement: None,
1311 at: now,
1312 }),
1313 Event::ApprovalGranted(ApprovalGrantedEvent {
1314 run_id,
1315 step_id: None,
1316 approved_by: "alice".to_string(),
1317 approvals_received: 1,
1318 approvals_required: 1,
1319 requirement: None,
1320 at: now,
1321 }),
1322 Event::ApprovalRejected(ApprovalRejectedEvent {
1323 run_id,
1324 step_id: None,
1325 rejected_by: "bob".to_string(),
1326 requirement: None,
1327 at: now,
1328 }),
1329 Event::LogLine(LogLineEvent {
1330 id: Uuid::now_v7(),
1331 run_id,
1332 step_id: Uuid::now_v7(),
1333 step_name: "s".to_string(),
1334 stream: LogStream::Stdout,
1335 line: "l".to_string(),
1336 at: now,
1337 }),
1338 ];
1339
1340 for event in &events {
1341 assert_eq!(
1342 event.run_id(),
1343 Some(run_id),
1344 "{} should carry a run_id",
1345 event.event_type()
1346 );
1347 }
1348 }
1349
1350 #[test]
1351 fn run_id_returns_none_for_auth_events() {
1352 let user_id = Uuid::now_v7();
1353 let now = Utc::now();
1354
1355 let events = vec![
1356 Event::UserSignedIn(UserSignedInEvent {
1357 user_id,
1358 username: "alice".to_string(),
1359 at: now,
1360 }),
1361 Event::UserSignedUp(UserSignedUpEvent {
1362 user_id,
1363 username: "alice".to_string(),
1364 at: now,
1365 }),
1366 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1367 ];
1368
1369 for event in &events {
1370 assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1371 }
1372 }
1373
1374 #[test]
1375 fn step_id_returns_some_only_for_step_events() {
1376 let step_id = Uuid::now_v7();
1377 let run_id = Uuid::now_v7();
1378 let now = Utc::now();
1379
1380 let with_step = vec![
1381 Event::StepCompleted(StepCompletedEvent {
1382 run_id,
1383 step_id,
1384 step_name: "s".to_string(),
1385 kind: StepKind::Shell,
1386 duration_ms: 0,
1387 cost_usd: Decimal::ZERO,
1388 at: now,
1389 }),
1390 Event::StepFailed(StepFailedEvent {
1391 run_id,
1392 step_id,
1393 step_name: "s".to_string(),
1394 kind: StepKind::Shell,
1395 error: "e".to_string(),
1396 at: now,
1397 }),
1398 Event::ApprovalRequested(ApprovalRequestedEvent {
1399 run_id,
1400 step_id,
1401 message: "ok?".to_string(),
1402 requirement: None,
1403 at: now,
1404 }),
1405 Event::ApprovalGranted(ApprovalGrantedEvent {
1406 run_id,
1407 step_id: Some(step_id),
1408 approved_by: "alice".to_string(),
1409 approvals_received: 1,
1410 approvals_required: 2,
1411 requirement: None,
1412 at: now,
1413 }),
1414 Event::ApprovalRejected(ApprovalRejectedEvent {
1415 run_id,
1416 step_id: Some(step_id),
1417 rejected_by: "bob".to_string(),
1418 requirement: None,
1419 at: now,
1420 }),
1421 ];
1422
1423 for event in &with_step {
1424 assert_eq!(
1425 event.step_id(),
1426 Some(step_id),
1427 "{} should carry a step_id",
1428 event.event_type()
1429 );
1430 }
1431
1432 let without_step = vec![
1433 Event::RunCreated(RunCreatedEvent {
1434 run_id,
1435 workflow_name: "w".to_string(),
1436 at: now,
1437 }),
1438 Event::ApprovalGranted(ApprovalGrantedEvent {
1440 run_id,
1441 step_id: None,
1442 approved_by: "alice".to_string(),
1443 approvals_received: 1,
1444 approvals_required: 1,
1445 requirement: None,
1446 at: now,
1447 }),
1448 Event::LogLine(LogLineEvent {
1451 id: Uuid::now_v7(),
1452 run_id,
1453 step_id,
1454 step_name: "s".to_string(),
1455 stream: LogStream::Stdout,
1456 line: "l".to_string(),
1457 at: now,
1458 }),
1459 Event::UserSignedOut(UserSignedOutEvent {
1460 user_id: Uuid::now_v7(),
1461 at: now,
1462 }),
1463 ];
1464
1465 for event in &without_step {
1466 assert_eq!(
1467 event.step_id(),
1468 None,
1469 "{} should not carry a step_id",
1470 event.event_type()
1471 );
1472 }
1473 }
1474
1475 #[test]
1476 fn user_id_returns_some_only_for_auth_events() {
1477 let user_id = Uuid::now_v7();
1478 let run_id = Uuid::now_v7();
1479 let now = Utc::now();
1480
1481 let auth = vec![
1482 Event::UserSignedIn(UserSignedInEvent {
1483 user_id,
1484 username: "alice".to_string(),
1485 at: now,
1486 }),
1487 Event::UserSignedUp(UserSignedUpEvent {
1488 user_id,
1489 username: "alice".to_string(),
1490 at: now,
1491 }),
1492 Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1493 ];
1494
1495 for event in &auth {
1496 assert_eq!(
1497 event.user_id(),
1498 Some(user_id),
1499 "{} should carry a user_id",
1500 event.event_type()
1501 );
1502 }
1503
1504 let non_auth = vec![
1505 Event::RunCreated(RunCreatedEvent {
1506 run_id,
1507 workflow_name: "w".to_string(),
1508 at: now,
1509 }),
1510 Event::StepFailed(StepFailedEvent {
1511 run_id,
1512 step_id: Uuid::now_v7(),
1513 step_name: "s".to_string(),
1514 kind: StepKind::Shell,
1515 error: "e".to_string(),
1516 at: now,
1517 }),
1518 ];
1519
1520 for event in &non_auth {
1521 assert_eq!(
1522 event.user_id(),
1523 None,
1524 "{} should not carry a user_id",
1525 event.event_type()
1526 );
1527 }
1528 }
1529
1530 #[test]
1531 fn event_type_all_variants() {
1532 let id = Uuid::now_v7();
1533 let now = Utc::now();
1534
1535 let cases: Vec<(Event, &str)> = vec![
1536 (
1537 Event::RunCreated(RunCreatedEvent {
1538 run_id: id,
1539 workflow_name: "w".to_string(),
1540 at: now,
1541 }),
1542 "run_created",
1543 ),
1544 (
1545 Event::RunStatusChanged(RunStatusChangedEvent {
1546 run_id: id,
1547 workflow_name: "w".to_string(),
1548 from: RunStatus::Pending,
1549 to: RunStatus::Running,
1550 error: None,
1551 cost_usd: Decimal::ZERO,
1552 duration_ms: 0,
1553 labels: HashMap::new(),
1554 at: now,
1555 }),
1556 "run_status_changed",
1557 ),
1558 (
1559 Event::RunFailed(RunFailedEvent {
1560 run_id: id,
1561 workflow_name: "w".to_string(),
1562 error: Some("boom".to_string()),
1563 cost_usd: Decimal::ZERO,
1564 duration_ms: 0,
1565 labels: HashMap::new(),
1566 at: now,
1567 }),
1568 "run_failed",
1569 ),
1570 (
1571 Event::RunBudgetExceeded(RunBudgetExceededEvent {
1572 run_id: id,
1573 workflow_name: "w".to_string(),
1574 limit_usd: Decimal::new(200, 2),
1575 spent_usd: Decimal::new(180, 2),
1576 step_budget_usd: Decimal::new(50, 2),
1577 at: now,
1578 }),
1579 "run_budget_exceeded",
1580 ),
1581 (
1582 Event::RetryForced(RetryForcedEvent {
1583 run_id: id,
1584 workflow_name: "w".to_string(),
1585 original_version: "1".to_string(),
1586 current_version: "2".to_string(),
1587 at: now,
1588 }),
1589 "retry_forced",
1590 ),
1591 (
1592 Event::StepCompleted(StepCompletedEvent {
1593 run_id: id,
1594 step_id: id,
1595 step_name: "s".to_string(),
1596 kind: StepKind::Shell,
1597 duration_ms: 0,
1598 cost_usd: Decimal::ZERO,
1599 at: now,
1600 }),
1601 "step_completed",
1602 ),
1603 (
1604 Event::StepFailed(StepFailedEvent {
1605 run_id: id,
1606 step_id: id,
1607 step_name: "s".to_string(),
1608 kind: StepKind::Shell,
1609 error: "err".to_string(),
1610 at: now,
1611 }),
1612 "step_failed",
1613 ),
1614 (
1615 Event::ApprovalRequested(ApprovalRequestedEvent {
1616 run_id: id,
1617 step_id: id,
1618 message: "ok?".to_string(),
1619 requirement: None,
1620 at: now,
1621 }),
1622 "approval_requested",
1623 ),
1624 (
1625 Event::ApprovalGranted(ApprovalGrantedEvent {
1626 run_id: id,
1627 step_id: Some(id),
1628 approved_by: "alice".to_string(),
1629 approvals_received: 1,
1630 approvals_required: 1,
1631 requirement: None,
1632 at: now,
1633 }),
1634 "approval_granted",
1635 ),
1636 (
1637 Event::ApprovalRejected(ApprovalRejectedEvent {
1638 run_id: id,
1639 step_id: Some(id),
1640 rejected_by: "bob".to_string(),
1641 requirement: None,
1642 at: now,
1643 }),
1644 "approval_rejected",
1645 ),
1646 (
1647 Event::ApprovalEscalated(ApprovalEscalatedEvent {
1648 run_id: id,
1649 step_id: id,
1650 step_name: "prod-gate".to_string(),
1651 stage: 0,
1652 policy: "auto_reject".to_string(),
1653 action: "rejected".to_string(),
1654 reason: "approval deadline of 3600s expired".to_string(),
1655 assignee: None,
1656 at: now,
1657 }),
1658 "approval_escalated",
1659 ),
1660 (
1661 Event::LogLine(LogLineEvent {
1662 id,
1663 run_id: id,
1664 step_id: id,
1665 step_name: "build".to_string(),
1666 stream: LogStream::Stdout,
1667 line: "Compiling ironflow v0.1.0".to_string(),
1668 at: now,
1669 }),
1670 "log_line",
1671 ),
1672 (
1673 Event::UserSignedIn(UserSignedInEvent {
1674 user_id: id,
1675 username: "u".to_string(),
1676 at: now,
1677 }),
1678 "user_signed_in",
1679 ),
1680 (
1681 Event::UserSignedUp(UserSignedUpEvent {
1682 user_id: id,
1683 username: "u".to_string(),
1684 at: now,
1685 }),
1686 "user_signed_up",
1687 ),
1688 (
1689 Event::UserSignedOut(UserSignedOutEvent {
1690 user_id: id,
1691 at: now,
1692 }),
1693 "user_signed_out",
1694 ),
1695 ];
1696
1697 assert_eq!(
1698 cases.len(),
1699 Event::ALL.len(),
1700 "every variant must be covered"
1701 );
1702
1703 for (event, expected_type) in cases {
1704 assert_eq!(event.event_type(), expected_type);
1705 }
1706 }
1707
1708 #[test]
1709 fn approval_escalated_serde_roundtrip() {
1710 let run_id = Uuid::now_v7();
1711 let step_id = Uuid::now_v7();
1712 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1713 run_id,
1714 step_id,
1715 step_name: "prod-gate".to_string(),
1716 stage: 1,
1717 policy: "escalate".to_string(),
1718 action: "reassigned to sre-oncall".to_string(),
1719 reason: "approval deadline of 3600s expired".to_string(),
1720 assignee: Some(Assignee::group("sre-oncall")),
1721 at: Utc::now(),
1722 });
1723
1724 let json = serde_json::to_string(&event).expect("serialize");
1725 assert!(
1726 json.contains("\"type\":\"approval_escalated\""),
1727 "got {json}"
1728 );
1729
1730 let back: Event = serde_json::from_str(&json).expect("deserialize");
1731 let Event::ApprovalEscalated(payload) = back else {
1732 panic!("expected an approval_escalated event");
1733 };
1734 assert_eq!(payload.run_id, run_id);
1735 assert_eq!(payload.stage, 1);
1736 assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1737 }
1738
1739 #[test]
1740 fn approval_escalated_carries_run_and_step_ids() {
1741 let run_id = Uuid::now_v7();
1742 let step_id = Uuid::now_v7();
1743 let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1744 run_id,
1745 step_id,
1746 step_name: "prod-gate".to_string(),
1747 stage: 0,
1748 policy: "notify".to_string(),
1749 action: "notified 1 target".to_string(),
1750 reason: "approval deadline of 60s expired".to_string(),
1751 assignee: None,
1752 at: Utc::now(),
1753 });
1754
1755 assert_eq!(event.run_id(), Some(run_id));
1756 assert_eq!(event.step_id(), Some(step_id));
1757 assert_eq!(event.user_id(), None);
1758 }
1759}