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, 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/// A domain event emitted by the ironflow system.
766///
767/// Covers the full lifecycle: runs, steps, approvals, and authentication.
768/// Subscribers receive these via [`EventPublisher`](super::EventPublisher)
769/// and pattern-match on the variants they care about.
770///
771/// Each variant wraps a dedicated payload struct. The serialized form stays
772/// flat: the `type` discriminant sits next to the payload fields, so
773/// `{"type":"run_created","run_id":...}` round-trips unchanged.
774///
775/// # Examples
776///
777/// ```
778/// use std::collections::HashMap;
779/// use ironflow_engine::notify::{Event, RunStatusChangedEvent};
780/// use ironflow_store::models::RunStatus;
781/// use uuid::Uuid;
782///
783/// let event = Event::RunStatusChanged(RunStatusChangedEvent {
784///     run_id: Uuid::now_v7(),
785///     workflow_name: "deploy".to_string(),
786///     from: RunStatus::Running,
787///     to: RunStatus::Completed,
788///     error: None,
789///     cost_usd: rust_decimal::Decimal::ZERO,
790///     duration_ms: 5000,
791///     labels: HashMap::new(),
792///     at: chrono::Utc::now(),
793/// });
794/// assert_eq!(event.event_type(), "run_status_changed");
795/// ```
796#[derive(Debug, Clone, Serialize, Deserialize)]
797#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
798#[serde(tag = "type", rename_all = "snake_case")]
799pub enum Event {
800    // -- Run lifecycle --
801    /// A new run was created (status: Pending).
802    RunCreated(RunCreatedEvent),
803
804    /// A run changed status.
805    RunStatusChanged(RunStatusChangedEvent),
806
807    /// A run transitioned to [`Failed`](ironflow_store::models::RunStatus::Failed).
808    ///
809    /// This is a convenience event emitted alongside [`RunStatusChanged`](Event::RunStatusChanged)
810    /// when the target status is `Failed`. Subscribe to this instead of
811    /// `RUN_STATUS_CHANGED` when you only care about failures.
812    RunFailed(RunFailedEvent),
813
814    /// A run was stopped because it reached its cumulative cost cap.
815    ///
816    /// Emitted when the engine refuses an agent step that would cross the run's
817    /// `max_cost_usd`. The run transitions to
818    /// [`Cancelled`](ironflow_store::models::RunStatus::Cancelled) and the step
819    /// is never launched, so the reported spend is what the run had already
820    /// consumed.
821    RunBudgetExceeded(RunBudgetExceededEvent),
822
823    /// A manual retry was forced despite a handler version mismatch.
824    ///
825    /// Emitted when a caller passes `force=true` on a retry where the
826    /// handler version differs from the run's recorded version. This is
827    /// an audit event: it means the new code will execute on the old
828    /// payload without the handler explicitly declaring compatibility.
829    RetryForced(RetryForcedEvent),
830
831    // -- Step lifecycle --
832    /// A step completed successfully.
833    StepCompleted(StepCompletedEvent),
834
835    /// A step failed.
836    StepFailed(StepFailedEvent),
837
838    // -- Approval --
839    /// A run is waiting for human approval.
840    ApprovalRequested(ApprovalRequestedEvent),
841
842    /// A run was approved by a human.
843    ApprovalGranted(ApprovalGrantedEvent),
844
845    /// A run was rejected by a human.
846    ApprovalRejected(ApprovalRejectedEvent),
847
848    /// An approval gate missed its SLA deadline and an escalation policy ran.
849    ApprovalEscalated(ApprovalEscalatedEvent),
850
851    // -- Log streaming --
852    /// A log line emitted during step execution.
853    ///
854    /// Pushed by the worker in real time so that SSE clients can stream
855    /// step output as it happens, without waiting for step completion.
856    LogLine(LogLineEvent),
857
858    // -- Authentication --
859    /// A user signed in.
860    UserSignedIn(UserSignedInEvent),
861
862    /// A new user signed up.
863    UserSignedUp(UserSignedUpEvent),
864
865    /// A user signed out.
866    UserSignedOut(UserSignedOutEvent),
867
868    // -- Provider Accounts --
869    /// A Provider Account was created, updated, deleted or had its token replaced.
870    #[serde(rename = "provider_account.updated")]
871    ProviderAccountUpdated(ProviderAccountUpdatedEvent),
872
873    /// New usage windows were recorded for a Provider Account.
874    #[serde(rename = "provider_account.usage_updated")]
875    ProviderAccountUsageUpdated(ProviderAccountUsageUpdatedEvent),
876
877    // -- Signals --
878    /// A run started waiting for a signal.
879    SignalAwaited(SignalAwaitedEvent),
880
881    /// A signal was received.
882    SignalReceived(SignalReceivedEvent),
883}
884
885impl Event {
886    /// Event type constant for [`RunCreated`](Event::RunCreated).
887    pub const RUN_CREATED: &'static str = "run_created";
888    /// Event type constant for [`RunStatusChanged`](Event::RunStatusChanged).
889    pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
890    /// Event type constant for [`RunFailed`](Event::RunFailed).
891    pub const RUN_FAILED: &'static str = "run_failed";
892    /// Event type constant for [`RunBudgetExceeded`](Event::RunBudgetExceeded).
893    pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
894    /// Event type constant for [`RetryForced`](Event::RetryForced).
895    pub const RETRY_FORCED: &'static str = "retry_forced";
896    /// Event type constant for [`StepCompleted`](Event::StepCompleted).
897    pub const STEP_COMPLETED: &'static str = "step_completed";
898    /// Event type constant for [`StepFailed`](Event::StepFailed).
899    pub const STEP_FAILED: &'static str = "step_failed";
900    /// Event type constant for [`ApprovalRequested`](Event::ApprovalRequested).
901    pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
902    /// Event type constant for [`ApprovalGranted`](Event::ApprovalGranted).
903    pub const APPROVAL_GRANTED: &'static str = "approval_granted";
904    /// Event type constant for [`ApprovalRejected`](Event::ApprovalRejected).
905    pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
906    /// Event type constant for [`ApprovalEscalated`](Event::ApprovalEscalated).
907    pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
908    /// Event type constant for [`LogLine`](Event::LogLine).
909    pub const LOG_LINE: &'static str = "log_line";
910    /// Event type constant for [`UserSignedIn`](Event::UserSignedIn).
911    pub const USER_SIGNED_IN: &'static str = "user_signed_in";
912    /// Event type constant for [`UserSignedUp`](Event::UserSignedUp).
913    pub const USER_SIGNED_UP: &'static str = "user_signed_up";
914    /// Event type constant for [`UserSignedOut`](Event::UserSignedOut).
915    pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
916    /// Event type constant for [`ProviderAccountUpdated`](Event::ProviderAccountUpdated).
917    pub const PROVIDER_ACCOUNT_UPDATED: &'static str = "provider_account.updated";
918    /// Event type constant for
919    /// [`ProviderAccountUsageUpdated`](Event::ProviderAccountUsageUpdated).
920    pub const PROVIDER_ACCOUNT_USAGE_UPDATED: &'static str = "provider_account.usage_updated";
921    /// Event type constant for [`SignalAwaited`](Event::SignalAwaited).
922    pub const SIGNAL_AWAITED: &'static str = "signal_awaited";
923    /// Event type constant for [`SignalReceived`](Event::SignalReceived).
924    pub const SIGNAL_RECEIVED: &'static str = "signal_received";
925
926    /// All event types. Pass this to
927    /// [`EventPublisher::subscribe`](super::EventPublisher::subscribe) to
928    /// receive every event.
929    ///
930    /// # Examples
931    ///
932    /// ```no_run
933    /// use ironflow_engine::notify::{Event, EventPublisher, WebhookSubscriber};
934    ///
935    /// let mut publisher = EventPublisher::new();
936    /// publisher.subscribe(
937    ///     WebhookSubscriber::new("https://example.com/all"),
938    ///     Event::ALL,
939    /// );
940    /// ```
941    pub const ALL: &'static [&'static str] = &[
942        Self::RUN_CREATED,
943        Self::RUN_STATUS_CHANGED,
944        Self::RUN_FAILED,
945        Self::RUN_BUDGET_EXCEEDED,
946        Self::STEP_COMPLETED,
947        Self::STEP_FAILED,
948        Self::APPROVAL_REQUESTED,
949        Self::APPROVAL_GRANTED,
950        Self::APPROVAL_REJECTED,
951        Self::APPROVAL_ESCALATED,
952        Self::LOG_LINE,
953        Self::USER_SIGNED_IN,
954        Self::USER_SIGNED_UP,
955        Self::USER_SIGNED_OUT,
956        Self::RETRY_FORCED,
957        Self::PROVIDER_ACCOUNT_UPDATED,
958        Self::PROVIDER_ACCOUNT_USAGE_UPDATED,
959        Self::SIGNAL_AWAITED,
960        Self::SIGNAL_RECEIVED,
961    ];
962
963    /// Returns the event type as a static string (e.g. `"run_status_changed"`).
964    ///
965    /// Useful for filtering and logging without deserializing.
966    ///
967    /// # Examples
968    ///
969    /// ```
970    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
971    /// use uuid::Uuid;
972    /// use chrono::Utc;
973    ///
974    /// let event = Event::UserSignedIn(UserSignedInEvent {
975    ///     user_id: Uuid::now_v7(),
976    ///     username: "alice".to_string(),
977    ///     at: Utc::now(),
978    /// });
979    /// assert_eq!(event.event_type(), "user_signed_in");
980    /// ```
981    #[deny(unreachable_patterns)]
982    pub fn event_type(&self) -> &'static str {
983        match self {
984            Event::RunCreated(_) => Self::RUN_CREATED,
985            Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
986            Event::RunFailed(_) => Self::RUN_FAILED,
987            Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
988            Event::RetryForced(_) => Self::RETRY_FORCED,
989            Event::StepCompleted(_) => Self::STEP_COMPLETED,
990            Event::StepFailed(_) => Self::STEP_FAILED,
991            Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
992            Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
993            Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
994            Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
995            Event::LogLine(_) => Self::LOG_LINE,
996            Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
997            Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
998            Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
999            Event::ProviderAccountUpdated(_) => Self::PROVIDER_ACCOUNT_UPDATED,
1000            Event::ProviderAccountUsageUpdated(_) => Self::PROVIDER_ACCOUNT_USAGE_UPDATED,
1001            Event::SignalAwaited(_) => Self::SIGNAL_AWAITED,
1002            Event::SignalReceived(_) => Self::SIGNAL_RECEIVED,
1003        }
1004    }
1005
1006    /// Returns the run this event belongs to, if any.
1007    ///
1008    /// Auth events ([`UserSignedIn`](Event::UserSignedIn),
1009    /// [`UserSignedUp`](Event::UserSignedUp),
1010    /// [`UserSignedOut`](Event::UserSignedOut)) are not tied to a run and
1011    /// return `None`.
1012    ///
1013    /// # Examples
1014    ///
1015    /// ```
1016    /// use ironflow_engine::notify::{Event, RunCreatedEvent};
1017    /// use uuid::Uuid;
1018    /// use chrono::Utc;
1019    ///
1020    /// let run_id = Uuid::now_v7();
1021    /// let event = Event::RunCreated(RunCreatedEvent {
1022    ///     run_id,
1023    ///     workflow_name: "deploy".to_string(),
1024    ///     at: Utc::now(),
1025    /// });
1026    /// assert_eq!(event.run_id(), Some(run_id));
1027    /// ```
1028    #[deny(unreachable_patterns)]
1029    pub fn run_id(&self) -> Option<Uuid> {
1030        match self {
1031            Event::RunCreated(e) => Some(e.run_id),
1032            Event::RunStatusChanged(e) => Some(e.run_id),
1033            Event::RunFailed(e) => Some(e.run_id),
1034            Event::RunBudgetExceeded(e) => Some(e.run_id),
1035            Event::RetryForced(e) => Some(e.run_id),
1036            Event::StepCompleted(e) => Some(e.run_id),
1037            Event::StepFailed(e) => Some(e.run_id),
1038            Event::ApprovalRequested(e) => Some(e.run_id),
1039            Event::ApprovalGranted(e) => Some(e.run_id),
1040            Event::ApprovalRejected(e) => Some(e.run_id),
1041            Event::ApprovalEscalated(e) => Some(e.run_id),
1042            Event::LogLine(e) => Some(e.run_id),
1043            Event::SignalAwaited(e) => Some(e.run_id),
1044            Event::UserSignedIn(_)
1045            | Event::UserSignedUp(_)
1046            | Event::UserSignedOut(_)
1047            | Event::ProviderAccountUpdated(_)
1048            | Event::ProviderAccountUsageUpdated(_)
1049            | Event::SignalReceived(_) => None,
1050        }
1051    }
1052
1053    /// Returns the step this event belongs to, if any.
1054    ///
1055    /// Only [`StepCompleted`](Event::StepCompleted),
1056    /// [`StepFailed`](Event::StepFailed),
1057    /// [`ApprovalRequested`](Event::ApprovalRequested) and
1058    /// [`ApprovalEscalated`](Event::ApprovalEscalated) always carry a step
1059    /// identifier; [`ApprovalGranted`](Event::ApprovalGranted) and
1060    /// [`ApprovalRejected`](Event::ApprovalRejected) carry one when it was
1061    /// recorded; every other variant returns `None`.
1062    ///
1063    /// # Examples
1064    ///
1065    /// ```
1066    /// use ironflow_engine::notify::{Event, StepFailedEvent};
1067    /// use ironflow_store::models::StepKind;
1068    /// use uuid::Uuid;
1069    /// use chrono::Utc;
1070    ///
1071    /// let step_id = Uuid::now_v7();
1072    /// let event = Event::StepFailed(StepFailedEvent {
1073    ///     run_id: Uuid::now_v7(),
1074    ///     step_id,
1075    ///     step_name: "build".to_string(),
1076    ///     kind: StepKind::Shell,
1077    ///     error: "exit code 1".to_string(),
1078    ///     at: Utc::now(),
1079    /// });
1080    /// assert_eq!(event.step_id(), Some(step_id));
1081    /// ```
1082    #[deny(unreachable_patterns)]
1083    pub fn step_id(&self) -> Option<Uuid> {
1084        match self {
1085            Event::StepCompleted(e) => Some(e.step_id),
1086            Event::StepFailed(e) => Some(e.step_id),
1087            Event::ApprovalRequested(e) => Some(e.step_id),
1088            Event::ApprovalEscalated(e) => Some(e.step_id),
1089            Event::ApprovalGranted(e) => e.step_id,
1090            Event::ApprovalRejected(e) => e.step_id,
1091            Event::SignalAwaited(e) => Some(e.step_id),
1092            Event::RunCreated(_)
1093            | Event::RunStatusChanged(_)
1094            | Event::RunFailed(_)
1095            | Event::RunBudgetExceeded(_)
1096            | Event::RetryForced(_)
1097            | Event::LogLine(_)
1098            | Event::UserSignedIn(_)
1099            | Event::UserSignedUp(_)
1100            | Event::UserSignedOut(_)
1101            | Event::ProviderAccountUpdated(_)
1102            | Event::ProviderAccountUsageUpdated(_)
1103            | Event::SignalReceived(_) => None,
1104        }
1105    }
1106
1107    /// Returns the user this event belongs to, if any.
1108    ///
1109    /// Only the auth events ([`UserSignedIn`](Event::UserSignedIn),
1110    /// [`UserSignedUp`](Event::UserSignedUp),
1111    /// [`UserSignedOut`](Event::UserSignedOut)) carry a user identifier;
1112    /// every other variant returns `None`.
1113    ///
1114    /// # Examples
1115    ///
1116    /// ```
1117    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
1118    /// use uuid::Uuid;
1119    /// use chrono::Utc;
1120    ///
1121    /// let user_id = Uuid::now_v7();
1122    /// let event = Event::UserSignedIn(UserSignedInEvent {
1123    ///     user_id,
1124    ///     username: "alice".to_string(),
1125    ///     at: Utc::now(),
1126    /// });
1127    /// assert_eq!(event.user_id(), Some(user_id));
1128    /// ```
1129    #[deny(unreachable_patterns)]
1130    pub fn user_id(&self) -> Option<Uuid> {
1131        match self {
1132            Event::UserSignedIn(e) => Some(e.user_id),
1133            Event::UserSignedUp(e) => Some(e.user_id),
1134            Event::UserSignedOut(e) => Some(e.user_id),
1135            Event::RunCreated(_)
1136            | Event::RunStatusChanged(_)
1137            | Event::RunFailed(_)
1138            | Event::RunBudgetExceeded(_)
1139            | Event::RetryForced(_)
1140            | Event::StepCompleted(_)
1141            | Event::StepFailed(_)
1142            | Event::ApprovalRequested(_)
1143            | Event::ApprovalGranted(_)
1144            | Event::ApprovalRejected(_)
1145            | Event::ApprovalEscalated(_)
1146            | Event::LogLine(_)
1147            | Event::ProviderAccountUpdated(_)
1148            | Event::ProviderAccountUsageUpdated(_)
1149            | Event::SignalAwaited(_)
1150            | Event::SignalReceived(_) => None,
1151        }
1152    }
1153}
1154
1155#[cfg(test)]
1156mod tests {
1157    use super::*;
1158
1159    #[test]
1160    fn run_status_changed_serde_roundtrip() {
1161        let event = Event::RunStatusChanged(RunStatusChangedEvent {
1162            run_id: Uuid::now_v7(),
1163            workflow_name: "deploy".to_string(),
1164            from: RunStatus::Running,
1165            to: RunStatus::Completed,
1166            error: None,
1167            cost_usd: Decimal::new(42, 2),
1168            duration_ms: 5000,
1169            labels: HashMap::new(),
1170            at: Utc::now(),
1171        });
1172
1173        let json = serde_json::to_string(&event).expect("serialize");
1174        let back: Event = serde_json::from_str(&json).expect("deserialize");
1175
1176        assert_eq!(back.event_type(), "run_status_changed");
1177        assert!(json.contains("\"type\":\"run_status_changed\""));
1178    }
1179
1180    #[test]
1181    fn run_failed_serde_roundtrip() {
1182        let event = Event::RunFailed(RunFailedEvent {
1183            run_id: Uuid::now_v7(),
1184            workflow_name: "deploy".to_string(),
1185            error: Some("step crashed".to_string()),
1186            cost_usd: Decimal::new(10, 2),
1187            duration_ms: 3000,
1188            labels: HashMap::new(),
1189            at: Utc::now(),
1190        });
1191
1192        let json = serde_json::to_string(&event).expect("serialize");
1193        let back: Event = serde_json::from_str(&json).expect("deserialize");
1194
1195        assert_eq!(back.event_type(), "run_failed");
1196        assert!(json.contains("\"type\":\"run_failed\""));
1197        assert!(json.contains("step crashed"));
1198    }
1199
1200    #[test]
1201    fn run_budget_exceeded_serde_roundtrip() {
1202        let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
1203            run_id: Uuid::now_v7(),
1204            workflow_name: "deploy".to_string(),
1205            limit_usd: Decimal::new(200, 2),
1206            spent_usd: Decimal::new(180, 2),
1207            step_budget_usd: Decimal::new(50, 2),
1208            at: Utc::now(),
1209        });
1210
1211        let json = serde_json::to_string(&event).expect("serialize");
1212        let back: Event = serde_json::from_str(&json).expect("deserialize");
1213
1214        assert_eq!(back.event_type(), "run_budget_exceeded");
1215        assert!(json.contains("\"type\":\"run_budget_exceeded\""));
1216        assert!(json.contains("limit_usd"));
1217        assert!(json.contains("step_budget_usd"));
1218    }
1219
1220    #[test]
1221    fn all_contains_run_budget_exceeded() {
1222        assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
1223    }
1224
1225    #[test]
1226    fn user_signed_in_serde_roundtrip() {
1227        let event = Event::UserSignedIn(UserSignedInEvent {
1228            user_id: Uuid::now_v7(),
1229            username: "alice".to_string(),
1230            at: Utc::now(),
1231        });
1232
1233        let json = serde_json::to_string(&event).expect("serialize");
1234        let back: Event = serde_json::from_str(&json).expect("deserialize");
1235
1236        assert_eq!(back.event_type(), "user_signed_in");
1237        assert!(json.contains("alice"));
1238    }
1239
1240    #[test]
1241    fn step_failed_serde_roundtrip() {
1242        let event = Event::StepFailed(StepFailedEvent {
1243            run_id: Uuid::now_v7(),
1244            step_id: Uuid::now_v7(),
1245            step_name: "build".to_string(),
1246            kind: StepKind::Shell,
1247            error: "exit code 1".to_string(),
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(), "step_failed");
1255    }
1256
1257    #[test]
1258    fn legacy_approval_granted_defaults_to_a_single_vote() {
1259        let raw = r#"{"type":"approval_granted","run_id":"01890000-0000-7000-8000-000000000000","approved_by":"alice","at":"2026-01-01T00:00:00Z"}"#;
1260        let event: Event = serde_json::from_str(raw).expect("deserialize");
1261        let Event::ApprovalGranted(event) = event else {
1262            panic!("expected approval_granted");
1263        };
1264
1265        assert_eq!(event.step_id, None);
1266        assert_eq!(event.approvals_received, 1);
1267        assert_eq!(event.approvals_required, 1);
1268        assert!(event.requirement.is_none());
1269    }
1270
1271    #[test]
1272    fn legacy_approval_requested_and_rejected_have_no_requirement() {
1273        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"}"#;
1274        let requested: Event = serde_json::from_str(raw).expect("deserialize");
1275        let Event::ApprovalRequested(requested) = requested else {
1276            panic!("expected approval_requested");
1277        };
1278        assert!(requested.requirement.is_none());
1279
1280        let raw = r#"{"type":"approval_rejected","run_id":"01890000-0000-7000-8000-000000000000","rejected_by":"bob","at":"2026-01-01T00:00:00Z"}"#;
1281        let rejected: Event = serde_json::from_str(raw).expect("deserialize");
1282        let Event::ApprovalRejected(rejected) = rejected else {
1283            panic!("expected approval_rejected");
1284        };
1285        assert_eq!(rejected.step_id, None);
1286        assert!(rejected.requirement.is_none());
1287    }
1288
1289    #[test]
1290    fn approval_granted_roundtrips_the_vote_counts() {
1291        let requirement = ApprovalRequirement {
1292            reason: Some("amount > 10k".to_string()),
1293            required_approvers: 2,
1294            approver_groups: vec!["finance".to_string()],
1295        };
1296        let event = Event::ApprovalGranted(ApprovalGrantedEvent {
1297            run_id: Uuid::now_v7(),
1298            step_id: Some(Uuid::now_v7()),
1299            approved_by: "alice".to_string(),
1300            approvals_received: 1,
1301            approvals_required: 2,
1302            requirement: Some(requirement.clone()),
1303            at: Utc::now(),
1304        });
1305
1306        let json = serde_json::to_string(&event).expect("serialize");
1307        let back: Event = serde_json::from_str(&json).expect("deserialize");
1308        let Event::ApprovalGranted(back) = back else {
1309            panic!("expected approval_granted");
1310        };
1311        assert_eq!(back.approvals_received, 1);
1312        assert_eq!(back.approvals_required, 2);
1313        assert_eq!(back.requirement, Some(requirement));
1314    }
1315
1316    #[test]
1317    fn approval_requested_serde_roundtrip() {
1318        let event = Event::ApprovalRequested(ApprovalRequestedEvent {
1319            run_id: Uuid::now_v7(),
1320            step_id: Uuid::now_v7(),
1321            message: "Deploy to prod?".to_string(),
1322            requirement: None,
1323            at: Utc::now(),
1324        });
1325
1326        let json = serde_json::to_string(&event).expect("serialize");
1327        assert!(json.contains("approval_requested"));
1328    }
1329
1330    #[test]
1331    fn log_line_serde_roundtrip() {
1332        let event = Event::LogLine(LogLineEvent {
1333            id: Uuid::now_v7(),
1334            run_id: Uuid::now_v7(),
1335            step_id: Uuid::now_v7(),
1336            step_name: "build".to_string(),
1337            stream: LogStream::Stdout,
1338            line: "Compiling ironflow v0.1.0".to_string(),
1339            at: Utc::now(),
1340        });
1341
1342        let json = serde_json::to_string(&event).expect("serialize");
1343        let back: Event = serde_json::from_str(&json).expect("deserialize");
1344
1345        assert_eq!(back.event_type(), "log_line");
1346        assert!(json.contains("\"type\":\"log_line\""));
1347        assert!(json.contains("Compiling ironflow"));
1348    }
1349
1350    /// The pre-refactor wire format used flat inline-struct variants. Newtype
1351    /// variants over named-field payloads produce and accept the exact same
1352    /// JSON, so audit rows and in-flight payloads written before the refactor
1353    /// still deserialize. No data migration is required.
1354    #[test]
1355    fn legacy_flat_json_deserializes_into_typed_payload() {
1356        let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1357            .parse()
1358            .expect("valid uuid");
1359
1360        let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1361        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1362        match event {
1363            Event::RunCreated(e) => {
1364                assert_eq!(e.run_id, run_id);
1365                assert_eq!(e.workflow_name, "deploy");
1366            }
1367            other => panic!("expected RunCreated, got {other:?}"),
1368        }
1369
1370        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"}"#;
1371        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1372        match event {
1373            Event::RunStatusChanged(e) => {
1374                assert_eq!(e.from, RunStatus::Running);
1375                assert_eq!(e.to, RunStatus::Completed);
1376                assert_eq!(e.cost_usd, Decimal::new(5, 1));
1377                assert_eq!(e.duration_ms, 5000);
1378                assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1379            }
1380            other => panic!("expected RunStatusChanged, got {other:?}"),
1381        }
1382
1383        // `labels` predates no payload: omitting it must still work via `#[serde(default)]`.
1384        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"}"#;
1385        let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1386        match event {
1387            Event::RunStatusChanged(e) => {
1388                assert!(e.labels.is_empty());
1389                assert_eq!(e.error.as_deref(), Some("boom"));
1390            }
1391            other => panic!("expected RunStatusChanged, got {other:?}"),
1392        }
1393
1394        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"}"#;
1395        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1396        match event {
1397            Event::RunFailed(e) => {
1398                assert_eq!(e.error.as_deref(), Some("boom"));
1399                assert!(e.labels.is_empty());
1400            }
1401            other => panic!("expected RunFailed, got {other:?}"),
1402        }
1403
1404        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"}"#;
1405        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1406        match event {
1407            Event::StepFailed(e) => {
1408                assert_eq!(e.kind, StepKind::Shell);
1409                assert_eq!(e.error, "exit code 1");
1410            }
1411            other => panic!("expected StepFailed, got {other:?}"),
1412        }
1413
1414        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"}"#;
1415        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1416        match event {
1417            Event::LogLine(e) => {
1418                assert_eq!(e.stream, LogStream::Stdout);
1419                assert_eq!(e.line, "hello");
1420                // A payload predating the `id` field degrades to the nil UUID
1421                // rather than failing to deserialize.
1422                assert_eq!(e.id, Uuid::nil());
1423            }
1424            other => panic!("expected LogLine, got {other:?}"),
1425        }
1426
1427        let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1428        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1429        match event {
1430            Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1431            other => panic!("expected UserSignedIn, got {other:?}"),
1432        }
1433    }
1434
1435    /// Guards the internally-tagged representation: payload fields must stay
1436    /// siblings of `type`, never nested under a variant key.
1437    #[test]
1438    fn serialized_event_is_flat_with_type_tag() {
1439        let run_id = Uuid::now_v7();
1440        let event = Event::RunCreated(RunCreatedEvent {
1441            run_id,
1442            workflow_name: "deploy".to_string(),
1443            at: Utc::now(),
1444        });
1445
1446        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1447        let object = value.as_object().expect("event serializes to an object");
1448
1449        assert_eq!(
1450            object.get("type").and_then(|v| v.as_str()),
1451            Some("run_created")
1452        );
1453        assert_eq!(
1454            object.get("workflow_name").and_then(|v| v.as_str()),
1455            Some("deploy")
1456        );
1457        assert_eq!(
1458            object.get("run_id").and_then(|v| v.as_str()),
1459            Some(run_id.to_string().as_str())
1460        );
1461        assert!(object.contains_key("at"));
1462        assert_eq!(object.len(), 4, "no nesting: {object:?}");
1463        assert!(!object.contains_key("RunCreated"));
1464    }
1465
1466    #[test]
1467    fn run_id_returns_some_for_run_events() {
1468        let run_id = Uuid::now_v7();
1469        let now = Utc::now();
1470
1471        let events = vec![
1472            Event::RunCreated(RunCreatedEvent {
1473                run_id,
1474                workflow_name: "w".to_string(),
1475                at: now,
1476            }),
1477            Event::RunStatusChanged(RunStatusChangedEvent {
1478                run_id,
1479                workflow_name: "w".to_string(),
1480                from: RunStatus::Pending,
1481                to: RunStatus::Running,
1482                error: None,
1483                cost_usd: Decimal::ZERO,
1484                duration_ms: 0,
1485                labels: HashMap::new(),
1486                at: now,
1487            }),
1488            Event::RunFailed(RunFailedEvent {
1489                run_id,
1490                workflow_name: "w".to_string(),
1491                error: None,
1492                cost_usd: Decimal::ZERO,
1493                duration_ms: 0,
1494                labels: HashMap::new(),
1495                at: now,
1496            }),
1497            Event::RunBudgetExceeded(RunBudgetExceededEvent {
1498                run_id,
1499                workflow_name: "w".to_string(),
1500                limit_usd: Decimal::ZERO,
1501                spent_usd: Decimal::ZERO,
1502                step_budget_usd: Decimal::ZERO,
1503                at: now,
1504            }),
1505            Event::RetryForced(RetryForcedEvent {
1506                run_id,
1507                workflow_name: "w".to_string(),
1508                original_version: "1".to_string(),
1509                current_version: "2".to_string(),
1510                at: now,
1511            }),
1512            Event::StepCompleted(StepCompletedEvent {
1513                run_id,
1514                step_id: Uuid::now_v7(),
1515                step_name: "s".to_string(),
1516                kind: StepKind::Shell,
1517                duration_ms: 0,
1518                cost_usd: Decimal::ZERO,
1519                at: now,
1520            }),
1521            Event::StepFailed(StepFailedEvent {
1522                run_id,
1523                step_id: Uuid::now_v7(),
1524                step_name: "s".to_string(),
1525                kind: StepKind::Shell,
1526                error: "e".to_string(),
1527                at: now,
1528            }),
1529            Event::ApprovalRequested(ApprovalRequestedEvent {
1530                run_id,
1531                step_id: Uuid::now_v7(),
1532                message: "ok?".to_string(),
1533                requirement: None,
1534                at: now,
1535            }),
1536            Event::ApprovalGranted(ApprovalGrantedEvent {
1537                run_id,
1538                step_id: None,
1539                approved_by: "alice".to_string(),
1540                approvals_received: 1,
1541                approvals_required: 1,
1542                requirement: None,
1543                at: now,
1544            }),
1545            Event::ApprovalRejected(ApprovalRejectedEvent {
1546                run_id,
1547                step_id: None,
1548                rejected_by: "bob".to_string(),
1549                requirement: None,
1550                at: now,
1551            }),
1552            Event::LogLine(LogLineEvent {
1553                id: Uuid::now_v7(),
1554                run_id,
1555                step_id: Uuid::now_v7(),
1556                step_name: "s".to_string(),
1557                stream: LogStream::Stdout,
1558                line: "l".to_string(),
1559                at: now,
1560            }),
1561        ];
1562
1563        for event in &events {
1564            assert_eq!(
1565                event.run_id(),
1566                Some(run_id),
1567                "{} should carry a run_id",
1568                event.event_type()
1569            );
1570        }
1571    }
1572
1573    #[test]
1574    fn run_id_returns_none_for_auth_events() {
1575        let user_id = Uuid::now_v7();
1576        let now = Utc::now();
1577
1578        let events = vec![
1579            Event::UserSignedIn(UserSignedInEvent {
1580                user_id,
1581                username: "alice".to_string(),
1582                at: now,
1583            }),
1584            Event::UserSignedUp(UserSignedUpEvent {
1585                user_id,
1586                username: "alice".to_string(),
1587                at: now,
1588            }),
1589            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1590        ];
1591
1592        for event in &events {
1593            assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1594        }
1595    }
1596
1597    #[test]
1598    fn step_id_returns_some_only_for_step_events() {
1599        let step_id = Uuid::now_v7();
1600        let run_id = Uuid::now_v7();
1601        let now = Utc::now();
1602
1603        let with_step = vec![
1604            Event::StepCompleted(StepCompletedEvent {
1605                run_id,
1606                step_id,
1607                step_name: "s".to_string(),
1608                kind: StepKind::Shell,
1609                duration_ms: 0,
1610                cost_usd: Decimal::ZERO,
1611                at: now,
1612            }),
1613            Event::StepFailed(StepFailedEvent {
1614                run_id,
1615                step_id,
1616                step_name: "s".to_string(),
1617                kind: StepKind::Shell,
1618                error: "e".to_string(),
1619                at: now,
1620            }),
1621            Event::ApprovalRequested(ApprovalRequestedEvent {
1622                run_id,
1623                step_id,
1624                message: "ok?".to_string(),
1625                requirement: None,
1626                at: now,
1627            }),
1628            Event::ApprovalGranted(ApprovalGrantedEvent {
1629                run_id,
1630                step_id: Some(step_id),
1631                approved_by: "alice".to_string(),
1632                approvals_received: 1,
1633                approvals_required: 2,
1634                requirement: None,
1635                at: now,
1636            }),
1637            Event::ApprovalRejected(ApprovalRejectedEvent {
1638                run_id,
1639                step_id: Some(step_id),
1640                rejected_by: "bob".to_string(),
1641                requirement: None,
1642                at: now,
1643            }),
1644        ];
1645
1646        for event in &with_step {
1647            assert_eq!(
1648                event.step_id(),
1649                Some(step_id),
1650                "{} should carry a step_id",
1651                event.event_type()
1652            );
1653        }
1654
1655        let without_step = vec![
1656            Event::RunCreated(RunCreatedEvent {
1657                run_id,
1658                workflow_name: "w".to_string(),
1659                at: now,
1660            }),
1661            // A granted event recorded before step ids were carried.
1662            Event::ApprovalGranted(ApprovalGrantedEvent {
1663                run_id,
1664                step_id: None,
1665                approved_by: "alice".to_string(),
1666                approvals_received: 1,
1667                approvals_required: 1,
1668                requirement: None,
1669                at: now,
1670            }),
1671            // LogLine carries a step_id field but is reported as a run-level
1672            // stream event, matching the pre-refactor behaviour.
1673            Event::LogLine(LogLineEvent {
1674                id: Uuid::now_v7(),
1675                run_id,
1676                step_id,
1677                step_name: "s".to_string(),
1678                stream: LogStream::Stdout,
1679                line: "l".to_string(),
1680                at: now,
1681            }),
1682            Event::UserSignedOut(UserSignedOutEvent {
1683                user_id: Uuid::now_v7(),
1684                at: now,
1685            }),
1686        ];
1687
1688        for event in &without_step {
1689            assert_eq!(
1690                event.step_id(),
1691                None,
1692                "{} should not carry a step_id",
1693                event.event_type()
1694            );
1695        }
1696    }
1697
1698    #[test]
1699    fn user_id_returns_some_only_for_auth_events() {
1700        let user_id = Uuid::now_v7();
1701        let run_id = Uuid::now_v7();
1702        let now = Utc::now();
1703
1704        let auth = vec![
1705            Event::UserSignedIn(UserSignedInEvent {
1706                user_id,
1707                username: "alice".to_string(),
1708                at: now,
1709            }),
1710            Event::UserSignedUp(UserSignedUpEvent {
1711                user_id,
1712                username: "alice".to_string(),
1713                at: now,
1714            }),
1715            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1716        ];
1717
1718        for event in &auth {
1719            assert_eq!(
1720                event.user_id(),
1721                Some(user_id),
1722                "{} should carry a user_id",
1723                event.event_type()
1724            );
1725        }
1726
1727        let non_auth = vec![
1728            Event::RunCreated(RunCreatedEvent {
1729                run_id,
1730                workflow_name: "w".to_string(),
1731                at: now,
1732            }),
1733            Event::StepFailed(StepFailedEvent {
1734                run_id,
1735                step_id: Uuid::now_v7(),
1736                step_name: "s".to_string(),
1737                kind: StepKind::Shell,
1738                error: "e".to_string(),
1739                at: now,
1740            }),
1741        ];
1742
1743        for event in &non_auth {
1744            assert_eq!(
1745                event.user_id(),
1746                None,
1747                "{} should not carry a user_id",
1748                event.event_type()
1749            );
1750        }
1751    }
1752
1753    #[test]
1754    fn event_type_all_variants() {
1755        let id = Uuid::now_v7();
1756        let now = Utc::now();
1757
1758        let cases: Vec<(Event, &str)> = vec![
1759            (
1760                Event::RunCreated(RunCreatedEvent {
1761                    run_id: id,
1762                    workflow_name: "w".to_string(),
1763                    at: now,
1764                }),
1765                "run_created",
1766            ),
1767            (
1768                Event::RunStatusChanged(RunStatusChangedEvent {
1769                    run_id: id,
1770                    workflow_name: "w".to_string(),
1771                    from: RunStatus::Pending,
1772                    to: RunStatus::Running,
1773                    error: None,
1774                    cost_usd: Decimal::ZERO,
1775                    duration_ms: 0,
1776                    labels: HashMap::new(),
1777                    at: now,
1778                }),
1779                "run_status_changed",
1780            ),
1781            (
1782                Event::RunFailed(RunFailedEvent {
1783                    run_id: id,
1784                    workflow_name: "w".to_string(),
1785                    error: Some("boom".to_string()),
1786                    cost_usd: Decimal::ZERO,
1787                    duration_ms: 0,
1788                    labels: HashMap::new(),
1789                    at: now,
1790                }),
1791                "run_failed",
1792            ),
1793            (
1794                Event::RunBudgetExceeded(RunBudgetExceededEvent {
1795                    run_id: id,
1796                    workflow_name: "w".to_string(),
1797                    limit_usd: Decimal::new(200, 2),
1798                    spent_usd: Decimal::new(180, 2),
1799                    step_budget_usd: Decimal::new(50, 2),
1800                    at: now,
1801                }),
1802                "run_budget_exceeded",
1803            ),
1804            (
1805                Event::RetryForced(RetryForcedEvent {
1806                    run_id: id,
1807                    workflow_name: "w".to_string(),
1808                    original_version: "1".to_string(),
1809                    current_version: "2".to_string(),
1810                    at: now,
1811                }),
1812                "retry_forced",
1813            ),
1814            (
1815                Event::StepCompleted(StepCompletedEvent {
1816                    run_id: id,
1817                    step_id: id,
1818                    step_name: "s".to_string(),
1819                    kind: StepKind::Shell,
1820                    duration_ms: 0,
1821                    cost_usd: Decimal::ZERO,
1822                    at: now,
1823                }),
1824                "step_completed",
1825            ),
1826            (
1827                Event::StepFailed(StepFailedEvent {
1828                    run_id: id,
1829                    step_id: id,
1830                    step_name: "s".to_string(),
1831                    kind: StepKind::Shell,
1832                    error: "err".to_string(),
1833                    at: now,
1834                }),
1835                "step_failed",
1836            ),
1837            (
1838                Event::ApprovalRequested(ApprovalRequestedEvent {
1839                    run_id: id,
1840                    step_id: id,
1841                    message: "ok?".to_string(),
1842                    requirement: None,
1843                    at: now,
1844                }),
1845                "approval_requested",
1846            ),
1847            (
1848                Event::ApprovalGranted(ApprovalGrantedEvent {
1849                    run_id: id,
1850                    step_id: Some(id),
1851                    approved_by: "alice".to_string(),
1852                    approvals_received: 1,
1853                    approvals_required: 1,
1854                    requirement: None,
1855                    at: now,
1856                }),
1857                "approval_granted",
1858            ),
1859            (
1860                Event::ApprovalRejected(ApprovalRejectedEvent {
1861                    run_id: id,
1862                    step_id: Some(id),
1863                    rejected_by: "bob".to_string(),
1864                    requirement: None,
1865                    at: now,
1866                }),
1867                "approval_rejected",
1868            ),
1869            (
1870                Event::ApprovalEscalated(ApprovalEscalatedEvent {
1871                    run_id: id,
1872                    step_id: id,
1873                    step_name: "prod-gate".to_string(),
1874                    stage: 0,
1875                    policy: "auto_reject".to_string(),
1876                    action: "rejected".to_string(),
1877                    reason: "approval deadline of 3600s expired".to_string(),
1878                    assignee: None,
1879                    at: now,
1880                }),
1881                "approval_escalated",
1882            ),
1883            (
1884                Event::LogLine(LogLineEvent {
1885                    id,
1886                    run_id: id,
1887                    step_id: id,
1888                    step_name: "build".to_string(),
1889                    stream: LogStream::Stdout,
1890                    line: "Compiling ironflow v0.1.0".to_string(),
1891                    at: now,
1892                }),
1893                "log_line",
1894            ),
1895            (
1896                Event::UserSignedIn(UserSignedInEvent {
1897                    user_id: id,
1898                    username: "u".to_string(),
1899                    at: now,
1900                }),
1901                "user_signed_in",
1902            ),
1903            (
1904                Event::UserSignedUp(UserSignedUpEvent {
1905                    user_id: id,
1906                    username: "u".to_string(),
1907                    at: now,
1908                }),
1909                "user_signed_up",
1910            ),
1911            (
1912                Event::UserSignedOut(UserSignedOutEvent {
1913                    user_id: id,
1914                    at: now,
1915                }),
1916                "user_signed_out",
1917            ),
1918            (
1919                Event::ProviderAccountUpdated(ProviderAccountUpdatedEvent {
1920                    account_id: id,
1921                    name: "perso".to_string(),
1922                    change: ProviderAccountChange::TokenReplaced,
1923                    at: now,
1924                }),
1925                "provider_account.updated",
1926            ),
1927            (
1928                Event::ProviderAccountUsageUpdated(ProviderAccountUsageUpdatedEvent {
1929                    account_id: id,
1930                    name: "perso".to_string(),
1931                    windows: Vec::new(),
1932                    at: now,
1933                }),
1934                "provider_account.usage_updated",
1935            ),
1936            (
1937                Event::SignalAwaited(SignalAwaitedEvent {
1938                    run_id: id,
1939                    step_id: id,
1940                    step_name: "wait-ci".to_string(),
1941                    name: "ci.done".to_string(),
1942                    key: "abc".to_string(),
1943                    deadline_at: now,
1944                    at: now,
1945                }),
1946                "signal_awaited",
1947            ),
1948            (
1949                Event::SignalReceived(SignalReceivedEvent {
1950                    signal_id: id,
1951                    name: "ci.done".to_string(),
1952                    key: "abc".to_string(),
1953                    resumed_runs: vec![id],
1954                    at: now,
1955                }),
1956                "signal_received",
1957            ),
1958        ];
1959
1960        assert_eq!(
1961            cases.len(),
1962            Event::ALL.len(),
1963            "every variant must be covered"
1964        );
1965
1966        for (event, expected_type) in cases {
1967            assert_eq!(event.event_type(), expected_type);
1968        }
1969    }
1970
1971    #[test]
1972    fn approval_escalated_serde_roundtrip() {
1973        let run_id = Uuid::now_v7();
1974        let step_id = Uuid::now_v7();
1975        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1976            run_id,
1977            step_id,
1978            step_name: "prod-gate".to_string(),
1979            stage: 1,
1980            policy: "escalate".to_string(),
1981            action: "reassigned to sre-oncall".to_string(),
1982            reason: "approval deadline of 3600s expired".to_string(),
1983            assignee: Some(Assignee::group("sre-oncall")),
1984            at: Utc::now(),
1985        });
1986
1987        let json = serde_json::to_string(&event).expect("serialize");
1988        assert!(
1989            json.contains("\"type\":\"approval_escalated\""),
1990            "got {json}"
1991        );
1992
1993        let back: Event = serde_json::from_str(&json).expect("deserialize");
1994        let Event::ApprovalEscalated(payload) = back else {
1995            panic!("expected an approval_escalated event");
1996        };
1997        assert_eq!(payload.run_id, run_id);
1998        assert_eq!(payload.stage, 1);
1999        assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
2000    }
2001
2002    #[test]
2003    fn approval_escalated_carries_run_and_step_ids() {
2004        let run_id = Uuid::now_v7();
2005        let step_id = Uuid::now_v7();
2006        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
2007            run_id,
2008            step_id,
2009            step_name: "prod-gate".to_string(),
2010            stage: 0,
2011            policy: "notify".to_string(),
2012            action: "notified 1 target".to_string(),
2013            reason: "approval deadline of 60s expired".to_string(),
2014            assignee: None,
2015            at: Utc::now(),
2016        });
2017
2018        assert_eq!(event.run_id(), Some(run_id));
2019        assert_eq!(event.step_id(), Some(step_id));
2020        assert_eq!(event.user_id(), None);
2021    }
2022
2023    #[test]
2024    fn signal_awaited_serde_roundtrip() {
2025        let run_id = Uuid::now_v7();
2026        let step_id = Uuid::now_v7();
2027        let event = Event::SignalAwaited(SignalAwaitedEvent {
2028            run_id,
2029            step_id,
2030            step_name: "wait-ci".to_string(),
2031            name: "ci.done".to_string(),
2032            key: "abc".to_string(),
2033            deadline_at: Utc::now(),
2034            at: Utc::now(),
2035        });
2036
2037        let json = serde_json::to_string(&event).expect("serialize");
2038        assert!(json.contains("\"type\":\"signal_awaited\""), "got {json}");
2039        let back: Event = serde_json::from_str(&json).expect("deserialize");
2040        assert_eq!(back.event_type(), Event::SIGNAL_AWAITED);
2041        assert_eq!(back.run_id(), Some(run_id));
2042        assert_eq!(back.step_id(), Some(step_id));
2043        assert_eq!(back.user_id(), None);
2044    }
2045
2046    #[test]
2047    fn signal_received_serde_roundtrip() {
2048        let resumed = Uuid::now_v7();
2049        let event = Event::SignalReceived(SignalReceivedEvent {
2050            signal_id: Uuid::now_v7(),
2051            name: "ci.done".to_string(),
2052            key: "abc".to_string(),
2053            resumed_runs: vec![resumed],
2054            at: Utc::now(),
2055        });
2056
2057        let json = serde_json::to_string(&event).expect("serialize");
2058        assert!(json.contains("\"type\":\"signal_received\""), "got {json}");
2059        let back: Event = serde_json::from_str(&json).expect("deserialize");
2060        let Event::SignalReceived(payload) = &back else {
2061            panic!("expected a signal_received event, got {back:?}");
2062        };
2063        assert_eq!(payload.resumed_runs, vec![resumed]);
2064        assert_eq!(back.run_id(), None);
2065        assert_eq!(back.step_id(), None);
2066        assert_eq!(back.user_id(), None);
2067    }
2068}