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).
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/// A granular step-level event for real-time workflow monitoring.
218///
219/// Unlike [`Event`](super::Event) which covers the full system lifecycle
220/// (runs, auth, audit), `WorkflowEvent` tracks individual step transitions
221/// within a single run. Serialized with a `type` discriminant for UI
222/// consumption.
223///
224/// Each variant wraps a dedicated payload struct; the serialized form stays
225/// flat, with `type` sitting next to the payload fields.
226///
227/// # Examples
228///
229/// ```
230/// use ironflow_engine::notify::{WorkflowEvent, WorkflowStepStartedEvent};
231/// use chrono::Utc;
232///
233/// let event = WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
234///     step_name: "deploy".to_string(),
235///     step_index: 0,
236///     timestamp: Utc::now(),
237/// });
238/// assert_eq!(event.event_type(), "step_started");
239///
240/// let json = serde_json::to_string(&event)?;
241/// assert!(json.contains("\"type\":\"step_started\""));
242/// # Ok::<(), serde_json::Error>(())
243/// ```
244#[derive(Debug, Clone, Serialize, Deserialize)]
245#[cfg_attr(feature = "openapi", derive(utoipa::ToSchema))]
246#[serde(tag = "type", rename_all = "snake_case")]
247pub enum WorkflowEvent {
248    /// A step began execution.
249    StepStarted(WorkflowStepStartedEvent),
250
251    /// A step completed successfully.
252    StepCompleted(WorkflowStepCompletedEvent),
253
254    /// A step failed.
255    StepFailed(WorkflowStepFailedEvent),
256
257    /// A step requires human approval before the run can continue.
258    ApprovalRequired(WorkflowApprovalRequiredEvent),
259
260    /// A step waits for a typed human input before the run can continue.
261    InputRequired(WorkflowInputRequiredEvent),
262
263    /// Token usage report for an agent step.
264    AgentStepTokensUsed(WorkflowAgentStepTokensUsedEvent),
265}
266
267impl WorkflowEvent {
268    /// Event type constant for [`StepStarted`](WorkflowEvent::StepStarted).
269    pub const STEP_STARTED: &'static str = "step_started";
270    /// Event type constant for [`StepCompleted`](WorkflowEvent::StepCompleted).
271    pub const STEP_COMPLETED: &'static str = "step_completed";
272    /// Event type constant for [`StepFailed`](WorkflowEvent::StepFailed).
273    pub const STEP_FAILED: &'static str = "step_failed";
274    /// Event type constant for [`ApprovalRequired`](WorkflowEvent::ApprovalRequired).
275    pub const APPROVAL_REQUIRED: &'static str = "approval_required";
276    /// Event type constant for [`InputRequired`](WorkflowEvent::InputRequired).
277    pub const INPUT_REQUIRED: &'static str = "input_required";
278    /// Event type constant for [`AgentStepTokensUsed`](WorkflowEvent::AgentStepTokensUsed).
279    pub const AGENT_STEP_TOKENS_USED: &'static str = "agent_step_tokens_used";
280
281    /// Returns the event type as a static string (e.g. `"step_started"`).
282    ///
283    /// # Examples
284    ///
285    /// ```
286    /// use ironflow_engine::notify::{WorkflowEvent, WorkflowStepStartedEvent};
287    /// use chrono::Utc;
288    ///
289    /// let event = WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
290    ///     step_name: "build".to_string(),
291    ///     step_index: 0,
292    ///     timestamp: Utc::now(),
293    /// });
294    /// assert_eq!(event.event_type(), "step_started");
295    /// ```
296    #[deny(unreachable_patterns)]
297    pub fn event_type(&self) -> &'static str {
298        match self {
299            WorkflowEvent::StepStarted(_) => Self::STEP_STARTED,
300            WorkflowEvent::StepCompleted(_) => Self::STEP_COMPLETED,
301            WorkflowEvent::StepFailed(_) => Self::STEP_FAILED,
302            WorkflowEvent::ApprovalRequired(_) => Self::APPROVAL_REQUIRED,
303            WorkflowEvent::InputRequired(_) => Self::INPUT_REQUIRED,
304            WorkflowEvent::AgentStepTokensUsed(_) => Self::AGENT_STEP_TOKENS_USED,
305        }
306    }
307}
308
309/// Per-workflow broadcast event bus for real-time monitoring.
310///
311/// Maintains one [`tokio::sync::broadcast`] channel per workflow run.
312/// Consumers call [`subscribe`](Self::subscribe) to receive events for a
313/// specific run; producers call [`publish`](Self::publish) to broadcast
314/// an event to all subscribers of that run.
315///
316/// Thread-safe and cheaply cloneable (`Clone` shares the same inner state).
317///
318/// # Examples
319///
320/// ```
321/// use ironflow_engine::notify::{WorkflowEvent, WorkflowEventBus, WorkflowStepStartedEvent};
322/// use uuid::Uuid;
323/// use chrono::Utc;
324///
325/// let bus = WorkflowEventBus::new();
326/// let run_id = Uuid::now_v7();
327///
328/// let mut rx = bus.subscribe(run_id);
329/// bus.publish(run_id, WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
330///     step_name: "build".to_string(),
331///     step_index: 0,
332///     timestamp: Utc::now(),
333/// }));
334/// ```
335#[derive(Clone)]
336pub struct WorkflowEventBus {
337    channels: std::sync::Arc<RwLock<HashMap<Uuid, broadcast::Sender<WorkflowEvent>>>>,
338}
339
340impl WorkflowEventBus {
341    /// Create a new empty event bus.
342    ///
343    /// # Examples
344    ///
345    /// ```
346    /// use ironflow_engine::notify::WorkflowEventBus;
347    ///
348    /// let bus = WorkflowEventBus::new();
349    /// ```
350    pub fn new() -> Self {
351        Self {
352            channels: std::sync::Arc::new(RwLock::new(HashMap::new())),
353        }
354    }
355
356    /// Subscribe to events for a specific workflow run.
357    ///
358    /// If no channel exists for this `run_id`, one is created on demand.
359    /// Returns a broadcast receiver that yields [`WorkflowEvent`]s for
360    /// that run only.
361    ///
362    /// # Examples
363    ///
364    /// ```
365    /// use ironflow_engine::notify::WorkflowEventBus;
366    /// use uuid::Uuid;
367    ///
368    /// let bus = WorkflowEventBus::new();
369    /// let run_id = Uuid::now_v7();
370    /// let _rx = bus.subscribe(run_id);
371    /// ```
372    pub fn subscribe(&self, run_id: Uuid) -> broadcast::Receiver<WorkflowEvent> {
373        let mut channels = self.channels.write().expect("event bus lock poisoned");
374        let sender = channels
375            .entry(run_id)
376            .or_insert_with(|| broadcast::channel(DEFAULT_BUFFER_SIZE).0);
377        sender.subscribe()
378    }
379
380    /// Broadcast an event to all subscribers of a specific workflow run.
381    ///
382    /// If no channel exists for `run_id` (no subscriber has called
383    /// [`subscribe`](Self::subscribe)), the event is silently dropped.
384    /// If subscribers exist but none are actively listening, the send
385    /// error is ignored.
386    ///
387    /// # Examples
388    ///
389    /// ```
390    /// use ironflow_engine::notify::{WorkflowEvent, WorkflowEventBus, WorkflowStepStartedEvent};
391    /// use uuid::Uuid;
392    /// use chrono::Utc;
393    ///
394    /// let bus = WorkflowEventBus::new();
395    /// let run_id = Uuid::now_v7();
396    ///
397    /// // No subscriber -- silently dropped.
398    /// bus.publish(run_id, WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
399    ///     step_name: "build".to_string(),
400    ///     step_index: 0,
401    ///     timestamp: Utc::now(),
402    /// }));
403    /// ```
404    pub fn publish(&self, run_id: Uuid, event: WorkflowEvent) {
405        let channels = self.channels.read().expect("event bus lock poisoned");
406        if let Some(sender) = channels.get(&run_id) {
407            let _ = sender.send(event);
408        }
409    }
410
411    /// Remove the channel for a workflow run.
412    ///
413    /// Call this when a run completes or is cleaned up to free resources.
414    /// If no channel exists for `run_id`, this is a no-op.
415    ///
416    /// # Examples
417    ///
418    /// ```
419    /// use ironflow_engine::notify::WorkflowEventBus;
420    /// use uuid::Uuid;
421    ///
422    /// let bus = WorkflowEventBus::new();
423    /// let run_id = Uuid::now_v7();
424    /// let _rx = bus.subscribe(run_id);
425    /// bus.remove(run_id);
426    /// ```
427    pub fn remove(&self, run_id: Uuid) {
428        let mut channels = self.channels.write().expect("event bus lock poisoned");
429        channels.remove(&run_id);
430    }
431}
432
433impl Default for WorkflowEventBus {
434    fn default() -> Self {
435        Self::new()
436    }
437}
438
439impl std::fmt::Debug for WorkflowEventBus {
440    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
441        let count = self.channels.read().map(|c| c.len()).unwrap_or(0);
442        f.debug_struct("WorkflowEventBus")
443            .field("active_channels", &count)
444            .finish()
445    }
446}
447
448#[cfg(test)]
449mod tests {
450    use serde_json::json;
451
452    use super::*;
453
454    fn step_started(step_name: &str) -> WorkflowEvent {
455        WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
456            step_name: step_name.to_string(),
457            step_index: 0,
458            timestamp: Utc::now(),
459        })
460    }
461
462    #[tokio::test]
463    async fn subscribe_receives_published_events() {
464        let bus = WorkflowEventBus::new();
465        let run_id = Uuid::now_v7();
466
467        let mut rx = bus.subscribe(run_id);
468
469        bus.publish(run_id, step_started("build"));
470
471        let received = rx.recv().await.expect("should receive event");
472        assert_eq!(received.event_type(), "step_started");
473        match received {
474            WorkflowEvent::StepStarted(e) => {
475                assert_eq!(e.step_name, "build");
476                assert_eq!(e.step_index, 0);
477            }
478            _ => panic!("expected StepStarted"),
479        }
480    }
481
482    #[test]
483    fn subscribe_creates_channel_on_demand() {
484        let bus = WorkflowEventBus::new();
485        let run_id = Uuid::now_v7();
486
487        let count_before = bus.channels.read().unwrap().len();
488        assert_eq!(count_before, 0);
489
490        let _rx = bus.subscribe(run_id);
491
492        let count_after = bus.channels.read().unwrap().len();
493        assert_eq!(count_after, 1);
494    }
495
496    #[test]
497    fn publish_unknown_run_is_noop() {
498        let bus = WorkflowEventBus::new();
499        let unknown_run = Uuid::now_v7();
500
501        bus.publish(unknown_run, step_started("build"));
502    }
503
504    #[test]
505    fn remove_cleans_up_channel() {
506        let bus = WorkflowEventBus::new();
507        let run_id = Uuid::now_v7();
508
509        let _rx = bus.subscribe(run_id);
510        assert_eq!(bus.channels.read().unwrap().len(), 1);
511
512        bus.remove(run_id);
513        assert_eq!(bus.channels.read().unwrap().len(), 0);
514    }
515
516    #[test]
517    fn remove_unknown_is_noop() {
518        let bus = WorkflowEventBus::new();
519        bus.remove(Uuid::now_v7());
520    }
521
522    #[test]
523    fn workflow_event_serde_roundtrip() {
524        let cases: Vec<WorkflowEvent> = vec![
525            WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
526                step_name: "build".to_string(),
527                step_index: 0,
528                timestamp: Utc::now(),
529            }),
530            WorkflowEvent::StepCompleted(WorkflowStepCompletedEvent {
531                step_name: "deploy".to_string(),
532                step_index: 1,
533                duration_ms: 5000,
534                output_summary: Some("deployed v1.2.3".to_string()),
535            }),
536            WorkflowEvent::StepFailed(WorkflowStepFailedEvent {
537                step_name: "test".to_string(),
538                step_index: 2,
539                error: "exit code 1".to_string(),
540                duration_ms: 3000,
541            }),
542            WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
543                step_name: "prod-gate".to_string(),
544                step_index: 3,
545                approval_id: Uuid::now_v7(),
546            }),
547            WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
548                run_id: Uuid::now_v7(),
549                step_id: Uuid::now_v7(),
550                step_name: "clarify".to_string(),
551                step_index: 4,
552                message: "Answer the questions".to_string(),
553                schema: json!({"type": "object"}),
554            }),
555            WorkflowEvent::AgentStepTokensUsed(WorkflowAgentStepTokensUsedEvent {
556                step_name: "review".to_string(),
557                tokens: 15000,
558                cost_usd: Decimal::new(42, 4),
559            }),
560        ];
561
562        for event in &cases {
563            let json = serde_json::to_string(event).expect("serialize");
564            let back: WorkflowEvent = serde_json::from_str(&json).expect("deserialize");
565
566            assert_eq!(back.event_type(), event.event_type());
567            assert!(json.contains(&format!("\"type\":\"{}\"", event.event_type())));
568        }
569    }
570
571    /// The pre-refactor wire format used flat inline-struct variants. Newtype
572    /// variants produce and accept the same JSON, so SSE consumers and stored
573    /// payloads need no migration.
574    #[test]
575    fn workflow_event_legacy_flat_json_deserializes() {
576        let approval_id: Uuid = "01890000-0000-7000-8000-000000000002"
577            .parse()
578            .expect("valid uuid");
579
580        let raw = r#"{"type":"step_started","step_name":"build","step_index":0,"timestamp":"2026-01-01T00:00:00Z"}"#;
581        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
582            WorkflowEvent::StepStarted(e) => {
583                assert_eq!(e.step_name, "build");
584                assert_eq!(e.step_index, 0);
585            }
586            other => panic!("expected StepStarted, got {other:?}"),
587        }
588
589        let raw = r#"{"type":"step_completed","step_name":"deploy","step_index":1,"duration_ms":5000,"output_summary":"deployed v1.2.3"}"#;
590        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
591            WorkflowEvent::StepCompleted(e) => {
592                assert_eq!(e.duration_ms, 5000);
593                assert_eq!(e.output_summary.as_deref(), Some("deployed v1.2.3"));
594            }
595            other => panic!("expected StepCompleted, got {other:?}"),
596        }
597
598        let raw = r#"{"type":"step_failed","step_name":"test","step_index":2,"error":"exit code 1","duration_ms":3000}"#;
599        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
600            WorkflowEvent::StepFailed(e) => {
601                assert_eq!(e.error, "exit code 1");
602                assert_eq!(e.duration_ms, 3000);
603            }
604            other => panic!("expected StepFailed, got {other:?}"),
605        }
606
607        let raw = r#"{"type":"approval_required","step_name":"prod-gate","step_index":3,"approval_id":"01890000-0000-7000-8000-000000000002"}"#;
608        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
609            WorkflowEvent::ApprovalRequired(e) => {
610                assert_eq!(e.approval_id, approval_id);
611            }
612            other => panic!("expected ApprovalRequired, got {other:?}"),
613        }
614
615        let raw = r#"{"type":"agent_step_tokens_used","step_name":"review","tokens":15000,"cost_usd":0.5}"#;
616        match serde_json::from_str::<WorkflowEvent>(raw).expect("legacy payload") {
617            WorkflowEvent::AgentStepTokensUsed(e) => {
618                assert_eq!(e.tokens, 15000);
619                assert_eq!(e.cost_usd, Decimal::new(5, 1));
620            }
621            other => panic!("expected AgentStepTokensUsed, got {other:?}"),
622        }
623    }
624
625    /// Guards the internally-tagged representation: payload fields must stay
626    /// siblings of `type`, never nested under a variant key.
627    #[test]
628    fn serialized_workflow_event_is_flat_with_type_tag() {
629        let event = WorkflowEvent::StepFailed(WorkflowStepFailedEvent {
630            step_name: "test".to_string(),
631            step_index: 2,
632            error: "exit code 1".to_string(),
633            duration_ms: 3000,
634        });
635
636        let value: serde_json::Value = serde_json::to_value(&event).expect("serialize");
637        let object = value.as_object().expect("event serializes to an object");
638
639        assert_eq!(
640            object.get("type").and_then(|v| v.as_str()),
641            Some("step_failed")
642        );
643        assert_eq!(
644            object.get("step_name").and_then(|v| v.as_str()),
645            Some("test")
646        );
647        assert_eq!(object.get("step_index").and_then(|v| v.as_u64()), Some(2));
648        assert_eq!(
649            object.get("error").and_then(|v| v.as_str()),
650            Some("exit code 1")
651        );
652        assert_eq!(
653            object.get("duration_ms").and_then(|v| v.as_u64()),
654            Some(3000)
655        );
656        assert_eq!(object.len(), 5, "no nesting: {object:?}");
657    }
658
659    #[test]
660    fn input_required_event_serializes_flat_with_its_schema() {
661        let step_id = Uuid::now_v7();
662        let event = WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
663            run_id: Uuid::now_v7(),
664            step_id,
665            step_name: "clarify".to_string(),
666            step_index: 1,
667            message: "Answer the questions".to_string(),
668            schema: json!({"type": "object", "required": ["answers"]}),
669        });
670
671        let value = serde_json::to_value(&event).expect("serialize");
672        assert_eq!(value["type"], "input_required");
673        assert_eq!(value["step_id"], step_id.to_string());
674        assert_eq!(value["message"], "Answer the questions");
675        assert_eq!(value["schema"]["required"][0], "answers");
676        assert_eq!(event.event_type(), WorkflowEvent::INPUT_REQUIRED);
677    }
678
679    #[test]
680    fn event_type_all_variants() {
681        let cases: Vec<(WorkflowEvent, &str)> = vec![
682            (
683                WorkflowEvent::StepStarted(WorkflowStepStartedEvent {
684                    step_name: "s".to_string(),
685                    step_index: 0,
686                    timestamp: Utc::now(),
687                }),
688                "step_started",
689            ),
690            (
691                WorkflowEvent::StepCompleted(WorkflowStepCompletedEvent {
692                    step_name: "s".to_string(),
693                    step_index: 0,
694                    duration_ms: 0,
695                    output_summary: None,
696                }),
697                "step_completed",
698            ),
699            (
700                WorkflowEvent::StepFailed(WorkflowStepFailedEvent {
701                    step_name: "s".to_string(),
702                    step_index: 0,
703                    error: "e".to_string(),
704                    duration_ms: 0,
705                }),
706                "step_failed",
707            ),
708            (
709                WorkflowEvent::ApprovalRequired(WorkflowApprovalRequiredEvent {
710                    step_name: "s".to_string(),
711                    step_index: 0,
712                    approval_id: Uuid::now_v7(),
713                }),
714                "approval_required",
715            ),
716            (
717                WorkflowEvent::InputRequired(WorkflowInputRequiredEvent {
718                    run_id: Uuid::now_v7(),
719                    step_id: Uuid::now_v7(),
720                    step_name: "s".to_string(),
721                    step_index: 0,
722                    message: "m".to_string(),
723                    schema: Value::Null,
724                }),
725                "input_required",
726            ),
727            (
728                WorkflowEvent::AgentStepTokensUsed(WorkflowAgentStepTokensUsedEvent {
729                    step_name: "s".to_string(),
730                    tokens: 0,
731                    cost_usd: Decimal::ZERO,
732                }),
733                "agent_step_tokens_used",
734            ),
735        ];
736
737        for (event, expected) in cases {
738            assert_eq!(event.event_type(), expected);
739        }
740    }
741
742    #[tokio::test]
743    async fn multiple_subscribers_receive_same_event() {
744        let bus = WorkflowEventBus::new();
745        let run_id = Uuid::now_v7();
746
747        let mut rx1 = bus.subscribe(run_id);
748        let mut rx2 = bus.subscribe(run_id);
749
750        bus.publish(run_id, step_started("build"));
751
752        let e1 = rx1.recv().await.expect("rx1 should receive");
753        let e2 = rx2.recv().await.expect("rx2 should receive");
754
755        assert_eq!(e1.event_type(), "step_started");
756        assert_eq!(e2.event_type(), "step_started");
757    }
758
759    #[tokio::test]
760    async fn events_isolated_between_runs() {
761        let bus = WorkflowEventBus::new();
762        let run_a = Uuid::now_v7();
763        let run_b = Uuid::now_v7();
764
765        let mut rx_a = bus.subscribe(run_a);
766        let mut rx_b = bus.subscribe(run_b);
767
768        bus.publish(run_a, step_started("only-for-a"));
769
770        let received = rx_a.recv().await.expect("rx_a should receive");
771        match received {
772            WorkflowEvent::StepStarted(e) => {
773                assert_eq!(e.step_name, "only-for-a");
774            }
775            _ => panic!("expected StepStarted"),
776        }
777
778        // rx_b should have nothing -- try_recv returns Empty.
779        assert!(rx_b.try_recv().is_err());
780    }
781}