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::{Assignee, RunStatus, StepKind};
12
13/// Payload of the `Event::RunCreated` event.
14///
15/// # Examples
16///
17/// ```
18/// use chrono::Utc;
19/// use ironflow_engine::notify::RunCreatedEvent;
20/// use uuid::Uuid;
21///
22/// let payload = RunCreatedEvent {
23///     run_id: Uuid::now_v7(),
24///     workflow_name: "deploy".to_string(),
25///     at: Utc::now(),
26/// };
27/// assert_eq!(payload.workflow_name, "deploy");
28/// ```
29#[derive(Debug, Clone, Serialize, Deserialize)]
30#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
31pub struct RunCreatedEvent {
32    /// Run identifier.
33    pub run_id: Uuid,
34    /// Workflow name.
35    pub workflow_name: String,
36    /// When the run was created.
37    pub at: DateTime<Utc>,
38}
39
40/// Payload of the `Event::RunStatusChanged` event.
41///
42/// # Examples
43///
44/// ```
45/// use std::collections::HashMap;
46///
47/// use chrono::Utc;
48/// use ironflow_engine::notify::RunStatusChangedEvent;
49/// use ironflow_store::models::RunStatus;
50/// use rust_decimal::Decimal;
51/// use uuid::Uuid;
52///
53/// let payload = RunStatusChangedEvent {
54///     run_id: Uuid::now_v7(),
55///     workflow_name: "deploy".to_string(),
56///     from: RunStatus::Running,
57///     to: RunStatus::Completed,
58///     error: None,
59///     cost_usd: Decimal::ZERO,
60///     duration_ms: 5000,
61///     labels: HashMap::new(),
62///     at: Utc::now(),
63/// };
64/// assert_eq!(payload.to, RunStatus::Completed);
65/// ```
66#[derive(Debug, Clone, Serialize, Deserialize)]
67#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
68pub struct RunStatusChangedEvent {
69    /// Run identifier.
70    pub run_id: Uuid,
71    /// Workflow name.
72    pub workflow_name: String,
73    /// Previous status.
74    pub from: RunStatus,
75    /// New status.
76    pub to: RunStatus,
77    /// Error message (when transitioning to Failed).
78    pub error: Option<String>,
79    /// Aggregated cost in USD at the time of transition.
80    pub cost_usd: Decimal,
81    /// Aggregated duration in milliseconds at the time of transition.
82    pub duration_ms: u64,
83    /// Labels of the run at the time of the transition.
84    #[serde(default)]
85    pub labels: HashMap<String, String>,
86    /// When the transition occurred.
87    pub at: DateTime<Utc>,
88}
89
90/// Payload of the `Event::RunFailed` event.
91///
92/// # Examples
93///
94/// ```
95/// use std::collections::HashMap;
96///
97/// use chrono::Utc;
98/// use ironflow_engine::notify::RunFailedEvent;
99/// use rust_decimal::Decimal;
100/// use uuid::Uuid;
101///
102/// let payload = RunFailedEvent {
103///     run_id: Uuid::now_v7(),
104///     workflow_name: "deploy".to_string(),
105///     error: Some("step crashed".to_string()),
106///     cost_usd: Decimal::ZERO,
107///     duration_ms: 3000,
108///     labels: HashMap::new(),
109///     at: Utc::now(),
110/// };
111/// assert_eq!(payload.error.as_deref(), Some("step crashed"));
112/// ```
113#[derive(Debug, Clone, Serialize, Deserialize)]
114#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
115pub struct RunFailedEvent {
116    /// Run identifier.
117    pub run_id: Uuid,
118    /// Workflow name.
119    pub workflow_name: String,
120    /// Error message.
121    pub error: Option<String>,
122    /// Aggregated cost in USD at the time of failure.
123    pub cost_usd: Decimal,
124    /// Aggregated duration in milliseconds at the time of failure.
125    pub duration_ms: u64,
126    /// Labels of the run at the time of the failure.
127    #[serde(default)]
128    pub labels: HashMap<String, String>,
129    /// When the failure occurred.
130    pub at: DateTime<Utc>,
131}
132
133/// Payload of the `Event::RunBudgetExceeded` event.
134///
135/// # Examples
136///
137/// ```
138/// use chrono::Utc;
139/// use ironflow_engine::notify::RunBudgetExceededEvent;
140/// use rust_decimal::Decimal;
141/// use uuid::Uuid;
142///
143/// let payload = RunBudgetExceededEvent {
144///     run_id: Uuid::now_v7(),
145///     workflow_name: "deploy".to_string(),
146///     limit_usd: Decimal::new(200, 2),
147///     spent_usd: Decimal::new(180, 2),
148///     step_budget_usd: Decimal::new(50, 2),
149///     at: Utc::now(),
150/// };
151/// assert_eq!(payload.limit_usd, Decimal::new(200, 2));
152/// ```
153#[derive(Debug, Clone, Serialize, Deserialize)]
154#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
155pub struct RunBudgetExceededEvent {
156    /// Run identifier.
157    pub run_id: Uuid,
158    /// Workflow name.
159    pub workflow_name: String,
160    /// The configured cost cap in USD.
161    pub limit_usd: Decimal,
162    /// Cost already consumed when the cap was reached, in USD.
163    pub spent_usd: Decimal,
164    /// Declared budget of the refused step, in USD.
165    pub step_budget_usd: Decimal,
166    /// When the refusal occurred.
167    pub at: DateTime<Utc>,
168}
169
170/// Payload of the `Event::RetryForced` event.
171///
172/// # Examples
173///
174/// ```
175/// use chrono::Utc;
176/// use ironflow_engine::notify::RetryForcedEvent;
177/// use uuid::Uuid;
178///
179/// let payload = RetryForcedEvent {
180///     run_id: Uuid::now_v7(),
181///     workflow_name: "deploy".to_string(),
182///     original_version: "1".to_string(),
183///     current_version: "2".to_string(),
184///     at: Utc::now(),
185/// };
186/// assert_eq!(payload.current_version, "2");
187/// ```
188#[derive(Debug, Clone, Serialize, Deserialize)]
189#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
190pub struct RetryForcedEvent {
191    /// The new run created by the forced retry.
192    pub run_id: Uuid,
193    /// Workflow name.
194    pub workflow_name: String,
195    /// Version stored on the original run.
196    pub original_version: String,
197    /// Current version of the handler.
198    pub current_version: String,
199    /// When the forced retry occurred.
200    pub at: DateTime<Utc>,
201}
202
203/// Payload of the `Event::StepCompleted` event.
204///
205/// # Examples
206///
207/// ```
208/// use chrono::Utc;
209/// use ironflow_engine::notify::StepCompletedEvent;
210/// use ironflow_store::models::StepKind;
211/// use rust_decimal::Decimal;
212/// use uuid::Uuid;
213///
214/// let payload = StepCompletedEvent {
215///     run_id: Uuid::now_v7(),
216///     step_id: Uuid::now_v7(),
217///     step_name: "build".to_string(),
218///     kind: StepKind::Shell,
219///     duration_ms: 1200,
220///     cost_usd: Decimal::ZERO,
221///     at: Utc::now(),
222/// };
223/// assert_eq!(payload.step_name, "build");
224/// ```
225#[derive(Debug, Clone, Serialize, Deserialize)]
226#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
227pub struct StepCompletedEvent {
228    /// Run identifier.
229    pub run_id: Uuid,
230    /// Step identifier.
231    pub step_id: Uuid,
232    /// Human-readable step name.
233    pub step_name: String,
234    /// Step operation kind.
235    #[cfg_attr(feature = "openapi", schema(value_type = String))]
236    pub kind: StepKind,
237    /// Step duration in milliseconds.
238    pub duration_ms: u64,
239    /// Step cost in USD.
240    pub cost_usd: Decimal,
241    /// When the step completed.
242    pub at: DateTime<Utc>,
243}
244
245/// Payload of the `Event::StepFailed` event.
246///
247/// # Examples
248///
249/// ```
250/// use chrono::Utc;
251/// use ironflow_engine::notify::StepFailedEvent;
252/// use ironflow_store::models::StepKind;
253/// use uuid::Uuid;
254///
255/// let payload = StepFailedEvent {
256///     run_id: Uuid::now_v7(),
257///     step_id: Uuid::now_v7(),
258///     step_name: "build".to_string(),
259///     kind: StepKind::Shell,
260///     error: "exit code 1".to_string(),
261///     at: Utc::now(),
262/// };
263/// assert_eq!(payload.error, "exit code 1");
264/// ```
265#[derive(Debug, Clone, Serialize, Deserialize)]
266#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
267pub struct StepFailedEvent {
268    /// Run identifier.
269    pub run_id: Uuid,
270    /// Step identifier.
271    pub step_id: Uuid,
272    /// Human-readable step name.
273    pub step_name: String,
274    /// Step operation kind.
275    #[cfg_attr(feature = "openapi", schema(value_type = String))]
276    pub kind: StepKind,
277    /// Error message.
278    pub error: String,
279    /// When the step failed.
280    pub at: DateTime<Utc>,
281}
282
283/// Payload of the `Event::ApprovalRequested` event.
284///
285/// # Examples
286///
287/// ```
288/// use chrono::Utc;
289/// use ironflow_engine::notify::ApprovalRequestedEvent;
290/// use uuid::Uuid;
291///
292/// let payload = ApprovalRequestedEvent {
293///     run_id: Uuid::now_v7(),
294///     step_id: Uuid::now_v7(),
295///     message: "Deploy to prod?".to_string(),
296///     at: Utc::now(),
297/// };
298/// assert_eq!(payload.message, "Deploy to prod?");
299/// ```
300#[derive(Debug, Clone, Serialize, Deserialize)]
301#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
302pub struct ApprovalRequestedEvent {
303    /// Run identifier.
304    pub run_id: Uuid,
305    /// Approval step identifier.
306    pub step_id: Uuid,
307    /// Message displayed to reviewers.
308    pub message: String,
309    /// When the approval was requested.
310    pub at: DateTime<Utc>,
311}
312
313/// Payload of the `Event::ApprovalGranted` event.
314///
315/// # Examples
316///
317/// ```
318/// use chrono::Utc;
319/// use ironflow_engine::notify::ApprovalGrantedEvent;
320/// use uuid::Uuid;
321///
322/// let payload = ApprovalGrantedEvent {
323///     run_id: Uuid::now_v7(),
324///     approved_by: "alice".to_string(),
325///     at: Utc::now(),
326/// };
327/// assert_eq!(payload.approved_by, "alice");
328/// ```
329#[derive(Debug, Clone, Serialize, Deserialize)]
330#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
331pub struct ApprovalGrantedEvent {
332    /// Run identifier.
333    pub run_id: Uuid,
334    /// User who approved (ID or username).
335    pub approved_by: String,
336    /// When the approval was granted.
337    pub at: DateTime<Utc>,
338}
339
340/// Payload of the `Event::ApprovalRejected` event.
341///
342/// # Examples
343///
344/// ```
345/// use chrono::Utc;
346/// use ironflow_engine::notify::ApprovalRejectedEvent;
347/// use uuid::Uuid;
348///
349/// let payload = ApprovalRejectedEvent {
350///     run_id: Uuid::now_v7(),
351///     rejected_by: "bob".to_string(),
352///     at: Utc::now(),
353/// };
354/// assert_eq!(payload.rejected_by, "bob");
355/// ```
356#[derive(Debug, Clone, Serialize, Deserialize)]
357#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
358pub struct ApprovalRejectedEvent {
359    /// Run identifier.
360    pub run_id: Uuid,
361    /// User who rejected (ID or username).
362    pub rejected_by: String,
363    /// When the rejection occurred.
364    pub at: DateTime<Utc>,
365}
366
367/// Payload of the `Event::ApprovalEscalated` event.
368///
369/// Emitted every time an approval gate misses its SLA deadline, including the
370/// repeated firings of a bare `Notify`/`Escalate` policy and the final
371/// "chain exhausted" notice. The audit log persists it verbatim, so the whole
372/// escalation history of a gate is reconstructable from it.
373///
374/// # Examples
375///
376/// ```
377/// use chrono::Utc;
378/// use ironflow_engine::notify::ApprovalEscalatedEvent;
379/// use uuid::Uuid;
380///
381/// let payload = ApprovalEscalatedEvent {
382///     run_id: Uuid::now_v7(),
383///     step_id: Uuid::now_v7(),
384///     step_name: "prod-gate".to_string(),
385///     stage: 0,
386///     policy: "auto_reject".to_string(),
387///     action: "rejected".to_string(),
388///     reason: "approval deadline of 3600s expired".to_string(),
389///     assignee: None,
390///     at: Utc::now(),
391/// };
392/// assert_eq!(payload.policy, "auto_reject");
393/// ```
394#[derive(Debug, Clone, Serialize, Deserialize)]
395#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
396pub struct ApprovalEscalatedEvent {
397    /// Run identifier.
398    pub run_id: Uuid,
399    /// Approval step identifier.
400    pub step_id: Uuid,
401    /// Human-readable step name.
402    pub step_name: String,
403    /// Escalation stage that fired (0-based index into the policy chain).
404    pub stage: u32,
405    /// Policy applied, e.g. `"auto_reject"`, `"notify"`, `"escalate"`.
406    pub policy: String,
407    /// What the escalation did, for the audit log.
408    pub action: String,
409    /// Why it fired, e.g. `"approval deadline of 3600s expired"`.
410    pub reason: String,
411    /// Assignee after the escalation, when it reassigned the gate.
412    #[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
413    pub assignee: Option<Assignee>,
414    /// When the escalation ran.
415    pub at: DateTime<Utc>,
416}
417
418/// Payload of the `Event::LogLine` event.
419///
420/// # Examples
421///
422/// ```
423/// use chrono::Utc;
424/// use ironflow_engine::notify::{LogLineEvent, LogStream};
425/// use uuid::Uuid;
426///
427/// let payload = LogLineEvent {
428///     run_id: Uuid::now_v7(),
429///     step_id: Uuid::now_v7(),
430///     step_name: "build".to_string(),
431///     stream: LogStream::Stdout,
432///     line: "Compiling ironflow v0.1.0".to_string(),
433///     at: Utc::now(),
434/// };
435/// assert_eq!(payload.line, "Compiling ironflow v0.1.0");
436/// ```
437#[derive(Debug, Clone, Serialize, Deserialize)]
438#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
439pub struct LogLineEvent {
440    /// Run identifier.
441    pub run_id: Uuid,
442    /// Step identifier.
443    pub step_id: Uuid,
444    /// Human-readable step name.
445    pub step_name: String,
446    /// Output stream.
447    pub stream: LogStream,
448    /// The log line content.
449    pub line: String,
450    /// When the line was emitted.
451    pub at: DateTime<Utc>,
452}
453
454/// Payload of the `Event::UserSignedIn` event.
455///
456/// # Examples
457///
458/// ```
459/// use chrono::Utc;
460/// use ironflow_engine::notify::UserSignedInEvent;
461/// use uuid::Uuid;
462///
463/// let payload = UserSignedInEvent {
464///     user_id: Uuid::now_v7(),
465///     username: "alice".to_string(),
466///     at: Utc::now(),
467/// };
468/// assert_eq!(payload.username, "alice");
469/// ```
470#[derive(Debug, Clone, Serialize, Deserialize)]
471#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
472pub struct UserSignedInEvent {
473    /// User identifier.
474    pub user_id: Uuid,
475    /// Username.
476    pub username: String,
477    /// When the sign-in occurred.
478    pub at: DateTime<Utc>,
479}
480
481/// Payload of the `Event::UserSignedUp` event.
482///
483/// # Examples
484///
485/// ```
486/// use chrono::Utc;
487/// use ironflow_engine::notify::UserSignedUpEvent;
488/// use uuid::Uuid;
489///
490/// let payload = UserSignedUpEvent {
491///     user_id: Uuid::now_v7(),
492///     username: "alice".to_string(),
493///     at: Utc::now(),
494/// };
495/// assert_eq!(payload.username, "alice");
496/// ```
497#[derive(Debug, Clone, Serialize, Deserialize)]
498#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
499pub struct UserSignedUpEvent {
500    /// User identifier.
501    pub user_id: Uuid,
502    /// Username.
503    pub username: String,
504    /// When the sign-up occurred.
505    pub at: DateTime<Utc>,
506}
507
508/// Payload of the `Event::UserSignedOut` event.
509///
510/// # Examples
511///
512/// ```
513/// use chrono::Utc;
514/// use ironflow_engine::notify::UserSignedOutEvent;
515/// use uuid::Uuid;
516///
517/// let user_id = Uuid::now_v7();
518/// let payload = UserSignedOutEvent {
519///     user_id,
520///     at: Utc::now(),
521/// };
522/// assert_eq!(payload.user_id, user_id);
523/// ```
524#[derive(Debug, Clone, Serialize, Deserialize)]
525#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
526pub struct UserSignedOutEvent {
527    /// User identifier.
528    pub user_id: Uuid,
529    /// When the sign-out occurred.
530    pub at: DateTime<Utc>,
531}
532
533/// A domain event emitted by the ironflow system.
534///
535/// Covers the full lifecycle: runs, steps, approvals, and authentication.
536/// Subscribers receive these via [`EventPublisher`](super::EventPublisher)
537/// and pattern-match on the variants they care about.
538///
539/// Each variant wraps a dedicated payload struct. The serialized form stays
540/// flat: the `type` discriminant sits next to the payload fields, so
541/// `{"type":"run_created","run_id":...}` round-trips unchanged.
542///
543/// # Examples
544///
545/// ```
546/// use std::collections::HashMap;
547/// use ironflow_engine::notify::{Event, RunStatusChangedEvent};
548/// use ironflow_store::models::RunStatus;
549/// use uuid::Uuid;
550///
551/// let event = Event::RunStatusChanged(RunStatusChangedEvent {
552///     run_id: Uuid::now_v7(),
553///     workflow_name: "deploy".to_string(),
554///     from: RunStatus::Running,
555///     to: RunStatus::Completed,
556///     error: None,
557///     cost_usd: rust_decimal::Decimal::ZERO,
558///     duration_ms: 5000,
559///     labels: HashMap::new(),
560///     at: chrono::Utc::now(),
561/// });
562/// assert_eq!(event.event_type(), "run_status_changed");
563/// ```
564#[derive(Debug, Clone, Serialize, Deserialize)]
565#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
566#[serde(tag = "type", rename_all = "snake_case")]
567pub enum Event {
568    // -- Run lifecycle --
569    /// A new run was created (status: Pending).
570    RunCreated(RunCreatedEvent),
571
572    /// A run changed status.
573    RunStatusChanged(RunStatusChangedEvent),
574
575    /// A run transitioned to [`Failed`](ironflow_store::models::RunStatus::Failed).
576    ///
577    /// This is a convenience event emitted alongside [`RunStatusChanged`](Event::RunStatusChanged)
578    /// when the target status is `Failed`. Subscribe to this instead of
579    /// `RUN_STATUS_CHANGED` when you only care about failures.
580    RunFailed(RunFailedEvent),
581
582    /// A run was stopped because it reached its cumulative cost cap.
583    ///
584    /// Emitted when the engine refuses an agent step that would cross the run's
585    /// `max_cost_usd`. The run transitions to
586    /// [`Cancelled`](ironflow_store::models::RunStatus::Cancelled) and the step
587    /// is never launched, so the reported spend is what the run had already
588    /// consumed.
589    RunBudgetExceeded(RunBudgetExceededEvent),
590
591    /// A manual retry was forced despite a handler version mismatch.
592    ///
593    /// Emitted when a caller passes `force=true` on a retry where the
594    /// handler version differs from the run's recorded version. This is
595    /// an audit event: it means the new code will execute on the old
596    /// payload without the handler explicitly declaring compatibility.
597    RetryForced(RetryForcedEvent),
598
599    // -- Step lifecycle --
600    /// A step completed successfully.
601    StepCompleted(StepCompletedEvent),
602
603    /// A step failed.
604    StepFailed(StepFailedEvent),
605
606    // -- Approval --
607    /// A run is waiting for human approval.
608    ApprovalRequested(ApprovalRequestedEvent),
609
610    /// A run was approved by a human.
611    ApprovalGranted(ApprovalGrantedEvent),
612
613    /// A run was rejected by a human.
614    ApprovalRejected(ApprovalRejectedEvent),
615
616    /// An approval gate missed its SLA deadline and an escalation policy ran.
617    ApprovalEscalated(ApprovalEscalatedEvent),
618
619    // -- Log streaming --
620    /// A log line emitted during step execution.
621    ///
622    /// Pushed by the worker in real time so that SSE clients can stream
623    /// step output as it happens, without waiting for step completion.
624    LogLine(LogLineEvent),
625
626    // -- Authentication --
627    /// A user signed in.
628    UserSignedIn(UserSignedInEvent),
629
630    /// A new user signed up.
631    UserSignedUp(UserSignedUpEvent),
632
633    /// A user signed out.
634    UserSignedOut(UserSignedOutEvent),
635}
636
637impl Event {
638    /// Event type constant for [`RunCreated`](Event::RunCreated).
639    pub const RUN_CREATED: &'static str = "run_created";
640    /// Event type constant for [`RunStatusChanged`](Event::RunStatusChanged).
641    pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
642    /// Event type constant for [`RunFailed`](Event::RunFailed).
643    pub const RUN_FAILED: &'static str = "run_failed";
644    /// Event type constant for [`RunBudgetExceeded`](Event::RunBudgetExceeded).
645    pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
646    /// Event type constant for [`RetryForced`](Event::RetryForced).
647    pub const RETRY_FORCED: &'static str = "retry_forced";
648    /// Event type constant for [`StepCompleted`](Event::StepCompleted).
649    pub const STEP_COMPLETED: &'static str = "step_completed";
650    /// Event type constant for [`StepFailed`](Event::StepFailed).
651    pub const STEP_FAILED: &'static str = "step_failed";
652    /// Event type constant for [`ApprovalRequested`](Event::ApprovalRequested).
653    pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
654    /// Event type constant for [`ApprovalGranted`](Event::ApprovalGranted).
655    pub const APPROVAL_GRANTED: &'static str = "approval_granted";
656    /// Event type constant for [`ApprovalRejected`](Event::ApprovalRejected).
657    pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
658    /// Event type constant for [`ApprovalEscalated`](Event::ApprovalEscalated).
659    pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
660    /// Event type constant for [`LogLine`](Event::LogLine).
661    pub const LOG_LINE: &'static str = "log_line";
662    /// Event type constant for [`UserSignedIn`](Event::UserSignedIn).
663    pub const USER_SIGNED_IN: &'static str = "user_signed_in";
664    /// Event type constant for [`UserSignedUp`](Event::UserSignedUp).
665    pub const USER_SIGNED_UP: &'static str = "user_signed_up";
666    /// Event type constant for [`UserSignedOut`](Event::UserSignedOut).
667    pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
668
669    /// All event types. Pass this to
670    /// [`EventPublisher::subscribe`](super::EventPublisher::subscribe) to
671    /// receive every event.
672    ///
673    /// # Examples
674    ///
675    /// ```no_run
676    /// use ironflow_engine::notify::{Event, EventPublisher, WebhookSubscriber};
677    ///
678    /// let mut publisher = EventPublisher::new();
679    /// publisher.subscribe(
680    ///     WebhookSubscriber::new("https://example.com/all"),
681    ///     Event::ALL,
682    /// );
683    /// ```
684    pub const ALL: &'static [&'static str] = &[
685        Self::RUN_CREATED,
686        Self::RUN_STATUS_CHANGED,
687        Self::RUN_FAILED,
688        Self::RUN_BUDGET_EXCEEDED,
689        Self::STEP_COMPLETED,
690        Self::STEP_FAILED,
691        Self::APPROVAL_REQUESTED,
692        Self::APPROVAL_GRANTED,
693        Self::APPROVAL_REJECTED,
694        Self::APPROVAL_ESCALATED,
695        Self::LOG_LINE,
696        Self::USER_SIGNED_IN,
697        Self::USER_SIGNED_UP,
698        Self::USER_SIGNED_OUT,
699        Self::RETRY_FORCED,
700    ];
701
702    /// Returns the event type as a static string (e.g. `"run_status_changed"`).
703    ///
704    /// Useful for filtering and logging without deserializing.
705    ///
706    /// # Examples
707    ///
708    /// ```
709    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
710    /// use uuid::Uuid;
711    /// use chrono::Utc;
712    ///
713    /// let event = Event::UserSignedIn(UserSignedInEvent {
714    ///     user_id: Uuid::now_v7(),
715    ///     username: "alice".to_string(),
716    ///     at: Utc::now(),
717    /// });
718    /// assert_eq!(event.event_type(), "user_signed_in");
719    /// ```
720    #[deny(unreachable_patterns)]
721    pub fn event_type(&self) -> &'static str {
722        match self {
723            Event::RunCreated(_) => Self::RUN_CREATED,
724            Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
725            Event::RunFailed(_) => Self::RUN_FAILED,
726            Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
727            Event::RetryForced(_) => Self::RETRY_FORCED,
728            Event::StepCompleted(_) => Self::STEP_COMPLETED,
729            Event::StepFailed(_) => Self::STEP_FAILED,
730            Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
731            Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
732            Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
733            Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
734            Event::LogLine(_) => Self::LOG_LINE,
735            Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
736            Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
737            Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
738        }
739    }
740
741    /// Returns the run this event belongs to, if any.
742    ///
743    /// Auth events ([`UserSignedIn`](Event::UserSignedIn),
744    /// [`UserSignedUp`](Event::UserSignedUp),
745    /// [`UserSignedOut`](Event::UserSignedOut)) are not tied to a run and
746    /// return `None`.
747    ///
748    /// # Examples
749    ///
750    /// ```
751    /// use ironflow_engine::notify::{Event, RunCreatedEvent};
752    /// use uuid::Uuid;
753    /// use chrono::Utc;
754    ///
755    /// let run_id = Uuid::now_v7();
756    /// let event = Event::RunCreated(RunCreatedEvent {
757    ///     run_id,
758    ///     workflow_name: "deploy".to_string(),
759    ///     at: Utc::now(),
760    /// });
761    /// assert_eq!(event.run_id(), Some(run_id));
762    /// ```
763    #[deny(unreachable_patterns)]
764    pub fn run_id(&self) -> Option<Uuid> {
765        match self {
766            Event::RunCreated(e) => Some(e.run_id),
767            Event::RunStatusChanged(e) => Some(e.run_id),
768            Event::RunFailed(e) => Some(e.run_id),
769            Event::RunBudgetExceeded(e) => Some(e.run_id),
770            Event::RetryForced(e) => Some(e.run_id),
771            Event::StepCompleted(e) => Some(e.run_id),
772            Event::StepFailed(e) => Some(e.run_id),
773            Event::ApprovalRequested(e) => Some(e.run_id),
774            Event::ApprovalGranted(e) => Some(e.run_id),
775            Event::ApprovalRejected(e) => Some(e.run_id),
776            Event::ApprovalEscalated(e) => Some(e.run_id),
777            Event::LogLine(e) => Some(e.run_id),
778            Event::UserSignedIn(_) | Event::UserSignedUp(_) | Event::UserSignedOut(_) => None,
779        }
780    }
781
782    /// Returns the step this event belongs to, if any.
783    ///
784    /// Only [`StepCompleted`](Event::StepCompleted),
785    /// [`StepFailed`](Event::StepFailed),
786    /// [`ApprovalRequested`](Event::ApprovalRequested) and
787    /// [`ApprovalEscalated`](Event::ApprovalEscalated) carry a step
788    /// identifier; every other variant returns `None`.
789    ///
790    /// # Examples
791    ///
792    /// ```
793    /// use ironflow_engine::notify::{Event, StepFailedEvent};
794    /// use ironflow_store::models::StepKind;
795    /// use uuid::Uuid;
796    /// use chrono::Utc;
797    ///
798    /// let step_id = Uuid::now_v7();
799    /// let event = Event::StepFailed(StepFailedEvent {
800    ///     run_id: Uuid::now_v7(),
801    ///     step_id,
802    ///     step_name: "build".to_string(),
803    ///     kind: StepKind::Shell,
804    ///     error: "exit code 1".to_string(),
805    ///     at: Utc::now(),
806    /// });
807    /// assert_eq!(event.step_id(), Some(step_id));
808    /// ```
809    #[deny(unreachable_patterns)]
810    pub fn step_id(&self) -> Option<Uuid> {
811        match self {
812            Event::StepCompleted(e) => Some(e.step_id),
813            Event::StepFailed(e) => Some(e.step_id),
814            Event::ApprovalRequested(e) => Some(e.step_id),
815            Event::ApprovalEscalated(e) => Some(e.step_id),
816            Event::RunCreated(_)
817            | Event::RunStatusChanged(_)
818            | Event::RunFailed(_)
819            | Event::RunBudgetExceeded(_)
820            | Event::RetryForced(_)
821            | Event::ApprovalGranted(_)
822            | Event::ApprovalRejected(_)
823            | Event::LogLine(_)
824            | Event::UserSignedIn(_)
825            | Event::UserSignedUp(_)
826            | Event::UserSignedOut(_) => None,
827        }
828    }
829
830    /// Returns the user this event belongs to, if any.
831    ///
832    /// Only the auth events ([`UserSignedIn`](Event::UserSignedIn),
833    /// [`UserSignedUp`](Event::UserSignedUp),
834    /// [`UserSignedOut`](Event::UserSignedOut)) carry a user identifier;
835    /// every other variant returns `None`.
836    ///
837    /// # Examples
838    ///
839    /// ```
840    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
841    /// use uuid::Uuid;
842    /// use chrono::Utc;
843    ///
844    /// let user_id = Uuid::now_v7();
845    /// let event = Event::UserSignedIn(UserSignedInEvent {
846    ///     user_id,
847    ///     username: "alice".to_string(),
848    ///     at: Utc::now(),
849    /// });
850    /// assert_eq!(event.user_id(), Some(user_id));
851    /// ```
852    #[deny(unreachable_patterns)]
853    pub fn user_id(&self) -> Option<Uuid> {
854        match self {
855            Event::UserSignedIn(e) => Some(e.user_id),
856            Event::UserSignedUp(e) => Some(e.user_id),
857            Event::UserSignedOut(e) => Some(e.user_id),
858            Event::RunCreated(_)
859            | Event::RunStatusChanged(_)
860            | Event::RunFailed(_)
861            | Event::RunBudgetExceeded(_)
862            | Event::RetryForced(_)
863            | Event::StepCompleted(_)
864            | Event::StepFailed(_)
865            | Event::ApprovalRequested(_)
866            | Event::ApprovalGranted(_)
867            | Event::ApprovalRejected(_)
868            | Event::ApprovalEscalated(_)
869            | Event::LogLine(_) => None,
870        }
871    }
872}
873
874#[cfg(test)]
875mod tests {
876    use super::*;
877
878    #[test]
879    fn run_status_changed_serde_roundtrip() {
880        let event = Event::RunStatusChanged(RunStatusChangedEvent {
881            run_id: Uuid::now_v7(),
882            workflow_name: "deploy".to_string(),
883            from: RunStatus::Running,
884            to: RunStatus::Completed,
885            error: None,
886            cost_usd: Decimal::new(42, 2),
887            duration_ms: 5000,
888            labels: HashMap::new(),
889            at: Utc::now(),
890        });
891
892        let json = serde_json::to_string(&event).expect("serialize");
893        let back: Event = serde_json::from_str(&json).expect("deserialize");
894
895        assert_eq!(back.event_type(), "run_status_changed");
896        assert!(json.contains("\"type\":\"run_status_changed\""));
897    }
898
899    #[test]
900    fn run_failed_serde_roundtrip() {
901        let event = Event::RunFailed(RunFailedEvent {
902            run_id: Uuid::now_v7(),
903            workflow_name: "deploy".to_string(),
904            error: Some("step crashed".to_string()),
905            cost_usd: Decimal::new(10, 2),
906            duration_ms: 3000,
907            labels: HashMap::new(),
908            at: Utc::now(),
909        });
910
911        let json = serde_json::to_string(&event).expect("serialize");
912        let back: Event = serde_json::from_str(&json).expect("deserialize");
913
914        assert_eq!(back.event_type(), "run_failed");
915        assert!(json.contains("\"type\":\"run_failed\""));
916        assert!(json.contains("step crashed"));
917    }
918
919    #[test]
920    fn run_budget_exceeded_serde_roundtrip() {
921        let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
922            run_id: Uuid::now_v7(),
923            workflow_name: "deploy".to_string(),
924            limit_usd: Decimal::new(200, 2),
925            spent_usd: Decimal::new(180, 2),
926            step_budget_usd: Decimal::new(50, 2),
927            at: Utc::now(),
928        });
929
930        let json = serde_json::to_string(&event).expect("serialize");
931        let back: Event = serde_json::from_str(&json).expect("deserialize");
932
933        assert_eq!(back.event_type(), "run_budget_exceeded");
934        assert!(json.contains("\"type\":\"run_budget_exceeded\""));
935        assert!(json.contains("limit_usd"));
936        assert!(json.contains("step_budget_usd"));
937    }
938
939    #[test]
940    fn all_contains_run_budget_exceeded() {
941        assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
942    }
943
944    #[test]
945    fn user_signed_in_serde_roundtrip() {
946        let event = Event::UserSignedIn(UserSignedInEvent {
947            user_id: Uuid::now_v7(),
948            username: "alice".to_string(),
949            at: Utc::now(),
950        });
951
952        let json = serde_json::to_string(&event).expect("serialize");
953        let back: Event = serde_json::from_str(&json).expect("deserialize");
954
955        assert_eq!(back.event_type(), "user_signed_in");
956        assert!(json.contains("alice"));
957    }
958
959    #[test]
960    fn step_failed_serde_roundtrip() {
961        let event = Event::StepFailed(StepFailedEvent {
962            run_id: Uuid::now_v7(),
963            step_id: Uuid::now_v7(),
964            step_name: "build".to_string(),
965            kind: StepKind::Shell,
966            error: "exit code 1".to_string(),
967            at: Utc::now(),
968        });
969
970        let json = serde_json::to_string(&event).expect("serialize");
971        let back: Event = serde_json::from_str(&json).expect("deserialize");
972
973        assert_eq!(back.event_type(), "step_failed");
974    }
975
976    #[test]
977    fn approval_requested_serde_roundtrip() {
978        let event = Event::ApprovalRequested(ApprovalRequestedEvent {
979            run_id: Uuid::now_v7(),
980            step_id: Uuid::now_v7(),
981            message: "Deploy to prod?".to_string(),
982            at: Utc::now(),
983        });
984
985        let json = serde_json::to_string(&event).expect("serialize");
986        assert!(json.contains("approval_requested"));
987    }
988
989    #[test]
990    fn log_line_serde_roundtrip() {
991        let event = Event::LogLine(LogLineEvent {
992            run_id: Uuid::now_v7(),
993            step_id: Uuid::now_v7(),
994            step_name: "build".to_string(),
995            stream: LogStream::Stdout,
996            line: "Compiling ironflow v0.1.0".to_string(),
997            at: Utc::now(),
998        });
999
1000        let json = serde_json::to_string(&event).expect("serialize");
1001        let back: Event = serde_json::from_str(&json).expect("deserialize");
1002
1003        assert_eq!(back.event_type(), "log_line");
1004        assert!(json.contains("\"type\":\"log_line\""));
1005        assert!(json.contains("Compiling ironflow"));
1006    }
1007
1008    /// The pre-refactor wire format used flat inline-struct variants. Newtype
1009    /// variants over named-field payloads produce and accept the exact same
1010    /// JSON, so audit rows and in-flight payloads written before the refactor
1011    /// still deserialize. No data migration is required.
1012    #[test]
1013    fn legacy_flat_json_deserializes_into_typed_payload() {
1014        let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1015            .parse()
1016            .expect("valid uuid");
1017
1018        let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1019        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1020        match event {
1021            Event::RunCreated(e) => {
1022                assert_eq!(e.run_id, run_id);
1023                assert_eq!(e.workflow_name, "deploy");
1024            }
1025            other => panic!("expected RunCreated, got {other:?}"),
1026        }
1027
1028        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"}"#;
1029        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1030        match event {
1031            Event::RunStatusChanged(e) => {
1032                assert_eq!(e.from, RunStatus::Running);
1033                assert_eq!(e.to, RunStatus::Completed);
1034                assert_eq!(e.cost_usd, Decimal::new(5, 1));
1035                assert_eq!(e.duration_ms, 5000);
1036                assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1037            }
1038            other => panic!("expected RunStatusChanged, got {other:?}"),
1039        }
1040
1041        // `labels` predates no payload: omitting it must still work via `#[serde(default)]`.
1042        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"}"#;
1043        let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1044        match event {
1045            Event::RunStatusChanged(e) => {
1046                assert!(e.labels.is_empty());
1047                assert_eq!(e.error.as_deref(), Some("boom"));
1048            }
1049            other => panic!("expected RunStatusChanged, got {other:?}"),
1050        }
1051
1052        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"}"#;
1053        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1054        match event {
1055            Event::RunFailed(e) => {
1056                assert_eq!(e.error.as_deref(), Some("boom"));
1057                assert!(e.labels.is_empty());
1058            }
1059            other => panic!("expected RunFailed, got {other:?}"),
1060        }
1061
1062        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"}"#;
1063        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1064        match event {
1065            Event::StepFailed(e) => {
1066                assert_eq!(e.kind, StepKind::Shell);
1067                assert_eq!(e.error, "exit code 1");
1068            }
1069            other => panic!("expected StepFailed, got {other:?}"),
1070        }
1071
1072        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"}"#;
1073        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1074        match event {
1075            Event::LogLine(e) => {
1076                assert_eq!(e.stream, LogStream::Stdout);
1077                assert_eq!(e.line, "hello");
1078            }
1079            other => panic!("expected LogLine, got {other:?}"),
1080        }
1081
1082        let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1083        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1084        match event {
1085            Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1086            other => panic!("expected UserSignedIn, got {other:?}"),
1087        }
1088    }
1089
1090    /// Guards the internally-tagged representation: payload fields must stay
1091    /// siblings of `type`, never nested under a variant key.
1092    #[test]
1093    fn serialized_event_is_flat_with_type_tag() {
1094        let run_id = Uuid::now_v7();
1095        let event = Event::RunCreated(RunCreatedEvent {
1096            run_id,
1097            workflow_name: "deploy".to_string(),
1098            at: Utc::now(),
1099        });
1100
1101        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1102        let object = value.as_object().expect("event serializes to an object");
1103
1104        assert_eq!(
1105            object.get("type").and_then(|v| v.as_str()),
1106            Some("run_created")
1107        );
1108        assert_eq!(
1109            object.get("workflow_name").and_then(|v| v.as_str()),
1110            Some("deploy")
1111        );
1112        assert_eq!(
1113            object.get("run_id").and_then(|v| v.as_str()),
1114            Some(run_id.to_string().as_str())
1115        );
1116        assert!(object.contains_key("at"));
1117        assert_eq!(object.len(), 4, "no nesting: {object:?}");
1118        assert!(!object.contains_key("RunCreated"));
1119    }
1120
1121    #[test]
1122    fn run_id_returns_some_for_run_events() {
1123        let run_id = Uuid::now_v7();
1124        let now = Utc::now();
1125
1126        let events = vec![
1127            Event::RunCreated(RunCreatedEvent {
1128                run_id,
1129                workflow_name: "w".to_string(),
1130                at: now,
1131            }),
1132            Event::RunStatusChanged(RunStatusChangedEvent {
1133                run_id,
1134                workflow_name: "w".to_string(),
1135                from: RunStatus::Pending,
1136                to: RunStatus::Running,
1137                error: None,
1138                cost_usd: Decimal::ZERO,
1139                duration_ms: 0,
1140                labels: HashMap::new(),
1141                at: now,
1142            }),
1143            Event::RunFailed(RunFailedEvent {
1144                run_id,
1145                workflow_name: "w".to_string(),
1146                error: None,
1147                cost_usd: Decimal::ZERO,
1148                duration_ms: 0,
1149                labels: HashMap::new(),
1150                at: now,
1151            }),
1152            Event::RunBudgetExceeded(RunBudgetExceededEvent {
1153                run_id,
1154                workflow_name: "w".to_string(),
1155                limit_usd: Decimal::ZERO,
1156                spent_usd: Decimal::ZERO,
1157                step_budget_usd: Decimal::ZERO,
1158                at: now,
1159            }),
1160            Event::RetryForced(RetryForcedEvent {
1161                run_id,
1162                workflow_name: "w".to_string(),
1163                original_version: "1".to_string(),
1164                current_version: "2".to_string(),
1165                at: now,
1166            }),
1167            Event::StepCompleted(StepCompletedEvent {
1168                run_id,
1169                step_id: Uuid::now_v7(),
1170                step_name: "s".to_string(),
1171                kind: StepKind::Shell,
1172                duration_ms: 0,
1173                cost_usd: Decimal::ZERO,
1174                at: now,
1175            }),
1176            Event::StepFailed(StepFailedEvent {
1177                run_id,
1178                step_id: Uuid::now_v7(),
1179                step_name: "s".to_string(),
1180                kind: StepKind::Shell,
1181                error: "e".to_string(),
1182                at: now,
1183            }),
1184            Event::ApprovalRequested(ApprovalRequestedEvent {
1185                run_id,
1186                step_id: Uuid::now_v7(),
1187                message: "ok?".to_string(),
1188                at: now,
1189            }),
1190            Event::ApprovalGranted(ApprovalGrantedEvent {
1191                run_id,
1192                approved_by: "alice".to_string(),
1193                at: now,
1194            }),
1195            Event::ApprovalRejected(ApprovalRejectedEvent {
1196                run_id,
1197                rejected_by: "bob".to_string(),
1198                at: now,
1199            }),
1200            Event::LogLine(LogLineEvent {
1201                run_id,
1202                step_id: Uuid::now_v7(),
1203                step_name: "s".to_string(),
1204                stream: LogStream::Stdout,
1205                line: "l".to_string(),
1206                at: now,
1207            }),
1208        ];
1209
1210        for event in &events {
1211            assert_eq!(
1212                event.run_id(),
1213                Some(run_id),
1214                "{} should carry a run_id",
1215                event.event_type()
1216            );
1217        }
1218    }
1219
1220    #[test]
1221    fn run_id_returns_none_for_auth_events() {
1222        let user_id = Uuid::now_v7();
1223        let now = Utc::now();
1224
1225        let events = vec![
1226            Event::UserSignedIn(UserSignedInEvent {
1227                user_id,
1228                username: "alice".to_string(),
1229                at: now,
1230            }),
1231            Event::UserSignedUp(UserSignedUpEvent {
1232                user_id,
1233                username: "alice".to_string(),
1234                at: now,
1235            }),
1236            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1237        ];
1238
1239        for event in &events {
1240            assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1241        }
1242    }
1243
1244    #[test]
1245    fn step_id_returns_some_only_for_step_events() {
1246        let step_id = Uuid::now_v7();
1247        let run_id = Uuid::now_v7();
1248        let now = Utc::now();
1249
1250        let with_step = vec![
1251            Event::StepCompleted(StepCompletedEvent {
1252                run_id,
1253                step_id,
1254                step_name: "s".to_string(),
1255                kind: StepKind::Shell,
1256                duration_ms: 0,
1257                cost_usd: Decimal::ZERO,
1258                at: now,
1259            }),
1260            Event::StepFailed(StepFailedEvent {
1261                run_id,
1262                step_id,
1263                step_name: "s".to_string(),
1264                kind: StepKind::Shell,
1265                error: "e".to_string(),
1266                at: now,
1267            }),
1268            Event::ApprovalRequested(ApprovalRequestedEvent {
1269                run_id,
1270                step_id,
1271                message: "ok?".to_string(),
1272                at: now,
1273            }),
1274        ];
1275
1276        for event in &with_step {
1277            assert_eq!(
1278                event.step_id(),
1279                Some(step_id),
1280                "{} should carry a step_id",
1281                event.event_type()
1282            );
1283        }
1284
1285        let without_step = vec![
1286            Event::RunCreated(RunCreatedEvent {
1287                run_id,
1288                workflow_name: "w".to_string(),
1289                at: now,
1290            }),
1291            Event::ApprovalGranted(ApprovalGrantedEvent {
1292                run_id,
1293                approved_by: "alice".to_string(),
1294                at: now,
1295            }),
1296            // LogLine carries a step_id field but is reported as a run-level
1297            // stream event, matching the pre-refactor behaviour.
1298            Event::LogLine(LogLineEvent {
1299                run_id,
1300                step_id,
1301                step_name: "s".to_string(),
1302                stream: LogStream::Stdout,
1303                line: "l".to_string(),
1304                at: now,
1305            }),
1306            Event::UserSignedOut(UserSignedOutEvent {
1307                user_id: Uuid::now_v7(),
1308                at: now,
1309            }),
1310        ];
1311
1312        for event in &without_step {
1313            assert_eq!(
1314                event.step_id(),
1315                None,
1316                "{} should not carry a step_id",
1317                event.event_type()
1318            );
1319        }
1320    }
1321
1322    #[test]
1323    fn user_id_returns_some_only_for_auth_events() {
1324        let user_id = Uuid::now_v7();
1325        let run_id = Uuid::now_v7();
1326        let now = Utc::now();
1327
1328        let auth = vec![
1329            Event::UserSignedIn(UserSignedInEvent {
1330                user_id,
1331                username: "alice".to_string(),
1332                at: now,
1333            }),
1334            Event::UserSignedUp(UserSignedUpEvent {
1335                user_id,
1336                username: "alice".to_string(),
1337                at: now,
1338            }),
1339            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1340        ];
1341
1342        for event in &auth {
1343            assert_eq!(
1344                event.user_id(),
1345                Some(user_id),
1346                "{} should carry a user_id",
1347                event.event_type()
1348            );
1349        }
1350
1351        let non_auth = vec![
1352            Event::RunCreated(RunCreatedEvent {
1353                run_id,
1354                workflow_name: "w".to_string(),
1355                at: now,
1356            }),
1357            Event::StepFailed(StepFailedEvent {
1358                run_id,
1359                step_id: Uuid::now_v7(),
1360                step_name: "s".to_string(),
1361                kind: StepKind::Shell,
1362                error: "e".to_string(),
1363                at: now,
1364            }),
1365        ];
1366
1367        for event in &non_auth {
1368            assert_eq!(
1369                event.user_id(),
1370                None,
1371                "{} should not carry a user_id",
1372                event.event_type()
1373            );
1374        }
1375    }
1376
1377    #[test]
1378    fn event_type_all_variants() {
1379        let id = Uuid::now_v7();
1380        let now = Utc::now();
1381
1382        let cases: Vec<(Event, &str)> = vec![
1383            (
1384                Event::RunCreated(RunCreatedEvent {
1385                    run_id: id,
1386                    workflow_name: "w".to_string(),
1387                    at: now,
1388                }),
1389                "run_created",
1390            ),
1391            (
1392                Event::RunStatusChanged(RunStatusChangedEvent {
1393                    run_id: id,
1394                    workflow_name: "w".to_string(),
1395                    from: RunStatus::Pending,
1396                    to: RunStatus::Running,
1397                    error: None,
1398                    cost_usd: Decimal::ZERO,
1399                    duration_ms: 0,
1400                    labels: HashMap::new(),
1401                    at: now,
1402                }),
1403                "run_status_changed",
1404            ),
1405            (
1406                Event::RunFailed(RunFailedEvent {
1407                    run_id: id,
1408                    workflow_name: "w".to_string(),
1409                    error: Some("boom".to_string()),
1410                    cost_usd: Decimal::ZERO,
1411                    duration_ms: 0,
1412                    labels: HashMap::new(),
1413                    at: now,
1414                }),
1415                "run_failed",
1416            ),
1417            (
1418                Event::RunBudgetExceeded(RunBudgetExceededEvent {
1419                    run_id: id,
1420                    workflow_name: "w".to_string(),
1421                    limit_usd: Decimal::new(200, 2),
1422                    spent_usd: Decimal::new(180, 2),
1423                    step_budget_usd: Decimal::new(50, 2),
1424                    at: now,
1425                }),
1426                "run_budget_exceeded",
1427            ),
1428            (
1429                Event::RetryForced(RetryForcedEvent {
1430                    run_id: id,
1431                    workflow_name: "w".to_string(),
1432                    original_version: "1".to_string(),
1433                    current_version: "2".to_string(),
1434                    at: now,
1435                }),
1436                "retry_forced",
1437            ),
1438            (
1439                Event::StepCompleted(StepCompletedEvent {
1440                    run_id: id,
1441                    step_id: id,
1442                    step_name: "s".to_string(),
1443                    kind: StepKind::Shell,
1444                    duration_ms: 0,
1445                    cost_usd: Decimal::ZERO,
1446                    at: now,
1447                }),
1448                "step_completed",
1449            ),
1450            (
1451                Event::StepFailed(StepFailedEvent {
1452                    run_id: id,
1453                    step_id: id,
1454                    step_name: "s".to_string(),
1455                    kind: StepKind::Shell,
1456                    error: "err".to_string(),
1457                    at: now,
1458                }),
1459                "step_failed",
1460            ),
1461            (
1462                Event::ApprovalRequested(ApprovalRequestedEvent {
1463                    run_id: id,
1464                    step_id: id,
1465                    message: "ok?".to_string(),
1466                    at: now,
1467                }),
1468                "approval_requested",
1469            ),
1470            (
1471                Event::ApprovalGranted(ApprovalGrantedEvent {
1472                    run_id: id,
1473                    approved_by: "alice".to_string(),
1474                    at: now,
1475                }),
1476                "approval_granted",
1477            ),
1478            (
1479                Event::ApprovalRejected(ApprovalRejectedEvent {
1480                    run_id: id,
1481                    rejected_by: "bob".to_string(),
1482                    at: now,
1483                }),
1484                "approval_rejected",
1485            ),
1486            (
1487                Event::ApprovalEscalated(ApprovalEscalatedEvent {
1488                    run_id: id,
1489                    step_id: id,
1490                    step_name: "prod-gate".to_string(),
1491                    stage: 0,
1492                    policy: "auto_reject".to_string(),
1493                    action: "rejected".to_string(),
1494                    reason: "approval deadline of 3600s expired".to_string(),
1495                    assignee: None,
1496                    at: now,
1497                }),
1498                "approval_escalated",
1499            ),
1500            (
1501                Event::LogLine(LogLineEvent {
1502                    run_id: id,
1503                    step_id: id,
1504                    step_name: "build".to_string(),
1505                    stream: LogStream::Stdout,
1506                    line: "Compiling ironflow v0.1.0".to_string(),
1507                    at: now,
1508                }),
1509                "log_line",
1510            ),
1511            (
1512                Event::UserSignedIn(UserSignedInEvent {
1513                    user_id: id,
1514                    username: "u".to_string(),
1515                    at: now,
1516                }),
1517                "user_signed_in",
1518            ),
1519            (
1520                Event::UserSignedUp(UserSignedUpEvent {
1521                    user_id: id,
1522                    username: "u".to_string(),
1523                    at: now,
1524                }),
1525                "user_signed_up",
1526            ),
1527            (
1528                Event::UserSignedOut(UserSignedOutEvent {
1529                    user_id: id,
1530                    at: now,
1531                }),
1532                "user_signed_out",
1533            ),
1534        ];
1535
1536        assert_eq!(
1537            cases.len(),
1538            Event::ALL.len(),
1539            "every variant must be covered"
1540        );
1541
1542        for (event, expected_type) in cases {
1543            assert_eq!(event.event_type(), expected_type);
1544        }
1545    }
1546
1547    #[test]
1548    fn approval_escalated_serde_roundtrip() {
1549        let run_id = Uuid::now_v7();
1550        let step_id = Uuid::now_v7();
1551        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1552            run_id,
1553            step_id,
1554            step_name: "prod-gate".to_string(),
1555            stage: 1,
1556            policy: "escalate".to_string(),
1557            action: "reassigned to sre-oncall".to_string(),
1558            reason: "approval deadline of 3600s expired".to_string(),
1559            assignee: Some(Assignee::group("sre-oncall")),
1560            at: Utc::now(),
1561        });
1562
1563        let json = serde_json::to_string(&event).expect("serialize");
1564        assert!(
1565            json.contains("\"type\":\"approval_escalated\""),
1566            "got {json}"
1567        );
1568
1569        let back: Event = serde_json::from_str(&json).expect("deserialize");
1570        let Event::ApprovalEscalated(payload) = back else {
1571            panic!("expected an approval_escalated event");
1572        };
1573        assert_eq!(payload.run_id, run_id);
1574        assert_eq!(payload.stage, 1);
1575        assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1576    }
1577
1578    #[test]
1579    fn approval_escalated_carries_run_and_step_ids() {
1580        let run_id = Uuid::now_v7();
1581        let step_id = Uuid::now_v7();
1582        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1583            run_id,
1584            step_id,
1585            step_name: "prod-gate".to_string(),
1586            stage: 0,
1587            policy: "notify".to_string(),
1588            action: "notified 1 target".to_string(),
1589            reason: "approval deadline of 60s expired".to_string(),
1590            assignee: None,
1591            at: Utc::now(),
1592        });
1593
1594        assert_eq!(event.run_id(), Some(run_id));
1595        assert_eq!(event.step_id(), Some(step_id));
1596        assert_eq!(event.user_id(), None);
1597    }
1598}