Skip to main content

ironflow_engine/notify/
event_bus.rs

1//! Per-workflow broadcast event bus for real-time monitoring.
2//!
3//! [`WorkflowEventBus`] maintains one [`tokio::sync::broadcast`] channel per
4//! workflow run. Consumers (dashboards, SSE routes) subscribe to a specific
5//! `run_id` and receive only the events for that run.
6//!
7//! # Architecture
8//!
9//! - [`WorkflowEvent`] -- granular step-level events (started, completed,
10//!   failed, approval, human input, token usage, agent session resume).
11//! - [`WorkflowEventBus`] -- per-run broadcast channels with subscribe /
12//!   publish / remove lifecycle.
13//!
14//! # Examples
15//!
16//! ```
17//! use ironflow_engine::notify::{WorkflowEvent, WorkflowEventBus, WorkflowStepStartedEvent};
18//! use uuid::Uuid;
19//! use chrono::Utc;
20//!
21//! let bus = WorkflowEventBus::new();
22//! let run_id = Uuid::now_v7();
23//!
24//! let mut rx = bus.subscribe(run_id);
25//!
26//! bus.publish(run_id, WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
27//!     step_name: "build".to_string(),
28//!     step_index: 0,
29//!     timestamp: Utc::now(),
30//! }));
31//! ```
32
33use std::collections::HashMap;
34use std::sync::RwLock;
35
36use chrono::{DateTime, Utc};
37use rust_decimal::Decimal;
38use serde::{Deserialize, Serialize};
39use serde_json::Value;
40use tokio::sync::broadcast;
41use uuid::Uuid;
42
43/// Default broadcast channel buffer size per run.
44const DEFAULT_BUFFER_SIZE: usize = 64;
45
46/// Payload of the `WorkflowEvent::StepStarted` workflow event.
47///
48/// # Examples
49///
50/// ```
51/// use chrono::Utc;
52/// use ironflow_engine::notify::WorkflowStepStartedEvent;
53///
54/// let payload = WorkflowStepStartedEvent {
55///     step_name: "build".to_string(),
56///     step_index: 0,
57///     timestamp: Utc::now(),
58/// };
59/// assert_eq!(payload.step_index, 0);
60/// ```
61#[derive(Debug, Clone, Serialize, Deserialize)]
62#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
63pub struct WorkflowStepStartedEvent {
64    /// Human-readable step name.
65    pub step_name: String,
66    /// Zero-based position in the workflow.
67    pub step_index: u32,
68    /// When the step started.
69    pub timestamp: DateTime<Utc>,
70}
71
72/// Payload of the `WorkflowEvent::StepCompleted` workflow event.
73///
74/// # Examples
75///
76/// ```
77/// use ironflow_engine::notify::WorkflowStepCompletedEvent;
78///
79/// let payload = WorkflowStepCompletedEvent {
80///     step_name: "deploy".to_string(),
81///     step_index: 1,
82///     duration_ms: 5000,
83///     output_summary: Some("deployed v1.2.3".to_string()),
84/// };
85/// assert_eq!(payload.duration_ms, 5000);
86/// ```
87#[derive(Debug, Clone, Serialize, Deserialize)]
88#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
89pub struct WorkflowStepCompletedEvent {
90    /// Human-readable step name.
91    pub step_name: String,
92    /// Zero-based position in the workflow.
93    pub step_index: u32,
94    /// Step duration in milliseconds.
95    pub duration_ms: u64,
96    /// Optional summary of the step output.
97    pub output_summary: Option<String>,
98}
99
100/// Payload of the `WorkflowEvent::StepFailed` workflow event.
101///
102/// # Examples
103///
104/// ```
105/// use ironflow_engine::notify::WorkflowStepFailedEvent;
106///
107/// let payload = WorkflowStepFailedEvent {
108///     step_name: "test".to_string(),
109///     step_index: 2,
110///     error: "exit code 1".to_string(),
111///     duration_ms: 3000,
112/// };
113/// assert_eq!(payload.error, "exit code 1");
114/// ```
115#[derive(Debug, Clone, Serialize, Deserialize)]
116#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
117pub struct WorkflowStepFailedEvent {
118    /// Human-readable step name.
119    pub step_name: String,
120    /// Zero-based position in the workflow.
121    pub step_index: u32,
122    /// Error description.
123    pub error: String,
124    /// Step duration in milliseconds.
125    pub duration_ms: u64,
126}
127
128/// Payload of the `WorkflowEvent::ApprovalRequired` workflow event.
129///
130/// # Examples
131///
132/// ```
133/// use ironflow_engine::notify::WorkflowApprovalRequiredEvent;
134/// use uuid::Uuid;
135///
136/// let payload = WorkflowApprovalRequiredEvent {
137///     step_name: "prod-gate".to_string(),
138///     step_index: 3,
139///     approval_id: Uuid::now_v7(),
140/// };
141/// assert_eq!(payload.step_name, "prod-gate");
142/// ```
143#[derive(Debug, Clone, Serialize, Deserialize)]
144#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
145pub struct WorkflowApprovalRequiredEvent {
146    /// Human-readable step name.
147    pub step_name: String,
148    /// Zero-based position in the workflow.
149    pub step_index: u32,
150    /// Identifier of the approval gate.
151    pub approval_id: Uuid,
152}
153
154/// Payload of the `WorkflowEvent::InputRequired` workflow event.
155///
156/// # Examples
157///
158/// ```
159/// use ironflow_engine::notify::WorkflowInputRequiredEvent;
160/// use serde_json::json;
161/// use uuid::Uuid;
162///
163/// let payload = WorkflowInputRequiredEvent {
164///     run_id: Uuid::now_v7(),
165///     step_id: Uuid::now_v7(),
166///     step_name: "clarify".to_string(),
167///     step_index: 2,
168///     message: "Answer the questions".to_string(),
169///     schema: json!({"type": "object"}),
170/// };
171/// assert_eq!(payload.step_name, "clarify");
172/// ```
173#[derive(Debug, Clone, Serialize, Deserialize)]
174#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
175pub struct WorkflowInputRequiredEvent {
176    /// The run waiting for the input.
177    pub run_id: Uuid,
178    /// Identifier of the human input step; the answer is posted to it.
179    pub step_id: Uuid,
180    /// Human-readable step name.
181    pub step_name: String,
182    /// Zero-based position in the workflow.
183    pub step_index: u32,
184    /// Message displayed to the person answering.
185    pub message: String,
186    /// JSON schema the answer must match.
187    #[cfg_attr(feature = "openapi", schema(value_type = Object))]
188    pub schema: Value,
189}
190
191/// Payload of the `WorkflowEvent::AgentStepTokensUsed` workflow event.
192///
193/// # Examples
194///
195/// ```
196/// use ironflow_engine::notify::WorkflowAgentStepTokensUsedEvent;
197/// use rust_decimal::Decimal;
198///
199/// let payload = WorkflowAgentStepTokensUsedEvent {
200///     step_name: "review".to_string(),
201///     tokens: 15_000,
202///     cost_usd: Decimal::new(42, 4),
203/// };
204/// assert_eq!(payload.tokens, 15_000);
205/// ```
206#[derive(Debug, Clone, Serialize, Deserialize)]
207#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
208pub struct WorkflowAgentStepTokensUsedEvent {
209    /// Human-readable step name.
210    pub step_name: String,
211    /// Total tokens consumed.
212    pub tokens: u64,
213    /// Estimated cost in USD.
214    pub cost_usd: Decimal,
215}
216
217/// Payload of the `WorkflowEvent::AgentStepResumed` workflow event.
218///
219/// Published when an agent step interrupted by a lost lease is re-executed
220/// on the Claude Code session of the interrupted attempt instead of starting
221/// from scratch.
222///
223/// # Examples
224///
225/// ```
226/// use ironflow_engine::notify::WorkflowAgentStepResumedEvent;
227/// use chrono::Utc;
228///
229/// let payload = WorkflowAgentStepResumedEvent {
230///     step_name: "review".to_string(),
231///     step_index: 2,
232///     session_id: "0192f0c1-7d2e-7a4b-9c3d-1e2f3a4b5c6d".to_string(),
233///     timestamp: Utc::now(),
234/// };
235/// assert_eq!(payload.step_index, 2);
236/// ```
237#[derive(Debug, Clone, Serialize, Deserialize)]
238#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
239pub struct WorkflowAgentStepResumedEvent {
240    /// Human-readable step name.
241    pub step_name: String,
242    /// Zero-based position in the workflow.
243    pub step_index: u32,
244    /// Claude Code session the step resumes.
245    pub session_id: String,
246    /// When the step resumed.
247    pub timestamp: DateTime<Utc>,
248}
249
250/// A granular step-level event for real-time workflow monitoring.
251///
252/// Unlike [`Event`](super::Event) which covers the full system lifecycle
253/// (runs, auth, audit), `WorkflowEvent` tracks individual step transitions
254/// within a single run. Serialized with a `type` discriminant for UI
255/// consumption.
256///
257/// Each variant wraps a dedicated payload struct; the serialized form stays
258/// flat, with `type` sitting next to the payload fields.
259///
260/// # Examples
261///
262/// ```
263/// use ironflow_engine::notify::{WorkflowEvent, WorkflowStepStartedEvent};
264/// use chrono::Utc;
265///
266/// let event = WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
267///     step_name: "deploy".to_string(),
268///     step_index: 0,
269///     timestamp: Utc::now(),
270/// });
271/// assert_eq!(event.event_type(), "step_started");
272///
273/// let json = serde_json::to_string(&event)?;
274/// assert!(json.contains("\"type\":\"step_started\""));
275/// # Ok::<(), serde_json::Error>(())
276/// ```
277#[derive(Debug, Clone, Serialize, Deserialize)]
278#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
279#[serde(tag = "type", rename_all = "snake_case")]
280pub enum WorkflowEvent {
281    /// A step began execution.
282    StepStarted(WorkflowStepStartedEvent),
283
284    /// A step completed successfully.
285    StepCompleted(WorkflowStepCompletedEvent),
286
287    /// A step failed.
288    StepFailed(WorkflowStepFailedEvent),
289
290    /// A step requires human approval before the run can continue.
291    ApprovalRequired(WorkflowApprovalRequiredEvent),
292
293    /// A step waits for a typed human input before the run can continue.
294    InputRequired(WorkflowInputRequiredEvent),
295
296    /// Token usage report for an agent step.
297    AgentStepTokensUsed(WorkflowAgentStepTokensUsedEvent),
298
299    /// An interrupted agent step resumed its Claude Code session.
300    AgentStepResumed(WorkflowAgentStepResumedEvent),
301}
302
303impl WorkflowEvent {
304    /// Event type constant for [`StepStarted`](WorkflowEvent::StepStarted).
305    pub const STEP_STARTED: &'static str = "step_started";
306    /// Event type constant for [`StepCompleted`](WorkflowEvent::StepCompleted).
307    pub const STEP_COMPLETED: &'static str = "step_completed";
308    /// Event type constant for [`StepFailed`](WorkflowEvent::StepFailed).
309    pub const STEP_FAILED: &'static str = "step_failed";
310    /// Event type constant for [`ApprovalRequired`](WorkflowEvent::ApprovalRequired).
311    pub const APPROVAL_REQUIRED: &'static str = "approval_required";
312    /// Event type constant for [`InputRequired`](WorkflowEvent::InputRequired).
313    pub const INPUT_REQUIRED: &'static str = "input_required";
314    /// Event type constant for [`AgentStepTokensUsed`](WorkflowEvent::AgentStepTokensUsed).
315    pub const AGENT_STEP_TOKENS_USED: &'static str = "agent_step_tokens_used";
316    /// Event type constant for [`AgentStepResumed`](WorkflowEvent::AgentStepResumed).
317    pub const AGENT_STEP_RESUMED: &'static str = "agent_step_resumed";
318
319    /// Returns the event type as a static string (e.g. `"step_started"`).
320    ///
321    /// # Examples
322    ///
323    /// ```
324    /// use ironflow_engine::notify::{WorkflowEvent, WorkflowStepStartedEvent};
325    /// use chrono::Utc;
326    ///
327    /// let event = WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
328    ///     step_name: "build".to_string(),
329    ///     step_index: 0,
330    ///     timestamp: Utc::now(),
331    /// });
332    /// assert_eq!(event.event_type(), "step_started");
333    /// ```
334    #[deny(unreachable_patterns)]
335    pub fn event_type(&self) -> &'static str {
336        match self {
337            WorkflowEvent::StepStarted(_) => Self::STEP_STARTED,
338            WorkflowEvent::StepCompleted(_) => Self::STEP_COMPLETED,
339            WorkflowEvent::StepFailed(_) => Self::STEP_FAILED,
340            WorkflowEvent::ApprovalRequired(_) => Self::APPROVAL_REQUIRED,
341            WorkflowEvent::InputRequired(_) => Self::INPUT_REQUIRED,
342            WorkflowEvent::AgentStepTokensUsed(_) => Self::AGENT_STEP_TOKENS_USED,
343            WorkflowEvent::AgentStepResumed(_) => Self::AGENT_STEP_RESUMED,
344        }
345    }
346}
347
348/// Per-workflow broadcast event bus for real-time monitoring.
349///
350/// Maintains one [`tokio::sync::broadcast`] channel per workflow run.
351/// Consumers call [`subscribe`](Self::subscribe) to receive events for a
352/// specific run; producers call [`publish`](Self::publish) to broadcast
353/// an event to all subscribers of that run.
354///
355/// Thread-safe and cheaply cloneable (`Clone` shares the same inner state).
356///
357/// # Examples
358///
359/// ```
360/// use ironflow_engine::notify::{WorkflowEvent, WorkflowEventBus, WorkflowStepStartedEvent};
361/// use uuid::Uuid;
362/// use chrono::Utc;
363///
364/// let bus = WorkflowEventBus::new();
365/// let run_id = Uuid::now_v7();
366///
367/// let mut rx = bus.subscribe(run_id);
368/// bus.publish(run_id, WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
369///     step_name: "build".to_string(),
370///     step_index: 0,
371///     timestamp: Utc::now(),
372/// }));
373/// ```
374#[derive(Clone)]
375pub struct WorkflowEventBus {
376    channels: std::sync::Arc<RwLock<HashMap<Uuid, broadcast::Sender<WorkflowEvent>>>>,
377}
378
379impl WorkflowEventBus {
380    /// Create a new empty event bus.
381    ///
382    /// # Examples
383    ///
384    /// ```
385    /// use ironflow_engine::notify::WorkflowEventBus;
386    ///
387    /// let bus = WorkflowEventBus::new();
388    /// ```
389    pub fn new() -> Self {
390        Self {
391            channels: std::sync::Arc::new(RwLock::new(HashMap::new())),
392        }
393    }
394
395    /// Subscribe to events for a specific workflow run.
396    ///
397    /// If no channel exists for this `run_id`, one is created on demand.
398    /// Returns a broadcast receiver that yields [`WorkflowEvent`]s for
399    /// that run only.
400    ///
401    /// # Examples
402    ///
403    /// ```
404    /// use ironflow_engine::notify::WorkflowEventBus;
405    /// use uuid::Uuid;
406    ///
407    /// let bus = WorkflowEventBus::new();
408    /// let run_id = Uuid::now_v7();
409    /// let _rx = bus.subscribe(run_id);
410    /// ```
411    pub fn subscribe(&self, run_id: Uuid) -> broadcast::Receiver<WorkflowEvent> {
412        let mut channels = self.channels.write().expect("event bus lock poisoned");
413        let sender = channels
414            .entry(run_id)
415            .or_insert_with(|| broadcast::channel(DEFAULT_BUFFER_SIZE).0);
416        sender.subscribe()
417    }
418
419    /// Broadcast an event to all subscribers of a specific workflow run.
420    ///
421    /// If no channel exists for `run_id` (no subscriber has called
422    /// [`subscribe`](Self::subscribe)), the event is silently dropped.
423    /// If subscribers exist but none are actively listening, the send
424    /// error is ignored.
425    ///
426    /// # Examples
427    ///
428    /// ```
429    /// use ironflow_engine::notify::{WorkflowEvent, WorkflowEventBus, WorkflowStepStartedEvent};
430    /// use uuid::Uuid;
431    /// use chrono::Utc;
432    ///
433    /// let bus = WorkflowEventBus::new();
434    /// let run_id = Uuid::now_v7();
435    ///
436    /// // No subscriber -- silently dropped.
437    /// bus.publish(run_id, WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
438    ///     step_name: "build".to_string(),
439    ///     step_index: 0,
440    ///     timestamp: Utc::now(),
441    /// }));
442    /// ```
443    pub fn publish(&self, run_id: Uuid, event: WorkflowEvent) {
444        let channels = self.channels.read().expect("event bus lock poisoned");
445        if let Some(sender) = channels.get(&run_id) {
446            let _ = sender.send(event);
447        }
448    }
449
450    /// Remove the channel for a workflow run.
451    ///
452    /// Call this when a run completes or is cleaned up to free resources.
453    /// If no channel exists for `run_id`, this is a no-op.
454    ///
455    /// # Examples
456    ///
457    /// ```
458    /// use ironflow_engine::notify::WorkflowEventBus;
459    /// use uuid::Uuid;
460    ///
461    /// let bus = WorkflowEventBus::new();
462    /// let run_id = Uuid::now_v7();
463    /// let _rx = bus.subscribe(run_id);
464    /// bus.remove(run_id);
465    /// ```
466    pub fn remove(&self, run_id: Uuid) {
467        let mut channels = self.channels.write().expect("event bus lock poisoned");
468        channels.remove(&run_id);
469    }
470}
471
472impl Default for WorkflowEventBus {
473    fn default() -> Self {
474        Self::new()
475    }
476}
477
478impl std::fmt::Debug for WorkflowEventBus {
479    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
480        let count = self.channels.read().map(|c| c.len()).unwrap_or(0);
481        f.debug_struct("WorkflowEventBus")
482            .field("active_channels", &count)
483            .finish()
484    }
485}
486
487#[cfg(test)]
488mod tests {
489    use serde_json::json;
490
491    use super::*;
492
493    fn step_started(step_name: &str) -> WorkflowEvent {
494        WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
495            step_name: step_name.to_string(),
496            step_index: 0,
497            timestamp: Utc::now(),
498        })
499    }
500
501    #[tokio::test]
502    async fn subscribe_receives_published_events() {
503        let bus = WorkflowEventBus::new();
504        let run_id = Uuid::now_v7();
505
506        let mut rx = bus.subscribe(run_id);
507
508        bus.publish(run_id, step_started("build"));
509
510        let received = rx.recv().await.expect("should receive event");
511        assert_eq!(received.event_type(), "step_started");
512        match received {
513            WorkflowEvent::StepStarted(e) => {
514                assert_eq!(e.step_name, "build");
515                assert_eq!(e.step_index, 0);
516            }
517            _ => panic!("expected StepStarted"),
518        }
519    }
520
521    #[test]
522    fn subscribe_creates_channel_on_demand() {
523        let bus = WorkflowEventBus::new();
524        let run_id = Uuid::now_v7();
525
526        let count_before = bus.channels.read().unwrap().len();
527        assert_eq!(count_before, 0);
528
529        let _rx = bus.subscribe(run_id);
530
531        let count_after = bus.channels.read().unwrap().len();
532        assert_eq!(count_after, 1);
533    }
534
535    #[test]
536    fn publish_unknown_run_is_noop() {
537        let bus = WorkflowEventBus::new();
538        let unknown_run = Uuid::now_v7();
539
540        bus.publish(unknown_run, step_started("build"));
541    }
542
543    #[test]
544    fn remove_cleans_up_channel() {
545        let bus = WorkflowEventBus::new();
546        let run_id = Uuid::now_v7();
547
548        let _rx = bus.subscribe(run_id);
549        assert_eq!(bus.channels.read().unwrap().len(), 1);
550
551        bus.remove(run_id);
552        assert_eq!(bus.channels.read().unwrap().len(), 0);
553    }
554
555    #[test]
556    fn remove_unknown_is_noop() {
557        let bus = WorkflowEventBus::new();
558        bus.remove(Uuid::now_v7());
559    }
560
561    #[test]
562    fn workflow_event_serde_roundtrip() {
563        let cases: Vec<WorkflowEvent> = vec![
564            WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
565                step_name: "build".to_string(),
566                step_index: 0,
567                timestamp: Utc::now(),
568            }),
569            WorkflowEvent::StepCompleted(WorkflowStepCompletedEvent {
570                step_name: "deploy".to_string(),
571                step_index: 1,
572                duration_ms: 5000,
573                output_summary: Some("deployed v1.2.3".to_string()),
574            }),
575            WorkflowEvent::StepFailed(WorkflowStepFailedEvent {
576                step_name: "test".to_string(),
577                step_index: 2,
578                error: "exit code 1".to_string(),
579                duration_ms: 3000,
580            }),
581            WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
582                step_name: "prod-gate".to_string(),
583                step_index: 3,
584                approval_id: Uuid::now_v7(),
585            }),
586            WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
587                run_id: Uuid::now_v7(),
588                step_id: Uuid::now_v7(),
589                step_name: "clarify".to_string(),
590                step_index: 4,
591                message: "Answer the questions".to_string(),
592                schema: json!({"type": "object"}),
593            }),
594            WorkflowEvent::AgentStepTokensUsed(WorkflowAgentStepTokensUsedEvent {
595                step_name: "review".to_string(),
596                tokens: 15000,
597                cost_usd: Decimal::new(42, 4),
598            }),
599            WorkflowEvent::AgentStepResumed(WorkflowAgentStepResumedEvent {
600                step_name: "review".to_string(),
601                step_index: 5,
602                session_id: "0192f0c1-7d2e-7a4b-9c3d-1e2f3a4b5c6d".to_string(),
603                timestamp: Utc::now(),
604            }),
605        ];
606
607        for event in &cases {
608            let json = serde_json::to_string(event).expect("serialize");
609            let back: WorkflowEvent = serde_json::from_str(&json).expect("deserialize");
610
611            assert_eq!(back.event_type(), event.event_type());
612            assert!(json.contains(&format!("\"type\":\"{}\"", event.event_type())));
613        }
614    }
615
616    /// The pre-refactor wire format used flat inline-struct variants. Newtype
617    /// variants produce and accept the same JSON, so SSE consumers and stored
618    /// payloads need no migration.
619    #[test]
620    fn workflow_event_legacy_flat_json_deserializes() {
621        let approval_id: Uuid = "01890000-0000-7000-8000-000000000002"
622            .parse()
623            .expect("valid uuid");
624
625        let raw = r#"{"type":"step_started","step_name":"build","step_index":0,"timestamp":"2026-01-01T00:00:00Z"}"#;
626        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
627            WorkflowEvent::StepStarted(e) => {
628                assert_eq!(e.step_name, "build");
629                assert_eq!(e.step_index, 0);
630            }
631            other => panic!("expected StepStarted, got {other:?}"),
632        }
633
634        let raw = r#"{"type":"step_completed","step_name":"deploy","step_index":1,"duration_ms":5000,"output_summary":"deployed v1.2.3"}"#;
635        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
636            WorkflowEvent::StepCompleted(e) => {
637                assert_eq!(e.duration_ms, 5000);
638                assert_eq!(e.output_summary.as_deref(), Some("deployed v1.2.3"));
639            }
640            other => panic!("expected StepCompleted, got {other:?}"),
641        }
642
643        let raw = r#"{"type":"step_failed","step_name":"test","step_index":2,"error":"exit code 1","duration_ms":3000}"#;
644        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
645            WorkflowEvent::StepFailed(e) => {
646                assert_eq!(e.error, "exit code 1");
647                assert_eq!(e.duration_ms, 3000);
648            }
649            other => panic!("expected StepFailed, got {other:?}"),
650        }
651
652        let raw = r#"{"type":"approval_required","step_name":"prod-gate","step_index":3,"approval_id":"01890000-0000-7000-8000-000000000002"}"#;
653        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
654            WorkflowEvent::ApprovalRequired(e) => {
655                assert_eq!(e.approval_id, approval_id);
656            }
657            other => panic!("expected ApprovalRequired, got {other:?}"),
658        }
659
660        let raw = r#"{"type":"agent_step_tokens_used","step_name":"review","tokens":15000,"cost_usd":0.5}"#;
661        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
662            WorkflowEvent::AgentStepTokensUsed(e) => {
663                assert_eq!(e.tokens, 15000);
664                assert_eq!(e.cost_usd, Decimal::new(5, 1));
665            }
666            other => panic!("expected AgentStepTokensUsed, got {other:?}"),
667        }
668    }
669
670    /// Guards the internally-tagged representation: payload fields must stay
671    /// siblings of `type`, never nested under a variant key.
672    #[test]
673    fn serialized_workflow_event_is_flat_with_type_tag() {
674        let event = WorkflowEvent::StepFailed(WorkflowStepFailedEvent {
675            step_name: "test".to_string(),
676            step_index: 2,
677            error: "exit code 1".to_string(),
678            duration_ms: 3000,
679        });
680
681        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
682        let object = value.as_object().expect("event serializes to an object");
683
684        assert_eq!(
685            object.get("type").and_then(|v| v.as_str()),
686            Some("step_failed")
687        );
688        assert_eq!(
689            object.get("step_name").and_then(|v| v.as_str()),
690            Some("test")
691        );
692        assert_eq!(object.get("step_index").and_then(|v| v.as_u64()), Some(2));
693        assert_eq!(
694            object.get("error").and_then(|v| v.as_str()),
695            Some("exit code 1")
696        );
697        assert_eq!(
698            object.get("duration_ms").and_then(|v| v.as_u64()),
699            Some(3000)
700        );
701        assert_eq!(object.len(), 5, "no nesting: {object:?}");
702    }
703
704    #[test]
705    fn input_required_event_serializes_flat_with_its_schema() {
706        let step_id = Uuid::now_v7();
707        let event = WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
708            run_id: Uuid::now_v7(),
709            step_id,
710            step_name: "clarify".to_string(),
711            step_index: 1,
712            message: "Answer the questions".to_string(),
713            schema: json!({"type": "object", "required": ["answers"]}),
714        });
715
716        let value = serde_json::to_value(&event).expect("serialize");
717        assert_eq!(value["type"], "input_required");
718        assert_eq!(value["step_id"], step_id.to_string());
719        assert_eq!(value["message"], "Answer the questions");
720        assert_eq!(value["schema"]["required"][0], "answers");
721        assert_eq!(event.event_type(), WorkflowEvent::INPUT_REQUIRED);
722    }
723
724    #[test]
725    fn event_type_all_variants() {
726        let cases: Vec<(WorkflowEvent, &str)> = vec![
727            (
728                WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
729                    step_name: "s".to_string(),
730                    step_index: 0,
731                    timestamp: Utc::now(),
732                }),
733                "step_started",
734            ),
735            (
736                WorkflowEvent::StepCompleted(WorkflowStepCompletedEvent {
737                    step_name: "s".to_string(),
738                    step_index: 0,
739                    duration_ms: 0,
740                    output_summary: None,
741                }),
742                "step_completed",
743            ),
744            (
745                WorkflowEvent::StepFailed(WorkflowStepFailedEvent {
746                    step_name: "s".to_string(),
747                    step_index: 0,
748                    error: "e".to_string(),
749                    duration_ms: 0,
750                }),
751                "step_failed",
752            ),
753            (
754                WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
755                    step_name: "s".to_string(),
756                    step_index: 0,
757                    approval_id: Uuid::now_v7(),
758                }),
759                "approval_required",
760            ),
761            (
762                WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
763                    run_id: Uuid::now_v7(),
764                    step_id: Uuid::now_v7(),
765                    step_name: "s".to_string(),
766                    step_index: 0,
767                    message: "m".to_string(),
768                    schema: Value::Null,
769                }),
770                "input_required",
771            ),
772            (
773                WorkflowEvent::AgentStepTokensUsed(WorkflowAgentStepTokensUsedEvent {
774                    step_name: "s".to_string(),
775                    tokens: 0,
776                    cost_usd: Decimal::ZERO,
777                }),
778                "agent_step_tokens_used",
779            ),
780            (
781                WorkflowEvent::AgentStepResumed(WorkflowAgentStepResumedEvent {
782                    step_name: "s".to_string(),
783                    step_index: 0,
784                    session_id: "sid".to_string(),
785                    timestamp: Utc::now(),
786                }),
787                "agent_step_resumed",
788            ),
789        ];
790
791        for (event, expected) in cases {
792            assert_eq!(event.event_type(), expected);
793        }
794    }
795
796    #[tokio::test]
797    async fn multiple_subscribers_receive_same_event() {
798        let bus = WorkflowEventBus::new();
799        let run_id = Uuid::now_v7();
800
801        let mut rx1 = bus.subscribe(run_id);
802        let mut rx2 = bus.subscribe(run_id);
803
804        bus.publish(run_id, step_started("build"));
805
806        let e1 = rx1.recv().await.expect("rx1 should receive");
807        let e2 = rx2.recv().await.expect("rx2 should receive");
808
809        assert_eq!(e1.event_type(), "step_started");
810        assert_eq!(e2.event_type(), "step_started");
811    }
812
813    #[tokio::test]
814    async fn events_isolated_between_runs() {
815        let bus = WorkflowEventBus::new();
816        let run_a = Uuid::now_v7();
817        let run_b = Uuid::now_v7();
818
819        let mut rx_a = bus.subscribe(run_a);
820        let mut rx_b = bus.subscribe(run_b);
821
822        bus.publish(run_a, step_started("only-for-a"));
823
824        let received = rx_a.recv().await.expect("rx_a should receive");
825        match received {
826            WorkflowEvent::StepStarted(e) => {
827                assert_eq!(e.step_name, "only-for-a");
828            }
829            _ => panic!("expected StepStarted"),
830        }
831
832        // rx_b should have nothing -- try_recv returns Empty.
833        assert!(rx_b.try_recv().is_err());
834    }
835}