use std::collections::HashMap;
use chrono::{DateTime, Utc};
use rust_decimal::Decimal;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
pub use ironflow_store::entities::LogStream;
use ironflow_store::models::{Assignee, RunStatus, StepKind};
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct RunCreatedEvent {
pub run_id: Uuid,
pub workflow_name: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct RunStatusChangedEvent {
pub run_id: Uuid,
pub workflow_name: String,
pub from: RunStatus,
pub to: RunStatus,
pub error: Option<String>,
pub cost_usd: Decimal,
pub duration_ms: u64,
#[serde(default)]
pub labels: HashMap<String, String>,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct RunFailedEvent {
pub run_id: Uuid,
pub workflow_name: String,
pub error: Option<String>,
pub cost_usd: Decimal,
pub duration_ms: u64,
#[serde(default)]
pub labels: HashMap<String, String>,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct RunBudgetExceededEvent {
pub run_id: Uuid,
pub workflow_name: String,
pub limit_usd: Decimal,
pub spent_usd: Decimal,
pub step_budget_usd: Decimal,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct RetryForcedEvent {
pub run_id: Uuid,
pub workflow_name: String,
pub original_version: String,
pub current_version: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct StepCompletedEvent {
pub run_id: Uuid,
pub step_id: Uuid,
pub step_name: String,
#[cfg_attr(feature = "openapi", schema(value_type = String))]
pub kind: StepKind,
pub duration_ms: u64,
pub cost_usd: Decimal,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct StepFailedEvent {
pub run_id: Uuid,
pub step_id: Uuid,
pub step_name: String,
#[cfg_attr(feature = "openapi", schema(value_type = String))]
pub kind: StepKind,
pub error: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct ApprovalRequestedEvent {
pub run_id: Uuid,
pub step_id: Uuid,
pub message: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct ApprovalGrantedEvent {
pub run_id: Uuid,
pub approved_by: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct ApprovalRejectedEvent {
pub run_id: Uuid,
pub rejected_by: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct ApprovalEscalatedEvent {
pub run_id: Uuid,
pub step_id: Uuid,
pub step_name: String,
pub stage: u32,
pub policy: String,
pub action: String,
pub reason: String,
#[cfg_attr(feature = "openapi", schema(value_type = Option<String>))]
pub assignee: Option<Assignee>,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct LogLineEvent {
pub run_id: Uuid,
pub step_id: Uuid,
pub step_name: String,
pub stream: LogStream,
pub line: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct UserSignedInEvent {
pub user_id: Uuid,
pub username: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct UserSignedUpEvent {
pub user_id: Uuid,
pub username: String,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
pub struct UserSignedOutEvent {
pub user_id: Uuid,
pub at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum Event {
RunCreated(RunCreatedEvent),
RunStatusChanged(RunStatusChangedEvent),
RunFailed(RunFailedEvent),
RunBudgetExceeded(RunBudgetExceededEvent),
RetryForced(RetryForcedEvent),
StepCompleted(StepCompletedEvent),
StepFailed(StepFailedEvent),
ApprovalRequested(ApprovalRequestedEvent),
ApprovalGranted(ApprovalGrantedEvent),
ApprovalRejected(ApprovalRejectedEvent),
ApprovalEscalated(ApprovalEscalatedEvent),
LogLine(LogLineEvent),
UserSignedIn(UserSignedInEvent),
UserSignedUp(UserSignedUpEvent),
UserSignedOut(UserSignedOutEvent),
}
impl Event {
pub const RUN_CREATED: &'static str = "run_created";
pub const RUN_STATUS_CHANGED: &'static str = "run_status_changed";
pub const RUN_FAILED: &'static str = "run_failed";
pub const RUN_BUDGET_EXCEEDED: &'static str = "run_budget_exceeded";
pub const RETRY_FORCED: &'static str = "retry_forced";
pub const STEP_COMPLETED: &'static str = "step_completed";
pub const STEP_FAILED: &'static str = "step_failed";
pub const APPROVAL_REQUESTED: &'static str = "approval_requested";
pub const APPROVAL_GRANTED: &'static str = "approval_granted";
pub const APPROVAL_REJECTED: &'static str = "approval_rejected";
pub const APPROVAL_ESCALATED: &'static str = "approval_escalated";
pub const LOG_LINE: &'static str = "log_line";
pub const USER_SIGNED_IN: &'static str = "user_signed_in";
pub const USER_SIGNED_UP: &'static str = "user_signed_up";
pub const USER_SIGNED_OUT: &'static str = "user_signed_out";
pub const ALL: &'static [&'static str] = &[
Self::RUN_CREATED,
Self::RUN_STATUS_CHANGED,
Self::RUN_FAILED,
Self::RUN_BUDGET_EXCEEDED,
Self::STEP_COMPLETED,
Self::STEP_FAILED,
Self::APPROVAL_REQUESTED,
Self::APPROVAL_GRANTED,
Self::APPROVAL_REJECTED,
Self::APPROVAL_ESCALATED,
Self::LOG_LINE,
Self::USER_SIGNED_IN,
Self::USER_SIGNED_UP,
Self::USER_SIGNED_OUT,
Self::RETRY_FORCED,
];
#[deny(unreachable_patterns)]
pub fn event_type(&self) -> &'static str {
match self {
Event::RunCreated(_) => Self::RUN_CREATED,
Event::RunStatusChanged(_) => Self::RUN_STATUS_CHANGED,
Event::RunFailed(_) => Self::RUN_FAILED,
Event::RunBudgetExceeded(_) => Self::RUN_BUDGET_EXCEEDED,
Event::RetryForced(_) => Self::RETRY_FORCED,
Event::StepCompleted(_) => Self::STEP_COMPLETED,
Event::StepFailed(_) => Self::STEP_FAILED,
Event::ApprovalRequested(_) => Self::APPROVAL_REQUESTED,
Event::ApprovalGranted(_) => Self::APPROVAL_GRANTED,
Event::ApprovalRejected(_) => Self::APPROVAL_REJECTED,
Event::ApprovalEscalated(_) => Self::APPROVAL_ESCALATED,
Event::LogLine(_) => Self::LOG_LINE,
Event::UserSignedIn(_) => Self::USER_SIGNED_IN,
Event::UserSignedUp(_) => Self::USER_SIGNED_UP,
Event::UserSignedOut(_) => Self::USER_SIGNED_OUT,
}
}
#[deny(unreachable_patterns)]
pub fn run_id(&self) -> Option<Uuid> {
match self {
Event::RunCreated(e) => Some(e.run_id),
Event::RunStatusChanged(e) => Some(e.run_id),
Event::RunFailed(e) => Some(e.run_id),
Event::RunBudgetExceeded(e) => Some(e.run_id),
Event::RetryForced(e) => Some(e.run_id),
Event::StepCompleted(e) => Some(e.run_id),
Event::StepFailed(e) => Some(e.run_id),
Event::ApprovalRequested(e) => Some(e.run_id),
Event::ApprovalGranted(e) => Some(e.run_id),
Event::ApprovalRejected(e) => Some(e.run_id),
Event::ApprovalEscalated(e) => Some(e.run_id),
Event::LogLine(e) => Some(e.run_id),
Event::UserSignedIn(_) | Event::UserSignedUp(_) | Event::UserSignedOut(_) => None,
}
}
#[deny(unreachable_patterns)]
pub fn step_id(&self) -> Option<Uuid> {
match self {
Event::StepCompleted(e) => Some(e.step_id),
Event::StepFailed(e) => Some(e.step_id),
Event::ApprovalRequested(e) => Some(e.step_id),
Event::ApprovalEscalated(e) => Some(e.step_id),
Event::RunCreated(_)
| Event::RunStatusChanged(_)
| Event::RunFailed(_)
| Event::RunBudgetExceeded(_)
| Event::RetryForced(_)
| Event::ApprovalGranted(_)
| Event::ApprovalRejected(_)
| Event::LogLine(_)
| Event::UserSignedIn(_)
| Event::UserSignedUp(_)
| Event::UserSignedOut(_) => None,
}
}
#[deny(unreachable_patterns)]
pub fn user_id(&self) -> Option<Uuid> {
match self {
Event::UserSignedIn(e) => Some(e.user_id),
Event::UserSignedUp(e) => Some(e.user_id),
Event::UserSignedOut(e) => Some(e.user_id),
Event::RunCreated(_)
| Event::RunStatusChanged(_)
| Event::RunFailed(_)
| Event::RunBudgetExceeded(_)
| Event::RetryForced(_)
| Event::StepCompleted(_)
| Event::StepFailed(_)
| Event::ApprovalRequested(_)
| Event::ApprovalGranted(_)
| Event::ApprovalRejected(_)
| Event::ApprovalEscalated(_)
| Event::LogLine(_) => None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn run_status_changed_serde_roundtrip() {
let event = Event::RunStatusChanged(RunStatusChangedEvent {
run_id: Uuid::now_v7(),
workflow_name: "deploy".to_string(),
from: RunStatus::Running,
to: RunStatus::Completed,
error: None,
cost_usd: Decimal::new(42, 2),
duration_ms: 5000,
labels: HashMap::new(),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
let back: Event = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.event_type(), "run_status_changed");
assert!(json.contains("\"type\":\"run_status_changed\""));
}
#[test]
fn run_failed_serde_roundtrip() {
let event = Event::RunFailed(RunFailedEvent {
run_id: Uuid::now_v7(),
workflow_name: "deploy".to_string(),
error: Some("step crashed".to_string()),
cost_usd: Decimal::new(10, 2),
duration_ms: 3000,
labels: HashMap::new(),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
let back: Event = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.event_type(), "run_failed");
assert!(json.contains("\"type\":\"run_failed\""));
assert!(json.contains("step crashed"));
}
#[test]
fn run_budget_exceeded_serde_roundtrip() {
let event = Event::RunBudgetExceeded(RunBudgetExceededEvent {
run_id: Uuid::now_v7(),
workflow_name: "deploy".to_string(),
limit_usd: Decimal::new(200, 2),
spent_usd: Decimal::new(180, 2),
step_budget_usd: Decimal::new(50, 2),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
let back: Event = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.event_type(), "run_budget_exceeded");
assert!(json.contains("\"type\":\"run_budget_exceeded\""));
assert!(json.contains("limit_usd"));
assert!(json.contains("step_budget_usd"));
}
#[test]
fn all_contains_run_budget_exceeded() {
assert!(Event::ALL.contains(&Event::RUN_BUDGET_EXCEEDED));
}
#[test]
fn user_signed_in_serde_roundtrip() {
let event = Event::UserSignedIn(UserSignedInEvent {
user_id: Uuid::now_v7(),
username: "alice".to_string(),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
let back: Event = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.event_type(), "user_signed_in");
assert!(json.contains("alice"));
}
#[test]
fn step_failed_serde_roundtrip() {
let event = Event::StepFailed(StepFailedEvent {
run_id: Uuid::now_v7(),
step_id: Uuid::now_v7(),
step_name: "build".to_string(),
kind: StepKind::Shell,
error: "exit code 1".to_string(),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
let back: Event = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.event_type(), "step_failed");
}
#[test]
fn approval_requested_serde_roundtrip() {
let event = Event::ApprovalRequested(ApprovalRequestedEvent {
run_id: Uuid::now_v7(),
step_id: Uuid::now_v7(),
message: "Deploy to prod?".to_string(),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
assert!(json.contains("approval_requested"));
}
#[test]
fn log_line_serde_roundtrip() {
let event = Event::LogLine(LogLineEvent {
run_id: Uuid::now_v7(),
step_id: Uuid::now_v7(),
step_name: "build".to_string(),
stream: LogStream::Stdout,
line: "Compiling ironflow v0.1.0".to_string(),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
let back: Event = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.event_type(), "log_line");
assert!(json.contains("\"type\":\"log_line\""));
assert!(json.contains("Compiling ironflow"));
}
#[test]
fn legacy_flat_json_deserializes_into_typed_payload() {
let run_id: Uuid = "01890000-0000-7000-8000-000000000000"
.parse()
.expect("valid uuid");
let raw = r#"{"type":"run_created","run_id":"01890000-0000-7000-8000-000000000000","workflow_name":"deploy","at":"2026-01-01T00:00:00Z"}"#;
let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
match event {
Event::RunCreated(e) => {
assert_eq!(e.run_id, run_id);
assert_eq!(e.workflow_name, "deploy");
}
other => panic!("expected RunCreated, got {other:?}"),
}
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"}"#;
let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
match event {
Event::RunStatusChanged(e) => {
assert_eq!(e.from, RunStatus::Running);
assert_eq!(e.to, RunStatus::Completed);
assert_eq!(e.cost_usd, Decimal::new(5, 1));
assert_eq!(e.duration_ms, 5000);
assert_eq!(e.labels.get("env").map(String::as_str), Some("prod"));
}
other => panic!("expected RunStatusChanged, got {other:?}"),
}
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"}"#;
let event: Event = serde_json::from_str(raw).expect("missing labels must default");
match event {
Event::RunStatusChanged(e) => {
assert!(e.labels.is_empty());
assert_eq!(e.error.as_deref(), Some("boom"));
}
other => panic!("expected RunStatusChanged, got {other:?}"),
}
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"}"#;
let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
match event {
Event::RunFailed(e) => {
assert_eq!(e.error.as_deref(), Some("boom"));
assert!(e.labels.is_empty());
}
other => panic!("expected RunFailed, got {other:?}"),
}
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"}"#;
let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
match event {
Event::StepFailed(e) => {
assert_eq!(e.kind, StepKind::Shell);
assert_eq!(e.error, "exit code 1");
}
other => panic!("expected StepFailed, got {other:?}"),
}
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"}"#;
let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
match event {
Event::LogLine(e) => {
assert_eq!(e.stream, LogStream::Stdout);
assert_eq!(e.line, "hello");
}
other => panic!("expected LogLine, got {other:?}"),
}
let raw = r#"{"type":"user_signed_in","user_id":"01890000-0000-7000-8000-000000000000","username":"alice","at":"2026-01-01T00:00:00Z"}"#;
let event: Event = serde_json::from_str(raw).expect("legacy payload must deserialize");
match event {
Event::UserSignedIn(e) => assert_eq!(e.username, "alice"),
other => panic!("expected UserSignedIn, got {other:?}"),
}
}
#[test]
fn serialized_event_is_flat_with_type_tag() {
let run_id = Uuid::now_v7();
let event = Event::RunCreated(RunCreatedEvent {
run_id,
workflow_name: "deploy".to_string(),
at: Utc::now(),
});
let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
let object = value.as_object().expect("event serializes to an object");
assert_eq!(
object.get("type").and_then(|v| v.as_str()),
Some("run_created")
);
assert_eq!(
object.get("workflow_name").and_then(|v| v.as_str()),
Some("deploy")
);
assert_eq!(
object.get("run_id").and_then(|v| v.as_str()),
Some(run_id.to_string().as_str())
);
assert!(object.contains_key("at"));
assert_eq!(object.len(), 4, "no nesting: {object:?}");
assert!(!object.contains_key("RunCreated"));
}
#[test]
fn run_id_returns_some_for_run_events() {
let run_id = Uuid::now_v7();
let now = Utc::now();
let events = vec![
Event::RunCreated(RunCreatedEvent {
run_id,
workflow_name: "w".to_string(),
at: now,
}),
Event::RunStatusChanged(RunStatusChangedEvent {
run_id,
workflow_name: "w".to_string(),
from: RunStatus::Pending,
to: RunStatus::Running,
error: None,
cost_usd: Decimal::ZERO,
duration_ms: 0,
labels: HashMap::new(),
at: now,
}),
Event::RunFailed(RunFailedEvent {
run_id,
workflow_name: "w".to_string(),
error: None,
cost_usd: Decimal::ZERO,
duration_ms: 0,
labels: HashMap::new(),
at: now,
}),
Event::RunBudgetExceeded(RunBudgetExceededEvent {
run_id,
workflow_name: "w".to_string(),
limit_usd: Decimal::ZERO,
spent_usd: Decimal::ZERO,
step_budget_usd: Decimal::ZERO,
at: now,
}),
Event::RetryForced(RetryForcedEvent {
run_id,
workflow_name: "w".to_string(),
original_version: "1".to_string(),
current_version: "2".to_string(),
at: now,
}),
Event::StepCompleted(StepCompletedEvent {
run_id,
step_id: Uuid::now_v7(),
step_name: "s".to_string(),
kind: StepKind::Shell,
duration_ms: 0,
cost_usd: Decimal::ZERO,
at: now,
}),
Event::StepFailed(StepFailedEvent {
run_id,
step_id: Uuid::now_v7(),
step_name: "s".to_string(),
kind: StepKind::Shell,
error: "e".to_string(),
at: now,
}),
Event::ApprovalRequested(ApprovalRequestedEvent {
run_id,
step_id: Uuid::now_v7(),
message: "ok?".to_string(),
at: now,
}),
Event::ApprovalGranted(ApprovalGrantedEvent {
run_id,
approved_by: "alice".to_string(),
at: now,
}),
Event::ApprovalRejected(ApprovalRejectedEvent {
run_id,
rejected_by: "bob".to_string(),
at: now,
}),
Event::LogLine(LogLineEvent {
run_id,
step_id: Uuid::now_v7(),
step_name: "s".to_string(),
stream: LogStream::Stdout,
line: "l".to_string(),
at: now,
}),
];
for event in &events {
assert_eq!(
event.run_id(),
Some(run_id),
"{} should carry a run_id",
event.event_type()
);
}
}
#[test]
fn run_id_returns_none_for_auth_events() {
let user_id = Uuid::now_v7();
let now = Utc::now();
let events = vec![
Event::UserSignedIn(UserSignedInEvent {
user_id,
username: "alice".to_string(),
at: now,
}),
Event::UserSignedUp(UserSignedUpEvent {
user_id,
username: "alice".to_string(),
at: now,
}),
Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
];
for event in &events {
assert_eq!(event.run_id(), None, "{} has no run", event.event_type());
}
}
#[test]
fn step_id_returns_some_only_for_step_events() {
let step_id = Uuid::now_v7();
let run_id = Uuid::now_v7();
let now = Utc::now();
let with_step = vec![
Event::StepCompleted(StepCompletedEvent {
run_id,
step_id,
step_name: "s".to_string(),
kind: StepKind::Shell,
duration_ms: 0,
cost_usd: Decimal::ZERO,
at: now,
}),
Event::StepFailed(StepFailedEvent {
run_id,
step_id,
step_name: "s".to_string(),
kind: StepKind::Shell,
error: "e".to_string(),
at: now,
}),
Event::ApprovalRequested(ApprovalRequestedEvent {
run_id,
step_id,
message: "ok?".to_string(),
at: now,
}),
];
for event in &with_step {
assert_eq!(
event.step_id(),
Some(step_id),
"{} should carry a step_id",
event.event_type()
);
}
let without_step = vec![
Event::RunCreated(RunCreatedEvent {
run_id,
workflow_name: "w".to_string(),
at: now,
}),
Event::ApprovalGranted(ApprovalGrantedEvent {
run_id,
approved_by: "alice".to_string(),
at: now,
}),
Event::LogLine(LogLineEvent {
run_id,
step_id,
step_name: "s".to_string(),
stream: LogStream::Stdout,
line: "l".to_string(),
at: now,
}),
Event::UserSignedOut(UserSignedOutEvent {
user_id: Uuid::now_v7(),
at: now,
}),
];
for event in &without_step {
assert_eq!(
event.step_id(),
None,
"{} should not carry a step_id",
event.event_type()
);
}
}
#[test]
fn user_id_returns_some_only_for_auth_events() {
let user_id = Uuid::now_v7();
let run_id = Uuid::now_v7();
let now = Utc::now();
let auth = vec![
Event::UserSignedIn(UserSignedInEvent {
user_id,
username: "alice".to_string(),
at: now,
}),
Event::UserSignedUp(UserSignedUpEvent {
user_id,
username: "alice".to_string(),
at: now,
}),
Event::UserSignedOut(UserSignedOutEvent { user_id, at: now }),
];
for event in &auth {
assert_eq!(
event.user_id(),
Some(user_id),
"{} should carry a user_id",
event.event_type()
);
}
let non_auth = vec![
Event::RunCreated(RunCreatedEvent {
run_id,
workflow_name: "w".to_string(),
at: now,
}),
Event::StepFailed(StepFailedEvent {
run_id,
step_id: Uuid::now_v7(),
step_name: "s".to_string(),
kind: StepKind::Shell,
error: "e".to_string(),
at: now,
}),
];
for event in &non_auth {
assert_eq!(
event.user_id(),
None,
"{} should not carry a user_id",
event.event_type()
);
}
}
#[test]
fn event_type_all_variants() {
let id = Uuid::now_v7();
let now = Utc::now();
let cases: Vec<(Event, &str)> = vec![
(
Event::RunCreated(RunCreatedEvent {
run_id: id,
workflow_name: "w".to_string(),
at: now,
}),
"run_created",
),
(
Event::RunStatusChanged(RunStatusChangedEvent {
run_id: id,
workflow_name: "w".to_string(),
from: RunStatus::Pending,
to: RunStatus::Running,
error: None,
cost_usd: Decimal::ZERO,
duration_ms: 0,
labels: HashMap::new(),
at: now,
}),
"run_status_changed",
),
(
Event::RunFailed(RunFailedEvent {
run_id: id,
workflow_name: "w".to_string(),
error: Some("boom".to_string()),
cost_usd: Decimal::ZERO,
duration_ms: 0,
labels: HashMap::new(),
at: now,
}),
"run_failed",
),
(
Event::RunBudgetExceeded(RunBudgetExceededEvent {
run_id: id,
workflow_name: "w".to_string(),
limit_usd: Decimal::new(200, 2),
spent_usd: Decimal::new(180, 2),
step_budget_usd: Decimal::new(50, 2),
at: now,
}),
"run_budget_exceeded",
),
(
Event::RetryForced(RetryForcedEvent {
run_id: id,
workflow_name: "w".to_string(),
original_version: "1".to_string(),
current_version: "2".to_string(),
at: now,
}),
"retry_forced",
),
(
Event::StepCompleted(StepCompletedEvent {
run_id: id,
step_id: id,
step_name: "s".to_string(),
kind: StepKind::Shell,
duration_ms: 0,
cost_usd: Decimal::ZERO,
at: now,
}),
"step_completed",
),
(
Event::StepFailed(StepFailedEvent {
run_id: id,
step_id: id,
step_name: "s".to_string(),
kind: StepKind::Shell,
error: "err".to_string(),
at: now,
}),
"step_failed",
),
(
Event::ApprovalRequested(ApprovalRequestedEvent {
run_id: id,
step_id: id,
message: "ok?".to_string(),
at: now,
}),
"approval_requested",
),
(
Event::ApprovalGranted(ApprovalGrantedEvent {
run_id: id,
approved_by: "alice".to_string(),
at: now,
}),
"approval_granted",
),
(
Event::ApprovalRejected(ApprovalRejectedEvent {
run_id: id,
rejected_by: "bob".to_string(),
at: now,
}),
"approval_rejected",
),
(
Event::ApprovalEscalated(ApprovalEscalatedEvent {
run_id: id,
step_id: id,
step_name: "prod-gate".to_string(),
stage: 0,
policy: "auto_reject".to_string(),
action: "rejected".to_string(),
reason: "approval deadline of 3600s expired".to_string(),
assignee: None,
at: now,
}),
"approval_escalated",
),
(
Event::LogLine(LogLineEvent {
run_id: id,
step_id: id,
step_name: "build".to_string(),
stream: LogStream::Stdout,
line: "Compiling ironflow v0.1.0".to_string(),
at: now,
}),
"log_line",
),
(
Event::UserSignedIn(UserSignedInEvent {
user_id: id,
username: "u".to_string(),
at: now,
}),
"user_signed_in",
),
(
Event::UserSignedUp(UserSignedUpEvent {
user_id: id,
username: "u".to_string(),
at: now,
}),
"user_signed_up",
),
(
Event::UserSignedOut(UserSignedOutEvent {
user_id: id,
at: now,
}),
"user_signed_out",
),
];
assert_eq!(
cases.len(),
Event::ALL.len(),
"every variant must be covered"
);
for (event, expected_type) in cases {
assert_eq!(event.event_type(), expected_type);
}
}
#[test]
fn approval_escalated_serde_roundtrip() {
let run_id = Uuid::now_v7();
let step_id = Uuid::now_v7();
let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
run_id,
step_id,
step_name: "prod-gate".to_string(),
stage: 1,
policy: "escalate".to_string(),
action: "reassigned to sre-oncall".to_string(),
reason: "approval deadline of 3600s expired".to_string(),
assignee: Some(Assignee::group("sre-oncall")),
at: Utc::now(),
});
let json = serde_json::to_string(&event).expect("serialize");
assert!(
json.contains("\"type\":\"approval_escalated\""),
"got {json}"
);
let back: Event = serde_json::from_str(&json).expect("deserialize");
let Event::ApprovalEscalated(payload) = back else {
panic!("expected an approval_escalated event");
};
assert_eq!(payload.run_id, run_id);
assert_eq!(payload.stage, 1);
assert_eq!(payload.assignee, Some(Assignee::group("sre-oncall")));
}
#[test]
fn approval_escalated_carries_run_and_step_ids() {
let run_id = Uuid::now_v7();
let step_id = Uuid::now_v7();
let event = Event::ApprovalEscalated(ApprovalEscalatedEvent {
run_id,
step_id,
step_name: "prod-gate".to_string(),
stage: 0,
policy: "notify".to_string(),
action: "notified 1 target".to_string(),
reason: "approval deadline of 60s expired".to_string(),
assignee: None,
at: Utc::now(),
});
assert_eq!(event.run_id(), Some(run_id));
assert_eq!(event.step_id(), Some(step_id));
assert_eq!(event.user_id(), None);
}
}