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