Skip to main content

ironflow_engine/notify/
event.rs

1//! Domain events emitted throughout the ironflow lifecycle.
2
3use 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::{
12    ApprovalRequirement, Assignee, ProviderAccountWindow, RunStatus, ScheduleMissReason, StepKind,
13};
14
15/// Vote counts assumed for approval events serialized before multi-approver
16/// gates existed: one approval was always enough.
17fn default_approval_count() -> u32 {
18    1
19}
20
21/// Payload of the `Event::RunCreated` event.
22///
23/// # Examples
24///
25/// ```
26/// use chrono::Utc;
27/// use ironflow_engine::notify::RunCreatedEvent;
28/// use uuid::Uuid;
29///
30/// let payload = RunCreatedEvent {
31///     run_id: Uuid::now_v7(),
32///     workflow_name: "deploy".to_string(),
33///     at: Utc::now(),
34/// };
35/// assert_eq!(payload.workflow_name, "deploy");
36/// ```
37#[derive(Debug, Clone, Serialize, Deserialize)]
38#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
39pub struct RunCreatedEvent {
40    /// Run identifier.
41    pub run_id: Uuid,
42    /// Workflow name.
43    pub workflow_name: String,
44    /// When the run was created.
45    pub at: DateTime<Utc>,
46}
47
48/// Payload of the `Event::RunStatusChanged` event.
49///
50/// # Examples
51///
52/// ```
53/// use std::collections::HashMap;
54///
55/// use chrono::Utc;
56/// use ironflow_engine::notify::RunStatusChangedEvent;
57/// use ironflow_store::models::RunStatus;
58/// use rust_decimal::Decimal;
59/// use uuid::Uuid;
60///
61/// let payload = RunStatusChangedEvent {
62///     run_id: Uuid::now_v7(),
63///     workflow_name: "deploy".to_string(),
64///     from: RunStatus::Running,
65///     to: RunStatus::Completed,
66///     error: None,
67///     cost_usd: Decimal::ZERO,
68///     duration_ms: 5000,
69///     labels: HashMap::new(),
70///     at: Utc::now(),
71/// };
72/// assert_eq!(payload.to, RunStatus::Completed);
73/// ```
74#[derive(Debug, Clone, Serialize, Deserialize)]
75#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
76pub struct RunStatusChangedEvent {
77    /// Run identifier.
78    pub run_id: Uuid,
79    /// Workflow name.
80    pub workflow_name: String,
81    /// Previous status.
82    pub from: RunStatus,
83    /// New status.
84    pub to: RunStatus,
85    /// Error message (when transitioning to Failed).
86    pub error: Option<String>,
87    /// Aggregated cost in USD at the time of transition.
88    pub cost_usd: Decimal,
89    /// Aggregated duration in milliseconds at the time of transition.
90    pub duration_ms: u64,
91    /// Labels of the run at the time of the transition.
92    #[serde(default)]
93    pub labels: HashMap<String, String>,
94    /// When the transition occurred.
95    pub at: DateTime<Utc>,
96}
97
98/// Payload of the `Event::RunFailed` event.
99///
100/// # Examples
101///
102/// ```
103/// use std::collections::HashMap;
104///
105/// use chrono::Utc;
106/// use ironflow_engine::notify::RunFailedEvent;
107/// use rust_decimal::Decimal;
108/// use uuid::Uuid;
109///
110/// let payload = RunFailedEvent {
111///     run_id: Uuid::now_v7(),
112///     workflow_name: "deploy".to_string(),
113///     error: Some("step crashed".to_string()),
114///     cost_usd: Decimal::ZERO,
115///     duration_ms: 3000,
116///     labels: HashMap::new(),
117///     at: Utc::now(),
118/// };
119/// assert_eq!(payload.error.as_deref(), Some("step crashed"));
120/// ```
121#[derive(Debug, Clone, Serialize, Deserialize)]
122#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
123pub struct RunFailedEvent {
124    /// Run identifier.
125    pub run_id: Uuid,
126    /// Workflow name.
127    pub workflow_name: String,
128    /// Error message.
129    pub error: Option<String>,
130    /// Aggregated cost in USD at the time of failure.
131    pub cost_usd: Decimal,
132    /// Aggregated duration in milliseconds at the time of failure.
133    pub duration_ms: u64,
134    /// Labels of the run at the time of the failure.
135    #[serde(default)]
136    pub labels: HashMap<String, String>,
137    /// When the failure occurred.
138    pub at: DateTime<Utc>,
139}
140
141/// Payload of the `Event::RunBudgetExceeded` event.
142///
143/// # Examples
144///
145/// ```
146/// use chrono::Utc;
147/// use ironflow_engine::notify::RunBudgetExceededEvent;
148/// use rust_decimal::Decimal;
149/// use uuid::Uuid;
150///
151/// let payload = RunBudgetExceededEvent {
152///     run_id: Uuid::now_v7(),
153///     workflow_name: "deploy".to_string(),
154///     limit_usd: Decimal::new(200, 2),
155///     spent_usd: Decimal::new(180, 2),
156///     step_budget_usd: Decimal::new(50, 2),
157///     at: Utc::now(),
158/// };
159/// assert_eq!(payload.limit_usd, Decimal::new(200, 2));
160/// ```
161#[derive(Debug, Clone, Serialize, Deserialize)]
162#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
163pub struct RunBudgetExceededEvent {
164    /// Run identifier.
165    pub run_id: Uuid,
166    /// Workflow name.
167    pub workflow_name: String,
168    /// The configured cost cap in USD.
169    pub limit_usd: Decimal,
170    /// Cost already consumed when the cap was reached, in USD.
171    pub spent_usd: Decimal,
172    /// Declared budget of the refused step, in USD.
173    pub step_budget_usd: Decimal,
174    /// When the refusal occurred.
175    pub at: DateTime<Utc>,
176}
177
178/// Payload of the `Event::RetryForced` event.
179///
180/// # Examples
181///
182/// ```
183/// use chrono::Utc;
184/// use ironflow_engine::notify::RetryForcedEvent;
185/// use uuid::Uuid;
186///
187/// let payload = RetryForcedEvent {
188///     run_id: Uuid::now_v7(),
189///     workflow_name: "deploy".to_string(),
190///     original_version: "1".to_string(),
191///     current_version: "2".to_string(),
192///     at: Utc::now(),
193/// };
194/// assert_eq!(payload.current_version, "2");
195/// ```
196#[derive(Debug, Clone, Serialize, Deserialize)]
197#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
198pub struct RetryForcedEvent {
199    /// The new run created by the forced retry.
200    pub run_id: Uuid,
201    /// Workflow name.
202    pub workflow_name: String,
203    /// Version stored on the original run.
204    pub original_version: String,
205    /// Current version of the handler.
206    pub current_version: String,
207    /// When the forced retry occurred.
208    pub at: DateTime<Utc>,
209}
210
211/// Payload of the `Event::StepCompleted` event.
212///
213/// # Examples
214///
215/// ```
216/// use chrono::Utc;
217/// use ironflow_engine::notify::StepCompletedEvent;
218/// use ironflow_store::models::StepKind;
219/// use rust_decimal::Decimal;
220/// use uuid::Uuid;
221///
222/// let payload = StepCompletedEvent {
223///     run_id: Uuid::now_v7(),
224///     step_id: Uuid::now_v7(),
225///     step_name: "build".to_string(),
226///     kind: StepKind::Shell,
227///     duration_ms: 1200,
228///     cost_usd: Decimal::ZERO,
229///     at: Utc::now(),
230/// };
231/// assert_eq!(payload.step_name, "build");
232/// ```
233#[derive(Debug, Clone, Serialize, Deserialize)]
234#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
235pub struct StepCompletedEvent {
236    /// Run identifier.
237    pub run_id: Uuid,
238    /// Step identifier.
239    pub step_id: Uuid,
240    /// Human-readable step name.
241    pub step_name: String,
242    /// Step operation kind.
243    #[cfg_attr(feature = "openapi", schema(value_type = String))]
244    pub kind: StepKind,
245    /// Step duration in milliseconds.
246    pub duration_ms: u64,
247    /// Step cost in USD.
248    pub cost_usd: Decimal,
249    /// When the step completed.
250    pub at: DateTime<Utc>,
251}
252
253/// Payload of the `Event::StepFailed` event.
254///
255/// # Examples
256///
257/// ```
258/// use chrono::Utc;
259/// use ironflow_engine::notify::StepFailedEvent;
260/// use ironflow_store::models::StepKind;
261/// use uuid::Uuid;
262///
263/// let payload = StepFailedEvent {
264///     run_id: Uuid::now_v7(),
265///     step_id: Uuid::now_v7(),
266///     step_name: "build".to_string(),
267///     kind: StepKind::Shell,
268///     error: "exit code 1".to_string(),
269///     at: Utc::now(),
270/// };
271/// assert_eq!(payload.error, "exit code 1");
272/// ```
273#[derive(Debug, Clone, Serialize, Deserialize)]
274#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
275pub struct StepFailedEvent {
276    /// Run identifier.
277    pub run_id: Uuid,
278    /// Step identifier.
279    pub step_id: Uuid,
280    /// Human-readable step name.
281    pub step_name: String,
282    /// Step operation kind.
283    #[cfg_attr(feature = "openapi", schema(value_type = String))]
284    pub kind: StepKind,
285    /// Error message.
286    pub error: String,
287    /// When the step failed.
288    pub at: DateTime<Utc>,
289}
290
291/// Payload of the `Event::ApprovalRequested` event.
292///
293/// Published when an approval gate opens. Carries the approvers the handler
294/// required, if it set any.
295///
296/// # Examples
297///
298/// ```
299/// use chrono::Utc;
300/// use ironflow_engine::notify::ApprovalRequestedEvent;
301/// use uuid::Uuid;
302///
303/// let payload = ApprovalRequestedEvent {
304///     run_id: Uuid::now_v7(),
305///     step_id: Uuid::now_v7(),
306///     message: "Deploy to prod?".to_string(),
307///     requirement: None,
308///     at: Utc::now(),
309/// };
310/// assert_eq!(payload.message, "Deploy to prod?");
311/// ```
312#[derive(Debug, Clone, Serialize, Deserialize)]
313#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
314pub struct ApprovalRequestedEvent {
315    /// Run identifier.
316    pub run_id: Uuid,
317    /// Approval step identifier.
318    pub step_id: Uuid,
319    /// Message displayed to reviewers.
320    pub message: String,
321    /// Approvers the handler required. `None` for a gate opened without
322    /// approvers: one approval resolves it.
323    #[serde(default)]
324    pub requirement: Option<ApprovalRequirement>,
325    /// When the approval was requested.
326    pub at: DateTime<Utc>,
327}
328
329/// Payload of the `Event::ApprovalGranted` event.
330///
331/// Published for every vote cast on a gate. A vote with
332/// `approvals_received < approvals_required` is recorded but does not resolve
333/// the gate: the run stays `AwaitingApproval` until enough distinct approvers
334/// voted.
335///
336/// # Examples
337///
338/// ```
339/// use chrono::Utc;
340/// use ironflow_engine::notify::ApprovalGrantedEvent;
341/// use uuid::Uuid;
342///
343/// let payload = ApprovalGrantedEvent {
344///     run_id: Uuid::now_v7(),
345///     step_id: Some(Uuid::now_v7()),
346///     approved_by: "alice".to_string(),
347///     approvals_received: 1,
348///     approvals_required: 2,
349///     requirement: None,
350///     at: Utc::now(),
351/// };
352/// assert!(payload.approvals_received < payload.approvals_required);
353/// ```
354#[derive(Debug, Clone, Serialize, Deserialize)]
355#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
356pub struct ApprovalGrantedEvent {
357    /// Run identifier.
358    pub run_id: Uuid,
359    /// Approval step identifier. `None` in events recorded before it existed.
360    #[serde(default)]
361    pub step_id: Option<Uuid>,
362    /// User who approved (ID or username).
363    pub approved_by: String,
364    /// Distinct approvals recorded on the gate, this one included.
365    #[serde(default = "default_approval_count")]
366    pub approvals_received: u32,
367    /// Distinct approvals needed to resolve the gate.
368    #[serde(default = "default_approval_count")]
369    pub approvals_required: u32,
370    /// Approvers the handler required, if any.
371    #[serde(default)]
372    pub requirement: Option<ApprovalRequirement>,
373    /// When the approval was granted.
374    pub at: DateTime<Utc>,
375}
376
377/// Payload of the `Event::ApprovalRejected` event.
378///
379/// # Examples
380///
381/// ```
382/// use chrono::Utc;
383/// use ironflow_engine::notify::ApprovalRejectedEvent;
384/// use uuid::Uuid;
385///
386/// let payload = ApprovalRejectedEvent {
387///     run_id: Uuid::now_v7(),
388///     step_id: Some(Uuid::now_v7()),
389///     rejected_by: "bob".to_string(),
390///     requirement: None,
391///     at: Utc::now(),
392/// };
393/// assert_eq!(payload.rejected_by, "bob");
394/// ```
395#[derive(Debug, Clone, Serialize, Deserialize)]
396#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
397pub struct ApprovalRejectedEvent {
398    /// Run identifier.
399    pub run_id: Uuid,
400    /// Approval step identifier. `None` in events recorded before it existed.
401    #[serde(default)]
402    pub step_id: Option<Uuid>,
403    /// User who rejected (ID or username).
404    pub rejected_by: String,
405    /// Approvers the handler required, if any.
406    #[serde(default)]
407    pub requirement: Option<ApprovalRequirement>,
408    /// When the rejection occurred.
409    pub at: DateTime<Utc>,
410}
411
412/// Payload of the `Event::ApprovalEscalated` event.
413///
414/// Emitted every time an approval gate misses its SLA deadline, including the
415/// repeated firings of a bare `Notify`/`Escalate` policy and the final
416/// "chain exhausted" notice. The audit log persists it verbatim, so the whole
417/// escalation history of a gate is reconstructable from it.
418///
419/// # Examples
420///
421/// ```
422/// use chrono::Utc;
423/// use ironflow_engine::notify::ApprovalEscalatedEvent;
424/// use uuid::Uuid;
425///
426/// let payload = ApprovalEscalatedEvent {
427///     run_id: Uuid::now_v7(),
428///     step_id: Uuid::now_v7(),
429///     step_name: "prod-gate".to_string(),
430///     stage: 0,
431///     policy: "auto_reject".to_string(),
432///     action: "rejected".to_string(),
433///     reason: "approval deadline of 3600s expired".to_string(),
434///     assignee: None,
435///     at: Utc::now(),
436/// };
437/// assert_eq!(payload.policy, "auto_reject");
438/// ```
439#[derive(Debug, Clone, Serialize, Deserialize)]
440#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
441pub struct ApprovalEscalatedEvent {
442    /// Run identifier.
443    pub run_id: Uuid,
444    /// Approval step identifier.
445    pub step_id: Uuid,
446    /// Human-readable step name.
447    pub step_name: String,
448    /// Escalation stage that fired (0-based index into the policy chain).
449    pub stage: u32,
450    /// Policy applied, e.g. `"auto_reject"`, `"notify"`, `"escalate"`.
451    pub policy: String,
452    /// What the escalation did, for the audit log.
453    pub action: String,
454    /// Why it fired, e.g. `"approval deadline of 3600s expired"`.
455    pub reason: String,
456    /// Assignee after the escalation, when it reassigned the gate.
457    #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
458    pub assignee: Option<Assignee>,
459    /// When the escalation ran.
460    pub at: DateTime<Utc>,
461}
462
463/// Payload of the `Event::LogLine` event.
464///
465/// # Examples
466///
467/// ```
468/// use chrono::Utc;
469/// use ironflow_engine::notify::{LogLineEvent, LogStream};
470/// use uuid::Uuid;
471///
472/// let payload = LogLineEvent {
473///     id: Uuid::now_v7(),
474///     run_id: Uuid::now_v7(),
475///     step_id: Uuid::now_v7(),
476///     step_name: "build".to_string(),
477///     stream: LogStream::Stdout,
478///     line: "Compiling ironflow v0.1.0".to_string(),
479///     at: Utc::now(),
480/// };
481/// assert_eq!(payload.line, "Compiling ironflow v0.1.0");
482/// ```
483#[derive(Debug, Clone, Serialize, Deserialize)]
484#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
485pub struct LogLineEvent {
486    /// Persisted entry identifier (UUID v7, time-ordered).
487    ///
488    /// Matches the `id` of the log entry later stored for this line, so clients
489    /// can de-duplicate the live SSE stream against the persisted history
490    /// fetched from `GET /runs/:id/logs`.
491    ///
492    /// Defaults to the nil UUID when absent, so a payload emitted by an older
493    /// producer still deserializes instead of dropping the whole event.
494    #[serde(default)]
495    pub id: Uuid,
496    /// Run identifier.
497    pub run_id: Uuid,
498    /// Step identifier.
499    pub step_id: Uuid,
500    /// Human-readable step name.
501    pub step_name: String,
502    /// Output stream.
503    pub stream: LogStream,
504    /// The log line content.
505    pub line: String,
506    /// When the line was emitted.
507    pub at: DateTime<Utc>,
508}
509
510/// Payload of the `Event::UserSignedIn` event.
511///
512/// # Examples
513///
514/// ```
515/// use chrono::Utc;
516/// use ironflow_engine::notify::UserSignedInEvent;
517/// use uuid::Uuid;
518///
519/// let payload = UserSignedInEvent {
520///     user_id: Uuid::now_v7(),
521///     username: "alice".to_string(),
522///     at: Utc::now(),
523/// };
524/// assert_eq!(payload.username, "alice");
525/// ```
526#[derive(Debug, Clone, Serialize, Deserialize)]
527#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
528pub struct UserSignedInEvent {
529    /// User identifier.
530    pub user_id: Uuid,
531    /// Username.
532    pub username: String,
533    /// When the sign-in occurred.
534    pub at: DateTime<Utc>,
535}
536
537/// Payload of the `Event::UserSignedUp` event.
538///
539/// # Examples
540///
541/// ```
542/// use chrono::Utc;
543/// use ironflow_engine::notify::UserSignedUpEvent;
544/// use uuid::Uuid;
545///
546/// let payload = UserSignedUpEvent {
547///     user_id: Uuid::now_v7(),
548///     username: "alice".to_string(),
549///     at: Utc::now(),
550/// };
551/// assert_eq!(payload.username, "alice");
552/// ```
553#[derive(Debug, Clone, Serialize, Deserialize)]
554#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
555pub struct UserSignedUpEvent {
556    /// User identifier.
557    pub user_id: Uuid,
558    /// Username.
559    pub username: String,
560    /// When the sign-up occurred.
561    pub at: DateTime<Utc>,
562}
563
564/// Payload of the `Event::UserSignedOut` event.
565///
566/// # Examples
567///
568/// ```
569/// use chrono::Utc;
570/// use ironflow_engine::notify::UserSignedOutEvent;
571/// use uuid::Uuid;
572///
573/// let user_id = Uuid::now_v7();
574/// let payload = UserSignedOutEvent {
575///     user_id,
576///     at: Utc::now(),
577/// };
578/// assert_eq!(payload.user_id, user_id);
579/// ```
580#[derive(Debug, Clone, Serialize, Deserialize)]
581#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
582pub struct UserSignedOutEvent {
583    /// User identifier.
584    pub user_id: Uuid,
585    /// When the sign-out occurred.
586    pub at: DateTime<Utc>,
587}
588
589/// What happened to a Provider Account in a
590/// [`ProviderAccountUpdatedEvent`].
591///
592/// # Examples
593///
594/// ```
595/// use ironflow_engine::notify::ProviderAccountChange;
596///
597/// assert_eq!(ProviderAccountChange::TokenReplaced.as_str(), "token_replaced");
598/// ```
599#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
600#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
601#[serde(rename_all = "snake_case")]
602pub enum ProviderAccountChange {
603    /// The account was created.
604    Created,
605    /// Settings of the account changed.
606    Updated,
607    /// The account was deleted.
608    Deleted,
609    /// The credential was replaced.
610    TokenReplaced,
611}
612
613impl ProviderAccountChange {
614    /// Wire name of the change, as stored in audit payloads.
615    pub fn as_str(&self) -> &'static str {
616        match self {
617            Self::Created => "created",
618            Self::Updated => "updated",
619            Self::Deleted => "deleted",
620            Self::TokenReplaced => "token_replaced",
621        }
622    }
623}
624
625/// Payload of the `Event::ProviderAccountUpdated` event.
626///
627/// Never carries the credential.
628///
629/// # Examples
630///
631/// ```
632/// use chrono::Utc;
633/// use ironflow_engine::notify::{ProviderAccountChange, ProviderAccountUpdatedEvent};
634/// use uuid::Uuid;
635///
636/// let payload = ProviderAccountUpdatedEvent {
637///     account_id: Uuid::now_v7(),
638///     name: "perso".to_string(),
639///     change: ProviderAccountChange::Created,
640///     at: Utc::now(),
641/// };
642/// assert_eq!(payload.name, "perso");
643/// ```
644#[derive(Debug, Clone, Serialize, Deserialize)]
645#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
646pub struct ProviderAccountUpdatedEvent {
647    /// Account identifier.
648    pub account_id: Uuid,
649    /// Account name.
650    pub name: String,
651    /// What changed.
652    pub change: ProviderAccountChange,
653    /// When the change happened.
654    pub at: DateTime<Utc>,
655}
656
657/// Payload of the `Event::ProviderAccountUsageUpdated` event.
658///
659/// # Examples
660///
661/// ```
662/// use chrono::Utc;
663/// use ironflow_engine::notify::ProviderAccountUsageUpdatedEvent;
664/// use uuid::Uuid;
665///
666/// let payload = ProviderAccountUsageUpdatedEvent {
667///     account_id: Uuid::now_v7(),
668///     name: "perso".to_string(),
669///     windows: Vec::new(),
670///     at: Utc::now(),
671/// };
672/// assert!(payload.windows.is_empty());
673/// ```
674#[derive(Debug, Clone, Serialize, Deserialize)]
675#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
676pub struct ProviderAccountUsageUpdatedEvent {
677    /// Account identifier.
678    pub account_id: Uuid,
679    /// Account name.
680    pub name: String,
681    /// Current windows of the account.
682    pub windows: Vec<ProviderAccountWindow>,
683    /// When the windows were recorded.
684    pub at: DateTime<Utc>,
685}
686
687/// Payload of the `Event::SignalAwaited` event.
688///
689/// Emitted when a run suspends on `ctx.wait_for_signal` and goes `Sleeping`.
690///
691/// # Examples
692///
693/// ```
694/// use chrono::{TimeDelta, Utc};
695/// use ironflow_engine::notify::SignalAwaitedEvent;
696/// use uuid::Uuid;
697///
698/// let now = Utc::now();
699/// let payload = SignalAwaitedEvent {
700///     run_id: Uuid::now_v7(),
701///     step_id: Uuid::now_v7(),
702///     step_name: "wait-ci".to_string(),
703///     name: "ci.pipeline_finished".to_string(),
704///     key: "4f2a9c1".to_string(),
705///     deadline_at: now + TimeDelta::hours(1),
706///     at: now,
707/// };
708/// assert_eq!(payload.key, "4f2a9c1");
709/// ```
710#[derive(Debug, Clone, Serialize, Deserialize)]
711#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
712pub struct SignalAwaitedEvent {
713    /// Run identifier.
714    pub run_id: Uuid,
715    /// The signal step the run waits on.
716    pub step_id: Uuid,
717    /// Name of the signal step.
718    pub step_name: String,
719    /// Awaited signal name.
720    pub name: String,
721    /// Awaited occurrence key.
722    pub key: String,
723    /// When the wait times out.
724    pub deadline_at: DateTime<Utc>,
725    /// When the run started waiting.
726    pub at: DateTime<Utc>,
727}
728
729/// Payload of the `Event::SignalReceived` event.
730///
731/// Emitted once per stored signal; a duplicate (same idempotency ID) emits
732/// nothing.
733///
734/// # Examples
735///
736/// ```
737/// use chrono::Utc;
738/// use ironflow_engine::notify::SignalReceivedEvent;
739/// use uuid::Uuid;
740///
741/// let payload = SignalReceivedEvent {
742///     signal_id: Uuid::now_v7(),
743///     name: "ci.pipeline_finished".to_string(),
744///     key: "4f2a9c1".to_string(),
745///     resumed_runs: vec![Uuid::now_v7()],
746///     at: Utc::now(),
747/// };
748/// assert_eq!(payload.resumed_runs.len(), 1);
749/// ```
750#[derive(Debug, Clone, Serialize, Deserialize)]
751#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
752pub struct SignalReceivedEvent {
753    /// Stored signal identifier.
754    pub signal_id: Uuid,
755    /// Signal name.
756    pub name: String,
757    /// Occurrence key.
758    pub key: String,
759    /// Runs whose waiting step the signal resolved.
760    pub resumed_runs: Vec<Uuid>,
761    /// When the signal was received.
762    pub at: DateTime<Utc>,
763}
764
765/// Payload of the `Event::ScheduleOccurrencesMissed` event.
766///
767/// A cron schedule skipped some of its occurrences: the server was down
768/// longer than the catch-up window allows, the catch-up policy dropped them,
769/// or the previous run was still active under an overlap `skip` policy.
770///
771/// # Examples
772///
773/// ```
774/// use chrono::{TimeDelta, Utc};
775/// use ironflow_engine::notify::ScheduleOccurrencesMissedEvent;
776/// use ironflow_store::models::ScheduleMissReason;
777/// use uuid::Uuid;
778///
779/// let now = Utc::now();
780/// let payload = ScheduleOccurrencesMissedEvent {
781///     schedule_id: Uuid::now_v7(),
782///     workflow_name: "nightly-report".to_string(),
783///     reason: ScheduleMissReason::OutsideWindow,
784///     count: 3,
785///     first: now - TimeDelta::hours(3),
786///     last: now - TimeDelta::hours(1),
787///     at: now,
788/// };
789/// assert_eq!(payload.count, 3);
790/// ```
791#[derive(Debug, Clone, Serialize, Deserialize)]
792#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
793pub struct ScheduleOccurrencesMissedEvent {
794    /// Schedule that missed the occurrences.
795    pub schedule_id: Uuid,
796    /// Workflow the schedule starts.
797    pub workflow_name: String,
798    /// Why the occurrences were not run.
799    pub reason: ScheduleMissReason,
800    /// Number of missed occurrences.
801    pub count: u64,
802    /// Earliest missed occurrence.
803    pub first: DateTime<Utc>,
804    /// Latest missed occurrence.
805    pub last: DateTime<Utc>,
806    /// When the ticker noticed the miss.
807    pub at: DateTime<Utc>,
808}
809
810/// A domain event emitted by the ironflow system.
811///
812/// Covers the full lifecycle: runs, steps, approvals, and authentication.
813/// Subscribers receive these via [`EventPublisher`](super::EventPublisher)
814/// and pattern-match on the variants they care about.
815///
816/// Each variant wraps a dedicated payload struct. The serialized form stays
817/// flat: the `type` discriminant sits next to the payload fields, so
818/// `{"type":"run_created","run_id":...}` round-trips unchanged.
819///
820/// # Examples
821///
822/// ```
823/// use std::collections::HashMap;
824/// use ironflow_engine::notify::{Event, RunStatusChangedEvent};
825/// use ironflow_store::models::RunStatus;
826/// use uuid::Uuid;
827///
828/// let event = Event::RunStatusChanged(RunStatusChangedEvent {
829///     run_id: Uuid::now_v7(),
830///     workflow_name: "deploy".to_string(),
831///     from: RunStatus::Running,
832///     to: RunStatus::Completed,
833///     error: None,
834///     cost_usd: rust_decimal::Decimal::ZERO,
835///     duration_ms: 5000,
836///     labels: HashMap::new(),
837///     at: chrono::Utc::now(),
838/// });
839/// assert_eq!(event.event_type(), "run_status_changed");
840/// ```
841#[derive(Debug, Clone, Serialize, Deserialize)]
842#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
843#[serde(tag = "type", rename_all = "snake_case")]
844pub enum Event {
845    // -- Run lifecycle --
846    /// A new run was created (status: Pending).
847    RunCreated(RunCreatedEvent),
848
849    /// A run changed status.
850    RunStatusChanged(RunStatusChangedEvent),
851
852    /// A run transitioned to [`Failed`](ironflow_store::models::RunStatus::Failed).
853    ///
854    /// This is a convenience event emitted alongside [`RunStatusChanged`](Event::RunStatusChanged)
855    /// when the target status is `Failed`. Subscribe to this instead of
856    /// `RUN_STATUS_CHANGED` when you only care about failures.
857    RunFailed(RunFailedEvent),
858
859    /// A run was stopped because it reached its cumulative cost cap.
860    ///
861    /// Emitted when the engine refuses an agent step that would cross the run's
862    /// `max_cost_usd`. The run transitions to
863    /// [`Cancelled`](ironflow_store::models::RunStatus::Cancelled) and the step
864    /// is never launched, so the reported spend is what the run had already
865    /// consumed.
866    RunBudgetExceeded(RunBudgetExceededEvent),
867
868    /// A manual retry was forced despite a handler version mismatch.
869    ///
870    /// Emitted when a caller passes `force=true` on a retry where the
871    /// handler version differs from the run's recorded version. This is
872    /// an audit event: it means the new code will execute on the old
873    /// payload without the handler explicitly declaring compatibility.
874    RetryForced(RetryForcedEvent),
875
876    // -- Step lifecycle --
877    /// A step completed successfully.
878    StepCompleted(StepCompletedEvent),
879
880    /// A step failed.
881    StepFailed(StepFailedEvent),
882
883    // -- Approval --
884    /// A run is waiting for human approval.
885    ApprovalRequested(ApprovalRequestedEvent),
886
887    /// A run was approved by a human.
888    ApprovalGranted(ApprovalGrantedEvent),
889
890    /// A run was rejected by a human.
891    ApprovalRejected(ApprovalRejectedEvent),
892
893    /// An approval gate missed its SLA deadline and an escalation policy ran.
894    ApprovalEscalated(ApprovalEscalatedEvent),
895
896    // -- Log streaming --
897    /// A log line emitted during step execution.
898    ///
899    /// Pushed by the worker in real time so that SSE clients can stream
900    /// step output as it happens, without waiting for step completion.
901    LogLine(LogLineEvent),
902
903    // -- Authentication --
904    /// A user signed in.
905    UserSignedIn(UserSignedInEvent),
906
907    /// A new user signed up.
908    UserSignedUp(UserSignedUpEvent),
909
910    /// A user signed out.
911    UserSignedOut(UserSignedOutEvent),
912
913    // -- Provider Accounts --
914    /// A Provider Account was created, updated, deleted or had its token replaced.
915    #[serde(rename = "provider_account.updated")]
916    ProviderAccountUpdated(ProviderAccountUpdatedEvent),
917
918    /// New usage windows were recorded for a Provider Account.
919    #[serde(rename = "provider_account.usage_updated")]
920    ProviderAccountUsageUpdated(ProviderAccountUsageUpdatedEvent),
921
922    // -- Signals --
923    /// A run started waiting for a signal.
924    SignalAwaited(SignalAwaitedEvent),
925
926    /// A signal was received.
927    SignalReceived(SignalReceivedEvent),
928
929    // -- Schedules --
930    /// A cron schedule dropped occurrences instead of starting a run for them.
931    ///
932    /// Emitted by the schedule ticker once per schedule, tick and reason.
933    ScheduleOccurrencesMissed(ScheduleOccurrencesMissedEvent),
934}
935
936impl Event {
937    /// Event type constant for [`RunCreated`](Event::RunCreated).
938    pub const RUN_CREATED: &'static str = "run_created";
939    /// Event type constant for [`RunStatusChanged`](Event::RunStatusChanged).
940    pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
941    /// Event type constant for [`RunFailed`](Event::RunFailed).
942    pub const RUN_FAILED: &'static str = "run_failed";
943    /// Event type constant for [`RunBudgetExceeded`](Event::RunBudgetExceeded).
944    pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
945    /// Event type constant for [`RetryForced`](Event::RetryForced).
946    pub const RETRY_FORCED: &'static str = "retry_forced";
947    /// Event type constant for [`StepCompleted`](Event::StepCompleted).
948    pub const STEP_COMPLETED: &'static str = "step_completed";
949    /// Event type constant for [`StepFailed`](Event::StepFailed).
950    pub const STEP_FAILED: &'static str = "step_failed";
951    /// Event type constant for [`ApprovalRequested`](Event::ApprovalRequested).
952    pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
953    /// Event type constant for [`ApprovalGranted`](Event::ApprovalGranted).
954    pub const APPROVAL_GRANTED: &'static str = "approval_granted";
955    /// Event type constant for [`ApprovalRejected`](Event::ApprovalRejected).
956    pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
957    /// Event type constant for [`ApprovalEscalated`](Event::ApprovalEscalated).
958    pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
959    /// Event type constant for [`LogLine`](Event::LogLine).
960    pub const LOG_LINE: &'static str = "log_line";
961    /// Event type constant for [`UserSignedIn`](Event::UserSignedIn).
962    pub const USER_SIGNED_IN: &'static str = "user_signed_in";
963    /// Event type constant for [`UserSignedUp`](Event::UserSignedUp).
964    pub const USER_SIGNED_UP: &'static str = "user_signed_up";
965    /// Event type constant for [`UserSignedOut`](Event::UserSignedOut).
966    pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
967    /// Event type constant for [`ProviderAccountUpdated`](Event::ProviderAccountUpdated).
968    pub const PROVIDER_ACCOUNT_UPDATED: &'static str = "provider_account.updated";
969    /// Event type constant for
970    /// [`ProviderAccountUsageUpdated`](Event::ProviderAccountUsageUpdated).
971    pub const PROVIDER_ACCOUNT_USAGE_UPDATED: &'static str = "provider_account.usage_updated";
972    /// Event type constant for [`SignalAwaited`](Event::SignalAwaited).
973    pub const SIGNAL_AWAITED: &'static str = "signal_awaited";
974    /// Event type constant for [`SignalReceived`](Event::SignalReceived).
975    pub const SIGNAL_RECEIVED: &'static str = "signal_received";
976    /// Event type constant for
977    /// [`ScheduleOccurrencesMissed`](Event::ScheduleOccurrencesMissed).
978    pub const SCHEDULE_OCCURRENCES_MISSED: &'static str = "schedule_occurrences_missed";
979
980    /// All event types. Pass this to
981    /// [`EventPublisher::subscribe`](super::EventPublisher::subscribe) to
982    /// receive every event.
983    ///
984    /// # Examples
985    ///
986    /// ```no_run
987    /// use ironflow_engine::notify::{Event, EventPublisher, WebhookSubscriber};
988    ///
989    /// let mut publisher = EventPublisher::new();
990    /// publisher.subscribe(
991    ///     WebhookSubscriber::new("https://example.com/all"),
992    ///     Event::ALL,
993    /// );
994    /// ```
995    pub const ALL: &'static [&'static str] = &[
996        Self::RUN_CREATED,
997        Self::RUN_STATUS_CHANGED,
998        Self::RUN_FAILED,
999        Self::RUN_BUDGET_EXCEEDED,
1000        Self::STEP_COMPLETED,
1001        Self::STEP_FAILED,
1002        Self::APPROVAL_REQUESTED,
1003        Self::APPROVAL_GRANTED,
1004        Self::APPROVAL_REJECTED,
1005        Self::APPROVAL_ESCALATED,
1006        Self::LOG_LINE,
1007        Self::USER_SIGNED_IN,
1008        Self::USER_SIGNED_UP,
1009        Self::USER_SIGNED_OUT,
1010        Self::RETRY_FORCED,
1011        Self::PROVIDER_ACCOUNT_UPDATED,
1012        Self::PROVIDER_ACCOUNT_USAGE_UPDATED,
1013        Self::SIGNAL_AWAITED,
1014        Self::SIGNAL_RECEIVED,
1015        Self::SCHEDULE_OCCURRENCES_MISSED,
1016    ];
1017
1018    /// Returns the event type as a static string (e.g. `"run_status_changed"`).
1019    ///
1020    /// Useful for filtering and logging without deserializing.
1021    ///
1022    /// # Examples
1023    ///
1024    /// ```
1025    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
1026    /// use uuid::Uuid;
1027    /// use chrono::Utc;
1028    ///
1029    /// let event = Event::UserSignedIn(UserSignedInEvent {
1030    ///     user_id: Uuid::now_v7(),
1031    ///     username: "alice".to_string(),
1032    ///     at: Utc::now(),
1033    /// });
1034    /// assert_eq!(event.event_type(), "user_signed_in");
1035    /// ```
1036    #[deny(unreachable_patterns)]
1037    pub fn event_type(&self) -> &'static str {
1038        match self {
1039            Event::RunCreated(_) => Self::RUN_CREATED,
1040            Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
1041            Event::RunFailed(_) => Self::RUN_FAILED,
1042            Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
1043            Event::RetryForced(_) => Self::RETRY_FORCED,
1044            Event::StepCompleted(_) => Self::STEP_COMPLETED,
1045            Event::StepFailed(_) => Self::STEP_FAILED,
1046            Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
1047            Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
1048            Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
1049            Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
1050            Event::LogLine(_) => Self::LOG_LINE,
1051            Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
1052            Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
1053            Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
1054            Event::ProviderAccountUpdated(_) => Self::PROVIDER_ACCOUNT_UPDATED,
1055            Event::ProviderAccountUsageUpdated(_) => Self::PROVIDER_ACCOUNT_USAGE_UPDATED,
1056            Event::SignalAwaited(_) => Self::SIGNAL_AWAITED,
1057            Event::SignalReceived(_) => Self::SIGNAL_RECEIVED,
1058            Event::ScheduleOccurrencesMissed(_) => Self::SCHEDULE_OCCURRENCES_MISSED,
1059        }
1060    }
1061
1062    /// Returns the run this event belongs to, if any.
1063    ///
1064    /// Auth events ([`UserSignedIn`](Event::UserSignedIn),
1065    /// [`UserSignedUp`](Event::UserSignedUp),
1066    /// [`UserSignedOut`](Event::UserSignedOut)) are not tied to a run and
1067    /// return `None`.
1068    ///
1069    /// # Examples
1070    ///
1071    /// ```
1072    /// use ironflow_engine::notify::{Event, RunCreatedEvent};
1073    /// use uuid::Uuid;
1074    /// use chrono::Utc;
1075    ///
1076    /// let run_id = Uuid::now_v7();
1077    /// let event = Event::RunCreated(RunCreatedEvent {
1078    ///     run_id,
1079    ///     workflow_name: "deploy".to_string(),
1080    ///     at: Utc::now(),
1081    /// });
1082    /// assert_eq!(event.run_id(), Some(run_id));
1083    /// ```
1084    #[deny(unreachable_patterns)]
1085    pub fn run_id(&self) -> Option<Uuid> {
1086        match self {
1087            Event::RunCreated(e) => Some(e.run_id),
1088            Event::RunStatusChanged(e) => Some(e.run_id),
1089            Event::RunFailed(e) => Some(e.run_id),
1090            Event::RunBudgetExceeded(e) => Some(e.run_id),
1091            Event::RetryForced(e) => Some(e.run_id),
1092            Event::StepCompleted(e) => Some(e.run_id),
1093            Event::StepFailed(e) => Some(e.run_id),
1094            Event::ApprovalRequested(e) => Some(e.run_id),
1095            Event::ApprovalGranted(e) => Some(e.run_id),
1096            Event::ApprovalRejected(e) => Some(e.run_id),
1097            Event::ApprovalEscalated(e) => Some(e.run_id),
1098            Event::LogLine(e) => Some(e.run_id),
1099            Event::SignalAwaited(e) => Some(e.run_id),
1100            Event::UserSignedIn(_)
1101            | Event::UserSignedUp(_)
1102            | Event::UserSignedOut(_)
1103            | Event::ProviderAccountUpdated(_)
1104            | Event::ProviderAccountUsageUpdated(_)
1105            | Event::SignalReceived(_)
1106            | Event::ScheduleOccurrencesMissed(_) => None,
1107        }
1108    }
1109
1110    /// Returns the step this event belongs to, if any.
1111    ///
1112    /// Only [`StepCompleted`](Event::StepCompleted),
1113    /// [`StepFailed`](Event::StepFailed),
1114    /// [`ApprovalRequested`](Event::ApprovalRequested) and
1115    /// [`ApprovalEscalated`](Event::ApprovalEscalated) always carry a step
1116    /// identifier; [`ApprovalGranted`](Event::ApprovalGranted) and
1117    /// [`ApprovalRejected`](Event::ApprovalRejected) carry one when it was
1118    /// recorded; every other variant returns `None`.
1119    ///
1120    /// # Examples
1121    ///
1122    /// ```
1123    /// use ironflow_engine::notify::{Event, StepFailedEvent};
1124    /// use ironflow_store::models::StepKind;
1125    /// use uuid::Uuid;
1126    /// use chrono::Utc;
1127    ///
1128    /// let step_id = Uuid::now_v7();
1129    /// let event = Event::StepFailed(StepFailedEvent {
1130    ///     run_id: Uuid::now_v7(),
1131    ///     step_id,
1132    ///     step_name: "build".to_string(),
1133    ///     kind: StepKind::Shell,
1134    ///     error: "exit code 1".to_string(),
1135    ///     at: Utc::now(),
1136    /// });
1137    /// assert_eq!(event.step_id(), Some(step_id));
1138    /// ```
1139    #[deny(unreachable_patterns)]
1140    pub fn step_id(&self) -> Option<Uuid> {
1141        match self {
1142            Event::StepCompleted(e) => Some(e.step_id),
1143            Event::StepFailed(e) => Some(e.step_id),
1144            Event::ApprovalRequested(e) => Some(e.step_id),
1145            Event::ApprovalEscalated(e) => Some(e.step_id),
1146            Event::ApprovalGranted(e) => e.step_id,
1147            Event::ApprovalRejected(e) => e.step_id,
1148            Event::SignalAwaited(e) => Some(e.step_id),
1149            Event::RunCreated(_)
1150            | Event::RunStatusChanged(_)
1151            | Event::RunFailed(_)
1152            | Event::RunBudgetExceeded(_)
1153            | Event::RetryForced(_)
1154            | Event::LogLine(_)
1155            | Event::UserSignedIn(_)
1156            | Event::UserSignedUp(_)
1157            | Event::UserSignedOut(_)
1158            | Event::ProviderAccountUpdated(_)
1159            | Event::ProviderAccountUsageUpdated(_)
1160            | Event::SignalReceived(_)
1161            | Event::ScheduleOccurrencesMissed(_) => None,
1162        }
1163    }
1164
1165    /// Returns the user this event belongs to, if any.
1166    ///
1167    /// Only the auth events ([`UserSignedIn`](Event::UserSignedIn),
1168    /// [`UserSignedUp`](Event::UserSignedUp),
1169    /// [`UserSignedOut`](Event::UserSignedOut)) carry a user identifier;
1170    /// every other variant returns `None`.
1171    ///
1172    /// # Examples
1173    ///
1174    /// ```
1175    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
1176    /// use uuid::Uuid;
1177    /// use chrono::Utc;
1178    ///
1179    /// let user_id = Uuid::now_v7();
1180    /// let event = Event::UserSignedIn(UserSignedInEvent {
1181    ///     user_id,
1182    ///     username: "alice".to_string(),
1183    ///     at: Utc::now(),
1184    /// });
1185    /// assert_eq!(event.user_id(), Some(user_id));
1186    /// ```
1187    #[deny(unreachable_patterns)]
1188    pub fn user_id(&self) -> Option<Uuid> {
1189        match self {
1190            Event::UserSignedIn(e) => Some(e.user_id),
1191            Event::UserSignedUp(e) => Some(e.user_id),
1192            Event::UserSignedOut(e) => Some(e.user_id),
1193            Event::RunCreated(_)
1194            | Event::RunStatusChanged(_)
1195            | Event::RunFailed(_)
1196            | Event::RunBudgetExceeded(_)
1197            | Event::RetryForced(_)
1198            | Event::StepCompleted(_)
1199            | Event::StepFailed(_)
1200            | Event::ApprovalRequested(_)
1201            | Event::ApprovalGranted(_)
1202            | Event::ApprovalRejected(_)
1203            | Event::ApprovalEscalated(_)
1204            | Event::LogLine(_)
1205            | Event::ProviderAccountUpdated(_)
1206            | Event::ProviderAccountUsageUpdated(_)
1207            | Event::SignalAwaited(_)
1208            | Event::SignalReceived(_)
1209            | Event::ScheduleOccurrencesMissed(_) => None,
1210        }
1211    }
1212}
1213
1214#[cfg(test)]
1215mod tests {
1216    use super::*;
1217
1218    #[test]
1219    fn run_status_changed_serde_roundtrip() {
1220        let event = Event::RunStatusChanged(RunStatusChangedEvent {
1221            run_id: Uuid::now_v7(),
1222            workflow_name: "deploy".to_string(),
1223            from: RunStatus::Running,
1224            to: RunStatus::Completed,
1225            error: None,
1226            cost_usd: Decimal::new(42, 2),
1227            duration_ms: 5000,
1228            labels: HashMap::new(),
1229            at: Utc::now(),
1230        });
1231
1232        let json = serde_json::to_string(&event).expect("serialize");
1233        let back: Event = serde_json::from_str(&json).expect("deserialize");
1234
1235        assert_eq!(back.event_type(), "run_status_changed");
1236        assert!(json.contains("\"type\":\"run_status_changed\""));
1237    }
1238
1239    #[test]
1240    fn run_failed_serde_roundtrip() {
1241        let event = Event::RunFailed(RunFailedEvent {
1242            run_id: Uuid::now_v7(),
1243            workflow_name: "deploy".to_string(),
1244            error: Some("step crashed".to_string()),
1245            cost_usd: Decimal::new(10, 2),
1246            duration_ms: 3000,
1247            labels: HashMap::new(),
1248            at: Utc::now(),
1249        });
1250
1251        let json = serde_json::to_string(&event).expect("serialize");
1252        let back: Event = serde_json::from_str(&json).expect("deserialize");
1253
1254        assert_eq!(back.event_type(), "run_failed");
1255        assert!(json.contains("\"type\":\"run_failed\""));
1256        assert!(json.contains("step crashed"));
1257    }
1258
1259    #[test]
1260    fn run_budget_exceeded_serde_roundtrip() {
1261        let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
1262            run_id: Uuid::now_v7(),
1263            workflow_name: "deploy".to_string(),
1264            limit_usd: Decimal::new(200, 2),
1265            spent_usd: Decimal::new(180, 2),
1266            step_budget_usd: Decimal::new(50, 2),
1267            at: Utc::now(),
1268        });
1269
1270        let json = serde_json::to_string(&event).expect("serialize");
1271        let back: Event = serde_json::from_str(&json).expect("deserialize");
1272
1273        assert_eq!(back.event_type(), "run_budget_exceeded");
1274        assert!(json.contains("\"type\":\"run_budget_exceeded\""));
1275        assert!(json.contains("limit_usd"));
1276        assert!(json.contains("step_budget_usd"));
1277    }
1278
1279    #[test]
1280    fn all_contains_run_budget_exceeded() {
1281        assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
1282    }
1283
1284    #[test]
1285    fn user_signed_in_serde_roundtrip() {
1286        let event = Event::UserSignedIn(UserSignedInEvent {
1287            user_id: Uuid::now_v7(),
1288            username: "alice".to_string(),
1289            at: Utc::now(),
1290        });
1291
1292        let json = serde_json::to_string(&event).expect("serialize");
1293        let back: Event = serde_json::from_str(&json).expect("deserialize");
1294
1295        assert_eq!(back.event_type(), "user_signed_in");
1296        assert!(json.contains("alice"));
1297    }
1298
1299    #[test]
1300    fn step_failed_serde_roundtrip() {
1301        let event = Event::StepFailed(StepFailedEvent {
1302            run_id: Uuid::now_v7(),
1303            step_id: Uuid::now_v7(),
1304            step_name: "build".to_string(),
1305            kind: StepKind::Shell,
1306            error: "exit code 1".to_string(),
1307            at: Utc::now(),
1308        });
1309
1310        let json = serde_json::to_string(&event).expect("serialize");
1311        let back: Event = serde_json::from_str(&json).expect("deserialize");
1312
1313        assert_eq!(back.event_type(), "step_failed");
1314    }
1315
1316    #[test]
1317    fn legacy_approval_granted_defaults_to_a_single_vote() {
1318        let raw = r#"{"type":"approval_granted","run_id":"01890000-0000-7000-8000-000000000000","approved_by":"alice","at":"2026-01-01T00:00:00Z"}"#;
1319        let event: Event = serde_json::from_str(raw).expect("deserialize");
1320        let Event::ApprovalGranted(event) = event else {
1321            panic!("expected approval_granted");
1322        };
1323
1324        assert_eq!(event.step_id, None);
1325        assert_eq!(event.approvals_received, 1);
1326        assert_eq!(event.approvals_required, 1);
1327        assert!(event.requirement.is_none());
1328    }
1329
1330    #[test]
1331    fn legacy_approval_requested_and_rejected_have_no_requirement() {
1332        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"}"#;
1333        let requested: Event = serde_json::from_str(raw).expect("deserialize");
1334        let Event::ApprovalRequested(requested) = requested else {
1335            panic!("expected approval_requested");
1336        };
1337        assert!(requested.requirement.is_none());
1338
1339        let raw = r#"{"type":"approval_rejected","run_id":"01890000-0000-7000-8000-000000000000","rejected_by":"bob","at":"2026-01-01T00:00:00Z"}"#;
1340        let rejected: Event = serde_json::from_str(raw).expect("deserialize");
1341        let Event::ApprovalRejected(rejected) = rejected else {
1342            panic!("expected approval_rejected");
1343        };
1344        assert_eq!(rejected.step_id, None);
1345        assert!(rejected.requirement.is_none());
1346    }
1347
1348    #[test]
1349    fn approval_granted_roundtrips_the_vote_counts() {
1350        let requirement = ApprovalRequirement {
1351            reason: Some("amount > 10k".to_string()),
1352            required_approvers: 2,
1353            approver_groups: vec!["finance".to_string()],
1354        };
1355        let event = Event::ApprovalGranted(ApprovalGrantedEvent {
1356            run_id: Uuid::now_v7(),
1357            step_id: Some(Uuid::now_v7()),
1358            approved_by: "alice".to_string(),
1359            approvals_received: 1,
1360            approvals_required: 2,
1361            requirement: Some(requirement.clone()),
1362            at: Utc::now(),
1363        });
1364
1365        let json = serde_json::to_string(&event).expect("serialize");
1366        let back: Event = serde_json::from_str(&json).expect("deserialize");
1367        let Event::ApprovalGranted(back) = back else {
1368            panic!("expected approval_granted");
1369        };
1370        assert_eq!(back.approvals_received, 1);
1371        assert_eq!(back.approvals_required, 2);
1372        assert_eq!(back.requirement, Some(requirement));
1373    }
1374
1375    #[test]
1376    fn approval_requested_serde_roundtrip() {
1377        let event = Event::ApprovalRequested(ApprovalRequestedEvent {
1378            run_id: Uuid::now_v7(),
1379            step_id: Uuid::now_v7(),
1380            message: "Deploy to prod?".to_string(),
1381            requirement: None,
1382            at: Utc::now(),
1383        });
1384
1385        let json = serde_json::to_string(&event).expect("serialize");
1386        assert!(json.contains("approval_requested"));
1387    }
1388
1389    #[test]
1390    fn log_line_serde_roundtrip() {
1391        let event = Event::LogLine(LogLineEvent {
1392            id: Uuid::now_v7(),
1393            run_id: Uuid::now_v7(),
1394            step_id: Uuid::now_v7(),
1395            step_name: "build".to_string(),
1396            stream: LogStream::Stdout,
1397            line: "Compiling ironflow v0.1.0".to_string(),
1398            at: Utc::now(),
1399        });
1400
1401        let json = serde_json::to_string(&event).expect("serialize");
1402        let back: Event = serde_json::from_str(&json).expect("deserialize");
1403
1404        assert_eq!(back.event_type(), "log_line");
1405        assert!(json.contains("\"type\":\"log_line\""));
1406        assert!(json.contains("Compiling ironflow"));
1407    }
1408
1409    /// The pre-refactor wire format used flat inline-struct variants. Newtype
1410    /// variants over named-field payloads produce and accept the exact same
1411    /// JSON, so audit rows and in-flight payloads written before the refactor
1412    /// still deserialize. No data migration is required.
1413    #[test]
1414    fn legacy_flat_json_deserializes_into_typed_payload() {
1415        let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1416            .parse()
1417            .expect("valid uuid");
1418
1419        let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1420        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1421        match event {
1422            Event::RunCreated(e) => {
1423                assert_eq!(e.run_id, run_id);
1424                assert_eq!(e.workflow_name, "deploy");
1425            }
1426            other => panic!("expected RunCreated, got {other:?}"),
1427        }
1428
1429        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"}"#;
1430        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1431        match event {
1432            Event::RunStatusChanged(e) => {
1433                assert_eq!(e.from, RunStatus::Running);
1434                assert_eq!(e.to, RunStatus::Completed);
1435                assert_eq!(e.cost_usd, Decimal::new(5, 1));
1436                assert_eq!(e.duration_ms, 5000);
1437                assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1438            }
1439            other => panic!("expected RunStatusChanged, got {other:?}"),
1440        }
1441
1442        // `labels` predates no payload: omitting it must still work via `#[serde(default)]`.
1443        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"}"#;
1444        let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1445        match event {
1446            Event::RunStatusChanged(e) => {
1447                assert!(e.labels.is_empty());
1448                assert_eq!(e.error.as_deref(), Some("boom"));
1449            }
1450            other => panic!("expected RunStatusChanged, got {other:?}"),
1451        }
1452
1453        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"}"#;
1454        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1455        match event {
1456            Event::RunFailed(e) => {
1457                assert_eq!(e.error.as_deref(), Some("boom"));
1458                assert!(e.labels.is_empty());
1459            }
1460            other => panic!("expected RunFailed, got {other:?}"),
1461        }
1462
1463        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"}"#;
1464        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1465        match event {
1466            Event::StepFailed(e) => {
1467                assert_eq!(e.kind, StepKind::Shell);
1468                assert_eq!(e.error, "exit code 1");
1469            }
1470            other => panic!("expected StepFailed, got {other:?}"),
1471        }
1472
1473        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"}"#;
1474        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1475        match event {
1476            Event::LogLine(e) => {
1477                assert_eq!(e.stream, LogStream::Stdout);
1478                assert_eq!(e.line, "hello");
1479                // A payload predating the `id` field degrades to the nil UUID
1480                // rather than failing to deserialize.
1481                assert_eq!(e.id, Uuid::nil());
1482            }
1483            other => panic!("expected LogLine, got {other:?}"),
1484        }
1485
1486        let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1487        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1488        match event {
1489            Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1490            other => panic!("expected UserSignedIn, got {other:?}"),
1491        }
1492    }
1493
1494    /// Guards the internally-tagged representation: payload fields must stay
1495    /// siblings of `type`, never nested under a variant key.
1496    #[test]
1497    fn serialized_event_is_flat_with_type_tag() {
1498        let run_id = Uuid::now_v7();
1499        let event = Event::RunCreated(RunCreatedEvent {
1500            run_id,
1501            workflow_name: "deploy".to_string(),
1502            at: Utc::now(),
1503        });
1504
1505        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1506        let object = value.as_object().expect("event serializes to an object");
1507
1508        assert_eq!(
1509            object.get("type").and_then(|v| v.as_str()),
1510            Some("run_created")
1511        );
1512        assert_eq!(
1513            object.get("workflow_name").and_then(|v| v.as_str()),
1514            Some("deploy")
1515        );
1516        assert_eq!(
1517            object.get("run_id").and_then(|v| v.as_str()),
1518            Some(run_id.to_string().as_str())
1519        );
1520        assert!(object.contains_key("at"));
1521        assert_eq!(object.len(), 4, "no nesting: {object:?}");
1522        assert!(!object.contains_key("RunCreated"));
1523    }
1524
1525    #[test]
1526    fn run_id_returns_some_for_run_events() {
1527        let run_id = Uuid::now_v7();
1528        let now = Utc::now();
1529
1530        let events = vec![
1531            Event::RunCreated(RunCreatedEvent {
1532                run_id,
1533                workflow_name: "w".to_string(),
1534                at: now,
1535            }),
1536            Event::RunStatusChanged(RunStatusChangedEvent {
1537                run_id,
1538                workflow_name: "w".to_string(),
1539                from: RunStatus::Pending,
1540                to: RunStatus::Running,
1541                error: None,
1542                cost_usd: Decimal::ZERO,
1543                duration_ms: 0,
1544                labels: HashMap::new(),
1545                at: now,
1546            }),
1547            Event::RunFailed(RunFailedEvent {
1548                run_id,
1549                workflow_name: "w".to_string(),
1550                error: None,
1551                cost_usd: Decimal::ZERO,
1552                duration_ms: 0,
1553                labels: HashMap::new(),
1554                at: now,
1555            }),
1556            Event::RunBudgetExceeded(RunBudgetExceededEvent {
1557                run_id,
1558                workflow_name: "w".to_string(),
1559                limit_usd: Decimal::ZERO,
1560                spent_usd: Decimal::ZERO,
1561                step_budget_usd: Decimal::ZERO,
1562                at: now,
1563            }),
1564            Event::RetryForced(RetryForcedEvent {
1565                run_id,
1566                workflow_name: "w".to_string(),
1567                original_version: "1".to_string(),
1568                current_version: "2".to_string(),
1569                at: now,
1570            }),
1571            Event::StepCompleted(StepCompletedEvent {
1572                run_id,
1573                step_id: Uuid::now_v7(),
1574                step_name: "s".to_string(),
1575                kind: StepKind::Shell,
1576                duration_ms: 0,
1577                cost_usd: Decimal::ZERO,
1578                at: now,
1579            }),
1580            Event::StepFailed(StepFailedEvent {
1581                run_id,
1582                step_id: Uuid::now_v7(),
1583                step_name: "s".to_string(),
1584                kind: StepKind::Shell,
1585                error: "e".to_string(),
1586                at: now,
1587            }),
1588            Event::ApprovalRequested(ApprovalRequestedEvent {
1589                run_id,
1590                step_id: Uuid::now_v7(),
1591                message: "ok?".to_string(),
1592                requirement: None,
1593                at: now,
1594            }),
1595            Event::ApprovalGranted(ApprovalGrantedEvent {
1596                run_id,
1597                step_id: None,
1598                approved_by: "alice".to_string(),
1599                approvals_received: 1,
1600                approvals_required: 1,
1601                requirement: None,
1602                at: now,
1603            }),
1604            Event::ApprovalRejected(ApprovalRejectedEvent {
1605                run_id,
1606                step_id: None,
1607                rejected_by: "bob".to_string(),
1608                requirement: None,
1609                at: now,
1610            }),
1611            Event::LogLine(LogLineEvent {
1612                id: Uuid::now_v7(),
1613                run_id,
1614                step_id: Uuid::now_v7(),
1615                step_name: "s".to_string(),
1616                stream: LogStream::Stdout,
1617                line: "l".to_string(),
1618                at: now,
1619            }),
1620        ];
1621
1622        for event in &events {
1623            assert_eq!(
1624                event.run_id(),
1625                Some(run_id),
1626                "{} should carry a run_id",
1627                event.event_type()
1628            );
1629        }
1630    }
1631
1632    #[test]
1633    fn run_id_returns_none_for_auth_events() {
1634        let user_id = Uuid::now_v7();
1635        let now = Utc::now();
1636
1637        let events = vec![
1638            Event::UserSignedIn(UserSignedInEvent {
1639                user_id,
1640                username: "alice".to_string(),
1641                at: now,
1642            }),
1643            Event::UserSignedUp(UserSignedUpEvent {
1644                user_id,
1645                username: "alice".to_string(),
1646                at: now,
1647            }),
1648            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1649        ];
1650
1651        for event in &events {
1652            assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1653        }
1654    }
1655
1656    #[test]
1657    fn step_id_returns_some_only_for_step_events() {
1658        let step_id = Uuid::now_v7();
1659        let run_id = Uuid::now_v7();
1660        let now = Utc::now();
1661
1662        let with_step = vec![
1663            Event::StepCompleted(StepCompletedEvent {
1664                run_id,
1665                step_id,
1666                step_name: "s".to_string(),
1667                kind: StepKind::Shell,
1668                duration_ms: 0,
1669                cost_usd: Decimal::ZERO,
1670                at: now,
1671            }),
1672            Event::StepFailed(StepFailedEvent {
1673                run_id,
1674                step_id,
1675                step_name: "s".to_string(),
1676                kind: StepKind::Shell,
1677                error: "e".to_string(),
1678                at: now,
1679            }),
1680            Event::ApprovalRequested(ApprovalRequestedEvent {
1681                run_id,
1682                step_id,
1683                message: "ok?".to_string(),
1684                requirement: None,
1685                at: now,
1686            }),
1687            Event::ApprovalGranted(ApprovalGrantedEvent {
1688                run_id,
1689                step_id: Some(step_id),
1690                approved_by: "alice".to_string(),
1691                approvals_received: 1,
1692                approvals_required: 2,
1693                requirement: None,
1694                at: now,
1695            }),
1696            Event::ApprovalRejected(ApprovalRejectedEvent {
1697                run_id,
1698                step_id: Some(step_id),
1699                rejected_by: "bob".to_string(),
1700                requirement: None,
1701                at: now,
1702            }),
1703        ];
1704
1705        for event in &with_step {
1706            assert_eq!(
1707                event.step_id(),
1708                Some(step_id),
1709                "{} should carry a step_id",
1710                event.event_type()
1711            );
1712        }
1713
1714        let without_step = vec![
1715            Event::RunCreated(RunCreatedEvent {
1716                run_id,
1717                workflow_name: "w".to_string(),
1718                at: now,
1719            }),
1720            // A granted event recorded before step ids were carried.
1721            Event::ApprovalGranted(ApprovalGrantedEvent {
1722                run_id,
1723                step_id: None,
1724                approved_by: "alice".to_string(),
1725                approvals_received: 1,
1726                approvals_required: 1,
1727                requirement: None,
1728                at: now,
1729            }),
1730            // LogLine carries a step_id field but is reported as a run-level
1731            // stream event, matching the pre-refactor behaviour.
1732            Event::LogLine(LogLineEvent {
1733                id: Uuid::now_v7(),
1734                run_id,
1735                step_id,
1736                step_name: "s".to_string(),
1737                stream: LogStream::Stdout,
1738                line: "l".to_string(),
1739                at: now,
1740            }),
1741            Event::UserSignedOut(UserSignedOutEvent {
1742                user_id: Uuid::now_v7(),
1743                at: now,
1744            }),
1745        ];
1746
1747        for event in &without_step {
1748            assert_eq!(
1749                event.step_id(),
1750                None,
1751                "{} should not carry a step_id",
1752                event.event_type()
1753            );
1754        }
1755    }
1756
1757    #[test]
1758    fn user_id_returns_some_only_for_auth_events() {
1759        let user_id = Uuid::now_v7();
1760        let run_id = Uuid::now_v7();
1761        let now = Utc::now();
1762
1763        let auth = vec![
1764            Event::UserSignedIn(UserSignedInEvent {
1765                user_id,
1766                username: "alice".to_string(),
1767                at: now,
1768            }),
1769            Event::UserSignedUp(UserSignedUpEvent {
1770                user_id,
1771                username: "alice".to_string(),
1772                at: now,
1773            }),
1774            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1775        ];
1776
1777        for event in &auth {
1778            assert_eq!(
1779                event.user_id(),
1780                Some(user_id),
1781                "{} should carry a user_id",
1782                event.event_type()
1783            );
1784        }
1785
1786        let non_auth = vec![
1787            Event::RunCreated(RunCreatedEvent {
1788                run_id,
1789                workflow_name: "w".to_string(),
1790                at: now,
1791            }),
1792            Event::StepFailed(StepFailedEvent {
1793                run_id,
1794                step_id: Uuid::now_v7(),
1795                step_name: "s".to_string(),
1796                kind: StepKind::Shell,
1797                error: "e".to_string(),
1798                at: now,
1799            }),
1800        ];
1801
1802        for event in &non_auth {
1803            assert_eq!(
1804                event.user_id(),
1805                None,
1806                "{} should not carry a user_id",
1807                event.event_type()
1808            );
1809        }
1810    }
1811
1812    #[test]
1813    fn event_type_all_variants() {
1814        let id = Uuid::now_v7();
1815        let now = Utc::now();
1816
1817        let cases: Vec<(Event, &str)> = vec![
1818            (
1819                Event::RunCreated(RunCreatedEvent {
1820                    run_id: id,
1821                    workflow_name: "w".to_string(),
1822                    at: now,
1823                }),
1824                "run_created",
1825            ),
1826            (
1827                Event::RunStatusChanged(RunStatusChangedEvent {
1828                    run_id: id,
1829                    workflow_name: "w".to_string(),
1830                    from: RunStatus::Pending,
1831                    to: RunStatus::Running,
1832                    error: None,
1833                    cost_usd: Decimal::ZERO,
1834                    duration_ms: 0,
1835                    labels: HashMap::new(),
1836                    at: now,
1837                }),
1838                "run_status_changed",
1839            ),
1840            (
1841                Event::RunFailed(RunFailedEvent {
1842                    run_id: id,
1843                    workflow_name: "w".to_string(),
1844                    error: Some("boom".to_string()),
1845                    cost_usd: Decimal::ZERO,
1846                    duration_ms: 0,
1847                    labels: HashMap::new(),
1848                    at: now,
1849                }),
1850                "run_failed",
1851            ),
1852            (
1853                Event::RunBudgetExceeded(RunBudgetExceededEvent {
1854                    run_id: id,
1855                    workflow_name: "w".to_string(),
1856                    limit_usd: Decimal::new(200, 2),
1857                    spent_usd: Decimal::new(180, 2),
1858                    step_budget_usd: Decimal::new(50, 2),
1859                    at: now,
1860                }),
1861                "run_budget_exceeded",
1862            ),
1863            (
1864                Event::RetryForced(RetryForcedEvent {
1865                    run_id: id,
1866                    workflow_name: "w".to_string(),
1867                    original_version: "1".to_string(),
1868                    current_version: "2".to_string(),
1869                    at: now,
1870                }),
1871                "retry_forced",
1872            ),
1873            (
1874                Event::StepCompleted(StepCompletedEvent {
1875                    run_id: id,
1876                    step_id: id,
1877                    step_name: "s".to_string(),
1878                    kind: StepKind::Shell,
1879                    duration_ms: 0,
1880                    cost_usd: Decimal::ZERO,
1881                    at: now,
1882                }),
1883                "step_completed",
1884            ),
1885            (
1886                Event::StepFailed(StepFailedEvent {
1887                    run_id: id,
1888                    step_id: id,
1889                    step_name: "s".to_string(),
1890                    kind: StepKind::Shell,
1891                    error: "err".to_string(),
1892                    at: now,
1893                }),
1894                "step_failed",
1895            ),
1896            (
1897                Event::ApprovalRequested(ApprovalRequestedEvent {
1898                    run_id: id,
1899                    step_id: id,
1900                    message: "ok?".to_string(),
1901                    requirement: None,
1902                    at: now,
1903                }),
1904                "approval_requested",
1905            ),
1906            (
1907                Event::ApprovalGranted(ApprovalGrantedEvent {
1908                    run_id: id,
1909                    step_id: Some(id),
1910                    approved_by: "alice".to_string(),
1911                    approvals_received: 1,
1912                    approvals_required: 1,
1913                    requirement: None,
1914                    at: now,
1915                }),
1916                "approval_granted",
1917            ),
1918            (
1919                Event::ApprovalRejected(ApprovalRejectedEvent {
1920                    run_id: id,
1921                    step_id: Some(id),
1922                    rejected_by: "bob".to_string(),
1923                    requirement: None,
1924                    at: now,
1925                }),
1926                "approval_rejected",
1927            ),
1928            (
1929                Event::ApprovalEscalated(ApprovalEscalatedEvent {
1930                    run_id: id,
1931                    step_id: id,
1932                    step_name: "prod-gate".to_string(),
1933                    stage: 0,
1934                    policy: "auto_reject".to_string(),
1935                    action: "rejected".to_string(),
1936                    reason: "approval deadline of 3600s expired".to_string(),
1937                    assignee: None,
1938                    at: now,
1939                }),
1940                "approval_escalated",
1941            ),
1942            (
1943                Event::LogLine(LogLineEvent {
1944                    id,
1945                    run_id: id,
1946                    step_id: id,
1947                    step_name: "build".to_string(),
1948                    stream: LogStream::Stdout,
1949                    line: "Compiling ironflow v0.1.0".to_string(),
1950                    at: now,
1951                }),
1952                "log_line",
1953            ),
1954            (
1955                Event::UserSignedIn(UserSignedInEvent {
1956                    user_id: id,
1957                    username: "u".to_string(),
1958                    at: now,
1959                }),
1960                "user_signed_in",
1961            ),
1962            (
1963                Event::UserSignedUp(UserSignedUpEvent {
1964                    user_id: id,
1965                    username: "u".to_string(),
1966                    at: now,
1967                }),
1968                "user_signed_up",
1969            ),
1970            (
1971                Event::UserSignedOut(UserSignedOutEvent {
1972                    user_id: id,
1973                    at: now,
1974                }),
1975                "user_signed_out",
1976            ),
1977            (
1978                Event::ProviderAccountUpdated(ProviderAccountUpdatedEvent {
1979                    account_id: id,
1980                    name: "perso".to_string(),
1981                    change: ProviderAccountChange::TokenReplaced,
1982                    at: now,
1983                }),
1984                "provider_account.updated",
1985            ),
1986            (
1987                Event::ProviderAccountUsageUpdated(ProviderAccountUsageUpdatedEvent {
1988                    account_id: id,
1989                    name: "perso".to_string(),
1990                    windows: Vec::new(),
1991                    at: now,
1992                }),
1993                "provider_account.usage_updated",
1994            ),
1995            (
1996                Event::SignalAwaited(SignalAwaitedEvent {
1997                    run_id: id,
1998                    step_id: id,
1999                    step_name: "wait-ci".to_string(),
2000                    name: "ci.done".to_string(),
2001                    key: "abc".to_string(),
2002                    deadline_at: now,
2003                    at: now,
2004                }),
2005                "signal_awaited",
2006            ),
2007            (
2008                Event::SignalReceived(SignalReceivedEvent {
2009                    signal_id: id,
2010                    name: "ci.done".to_string(),
2011                    key: "abc".to_string(),
2012                    resumed_runs: vec![id],
2013                    at: now,
2014                }),
2015                "signal_received",
2016            ),
2017            (
2018                Event::ScheduleOccurrencesMissed(ScheduleOccurrencesMissedEvent {
2019                    schedule_id: id,
2020                    workflow_name: "nightly".to_string(),
2021                    reason: ScheduleMissReason::Overlap,
2022                    count: 1,
2023                    first: now,
2024                    last: now,
2025                    at: now,
2026                }),
2027                "schedule_occurrences_missed",
2028            ),
2029        ];
2030
2031        assert_eq!(
2032            cases.len(),
2033            Event::ALL.len(),
2034            "every variant must be covered"
2035        );
2036
2037        for (event, expected_type) in cases {
2038            assert_eq!(event.event_type(), expected_type);
2039        }
2040    }
2041
2042    #[test]
2043    fn approval_escalated_serde_roundtrip() {
2044        let run_id = Uuid::now_v7();
2045        let step_id = Uuid::now_v7();
2046        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
2047            run_id,
2048            step_id,
2049            step_name: "prod-gate".to_string(),
2050            stage: 1,
2051            policy: "escalate".to_string(),
2052            action: "reassigned to sre-oncall".to_string(),
2053            reason: "approval deadline of 3600s expired".to_string(),
2054            assignee: Some(Assignee::group("sre-oncall")),
2055            at: Utc::now(),
2056        });
2057
2058        let json = serde_json::to_string(&event).expect("serialize");
2059        assert!(
2060            json.contains("\"type\":\"approval_escalated\""),
2061            "got {json}"
2062        );
2063
2064        let back: Event = serde_json::from_str(&json).expect("deserialize");
2065        let Event::ApprovalEscalated(payload) = back else {
2066            panic!("expected an approval_escalated event");
2067        };
2068        assert_eq!(payload.run_id, run_id);
2069        assert_eq!(payload.stage, 1);
2070        assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
2071    }
2072
2073    #[test]
2074    fn approval_escalated_carries_run_and_step_ids() {
2075        let run_id = Uuid::now_v7();
2076        let step_id = Uuid::now_v7();
2077        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
2078            run_id,
2079            step_id,
2080            step_name: "prod-gate".to_string(),
2081            stage: 0,
2082            policy: "notify".to_string(),
2083            action: "notified 1 target".to_string(),
2084            reason: "approval deadline of 60s expired".to_string(),
2085            assignee: None,
2086            at: Utc::now(),
2087        });
2088
2089        assert_eq!(event.run_id(), Some(run_id));
2090        assert_eq!(event.step_id(), Some(step_id));
2091        assert_eq!(event.user_id(), None);
2092    }
2093
2094    #[test]
2095    fn signal_awaited_serde_roundtrip() {
2096        let run_id = Uuid::now_v7();
2097        let step_id = Uuid::now_v7();
2098        let event = Event::SignalAwaited(SignalAwaitedEvent {
2099            run_id,
2100            step_id,
2101            step_name: "wait-ci".to_string(),
2102            name: "ci.done".to_string(),
2103            key: "abc".to_string(),
2104            deadline_at: Utc::now(),
2105            at: Utc::now(),
2106        });
2107
2108        let json = serde_json::to_string(&event).expect("serialize");
2109        assert!(json.contains("\"type\":\"signal_awaited\""), "got {json}");
2110        let back: Event = serde_json::from_str(&json).expect("deserialize");
2111        assert_eq!(back.event_type(), Event::SIGNAL_AWAITED);
2112        assert_eq!(back.run_id(), Some(run_id));
2113        assert_eq!(back.step_id(), Some(step_id));
2114        assert_eq!(back.user_id(), None);
2115    }
2116
2117    #[test]
2118    fn signal_received_serde_roundtrip() {
2119        let resumed = Uuid::now_v7();
2120        let event = Event::SignalReceived(SignalReceivedEvent {
2121            signal_id: Uuid::now_v7(),
2122            name: "ci.done".to_string(),
2123            key: "abc".to_string(),
2124            resumed_runs: vec![resumed],
2125            at: Utc::now(),
2126        });
2127
2128        let json = serde_json::to_string(&event).expect("serialize");
2129        assert!(json.contains("\"type\":\"signal_received\""), "got {json}");
2130        let back: Event = serde_json::from_str(&json).expect("deserialize");
2131        let Event::SignalReceived(payload) = &back else {
2132            panic!("expected a signal_received event, got {back:?}");
2133        };
2134        assert_eq!(payload.resumed_runs, vec![resumed]);
2135        assert_eq!(back.run_id(), None);
2136        assert_eq!(back.step_id(), None);
2137        assert_eq!(back.user_id(), None);
2138    }
2139
2140    #[test]
2141    fn schedule_occurrences_missed_serde_roundtrip() {
2142        let schedule_id = Uuid::now_v7();
2143        let first = Utc::now();
2144        let event = Event::ScheduleOccurrencesMissed(ScheduleOccurrencesMissedEvent {
2145            schedule_id,
2146            workflow_name: "nightly".to_string(),
2147            reason: ScheduleMissReason::OutsideWindow,
2148            count: 4,
2149            first,
2150            last: first,
2151            at: first,
2152        });
2153
2154        let json = serde_json::to_string(&event).expect("serialize");
2155        assert!(
2156            json.contains("\"type\":\"schedule_occurrences_missed\""),
2157            "got {json}"
2158        );
2159        assert!(json.contains("\"reason\":\"outside_window\""), "got {json}");
2160        let back: Event = serde_json::from_str(&json).expect("deserialize");
2161        let Event::ScheduleOccurrencesMissed(payload) = &back else {
2162            panic!("expected a schedule_occurrences_missed event, got {back:?}");
2163        };
2164        assert_eq!(payload.schedule_id, schedule_id);
2165        assert_eq!(payload.reason, ScheduleMissReason::OutsideWindow);
2166        assert_eq!(payload.count, 4);
2167        assert_eq!(payload.first, first);
2168        assert!(Event::ALL.contains(&Event::SCHEDULE_OCCURRENCES_MISSED));
2169        assert_eq!(back.run_id(), None);
2170        assert_eq!(back.step_id(), None);
2171        assert_eq!(back.user_id(), None);
2172    }
2173}