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 requirement evaluated
292/// from the gate's approval rules, if it has 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    /// Requirement evaluated from the gate's approval rules. `None` for a
320    /// gate without rules: 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    /// Requirement evaluated from the gate's approval rules, 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    /// Requirement evaluated from the gate's approval rules, 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            rule_index: Some(0),
1068            condition: Some("payload.amount > 10000".to_string()),
1069            required_approvers: 2,
1070            approver_groups: vec!["finance".to_string()],
1071            evaluated: Vec::new(),
1072        };
1073        let event = Event::ApprovalGranted(ApprovalGrantedEvent {
1074            run_id: Uuid::now_v7(),
1075            step_id: Some(Uuid::now_v7()),
1076            approved_by: "alice".to_string(),
1077            approvals_received: 1,
1078            approvals_required: 2,
1079            requirement: Some(requirement.clone()),
1080            at: Utc::now(),
1081        });
1082
1083        let json = serde_json::to_string(&event).expect("serialize");
1084        let back: Event = serde_json::from_str(&json).expect("deserialize");
1085        let Event::ApprovalGranted(back) = back else {
1086            panic!("expected approval_granted");
1087        };
1088        assert_eq!(back.approvals_received, 1);
1089        assert_eq!(back.approvals_required, 2);
1090        assert_eq!(back.requirement, Some(requirement));
1091    }
1092
1093    #[test]
1094    fn approval_requested_serde_roundtrip() {
1095        let event = Event::ApprovalRequested(ApprovalRequestedEvent {
1096            run_id: Uuid::now_v7(),
1097            step_id: Uuid::now_v7(),
1098            message: "Deploy to prod?".to_string(),
1099            requirement: None,
1100            at: Utc::now(),
1101        });
1102
1103        let json = serde_json::to_string(&event).expect("serialize");
1104        assert!(json.contains("approval_requested"));
1105    }
1106
1107    #[test]
1108    fn log_line_serde_roundtrip() {
1109        let event = Event::LogLine(LogLineEvent {
1110            id: Uuid::now_v7(),
1111            run_id: Uuid::now_v7(),
1112            step_id: Uuid::now_v7(),
1113            step_name: "build".to_string(),
1114            stream: LogStream::Stdout,
1115            line: "Compiling ironflow v0.1.0".to_string(),
1116            at: Utc::now(),
1117        });
1118
1119        let json = serde_json::to_string(&event).expect("serialize");
1120        let back: Event = serde_json::from_str(&json).expect("deserialize");
1121
1122        assert_eq!(back.event_type(), "log_line");
1123        assert!(json.contains("\"type\":\"log_line\""));
1124        assert!(json.contains("Compiling ironflow"));
1125    }
1126
1127    /// The pre-refactor wire format used flat inline-struct variants. Newtype
1128    /// variants over named-field payloads produce and accept the exact same
1129    /// JSON, so audit rows and in-flight payloads written before the refactor
1130    /// still deserialize. No data migration is required.
1131    #[test]
1132    fn legacy_flat_json_deserializes_into_typed_payload() {
1133        let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
1134            .parse()
1135            .expect("valid uuid");
1136
1137        let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
1138        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1139        match event {
1140            Event::RunCreated(e) => {
1141                assert_eq!(e.run_id, run_id);
1142                assert_eq!(e.workflow_name, "deploy");
1143            }
1144            other => panic!("expected RunCreated, got {other:?}"),
1145        }
1146
1147        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"}"#;
1148        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1149        match event {
1150            Event::RunStatusChanged(e) => {
1151                assert_eq!(e.from, RunStatus::Running);
1152                assert_eq!(e.to, RunStatus::Completed);
1153                assert_eq!(e.cost_usd, Decimal::new(5, 1));
1154                assert_eq!(e.duration_ms, 5000);
1155                assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
1156            }
1157            other => panic!("expected RunStatusChanged, got {other:?}"),
1158        }
1159
1160        // `labels` predates no payload: omitting it must still work via `#[serde(default)]`.
1161        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"}"#;
1162        let event: Event = serde_json::from_str(raw).expect("missing labels must default");
1163        match event {
1164            Event::RunStatusChanged(e) => {
1165                assert!(e.labels.is_empty());
1166                assert_eq!(e.error.as_deref(), Some("boom"));
1167            }
1168            other => panic!("expected RunStatusChanged, got {other:?}"),
1169        }
1170
1171        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"}"#;
1172        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1173        match event {
1174            Event::RunFailed(e) => {
1175                assert_eq!(e.error.as_deref(), Some("boom"));
1176                assert!(e.labels.is_empty());
1177            }
1178            other => panic!("expected RunFailed, got {other:?}"),
1179        }
1180
1181        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"}"#;
1182        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1183        match event {
1184            Event::StepFailed(e) => {
1185                assert_eq!(e.kind, StepKind::Shell);
1186                assert_eq!(e.error, "exit code 1");
1187            }
1188            other => panic!("expected StepFailed, got {other:?}"),
1189        }
1190
1191        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"}"#;
1192        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1193        match event {
1194            Event::LogLine(e) => {
1195                assert_eq!(e.stream, LogStream::Stdout);
1196                assert_eq!(e.line, "hello");
1197                // A payload predating the `id` field degrades to the nil UUID
1198                // rather than failing to deserialize.
1199                assert_eq!(e.id, Uuid::nil());
1200            }
1201            other => panic!("expected LogLine, got {other:?}"),
1202        }
1203
1204        let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
1205        let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
1206        match event {
1207            Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
1208            other => panic!("expected UserSignedIn, got {other:?}"),
1209        }
1210    }
1211
1212    /// Guards the internally-tagged representation: payload fields must stay
1213    /// siblings of `type`, never nested under a variant key.
1214    #[test]
1215    fn serialized_event_is_flat_with_type_tag() {
1216        let run_id = Uuid::now_v7();
1217        let event = Event::RunCreated(RunCreatedEvent {
1218            run_id,
1219            workflow_name: "deploy".to_string(),
1220            at: Utc::now(),
1221        });
1222
1223        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
1224        let object = value.as_object().expect("event serializes to an object");
1225
1226        assert_eq!(
1227            object.get("type").and_then(|v| v.as_str()),
1228            Some("run_created")
1229        );
1230        assert_eq!(
1231            object.get("workflow_name").and_then(|v| v.as_str()),
1232            Some("deploy")
1233        );
1234        assert_eq!(
1235            object.get("run_id").and_then(|v| v.as_str()),
1236            Some(run_id.to_string().as_str())
1237        );
1238        assert!(object.contains_key("at"));
1239        assert_eq!(object.len(), 4, "no nesting: {object:?}");
1240        assert!(!object.contains_key("RunCreated"));
1241    }
1242
1243    #[test]
1244    fn run_id_returns_some_for_run_events() {
1245        let run_id = Uuid::now_v7();
1246        let now = Utc::now();
1247
1248        let events = vec![
1249            Event::RunCreated(RunCreatedEvent {
1250                run_id,
1251                workflow_name: "w".to_string(),
1252                at: now,
1253            }),
1254            Event::RunStatusChanged(RunStatusChangedEvent {
1255                run_id,
1256                workflow_name: "w".to_string(),
1257                from: RunStatus::Pending,
1258                to: RunStatus::Running,
1259                error: None,
1260                cost_usd: Decimal::ZERO,
1261                duration_ms: 0,
1262                labels: HashMap::new(),
1263                at: now,
1264            }),
1265            Event::RunFailed(RunFailedEvent {
1266                run_id,
1267                workflow_name: "w".to_string(),
1268                error: None,
1269                cost_usd: Decimal::ZERO,
1270                duration_ms: 0,
1271                labels: HashMap::new(),
1272                at: now,
1273            }),
1274            Event::RunBudgetExceeded(RunBudgetExceededEvent {
1275                run_id,
1276                workflow_name: "w".to_string(),
1277                limit_usd: Decimal::ZERO,
1278                spent_usd: Decimal::ZERO,
1279                step_budget_usd: Decimal::ZERO,
1280                at: now,
1281            }),
1282            Event::RetryForced(RetryForcedEvent {
1283                run_id,
1284                workflow_name: "w".to_string(),
1285                original_version: "1".to_string(),
1286                current_version: "2".to_string(),
1287                at: now,
1288            }),
1289            Event::StepCompleted(StepCompletedEvent {
1290                run_id,
1291                step_id: Uuid::now_v7(),
1292                step_name: "s".to_string(),
1293                kind: StepKind::Shell,
1294                duration_ms: 0,
1295                cost_usd: Decimal::ZERO,
1296                at: now,
1297            }),
1298            Event::StepFailed(StepFailedEvent {
1299                run_id,
1300                step_id: Uuid::now_v7(),
1301                step_name: "s".to_string(),
1302                kind: StepKind::Shell,
1303                error: "e".to_string(),
1304                at: now,
1305            }),
1306            Event::ApprovalRequested(ApprovalRequestedEvent {
1307                run_id,
1308                step_id: Uuid::now_v7(),
1309                message: "ok?".to_string(),
1310                requirement: None,
1311                at: now,
1312            }),
1313            Event::ApprovalGranted(ApprovalGrantedEvent {
1314                run_id,
1315                step_id: None,
1316                approved_by: "alice".to_string(),
1317                approvals_received: 1,
1318                approvals_required: 1,
1319                requirement: None,
1320                at: now,
1321            }),
1322            Event::ApprovalRejected(ApprovalRejectedEvent {
1323                run_id,
1324                step_id: None,
1325                rejected_by: "bob".to_string(),
1326                requirement: None,
1327                at: now,
1328            }),
1329            Event::LogLine(LogLineEvent {
1330                id: Uuid::now_v7(),
1331                run_id,
1332                step_id: Uuid::now_v7(),
1333                step_name: "s".to_string(),
1334                stream: LogStream::Stdout,
1335                line: "l".to_string(),
1336                at: now,
1337            }),
1338        ];
1339
1340        for event in &events {
1341            assert_eq!(
1342                event.run_id(),
1343                Some(run_id),
1344                "{} should carry a run_id",
1345                event.event_type()
1346            );
1347        }
1348    }
1349
1350    #[test]
1351    fn run_id_returns_none_for_auth_events() {
1352        let user_id = Uuid::now_v7();
1353        let now = Utc::now();
1354
1355        let events = vec![
1356            Event::UserSignedIn(UserSignedInEvent {
1357                user_id,
1358                username: "alice".to_string(),
1359                at: now,
1360            }),
1361            Event::UserSignedUp(UserSignedUpEvent {
1362                user_id,
1363                username: "alice".to_string(),
1364                at: now,
1365            }),
1366            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1367        ];
1368
1369        for event in &events {
1370            assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
1371        }
1372    }
1373
1374    #[test]
1375    fn step_id_returns_some_only_for_step_events() {
1376        let step_id = Uuid::now_v7();
1377        let run_id = Uuid::now_v7();
1378        let now = Utc::now();
1379
1380        let with_step = vec![
1381            Event::StepCompleted(StepCompletedEvent {
1382                run_id,
1383                step_id,
1384                step_name: "s".to_string(),
1385                kind: StepKind::Shell,
1386                duration_ms: 0,
1387                cost_usd: Decimal::ZERO,
1388                at: now,
1389            }),
1390            Event::StepFailed(StepFailedEvent {
1391                run_id,
1392                step_id,
1393                step_name: "s".to_string(),
1394                kind: StepKind::Shell,
1395                error: "e".to_string(),
1396                at: now,
1397            }),
1398            Event::ApprovalRequested(ApprovalRequestedEvent {
1399                run_id,
1400                step_id,
1401                message: "ok?".to_string(),
1402                requirement: None,
1403                at: now,
1404            }),
1405            Event::ApprovalGranted(ApprovalGrantedEvent {
1406                run_id,
1407                step_id: Some(step_id),
1408                approved_by: "alice".to_string(),
1409                approvals_received: 1,
1410                approvals_required: 2,
1411                requirement: None,
1412                at: now,
1413            }),
1414            Event::ApprovalRejected(ApprovalRejectedEvent {
1415                run_id,
1416                step_id: Some(step_id),
1417                rejected_by: "bob".to_string(),
1418                requirement: None,
1419                at: now,
1420            }),
1421        ];
1422
1423        for event in &with_step {
1424            assert_eq!(
1425                event.step_id(),
1426                Some(step_id),
1427                "{} should carry a step_id",
1428                event.event_type()
1429            );
1430        }
1431
1432        let without_step = vec![
1433            Event::RunCreated(RunCreatedEvent {
1434                run_id,
1435                workflow_name: "w".to_string(),
1436                at: now,
1437            }),
1438            // A granted event recorded before step ids were carried.
1439            Event::ApprovalGranted(ApprovalGrantedEvent {
1440                run_id,
1441                step_id: None,
1442                approved_by: "alice".to_string(),
1443                approvals_received: 1,
1444                approvals_required: 1,
1445                requirement: None,
1446                at: now,
1447            }),
1448            // LogLine carries a step_id field but is reported as a run-level
1449            // stream event, matching the pre-refactor behaviour.
1450            Event::LogLine(LogLineEvent {
1451                id: Uuid::now_v7(),
1452                run_id,
1453                step_id,
1454                step_name: "s".to_string(),
1455                stream: LogStream::Stdout,
1456                line: "l".to_string(),
1457                at: now,
1458            }),
1459            Event::UserSignedOut(UserSignedOutEvent {
1460                user_id: Uuid::now_v7(),
1461                at: now,
1462            }),
1463        ];
1464
1465        for event in &without_step {
1466            assert_eq!(
1467                event.step_id(),
1468                None,
1469                "{} should not carry a step_id",
1470                event.event_type()
1471            );
1472        }
1473    }
1474
1475    #[test]
1476    fn user_id_returns_some_only_for_auth_events() {
1477        let user_id = Uuid::now_v7();
1478        let run_id = Uuid::now_v7();
1479        let now = Utc::now();
1480
1481        let auth = vec![
1482            Event::UserSignedIn(UserSignedInEvent {
1483                user_id,
1484                username: "alice".to_string(),
1485                at: now,
1486            }),
1487            Event::UserSignedUp(UserSignedUpEvent {
1488                user_id,
1489                username: "alice".to_string(),
1490                at: now,
1491            }),
1492            Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
1493        ];
1494
1495        for event in &auth {
1496            assert_eq!(
1497                event.user_id(),
1498                Some(user_id),
1499                "{} should carry a user_id",
1500                event.event_type()
1501            );
1502        }
1503
1504        let non_auth = vec![
1505            Event::RunCreated(RunCreatedEvent {
1506                run_id,
1507                workflow_name: "w".to_string(),
1508                at: now,
1509            }),
1510            Event::StepFailed(StepFailedEvent {
1511                run_id,
1512                step_id: Uuid::now_v7(),
1513                step_name: "s".to_string(),
1514                kind: StepKind::Shell,
1515                error: "e".to_string(),
1516                at: now,
1517            }),
1518        ];
1519
1520        for event in &non_auth {
1521            assert_eq!(
1522                event.user_id(),
1523                None,
1524                "{} should not carry a user_id",
1525                event.event_type()
1526            );
1527        }
1528    }
1529
1530    #[test]
1531    fn event_type_all_variants() {
1532        let id = Uuid::now_v7();
1533        let now = Utc::now();
1534
1535        let cases: Vec<(Event, &str)> = vec![
1536            (
1537                Event::RunCreated(RunCreatedEvent {
1538                    run_id: id,
1539                    workflow_name: "w".to_string(),
1540                    at: now,
1541                }),
1542                "run_created",
1543            ),
1544            (
1545                Event::RunStatusChanged(RunStatusChangedEvent {
1546                    run_id: id,
1547                    workflow_name: "w".to_string(),
1548                    from: RunStatus::Pending,
1549                    to: RunStatus::Running,
1550                    error: None,
1551                    cost_usd: Decimal::ZERO,
1552                    duration_ms: 0,
1553                    labels: HashMap::new(),
1554                    at: now,
1555                }),
1556                "run_status_changed",
1557            ),
1558            (
1559                Event::RunFailed(RunFailedEvent {
1560                    run_id: id,
1561                    workflow_name: "w".to_string(),
1562                    error: Some("boom".to_string()),
1563                    cost_usd: Decimal::ZERO,
1564                    duration_ms: 0,
1565                    labels: HashMap::new(),
1566                    at: now,
1567                }),
1568                "run_failed",
1569            ),
1570            (
1571                Event::RunBudgetExceeded(RunBudgetExceededEvent {
1572                    run_id: id,
1573                    workflow_name: "w".to_string(),
1574                    limit_usd: Decimal::new(200, 2),
1575                    spent_usd: Decimal::new(180, 2),
1576                    step_budget_usd: Decimal::new(50, 2),
1577                    at: now,
1578                }),
1579                "run_budget_exceeded",
1580            ),
1581            (
1582                Event::RetryForced(RetryForcedEvent {
1583                    run_id: id,
1584                    workflow_name: "w".to_string(),
1585                    original_version: "1".to_string(),
1586                    current_version: "2".to_string(),
1587                    at: now,
1588                }),
1589                "retry_forced",
1590            ),
1591            (
1592                Event::StepCompleted(StepCompletedEvent {
1593                    run_id: id,
1594                    step_id: id,
1595                    step_name: "s".to_string(),
1596                    kind: StepKind::Shell,
1597                    duration_ms: 0,
1598                    cost_usd: Decimal::ZERO,
1599                    at: now,
1600                }),
1601                "step_completed",
1602            ),
1603            (
1604                Event::StepFailed(StepFailedEvent {
1605                    run_id: id,
1606                    step_id: id,
1607                    step_name: "s".to_string(),
1608                    kind: StepKind::Shell,
1609                    error: "err".to_string(),
1610                    at: now,
1611                }),
1612                "step_failed",
1613            ),
1614            (
1615                Event::ApprovalRequested(ApprovalRequestedEvent {
1616                    run_id: id,
1617                    step_id: id,
1618                    message: "ok?".to_string(),
1619                    requirement: None,
1620                    at: now,
1621                }),
1622                "approval_requested",
1623            ),
1624            (
1625                Event::ApprovalGranted(ApprovalGrantedEvent {
1626                    run_id: id,
1627                    step_id: Some(id),
1628                    approved_by: "alice".to_string(),
1629                    approvals_received: 1,
1630                    approvals_required: 1,
1631                    requirement: None,
1632                    at: now,
1633                }),
1634                "approval_granted",
1635            ),
1636            (
1637                Event::ApprovalRejected(ApprovalRejectedEvent {
1638                    run_id: id,
1639                    step_id: Some(id),
1640                    rejected_by: "bob".to_string(),
1641                    requirement: None,
1642                    at: now,
1643                }),
1644                "approval_rejected",
1645            ),
1646            (
1647                Event::ApprovalEscalated(ApprovalEscalatedEvent {
1648                    run_id: id,
1649                    step_id: id,
1650                    step_name: "prod-gate".to_string(),
1651                    stage: 0,
1652                    policy: "auto_reject".to_string(),
1653                    action: "rejected".to_string(),
1654                    reason: "approval deadline of 3600s expired".to_string(),
1655                    assignee: None,
1656                    at: now,
1657                }),
1658                "approval_escalated",
1659            ),
1660            (
1661                Event::LogLine(LogLineEvent {
1662                    id,
1663                    run_id: id,
1664                    step_id: id,
1665                    step_name: "build".to_string(),
1666                    stream: LogStream::Stdout,
1667                    line: "Compiling ironflow v0.1.0".to_string(),
1668                    at: now,
1669                }),
1670                "log_line",
1671            ),
1672            (
1673                Event::UserSignedIn(UserSignedInEvent {
1674                    user_id: id,
1675                    username: "u".to_string(),
1676                    at: now,
1677                }),
1678                "user_signed_in",
1679            ),
1680            (
1681                Event::UserSignedUp(UserSignedUpEvent {
1682                    user_id: id,
1683                    username: "u".to_string(),
1684                    at: now,
1685                }),
1686                "user_signed_up",
1687            ),
1688            (
1689                Event::UserSignedOut(UserSignedOutEvent {
1690                    user_id: id,
1691                    at: now,
1692                }),
1693                "user_signed_out",
1694            ),
1695        ];
1696
1697        assert_eq!(
1698            cases.len(),
1699            Event::ALL.len(),
1700            "every variant must be covered"
1701        );
1702
1703        for (event, expected_type) in cases {
1704            assert_eq!(event.event_type(), expected_type);
1705        }
1706    }
1707
1708    #[test]
1709    fn approval_escalated_serde_roundtrip() {
1710        let run_id = Uuid::now_v7();
1711        let step_id = Uuid::now_v7();
1712        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1713            run_id,
1714            step_id,
1715            step_name: "prod-gate".to_string(),
1716            stage: 1,
1717            policy: "escalate".to_string(),
1718            action: "reassigned to sre-oncall".to_string(),
1719            reason: "approval deadline of 3600s expired".to_string(),
1720            assignee: Some(Assignee::group("sre-oncall")),
1721            at: Utc::now(),
1722        });
1723
1724        let json = serde_json::to_string(&event).expect("serialize");
1725        assert!(
1726            json.contains("\"type\":\"approval_escalated\""),
1727            "got {json}"
1728        );
1729
1730        let back: Event = serde_json::from_str(&json).expect("deserialize");
1731        let Event::ApprovalEscalated(payload) = back else {
1732            panic!("expected an approval_escalated event");
1733        };
1734        assert_eq!(payload.run_id, run_id);
1735        assert_eq!(payload.stage, 1);
1736        assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
1737    }
1738
1739    #[test]
1740    fn approval_escalated_carries_run_and_step_ids() {
1741        let run_id = Uuid::now_v7();
1742        let step_id = Uuid::now_v7();
1743        let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
1744            run_id,
1745            step_id,
1746            step_name: "prod-gate".to_string(),
1747            stage: 0,
1748            policy: "notify".to_string(),
1749            action: "notified 1 target".to_string(),
1750            reason: "approval deadline of 60s expired".to_string(),
1751            assignee: None,
1752            at: Utc::now(),
1753        });
1754
1755        assert_eq!(event.run_id(), Some(run_id));
1756        assert_eq!(event.step_id(), Some(step_id));
1757        assert_eq!(event.user_id(), None);
1758    }
1759}