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///     id: Uuid::now_v7(),
429///     run_id: Uuid::now_v7(),
430///     step_id: Uuid::now_v7(),
431///     step_name: "build".to_string(),
432///     stream: LogStream::Stdout,
433///     line: "Compiling ironflow v0.1.0".to_string(),
434///     at: Utc::now(),
435/// };
436/// assert_eq!(payload.line, "Compiling ironflow v0.1.0");
437/// ```
438#[derive(Debug, Clone, Serialize, Deserialize)]
439#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
440pub struct LogLineEvent {
441    /// Persisted entry identifier (UUID v7, time-ordered).
442    ///
443    /// Matches the `id` of the log entry later stored for this line, so clients
444    /// can de-duplicate the live SSE stream against the persisted history
445    /// fetched from `GET /runs/:id/logs`.
446    ///
447    /// Defaults to the nil UUID when absent, so a payload emitted by an older
448    /// producer still deserializes instead of dropping the whole event.
449    #[serde(default)]
450    pub id: Uuid,
451    /// Run identifier.
452    pub run_id: Uuid,
453    /// Step identifier.
454    pub step_id: Uuid,
455    /// Human-readable step name.
456    pub step_name: String,
457    /// Output stream.
458    pub stream: LogStream,
459    /// The log line content.
460    pub line: String,
461    /// When the line was emitted.
462    pub at: DateTime<Utc>,
463}
464
465/// Payload of the `Event::UserSignedIn` event.
466///
467/// # Examples
468///
469/// ```
470/// use chrono::Utc;
471/// use ironflow_engine::notify::UserSignedInEvent;
472/// use uuid::Uuid;
473///
474/// let payload = UserSignedInEvent {
475///     user_id: Uuid::now_v7(),
476///     username: "alice".to_string(),
477///     at: Utc::now(),
478/// };
479/// assert_eq!(payload.username, "alice");
480/// ```
481#[derive(Debug, Clone, Serialize, Deserialize)]
482#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
483pub struct UserSignedInEvent {
484    /// User identifier.
485    pub user_id: Uuid,
486    /// Username.
487    pub username: String,
488    /// When the sign-in occurred.
489    pub at: DateTime<Utc>,
490}
491
492/// Payload of the `Event::UserSignedUp` event.
493///
494/// # Examples
495///
496/// ```
497/// use chrono::Utc;
498/// use ironflow_engine::notify::UserSignedUpEvent;
499/// use uuid::Uuid;
500///
501/// let payload = UserSignedUpEvent {
502///     user_id: Uuid::now_v7(),
503///     username: "alice".to_string(),
504///     at: Utc::now(),
505/// };
506/// assert_eq!(payload.username, "alice");
507/// ```
508#[derive(Debug, Clone, Serialize, Deserialize)]
509#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
510pub struct UserSignedUpEvent {
511    /// User identifier.
512    pub user_id: Uuid,
513    /// Username.
514    pub username: String,
515    /// When the sign-up occurred.
516    pub at: DateTime<Utc>,
517}
518
519/// Payload of the `Event::UserSignedOut` event.
520///
521/// # Examples
522///
523/// ```
524/// use chrono::Utc;
525/// use ironflow_engine::notify::UserSignedOutEvent;
526/// use uuid::Uuid;
527///
528/// let user_id = Uuid::now_v7();
529/// let payload = UserSignedOutEvent {
530///     user_id,
531///     at: Utc::now(),
532/// };
533/// assert_eq!(payload.user_id, user_id);
534/// ```
535#[derive(Debug, Clone, Serialize, Deserialize)]
536#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
537pub struct UserSignedOutEvent {
538    /// User identifier.
539    pub user_id: Uuid,
540    /// When the sign-out occurred.
541    pub at: DateTime<Utc>,
542}
543
544/// A domain event emitted by the ironflow system.
545///
546/// Covers the full lifecycle: runs, steps, approvals, and authentication.
547/// Subscribers receive these via [`EventPublisher`](super::EventPublisher)
548/// and pattern-match on the variants they care about.
549///
550/// Each variant wraps a dedicated payload struct. The serialized form stays
551/// flat: the `type` discriminant sits next to the payload fields, so
552/// `{"type":"run_created","run_id":...}` round-trips unchanged.
553///
554/// # Examples
555///
556/// ```
557/// use std::collections::HashMap;
558/// use ironflow_engine::notify::{Event, RunStatusChangedEvent};
559/// use ironflow_store::models::RunStatus;
560/// use uuid::Uuid;
561///
562/// let event = Event::RunStatusChanged(RunStatusChangedEvent {
563///     run_id: Uuid::now_v7(),
564///     workflow_name: "deploy".to_string(),
565///     from: RunStatus::Running,
566///     to: RunStatus::Completed,
567///     error: None,
568///     cost_usd: rust_decimal::Decimal::ZERO,
569///     duration_ms: 5000,
570///     labels: HashMap::new(),
571///     at: chrono::Utc::now(),
572/// });
573/// assert_eq!(event.event_type(), "run_status_changed");
574/// ```
575#[derive(Debug, Clone, Serialize, Deserialize)]
576#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
577#[serde(tag = "type", rename_all = "snake_case")]
578pub enum Event {
579    // -- Run lifecycle --
580    /// A new run was created (status: Pending).
581    RunCreated(RunCreatedEvent),
582
583    /// A run changed status.
584    RunStatusChanged(RunStatusChangedEvent),
585
586    /// A run transitioned to [`Failed`](ironflow_store::models::RunStatus::Failed).
587    ///
588    /// This is a convenience event emitted alongside [`RunStatusChanged`](Event::RunStatusChanged)
589    /// when the target status is `Failed`. Subscribe to this instead of
590    /// `RUN_STATUS_CHANGED` when you only care about failures.
591    RunFailed(RunFailedEvent),
592
593    /// A run was stopped because it reached its cumulative cost cap.
594    ///
595    /// Emitted when the engine refuses an agent step that would cross the run's
596    /// `max_cost_usd`. The run transitions to
597    /// [`Cancelled`](ironflow_store::models::RunStatus::Cancelled) and the step
598    /// is never launched, so the reported spend is what the run had already
599    /// consumed.
600    RunBudgetExceeded(RunBudgetExceededEvent),
601
602    /// A manual retry was forced despite a handler version mismatch.
603    ///
604    /// Emitted when a caller passes `force=true` on a retry where the
605    /// handler version differs from the run's recorded version. This is
606    /// an audit event: it means the new code will execute on the old
607    /// payload without the handler explicitly declaring compatibility.
608    RetryForced(RetryForcedEvent),
609
610    // -- Step lifecycle --
611    /// A step completed successfully.
612    StepCompleted(StepCompletedEvent),
613
614    /// A step failed.
615    StepFailed(StepFailedEvent),
616
617    // -- Approval --
618    /// A run is waiting for human approval.
619    ApprovalRequested(ApprovalRequestedEvent),
620
621    /// A run was approved by a human.
622    ApprovalGranted(ApprovalGrantedEvent),
623
624    /// A run was rejected by a human.
625    ApprovalRejected(ApprovalRejectedEvent),
626
627    /// An approval gate missed its SLA deadline and an escalation policy ran.
628    ApprovalEscalated(ApprovalEscalatedEvent),
629
630    // -- Log streaming --
631    /// A log line emitted during step execution.
632    ///
633    /// Pushed by the worker in real time so that SSE clients can stream
634    /// step output as it happens, without waiting for step completion.
635    LogLine(LogLineEvent),
636
637    // -- Authentication --
638    /// A user signed in.
639    UserSignedIn(UserSignedInEvent),
640
641    /// A new user signed up.
642    UserSignedUp(UserSignedUpEvent),
643
644    /// A user signed out.
645    UserSignedOut(UserSignedOutEvent),
646}
647
648impl Event {
649    /// Event type constant for [`RunCreated`](Event::RunCreated).
650    pub const RUN_CREATED: &'static str = "run_created";
651    /// Event type constant for [`RunStatusChanged`](Event::RunStatusChanged).
652    pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
653    /// Event type constant for [`RunFailed`](Event::RunFailed).
654    pub const RUN_FAILED: &'static str = "run_failed";
655    /// Event type constant for [`RunBudgetExceeded`](Event::RunBudgetExceeded).
656    pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
657    /// Event type constant for [`RetryForced`](Event::RetryForced).
658    pub const RETRY_FORCED: &'static str = "retry_forced";
659    /// Event type constant for [`StepCompleted`](Event::StepCompleted).
660    pub const STEP_COMPLETED: &'static str = "step_completed";
661    /// Event type constant for [`StepFailed`](Event::StepFailed).
662    pub const STEP_FAILED: &'static str = "step_failed";
663    /// Event type constant for [`ApprovalRequested`](Event::ApprovalRequested).
664    pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
665    /// Event type constant for [`ApprovalGranted`](Event::ApprovalGranted).
666    pub const APPROVAL_GRANTED: &'static str = "approval_granted";
667    /// Event type constant for [`ApprovalRejected`](Event::ApprovalRejected).
668    pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
669    /// Event type constant for [`ApprovalEscalated`](Event::ApprovalEscalated).
670    pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
671    /// Event type constant for [`LogLine`](Event::LogLine).
672    pub const LOG_LINE: &'static str = "log_line";
673    /// Event type constant for [`UserSignedIn`](Event::UserSignedIn).
674    pub const USER_SIGNED_IN: &'static str = "user_signed_in";
675    /// Event type constant for [`UserSignedUp`](Event::UserSignedUp).
676    pub const USER_SIGNED_UP: &'static str = "user_signed_up";
677    /// Event type constant for [`UserSignedOut`](Event::UserSignedOut).
678    pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
679
680    /// All event types. Pass this to
681    /// [`EventPublisher::subscribe`](super::EventPublisher::subscribe) to
682    /// receive every event.
683    ///
684    /// # Examples
685    ///
686    /// ```no_run
687    /// use ironflow_engine::notify::{Event, EventPublisher, WebhookSubscriber};
688    ///
689    /// let mut publisher = EventPublisher::new();
690    /// publisher.subscribe(
691    ///     WebhookSubscriber::new("https://example.com/all"),
692    ///     Event::ALL,
693    /// );
694    /// ```
695    pub const ALL: &'static [&'static str] = &[
696        Self::RUN_CREATED,
697        Self::RUN_STATUS_CHANGED,
698        Self::RUN_FAILED,
699        Self::RUN_BUDGET_EXCEEDED,
700        Self::STEP_COMPLETED,
701        Self::STEP_FAILED,
702        Self::APPROVAL_REQUESTED,
703        Self::APPROVAL_GRANTED,
704        Self::APPROVAL_REJECTED,
705        Self::APPROVAL_ESCALATED,
706        Self::LOG_LINE,
707        Self::USER_SIGNED_IN,
708        Self::USER_SIGNED_UP,
709        Self::USER_SIGNED_OUT,
710        Self::RETRY_FORCED,
711    ];
712
713    /// Returns the event type as a static string (e.g. `"run_status_changed"`).
714    ///
715    /// Useful for filtering and logging without deserializing.
716    ///
717    /// # Examples
718    ///
719    /// ```
720    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
721    /// use uuid::Uuid;
722    /// use chrono::Utc;
723    ///
724    /// let event = Event::UserSignedIn(UserSignedInEvent {
725    ///     user_id: Uuid::now_v7(),
726    ///     username: "alice".to_string(),
727    ///     at: Utc::now(),
728    /// });
729    /// assert_eq!(event.event_type(), "user_signed_in");
730    /// ```
731    #[deny(unreachable_patterns)]
732    pub fn event_type(&self) -> &'static str {
733        match self {
734            Event::RunCreated(_) => Self::RUN_CREATED,
735            Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
736            Event::RunFailed(_) => Self::RUN_FAILED,
737            Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
738            Event::RetryForced(_) => Self::RETRY_FORCED,
739            Event::StepCompleted(_) => Self::STEP_COMPLETED,
740            Event::StepFailed(_) => Self::STEP_FAILED,
741            Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
742            Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
743            Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
744            Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
745            Event::LogLine(_) => Self::LOG_LINE,
746            Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
747            Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
748            Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
749        }
750    }
751
752    /// Returns the run this event belongs to, if any.
753    ///
754    /// Auth events ([`UserSignedIn`](Event::UserSignedIn),
755    /// [`UserSignedUp`](Event::UserSignedUp),
756    /// [`UserSignedOut`](Event::UserSignedOut)) are not tied to a run and
757    /// return `None`.
758    ///
759    /// # Examples
760    ///
761    /// ```
762    /// use ironflow_engine::notify::{Event, RunCreatedEvent};
763    /// use uuid::Uuid;
764    /// use chrono::Utc;
765    ///
766    /// let run_id = Uuid::now_v7();
767    /// let event = Event::RunCreated(RunCreatedEvent {
768    ///     run_id,
769    ///     workflow_name: "deploy".to_string(),
770    ///     at: Utc::now(),
771    /// });
772    /// assert_eq!(event.run_id(), Some(run_id));
773    /// ```
774    #[deny(unreachable_patterns)]
775    pub fn run_id(&self) -> Option<Uuid> {
776        match self {
777            Event::RunCreated(e) => Some(e.run_id),
778            Event::RunStatusChanged(e) => Some(e.run_id),
779            Event::RunFailed(e) => Some(e.run_id),
780            Event::RunBudgetExceeded(e) => Some(e.run_id),
781            Event::RetryForced(e) => Some(e.run_id),
782            Event::StepCompleted(e) => Some(e.run_id),
783            Event::StepFailed(e) => Some(e.run_id),
784            Event::ApprovalRequested(e) => Some(e.run_id),
785            Event::ApprovalGranted(e) => Some(e.run_id),
786            Event::ApprovalRejected(e) => Some(e.run_id),
787            Event::ApprovalEscalated(e) => Some(e.run_id),
788            Event::LogLine(e) => Some(e.run_id),
789            Event::UserSignedIn(_) | Event::UserSignedUp(_) | Event::UserSignedOut(_) => None,
790        }
791    }
792
793    /// Returns the step this event belongs to, if any.
794    ///
795    /// Only [`StepCompleted`](Event::StepCompleted),
796    /// [`StepFailed`](Event::StepFailed),
797    /// [`ApprovalRequested`](Event::ApprovalRequested) and
798    /// [`ApprovalEscalated`](Event::ApprovalEscalated) carry a step
799    /// identifier; every other variant returns `None`.
800    ///
801    /// # Examples
802    ///
803    /// ```
804    /// use ironflow_engine::notify::{Event, StepFailedEvent};
805    /// use ironflow_store::models::StepKind;
806    /// use uuid::Uuid;
807    /// use chrono::Utc;
808    ///
809    /// let step_id = Uuid::now_v7();
810    /// let event = Event::StepFailed(StepFailedEvent {
811    ///     run_id: Uuid::now_v7(),
812    ///     step_id,
813    ///     step_name: "build".to_string(),
814    ///     kind: StepKind::Shell,
815    ///     error: "exit code 1".to_string(),
816    ///     at: Utc::now(),
817    /// });
818    /// assert_eq!(event.step_id(), Some(step_id));
819    /// ```
820    #[deny(unreachable_patterns)]
821    pub fn step_id(&self) -> Option<Uuid> {
822        match self {
823            Event::StepCompleted(e) => Some(e.step_id),
824            Event::StepFailed(e) => Some(e.step_id),
825            Event::ApprovalRequested(e) => Some(e.step_id),
826            Event::ApprovalEscalated(e) => Some(e.step_id),
827            Event::RunCreated(_)
828            | Event::RunStatusChanged(_)
829            | Event::RunFailed(_)
830            | Event::RunBudgetExceeded(_)
831            | Event::RetryForced(_)
832            | Event::ApprovalGranted(_)
833            | Event::ApprovalRejected(_)
834            | Event::LogLine(_)
835            | Event::UserSignedIn(_)
836            | Event::UserSignedUp(_)
837            | Event::UserSignedOut(_) => None,
838        }
839    }
840
841    /// Returns the user this event belongs to, if any.
842    ///
843    /// Only the auth events ([`UserSignedIn`](Event::UserSignedIn),
844    /// [`UserSignedUp`](Event::UserSignedUp),
845    /// [`UserSignedOut`](Event::UserSignedOut)) carry a user identifier;
846    /// every other variant returns `None`.
847    ///
848    /// # Examples
849    ///
850    /// ```
851    /// use ironflow_engine::notify::{Event, UserSignedInEvent};
852    /// use uuid::Uuid;
853    /// use chrono::Utc;
854    ///
855    /// let user_id = Uuid::now_v7();
856    /// let event = Event::UserSignedIn(UserSignedInEvent {
857    ///     user_id,
858    ///     username: "alice".to_string(),
859    ///     at: Utc::now(),
860    /// });
861    /// assert_eq!(event.user_id(), Some(user_id));
862    /// ```
863    #[deny(unreachable_patterns)]
864    pub fn user_id(&self) -> Option<Uuid> {
865        match self {
866            Event::UserSignedIn(e) => Some(e.user_id),
867            Event::UserSignedUp(e) => Some(e.user_id),
868            Event::UserSignedOut(e) => Some(e.user_id),
869            Event::RunCreated(_)
870            | Event::RunStatusChanged(_)
871            | Event::RunFailed(_)
872            | Event::RunBudgetExceeded(_)
873            | Event::RetryForced(_)
874            | Event::StepCompleted(_)
875            | Event::StepFailed(_)
876            | Event::ApprovalRequested(_)
877            | Event::ApprovalGranted(_)
878            | Event::ApprovalRejected(_)
879            | Event::ApprovalEscalated(_)
880            | Event::LogLine(_) => None,
881        }
882    }
883}
884
885#[cfg(test)]
886mod tests {
887    use super::*;
888
889    #[test]
890    fn run_status_changed_serde_roundtrip() {
891        let event = Event::RunStatusChanged(RunStatusChangedEvent {
892            run_id: Uuid::now_v7(),
893            workflow_name: "deploy".to_string(),
894            from: RunStatus::Running,
895            to: RunStatus::Completed,
896            error: None,
897            cost_usd: Decimal::new(42, 2),
898            duration_ms: 5000,
899            labels: HashMap::new(),
900            at: Utc::now(),
901        });
902
903        let json = serde_json::to_string(&event).expect("serialize");
904        let back: Event = serde_json::from_str(&json).expect("deserialize");
905
906        assert_eq!(back.event_type(), "run_status_changed");
907        assert!(json.contains("\"type\":\"run_status_changed\""));
908    }
909
910    #[test]
911    fn run_failed_serde_roundtrip() {
912        let event = Event::RunFailed(RunFailedEvent {
913            run_id: Uuid::now_v7(),
914            workflow_name: "deploy".to_string(),
915            error: Some("step crashed".to_string()),
916            cost_usd: Decimal::new(10, 2),
917            duration_ms: 3000,
918            labels: HashMap::new(),
919            at: Utc::now(),
920        });
921
922        let json = serde_json::to_string(&event).expect("serialize");
923        let back: Event = serde_json::from_str(&json).expect("deserialize");
924
925        assert_eq!(back.event_type(), "run_failed");
926        assert!(json.contains("\"type\":\"run_failed\""));
927        assert!(json.contains("step crashed"));
928    }
929
930    #[test]
931    fn run_budget_exceeded_serde_roundtrip() {
932        let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
933            run_id: Uuid::now_v7(),
934            workflow_name: "deploy".to_string(),
935            limit_usd: Decimal::new(200, 2),
936            spent_usd: Decimal::new(180, 2),
937            step_budget_usd: Decimal::new(50, 2),
938            at: Utc::now(),
939        });
940
941        let json = serde_json::to_string(&event).expect("serialize");
942        let back: Event = serde_json::from_str(&json).expect("deserialize");
943
944        assert_eq!(back.event_type(), "run_budget_exceeded");
945        assert!(json.contains("\"type\":\"run_budget_exceeded\""));
946        assert!(json.contains("limit_usd"));
947        assert!(json.contains("step_budget_usd"));
948    }
949
950    #[test]
951    fn all_contains_run_budget_exceeded() {
952        assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
953    }
954
955    #[test]
956    fn user_signed_in_serde_roundtrip() {
957        let event = Event::UserSignedIn(UserSignedInEvent {
958            user_id: Uuid::now_v7(),
959            username: "alice".to_string(),
960            at: Utc::now(),
961        });
962
963        let json = serde_json::to_string(&event).expect("serialize");
964        let back: Event = serde_json::from_str(&json).expect("deserialize");
965
966        assert_eq!(back.event_type(), "user_signed_in");
967        assert!(json.contains("alice"));
968    }
969
970    #[test]
971    fn step_failed_serde_roundtrip() {
972        let event = Event::StepFailed(StepFailedEvent {
973            run_id: Uuid::now_v7(),
974            step_id: Uuid::now_v7(),
975            step_name: "build".to_string(),
976            kind: StepKind::Shell,
977            error: "exit code 1".to_string(),
978            at: Utc::now(),
979        });
980
981        let json = serde_json::to_string(&event).expect("serialize");
982        let back: Event = serde_json::from_str(&json).expect("deserialize");
983
984        assert_eq!(back.event_type(), "step_failed");
985    }
986
987    #[test]
988    fn approval_requested_serde_roundtrip() {
989        let event = Event::ApprovalRequested(ApprovalRequestedEvent {
990            run_id: Uuid::now_v7(),
991            step_id: Uuid::now_v7(),
992            message: "Deploy to prod?".to_string(),
993            at: Utc::now(),
994        });
995
996        let json = serde_json::to_string(&event).expect("serialize");
997        assert!(json.contains("approval_requested"));
998    }
999
1000    #[test]
1001    fn log_line_serde_roundtrip() {
1002        let event = Event::LogLine(LogLineEvent {
1003            id: Uuid::now_v7(),
1004            run_id: Uuid::now_v7(),
1005            step_id: Uuid::now_v7(),
1006            step_name: "build".to_string(),
1007            stream: LogStream::Stdout,
1008            line: "Compiling ironflow v0.1.0".to_string(),
1009            at: Utc::now(),
1010        });
1011
1012        let json = serde_json::to_string(&event).expect("serialize");
1013        let back: Event = serde_json::from_str(&json).expect("deserialize");
1014
1015        assert_eq!(back.event_type(), "log_line");
1016        assert!(json.contains("\"type\":\"log_line\""));
1017        assert!(json.contains("Compiling ironflow"));
1018    }
1019
1020    /// The pre-refactor wire format used flat inline-struct variants. Newtype
1021    /// variants over named-field payloads produce and accept the exact same
1022    /// JSON, so audit rows and in-flight payloads written before the refactor
1023    /// still deserialize. No data migration is required.
1024    #[test]
1025    fn legacy_flat_json_deserializes_into_typed_payload() {
1026        let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1027            .parse()
1028            .expect("valid uuid");
1029
1030        let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1031        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1032        match event {
1033            Event::RunCreated(e) => {
1034                assert_eq!(e.run_id, run_id);
1035                assert_eq!(e.workflow_name, "deploy");
1036            }
1037            other => panic!("expected RunCreated, got {other:?}"),
1038        }
1039
1040        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"}"#;
1041        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1042        match event {
1043            Event::RunStatusChanged(e) => {
1044                assert_eq!(e.from, RunStatus::Running);
1045                assert_eq!(e.to, RunStatus::Completed);
1046                assert_eq!(e.cost_usd, Decimal::new(5, 1));
1047                assert_eq!(e.duration_ms, 5000);
1048                assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1049            }
1050            other => panic!("expected RunStatusChanged, got {other:?}"),
1051        }
1052
1053        // `labels` predates no payload: omitting it must still work via `#[serde(default)]`.
1054        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"}"#;
1055        let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1056        match event {
1057            Event::RunStatusChanged(e) => {
1058                assert!(e.labels.is_empty());
1059                assert_eq!(e.error.as_deref(), Some("boom"));
1060            }
1061            other => panic!("expected RunStatusChanged, got {other:?}"),
1062        }
1063
1064        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"}"#;
1065        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1066        match event {
1067            Event::RunFailed(e) => {
1068                assert_eq!(e.error.as_deref(), Some("boom"));
1069                assert!(e.labels.is_empty());
1070            }
1071            other => panic!("expected RunFailed, got {other:?}"),
1072        }
1073
1074        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"}"#;
1075        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1076        match event {
1077            Event::StepFailed(e) => {
1078                assert_eq!(e.kind, StepKind::Shell);
1079                assert_eq!(e.error, "exit code 1");
1080            }
1081            other => panic!("expected StepFailed, got {other:?}"),
1082        }
1083
1084        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"}"#;
1085        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1086        match event {
1087            Event::LogLine(e) => {
1088                assert_eq!(e.stream, LogStream::Stdout);
1089                assert_eq!(e.line, "hello");
1090                // A payload predating the `id` field degrades to the nil UUID
1091                // rather than failing to deserialize.
1092                assert_eq!(e.id, Uuid::nil());
1093            }
1094            other => panic!("expected LogLine, got {other:?}"),
1095        }
1096
1097        let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1098        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1099        match event {
1100            Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1101            other => panic!("expected UserSignedIn, got {other:?}"),
1102        }
1103    }
1104
1105    /// Guards the internally-tagged representation: payload fields must stay
1106    /// siblings of `type`, never nested under a variant key.
1107    #[test]
1108    fn serialized_event_is_flat_with_type_tag() {
1109        let run_id = Uuid::now_v7();
1110        let event = Event::RunCreated(RunCreatedEvent {
1111            run_id,
1112            workflow_name: "deploy".to_string(),
1113            at: Utc::now(),
1114        });
1115
1116        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1117        let object = value.as_object().expect("event serializes to an object");
1118
1119        assert_eq!(
1120            object.get("type").and_then(|v| v.as_str()),
1121            Some("run_created")
1122        );
1123        assert_eq!(
1124            object.get("workflow_name").and_then(|v| v.as_str()),
1125            Some("deploy")
1126        );
1127        assert_eq!(
1128            object.get("run_id").and_then(|v| v.as_str()),
1129            Some(run_id.to_string().as_str())
1130        );
1131        assert!(object.contains_key("at"));
1132        assert_eq!(object.len(), 4, "no nesting: {object:?}");
1133        assert!(!object.contains_key("RunCreated"));
1134    }
1135
1136    #[test]
1137    fn run_id_returns_some_for_run_events() {
1138        let run_id = Uuid::now_v7();
1139        let now = Utc::now();
1140
1141        let events = vec![
1142            Event::RunCreated(RunCreatedEvent {
1143                run_id,
1144                workflow_name: "w".to_string(),
1145                at: now,
1146            }),
1147            Event::RunStatusChanged(RunStatusChangedEvent {
1148                run_id,
1149                workflow_name: "w".to_string(),
1150                from: RunStatus::Pending,
1151                to: RunStatus::Running,
1152                error: None,
1153                cost_usd: Decimal::ZERO,
1154                duration_ms: 0,
1155                labels: HashMap::new(),
1156                at: now,
1157            }),
1158            Event::RunFailed(RunFailedEvent {
1159                run_id,
1160                workflow_name: "w".to_string(),
1161                error: None,
1162                cost_usd: Decimal::ZERO,
1163                duration_ms: 0,
1164                labels: HashMap::new(),
1165                at: now,
1166            }),
1167            Event::RunBudgetExceeded(RunBudgetExceededEvent {
1168                run_id,
1169                workflow_name: "w".to_string(),
1170                limit_usd: Decimal::ZERO,
1171                spent_usd: Decimal::ZERO,
1172                step_budget_usd: Decimal::ZERO,
1173                at: now,
1174            }),
1175            Event::RetryForced(RetryForcedEvent {
1176                run_id,
1177                workflow_name: "w".to_string(),
1178                original_version: "1".to_string(),
1179                current_version: "2".to_string(),
1180                at: now,
1181            }),
1182            Event::StepCompleted(StepCompletedEvent {
1183                run_id,
1184                step_id: Uuid::now_v7(),
1185                step_name: "s".to_string(),
1186                kind: StepKind::Shell,
1187                duration_ms: 0,
1188                cost_usd: Decimal::ZERO,
1189                at: now,
1190            }),
1191            Event::StepFailed(StepFailedEvent {
1192                run_id,
1193                step_id: Uuid::now_v7(),
1194                step_name: "s".to_string(),
1195                kind: StepKind::Shell,
1196                error: "e".to_string(),
1197                at: now,
1198            }),
1199            Event::ApprovalRequested(ApprovalRequestedEvent {
1200                run_id,
1201                step_id: Uuid::now_v7(),
1202                message: "ok?".to_string(),
1203                at: now,
1204            }),
1205            Event::ApprovalGranted(ApprovalGrantedEvent {
1206                run_id,
1207                approved_by: "alice".to_string(),
1208                at: now,
1209            }),
1210            Event::ApprovalRejected(ApprovalRejectedEvent {
1211                run_id,
1212                rejected_by: "bob".to_string(),
1213                at: now,
1214            }),
1215            Event::LogLine(LogLineEvent {
1216                id: Uuid::now_v7(),
1217                run_id,
1218                step_id: Uuid::now_v7(),
1219                step_name: "s".to_string(),
1220                stream: LogStream::Stdout,
1221                line: "l".to_string(),
1222                at: now,
1223            }),
1224        ];
1225
1226        for event in &events {
1227            assert_eq!(
1228                event.run_id(),
1229                Some(run_id),
1230                "{} should carry a run_id",
1231                event.event_type()
1232            );
1233        }
1234    }
1235
1236    #[test]
1237    fn run_id_returns_none_for_auth_events() {
1238        let user_id = Uuid::now_v7();
1239        let now = Utc::now();
1240
1241        let events = vec![
1242            Event::UserSignedIn(UserSignedInEvent {
1243                user_id,
1244                username: "alice".to_string(),
1245                at: now,
1246            }),
1247            Event::UserSignedUp(UserSignedUpEvent {
1248                user_id,
1249                username: "alice".to_string(),
1250                at: now,
1251            }),
1252            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1253        ];
1254
1255        for event in &events {
1256            assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1257        }
1258    }
1259
1260    #[test]
1261    fn step_id_returns_some_only_for_step_events() {
1262        let step_id = Uuid::now_v7();
1263        let run_id = Uuid::now_v7();
1264        let now = Utc::now();
1265
1266        let with_step = vec![
1267            Event::StepCompleted(StepCompletedEvent {
1268                run_id,
1269                step_id,
1270                step_name: "s".to_string(),
1271                kind: StepKind::Shell,
1272                duration_ms: 0,
1273                cost_usd: Decimal::ZERO,
1274                at: now,
1275            }),
1276            Event::StepFailed(StepFailedEvent {
1277                run_id,
1278                step_id,
1279                step_name: "s".to_string(),
1280                kind: StepKind::Shell,
1281                error: "e".to_string(),
1282                at: now,
1283            }),
1284            Event::ApprovalRequested(ApprovalRequestedEvent {
1285                run_id,
1286                step_id,
1287                message: "ok?".to_string(),
1288                at: now,
1289            }),
1290        ];
1291
1292        for event in &with_step {
1293            assert_eq!(
1294                event.step_id(),
1295                Some(step_id),
1296                "{} should carry a step_id",
1297                event.event_type()
1298            );
1299        }
1300
1301        let without_step = vec![
1302            Event::RunCreated(RunCreatedEvent {
1303                run_id,
1304                workflow_name: "w".to_string(),
1305                at: now,
1306            }),
1307            Event::ApprovalGranted(ApprovalGrantedEvent {
1308                run_id,
1309                approved_by: "alice".to_string(),
1310                at: now,
1311            }),
1312            // LogLine carries a step_id field but is reported as a run-level
1313            // stream event, matching the pre-refactor behaviour.
1314            Event::LogLine(LogLineEvent {
1315                id: Uuid::now_v7(),
1316                run_id,
1317                step_id,
1318                step_name: "s".to_string(),
1319                stream: LogStream::Stdout,
1320                line: "l".to_string(),
1321                at: now,
1322            }),
1323            Event::UserSignedOut(UserSignedOutEvent {
1324                user_id: Uuid::now_v7(),
1325                at: now,
1326            }),
1327        ];
1328
1329        for event in &without_step {
1330            assert_eq!(
1331                event.step_id(),
1332                None,
1333                "{} should not carry a step_id",
1334                event.event_type()
1335            );
1336        }
1337    }
1338
1339    #[test]
1340    fn user_id_returns_some_only_for_auth_events() {
1341        let user_id = Uuid::now_v7();
1342        let run_id = Uuid::now_v7();
1343        let now = Utc::now();
1344
1345        let auth = vec![
1346            Event::UserSignedIn(UserSignedInEvent {
1347                user_id,
1348                username: "alice".to_string(),
1349                at: now,
1350            }),
1351            Event::UserSignedUp(UserSignedUpEvent {
1352                user_id,
1353                username: "alice".to_string(),
1354                at: now,
1355            }),
1356            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1357        ];
1358
1359        for event in &auth {
1360            assert_eq!(
1361                event.user_id(),
1362                Some(user_id),
1363                "{} should carry a user_id",
1364                event.event_type()
1365            );
1366        }
1367
1368        let non_auth = vec![
1369            Event::RunCreated(RunCreatedEvent {
1370                run_id,
1371                workflow_name: "w".to_string(),
1372                at: now,
1373            }),
1374            Event::StepFailed(StepFailedEvent {
1375                run_id,
1376                step_id: Uuid::now_v7(),
1377                step_name: "s".to_string(),
1378                kind: StepKind::Shell,
1379                error: "e".to_string(),
1380                at: now,
1381            }),
1382        ];
1383
1384        for event in &non_auth {
1385            assert_eq!(
1386                event.user_id(),
1387                None,
1388                "{} should not carry a user_id",
1389                event.event_type()
1390            );
1391        }
1392    }
1393
1394    #[test]
1395    fn event_type_all_variants() {
1396        let id = Uuid::now_v7();
1397        let now = Utc::now();
1398
1399        let cases: Vec<(Event, &str)> = vec![
1400            (
1401                Event::RunCreated(RunCreatedEvent {
1402                    run_id: id,
1403                    workflow_name: "w".to_string(),
1404                    at: now,
1405                }),
1406                "run_created",
1407            ),
1408            (
1409                Event::RunStatusChanged(RunStatusChangedEvent {
1410                    run_id: id,
1411                    workflow_name: "w".to_string(),
1412                    from: RunStatus::Pending,
1413                    to: RunStatus::Running,
1414                    error: None,
1415                    cost_usd: Decimal::ZERO,
1416                    duration_ms: 0,
1417                    labels: HashMap::new(),
1418                    at: now,
1419                }),
1420                "run_status_changed",
1421            ),
1422            (
1423                Event::RunFailed(RunFailedEvent {
1424                    run_id: id,
1425                    workflow_name: "w".to_string(),
1426                    error: Some("boom".to_string()),
1427                    cost_usd: Decimal::ZERO,
1428                    duration_ms: 0,
1429                    labels: HashMap::new(),
1430                    at: now,
1431                }),
1432                "run_failed",
1433            ),
1434            (
1435                Event::RunBudgetExceeded(RunBudgetExceededEvent {
1436                    run_id: id,
1437                    workflow_name: "w".to_string(),
1438                    limit_usd: Decimal::new(200, 2),
1439                    spent_usd: Decimal::new(180, 2),
1440                    step_budget_usd: Decimal::new(50, 2),
1441                    at: now,
1442                }),
1443                "run_budget_exceeded",
1444            ),
1445            (
1446                Event::RetryForced(RetryForcedEvent {
1447                    run_id: id,
1448                    workflow_name: "w".to_string(),
1449                    original_version: "1".to_string(),
1450                    current_version: "2".to_string(),
1451                    at: now,
1452                }),
1453                "retry_forced",
1454            ),
1455            (
1456                Event::StepCompleted(StepCompletedEvent {
1457                    run_id: id,
1458                    step_id: id,
1459                    step_name: "s".to_string(),
1460                    kind: StepKind::Shell,
1461                    duration_ms: 0,
1462                    cost_usd: Decimal::ZERO,
1463                    at: now,
1464                }),
1465                "step_completed",
1466            ),
1467            (
1468                Event::StepFailed(StepFailedEvent {
1469                    run_id: id,
1470                    step_id: id,
1471                    step_name: "s".to_string(),
1472                    kind: StepKind::Shell,
1473                    error: "err".to_string(),
1474                    at: now,
1475                }),
1476                "step_failed",
1477            ),
1478            (
1479                Event::ApprovalRequested(ApprovalRequestedEvent {
1480                    run_id: id,
1481                    step_id: id,
1482                    message: "ok?".to_string(),
1483                    at: now,
1484                }),
1485                "approval_requested",
1486            ),
1487            (
1488                Event::ApprovalGranted(ApprovalGrantedEvent {
1489                    run_id: id,
1490                    approved_by: "alice".to_string(),
1491                    at: now,
1492                }),
1493                "approval_granted",
1494            ),
1495            (
1496                Event::ApprovalRejected(ApprovalRejectedEvent {
1497                    run_id: id,
1498                    rejected_by: "bob".to_string(),
1499                    at: now,
1500                }),
1501                "approval_rejected",
1502            ),
1503            (
1504                Event::ApprovalEscalated(ApprovalEscalatedEvent {
1505                    run_id: id,
1506                    step_id: id,
1507                    step_name: "prod-gate".to_string(),
1508                    stage: 0,
1509                    policy: "auto_reject".to_string(),
1510                    action: "rejected".to_string(),
1511                    reason: "approval deadline of 3600s expired".to_string(),
1512                    assignee: None,
1513                    at: now,
1514                }),
1515                "approval_escalated",
1516            ),
1517            (
1518                Event::LogLine(LogLineEvent {
1519                    id,
1520                    run_id: id,
1521                    step_id: id,
1522                    step_name: "build".to_string(),
1523                    stream: LogStream::Stdout,
1524                    line: "Compiling ironflow v0.1.0".to_string(),
1525                    at: now,
1526                }),
1527                "log_line",
1528            ),
1529            (
1530                Event::UserSignedIn(UserSignedInEvent {
1531                    user_id: id,
1532                    username: "u".to_string(),
1533                    at: now,
1534                }),
1535                "user_signed_in",
1536            ),
1537            (
1538                Event::UserSignedUp(UserSignedUpEvent {
1539                    user_id: id,
1540                    username: "u".to_string(),
1541                    at: now,
1542                }),
1543                "user_signed_up",
1544            ),
1545            (
1546                Event::UserSignedOut(UserSignedOutEvent {
1547                    user_id: id,
1548                    at: now,
1549                }),
1550                "user_signed_out",
1551            ),
1552        ];
1553
1554        assert_eq!(
1555            cases.len(),
1556            Event::ALL.len(),
1557            "every variant must be covered"
1558        );
1559
1560        for (event, expected_type) in cases {
1561            assert_eq!(event.event_type(), expected_type);
1562        }
1563    }
1564
1565    #[test]
1566    fn approval_escalated_serde_roundtrip() {
1567        let run_id = Uuid::now_v7();
1568        let step_id = Uuid::now_v7();
1569        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1570            run_id,
1571            step_id,
1572            step_name: "prod-gate".to_string(),
1573            stage: 1,
1574            policy: "escalate".to_string(),
1575            action: "reassigned to sre-oncall".to_string(),
1576            reason: "approval deadline of 3600s expired".to_string(),
1577            assignee: Some(Assignee::group("sre-oncall")),
1578            at: Utc::now(),
1579        });
1580
1581        let json = serde_json::to_string(&event).expect("serialize");
1582        assert!(
1583            json.contains("\"type\":\"approval_escalated\""),
1584            "got {json}"
1585        );
1586
1587        let back: Event = serde_json::from_str(&json).expect("deserialize");
1588        let Event::ApprovalEscalated(payload) = back else {
1589            panic!("expected an approval_escalated event");
1590        };
1591        assert_eq!(payload.run_id, run_id);
1592        assert_eq!(payload.stage, 1);
1593        assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1594    }
1595
1596    #[test]
1597    fn approval_escalated_carries_run_and_step_ids() {
1598        let run_id = Uuid::now_v7();
1599        let step_id = Uuid::now_v7();
1600        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1601            run_id,
1602            step_id,
1603            step_name: "prod-gate".to_string(),
1604            stage: 0,
1605            policy: "notify".to_string(),
1606            action: "notified 1 target".to_string(),
1607            reason: "approval deadline of 60s expired".to_string(),
1608            assignee: None,
1609            at: Utc::now(),
1610        });
1611
1612        assert_eq!(event.run_id(), Some(run_id));
1613        assert_eq!(event.step_id(), Some(step_id));
1614        assert_eq!(event.user_id(), None);
1615    }
1616}