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