Skip to main content

luft_core/contract/
event.rs

1//! Event bus contract (§1.4) — the single observability data source.
2//! Each [`AgentEvent`] is consumed by the event bus subscribers; the state store persists it.
3
4use crate::contract::backend::AgentStatus;
5use crate::contract::finding::Finding;
6use crate::contract::ids::{AgentId, PhaseId, RunId, TokenUsage};
7use chrono::{DateTime, Utc};
8use serde::{Deserialize, Serialize};
9use std::path::PathBuf;
10
11/// Broadcast sender shared by the scheduler and every event producer.
12pub type EventSender = tokio::sync::broadcast::Sender<AgentEvent>;
13
14#[derive(Debug, Clone, Serialize, Deserialize)]
15#[serde(tag = "type", rename_all = "snake_case")]
16pub enum AgentEvent {
17    RunStarted {
18        run_id: RunId,
19        task: String,
20        ts: DateTime<Utc>,
21    },
22    PhaseStarted {
23        run_id: RunId,
24        phase_id: PhaseId,
25        label: String,
26        planned: usize,
27        #[serde(default)]
28        parent_span_id: Option<u32>,
29        #[serde(default)]
30        description: Option<String>,
31        #[serde(default)]
32        role: Option<String>,
33        #[serde(default)]
34        ts: DateTime<Utc>,
35    },
36    AgentStarted {
37        run_id: RunId,
38        phase_id: PhaseId,
39        agent_id: AgentId,
40        prompt_preview: String,
41        model: Option<String>,
42        #[serde(default)]
43        description: Option<String>,
44        #[serde(default)]
45        role: Option<String>,
46        #[serde(default)]
47        name: Option<String>,
48        #[serde(default)]
49        agent_seq: u32,
50        #[serde(default)]
51        ts: DateTime<Utc>,
52    },
53    AgentProgress {
54        run_id: RunId,
55        agent_id: AgentId,
56        delta: ProgressDelta,
57    },
58    /// Raw ACP `session/update` passthrough — the verbatim notification, surfaced
59    /// for observability. Produced only when the ACP backend has raw events
60    /// enabled. Not persisted to
61    /// the journal (see `acp-raw-events.md`).
62    AcpRaw {
63        run_id: RunId,
64        agent_id: AgentId,
65        /// `SessionUpdate` discriminator (the `sessionUpdate` tag), e.g.
66        /// `"agent_message_chunk"`, `"plan"` — lets consumers filter without
67        /// parsing `raw`.
68        kind: String,
69        /// The ACP `SessionUpdate`, serialized verbatim.
70        raw: serde_json::Value,
71    },
72    AgentDone {
73        run_id: RunId,
74        agent_id: AgentId,
75        status: AgentStatus,
76        tokens: TokenUsage,
77        elapsed_ms: u64,
78        #[serde(default)]
79        name: Option<String>,
80        #[serde(default)]
81        agent_seq: u32,
82        #[serde(default)]
83        ts: DateTime<Utc>,
84        #[serde(default)]
85        output: serde_json::Value,
86        #[serde(default)]
87        findings: Vec<Finding>,
88        #[serde(default)]
89        prompt: String,
90        #[serde(default)]
91        retry_count: u32,
92    },
93    PhaseDone {
94        run_id: RunId,
95        phase_id: PhaseId,
96        ok: usize,
97        failed: usize,
98        #[serde(default)]
99        ts: DateTime<Utc>,
100    },
101    RunDone {
102        run_id: RunId,
103        status: RunStatus,
104        total_tokens: TokenUsage,
105        report: serde_json::Value,
106        #[serde(default)]
107        ts: DateTime<Utc>,
108    },
109    Log {
110        run_id: RunId,
111        agent_id: Option<AgentId>,
112        level: LogLevel,
113        msg: String,
114    },
115    // SDK primitive events (§ sdk-events.md) — DSL-granularity observability for
116    // the orchestration script. Blocking primitives emit a Started/Done span
117    // pair correlated by `span_id`; instantaneous ones emit a single event.
118    BudgetSet {
119        run_id: RunId,
120        time_limit_ms: Option<u64>,
121        max_rounds: Option<u32>,
122    },
123    ReportEmitted {
124        run_id: RunId,
125        phase_id: PhaseId,
126        report: serde_json::Value,
127    },
128    ParallelStarted {
129        run_id: RunId,
130        phase_id: PhaseId,
131        span_id: u64,
132        count: usize,
133    },
134    ParallelDone {
135        run_id: RunId,
136        phase_id: PhaseId,
137        span_id: u64,
138        ok: usize,
139        failed: usize,
140        results: serde_json::Value,
141        elapsed_ms: u64,
142    },
143    WorkflowStarted {
144        run_id: RunId,
145        span_id: u64,
146        path: String,
147        args: serde_json::Value,
148    },
149    WorkflowDone {
150        run_id: RunId,
151        span_id: u64,
152        path: String,
153        report: serde_json::Value,
154        elapsed_ms: u64,
155        error: Option<String>,
156    },
157    ConvergeStarted {
158        run_id: RunId,
159        phase_id: PhaseId,
160        span_id: u64,
161        items: usize,
162        max_rounds: u32,
163    },
164    ConvergeDone {
165        run_id: RunId,
166        phase_id: PhaseId,
167        span_id: u64,
168        rounds: u32,
169        converged: bool,
170        surviving: usize,
171        result: serde_json::Value,
172        elapsed_ms: u64,
173        error: Option<String>,
174    },
175    // M2 Pipeline events
176    PipelineStarted {
177        run_id: RunId,
178        total_stages: usize,
179        items: usize,
180    },
181    PipelineStageStarted {
182        run_id: RunId,
183        stage_index: usize,
184        label: String,
185        agents_in_stage: usize,
186    },
187    PipelineItemDone {
188        run_id: RunId,
189        stage_index: usize,
190        item_index: usize,
191        status: AgentStatus,
192        tokens: TokenUsage,
193        elapsed_ms: u64,
194    },
195    PipelineDone {
196        run_id: RunId,
197        stages_completed: usize,
198        total_ok: usize,
199        total_failed: usize,
200    },
201    /// Structural phase span started — emitted by `phase_begin()`.
202    PhaseSpanStarted {
203        run_id: RunId,
204        span_id: u32,
205        name: String,
206        parent_id: Option<u32>,
207        depth: u32,
208        planned: usize,
209    },
210    /// Structural phase span done — emitted by `phase_end()`.
211    PhaseSpanDone {
212        run_id: RunId,
213        span_id: u32,
214        name: String,
215        parent_id: Option<u32>,
216        depth: u32,
217        elapsed_ms: u64,
218        status: String,
219    },
220    /// Agent output failed schema validation and is being retried with corrective
221    /// feedback injected into the prompt. Consumers (CLI, event log) can use this
222    /// to inform users that an extra LLM round-trip is underway.
223    SchemaRetry {
224        run_id: RunId,
225        agent_id: AgentId,
226        attempt: u32,
227        max: u32,
228    },
229    /// Plan preview — emitted before execution starts, from the `meta` table.
230    /// Lists the declared phases so the CLI can render a plan overview before
231    /// real-time execution output begins.
232    PlanPreview {
233        run_id: RunId,
234        reasoning: String,
235        phases: Vec<PlanPhase>,
236    },
237    /// OS signal received (SIGINT / SIGTERM / Ctrl+C). `run_id` is `None` if
238    /// the signal arrived before a run had started. Emitted by the process-
239    /// level signal handler in [`crate::signal`].
240    SignalReceived {
241        run_id: Option<RunId>,
242        signal: String,
243        ts: DateTime<Utc>,
244    },
245}
246
247/// A single phase entry in the plan preview `meta.phases` array.
248#[derive(Debug, Clone, Serialize, Deserialize)]
249pub struct PlanPhase {
250    pub label: String,
251    #[serde(default)]
252    pub dynamic: bool,
253    #[serde(default)]
254    pub description: Option<String>,
255}
256
257#[derive(Debug, Clone, Serialize, Deserialize)]
258#[serde(tag = "kind", rename_all = "snake_case")]
259pub enum ProgressDelta {
260    Message { text: String },
261    ToolCall { name: String, summary: String },
262    FileEdit { path: PathBuf },
263    Tokens { usage: TokenUsage },
264}
265
266#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
267#[serde(rename_all = "snake_case")]
268pub enum RunStatus {
269    Completed,
270    Failed,
271    Cancelled,
272    Partial,
273}
274
275#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
276#[serde(rename_all = "lowercase")]
277pub enum LogLevel {
278    Trace,
279    Debug,
280    Info,
281    Warn,
282    Error,
283}
284
285#[cfg(test)]
286mod tests {
287    use super::*;
288    use crate::contract::finding::{Finding, Location, Severity};
289    use chrono::TimeZone;
290    use serde_json::json;
291    use uuid::Uuid;
292
293    fn ts() -> DateTime<Utc> {
294        Utc.with_ymd_and_hms(2025, 1, 2, 3, 4, 5).unwrap()
295    }
296
297    fn run_id() -> RunId {
298        Uuid::nil()
299    }
300
301    fn agent_id() -> AgentId {
302        Uuid::nil()
303    }
304
305    // ── RunStatus ────────────────────────────────────────────────
306
307    #[test]
308    fn run_status_serializes_as_snake_case() {
309        assert_eq!(
310            serde_json::to_string(&RunStatus::Completed).unwrap(),
311            "\"completed\""
312        );
313        assert_eq!(
314            serde_json::to_string(&RunStatus::Failed).unwrap(),
315            "\"failed\""
316        );
317        assert_eq!(
318            serde_json::to_string(&RunStatus::Cancelled).unwrap(),
319            "\"cancelled\""
320        );
321        assert_eq!(
322            serde_json::to_string(&RunStatus::Partial).unwrap(),
323            "\"partial\""
324        );
325    }
326
327    #[test]
328    fn run_status_deserializes_from_snake_case() {
329        assert_eq!(
330            serde_json::from_str::<RunStatus>("\"completed\"").unwrap(),
331            RunStatus::Completed
332        );
333        assert_eq!(
334            serde_json::from_str::<RunStatus>("\"failed\"").unwrap(),
335            RunStatus::Failed
336        );
337        assert_eq!(
338            serde_json::from_str::<RunStatus>("\"cancelled\"").unwrap(),
339            RunStatus::Cancelled
340        );
341        assert_eq!(
342            serde_json::from_str::<RunStatus>("\"partial\"").unwrap(),
343            RunStatus::Partial
344        );
345    }
346
347    #[test]
348    fn run_status_equality_and_copy() {
349        let a = RunStatus::Completed;
350        let b = a; // Copy semantics
351        let c = a;
352        assert_eq!(a, b);
353        assert_eq!(a, c);
354    }
355
356    #[test]
357    fn run_status_unknown_variant_fails() {
358        let r: Result<RunStatus, _> = serde_json::from_str("\"unknown\"");
359        assert!(r.is_err());
360    }
361
362    // ── LogLevel ─────────────────────────────────────────────────
363
364    #[test]
365    fn log_level_serializes_as_lowercase() {
366        assert_eq!(
367            serde_json::to_string(&LogLevel::Trace).unwrap(),
368            "\"trace\""
369        );
370        assert_eq!(
371            serde_json::to_string(&LogLevel::Debug).unwrap(),
372            "\"debug\""
373        );
374        assert_eq!(serde_json::to_string(&LogLevel::Info).unwrap(), "\"info\"");
375        assert_eq!(serde_json::to_string(&LogLevel::Warn).unwrap(), "\"warn\"");
376        assert_eq!(
377            serde_json::to_string(&LogLevel::Error).unwrap(),
378            "\"error\""
379        );
380    }
381
382    #[test]
383    fn log_level_deserializes_from_lowercase() {
384        assert_eq!(
385            serde_json::from_str::<LogLevel>("\"warn\"").unwrap(),
386            LogLevel::Warn
387        );
388        assert_eq!(
389            serde_json::from_str::<LogLevel>("\"error\"").unwrap(),
390            LogLevel::Error
391        );
392    }
393
394    #[test]
395    fn log_level_equality_and_copy() {
396        let a = LogLevel::Info;
397        let b = a;
398        assert_eq!(a, b);
399    }
400
401    // ── ProgressDelta ────────────────────────────────────────────
402
403    #[test]
404    fn progress_delta_message_roundtrip() {
405        let d = ProgressDelta::Message {
406            text: "hello".into(),
407        };
408        let s = serde_json::to_string(&d).unwrap();
409        assert!(s.contains("\"kind\":\"message\""));
410        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
411        match back {
412            ProgressDelta::Message { text } => assert_eq!(text, "hello"),
413            _ => panic!("wrong variant"),
414        }
415    }
416
417    #[test]
418    fn progress_delta_tool_call_roundtrip() {
419        let d = ProgressDelta::ToolCall {
420            name: "bash".into(),
421            summary: "run ls".into(),
422        };
423        let s = serde_json::to_string(&d).unwrap();
424        assert!(s.contains("\"kind\":\"tool_call\""));
425        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
426        match back {
427            ProgressDelta::ToolCall { name, summary } => {
428                assert_eq!(name, "bash");
429                assert_eq!(summary, "run ls");
430            }
431            _ => panic!("wrong variant"),
432        }
433    }
434
435    #[test]
436    fn progress_delta_file_edit_roundtrip() {
437        let d = ProgressDelta::FileEdit {
438            path: PathBuf::from("src/main.rs"),
439        };
440        let s = serde_json::to_string(&d).unwrap();
441        assert!(s.contains("\"kind\":\"file_edit\""));
442        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
443        match back {
444            ProgressDelta::FileEdit { path } => assert_eq!(path, PathBuf::from("src/main.rs")),
445            _ => panic!("wrong variant"),
446        }
447    }
448
449    #[test]
450    fn progress_delta_tokens_roundtrip() {
451        let d = ProgressDelta::Tokens {
452            usage: TokenUsage {
453                input: 10,
454                output: 20,
455                cache_read: 1,
456                cache_write: 2,
457            },
458        };
459        let s = serde_json::to_string(&d).unwrap();
460        assert!(s.contains("\"kind\":\"tokens\""));
461        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
462        match back {
463            ProgressDelta::Tokens { usage } => {
464                assert_eq!(usage.input, 10);
465                assert_eq!(usage.output, 20);
466            }
467            _ => panic!("wrong variant"),
468        }
469    }
470
471    // ── PlanPhase ────────────────────────────────────────────────
472
473    #[test]
474    fn plan_phase_minimal_roundtrip() {
475        let raw = json!({"label": "review"});
476        let p: PlanPhase = serde_json::from_value(raw.clone()).unwrap();
477        assert_eq!(p.label, "review");
478        assert!(!p.dynamic);
479        assert!(p.description.is_none());
480        let back = serde_json::to_value(&p).unwrap();
481        assert_eq!(
482            back,
483            json!({"label": "review", "dynamic": false, "description": null})
484        );
485    }
486
487    #[test]
488    fn plan_phase_full_roundtrip() {
489        let p = PlanPhase {
490            label: "implement".into(),
491            dynamic: true,
492            description: Some("build the thing".into()),
493        };
494        let s = serde_json::to_string(&p).unwrap();
495        let back: PlanPhase = serde_json::from_str(&s).unwrap();
496        assert_eq!(back.label, "implement");
497        assert!(back.dynamic);
498        assert_eq!(back.description.as_deref(), Some("build the thing"));
499    }
500
501    // ── AgentEvent variants ──────────────────────────────────────
502
503    #[test]
504    fn event_run_started_roundtrip() {
505        let ev = AgentEvent::RunStarted {
506            run_id: run_id(),
507            task: "investigate".into(),
508            ts: ts(),
509        };
510        let s = serde_json::to_string(&ev).unwrap();
511        assert!(s.contains("\"type\":\"run_started\""));
512        let back: AgentEvent = serde_json::from_str(&s).unwrap();
513        match back {
514            AgentEvent::RunStarted { task, .. } => assert_eq!(task, "investigate"),
515            _ => panic!("wrong variant"),
516        }
517    }
518
519    #[test]
520    fn event_phase_started_minimal_roundtrip() {
521        let ev = AgentEvent::PhaseStarted {
522            run_id: run_id(),
523            phase_id: 3,
524            label: "scan".into(),
525            planned: 5,
526            parent_span_id: None,
527            description: None,
528            role: None,
529            ts: DateTime::<Utc>::from_timestamp(0, 0).unwrap(),
530        };
531        let s = serde_json::to_string(&ev).unwrap();
532        assert!(s.contains("\"type\":\"phase_started\""));
533        let back: AgentEvent = serde_json::from_str(&s).unwrap();
534        match back {
535            AgentEvent::PhaseStarted {
536                phase_id,
537                label,
538                planned,
539                ..
540            } => {
541                assert_eq!(phase_id, 3);
542                assert_eq!(label, "scan");
543                assert_eq!(planned, 5);
544            }
545            _ => panic!("wrong variant"),
546        }
547    }
548
549    #[test]
550    fn event_phase_started_omits_optional_fields() {
551        let ev = AgentEvent::PhaseStarted {
552            run_id: run_id(),
553            phase_id: 0,
554            label: "x".into(),
555            planned: 0,
556            parent_span_id: None,
557            description: None,
558            role: None,
559            ts: ts(),
560        };
561        let s = serde_json::to_string(&ev).unwrap();
562        // Optional fields with #[serde(default)] deserialize when omitted but
563        // serialize as JSON null when None. Verify they can be omitted on the
564        // input side by deserializing a payload that drops them.
565        let back: AgentEvent = serde_json::from_str(&s).unwrap();
566        match back {
567            AgentEvent::PhaseStarted {
568                parent_span_id,
569                description,
570                role,
571                ..
572            } => {
573                assert!(parent_span_id.is_none());
574                assert!(description.is_none());
575                assert!(role.is_none());
576            }
577            _ => panic!("wrong variant"),
578        }
579        assert!(s.contains("\"parent_span_id\":null"));
580        assert!(s.contains("\"description\":null"));
581        assert!(s.contains("\"role\":null"));
582    }
583
584    #[test]
585    fn event_agent_started_roundtrip() {
586        let ev = AgentEvent::AgentStarted {
587            run_id: run_id(),
588            phase_id: 1,
589            agent_id: agent_id(),
590            prompt_preview: "find bugs".into(),
591            model: Some("gpt-4".into()),
592            description: Some("audit".into()),
593            role: Some("reviewer".into()),
594            name: Some("agent-1".into()),
595            agent_seq: 7,
596            ts: Default::default(),
597        };
598        let s = serde_json::to_string(&ev).unwrap();
599        assert!(s.contains("\"type\":\"agent_started\""));
600        let back: AgentEvent = serde_json::from_str(&s).unwrap();
601        match back {
602            AgentEvent::AgentStarted {
603                agent_seq,
604                model,
605                name,
606                ..
607            } => {
608                assert_eq!(agent_seq, 7);
609                assert_eq!(model.as_deref(), Some("gpt-4"));
610                assert_eq!(name.as_deref(), Some("agent-1"));
611            }
612            _ => panic!("wrong variant"),
613        }
614    }
615
616    #[test]
617    fn event_agent_progress_roundtrip() {
618        let ev = AgentEvent::AgentProgress {
619            run_id: run_id(),
620            agent_id: agent_id(),
621            delta: ProgressDelta::Message { text: "ok".into() },
622        };
623        let s = serde_json::to_string(&ev).unwrap();
624        assert!(s.contains("\"type\":\"agent_progress\""));
625        let back: AgentEvent = serde_json::from_str(&s).unwrap();
626        match back {
627            AgentEvent::AgentProgress { delta, .. } => match delta {
628                ProgressDelta::Message { text } => assert_eq!(text, "ok"),
629                _ => panic!("wrong delta"),
630            },
631            _ => panic!("wrong variant"),
632        }
633    }
634
635    #[test]
636    fn event_acp_raw_roundtrip() {
637        let ev = AgentEvent::AcpRaw {
638            run_id: run_id(),
639            agent_id: agent_id(),
640            kind: "agent_message_chunk".into(),
641            raw: json!({"chunk": "hi"}),
642        };
643        let s = serde_json::to_string(&ev).unwrap();
644        assert!(s.contains("\"type\":\"acp_raw\""));
645        assert!(s.contains("\"kind\":\"agent_message_chunk\""));
646        let back: AgentEvent = serde_json::from_str(&s).unwrap();
647        match back {
648            AgentEvent::AcpRaw { kind, raw, .. } => {
649                assert_eq!(kind, "agent_message_chunk");
650                assert_eq!(raw, json!({"chunk": "hi"}));
651            }
652            _ => panic!("wrong variant"),
653        }
654    }
655
656    #[test]
657    fn event_agent_done_minimal_roundtrip() {
658        // Required fields only; optional fields should default.
659        let raw = json!({
660            "type": "agent_done",
661            "run_id": run_id(),
662            "agent_id": agent_id(),
663            "status": "Ok",
664            "tokens": {"input": 0, "output": 0, "cache_read": 0, "cache_write": 0},
665            "elapsed_ms": 12,
666        });
667        let ev: AgentEvent = serde_json::from_value(raw).unwrap();
668        match ev {
669            AgentEvent::AgentDone {
670                status,
671                elapsed_ms,
672                name,
673                agent_seq,
674                output,
675                findings,
676                prompt,
677                retry_count,
678                ..
679            } => {
680                assert_eq!(status, AgentStatus::Ok);
681                assert_eq!(elapsed_ms, 12);
682                assert!(name.is_none());
683                assert_eq!(agent_seq, 0);
684                assert_eq!(output, serde_json::Value::Null);
685                assert!(findings.is_empty());
686                assert_eq!(prompt, "");
687                assert_eq!(retry_count, 0);
688            }
689            _ => panic!("wrong variant"),
690        }
691    }
692
693    #[test]
694    fn event_agent_done_full_roundtrip() {
695        let finding = Finding {
696            kind: "x".into(),
697            severity: Severity::Low,
698            title: "t".into(),
699            detail: "d".into(),
700            location: Some(Location {
701                file: PathBuf::from("f.rs"),
702                line: Some(1),
703            }),
704            evidence: vec!["e1".into()],
705            data: json!({"k": 1}),
706        };
707        let ev = AgentEvent::AgentDone {
708            run_id: run_id(),
709            agent_id: agent_id(),
710            status: AgentStatus::Ok,
711            tokens: TokenUsage {
712                input: 1,
713                output: 2,
714                cache_read: 0,
715                cache_write: 0,
716            },
717            elapsed_ms: 42,
718            name: Some("agent-A".into()),
719            agent_seq: 3,
720            output: json!({"answer": "yes"}),
721            findings: vec![finding],
722            prompt: "do it".into(),
723            retry_count: 1,
724            ts: Default::default(),
725        };
726        let s = serde_json::to_string(&ev).unwrap();
727        assert!(s.contains("\"type\":\"agent_done\""));
728        let back: AgentEvent = serde_json::from_str(&s).unwrap();
729        match back {
730            AgentEvent::AgentDone {
731                status,
732                elapsed_ms,
733                name,
734                agent_seq,
735                output,
736                findings,
737                prompt,
738                retry_count,
739                ..
740            } => {
741                assert_eq!(status, AgentStatus::Ok);
742                assert_eq!(elapsed_ms, 42);
743                assert_eq!(name.as_deref(), Some("agent-A"));
744                assert_eq!(agent_seq, 3);
745                assert_eq!(output, json!({"answer": "yes"}));
746                assert_eq!(findings.len(), 1);
747                assert_eq!(prompt, "do it");
748                assert_eq!(retry_count, 1);
749            }
750            _ => panic!("wrong variant"),
751        }
752    }
753
754    #[test]
755    fn event_phase_done_minimal_roundtrip() {
756        let raw = json!({
757            "type": "phase_done",
758            "run_id": run_id(),
759            "phase_id": 0,
760            "ok": 4,
761            "failed": 1,
762        });
763        let ev: AgentEvent = serde_json::from_value(raw).unwrap();
764        match ev {
765            AgentEvent::PhaseDone { ok, failed, ts, .. } => {
766                assert_eq!(ok, 4);
767                assert_eq!(failed, 1);
768                // Default timestamp on the optional field
769                let _ = ts;
770            }
771            _ => panic!("wrong variant"),
772        }
773    }
774
775    #[test]
776    fn event_run_done_roundtrip() {
777        let ev = AgentEvent::RunDone {
778            run_id: run_id(),
779            status: RunStatus::Completed,
780            total_tokens: TokenUsage {
781                input: 100,
782                output: 50,
783                cache_read: 0,
784                cache_write: 0,
785            },
786            report: json!({"summary": "ok"}),
787            ts: ts(),
788        };
789        let s = serde_json::to_string(&ev).unwrap();
790        assert!(s.contains("\"type\":\"run_done\""));
791        let back: AgentEvent = serde_json::from_str(&s).unwrap();
792        match back {
793            AgentEvent::RunDone { status, report, .. } => {
794                assert_eq!(status, RunStatus::Completed);
795                assert_eq!(report, json!({"summary": "ok"}));
796            }
797            _ => panic!("wrong variant"),
798        }
799    }
800
801    #[test]
802    fn event_log_roundtrip() {
803        let ev = AgentEvent::Log {
804            run_id: run_id(),
805            agent_id: None,
806            level: LogLevel::Warn,
807            msg: "watch out".into(),
808        };
809        let s = serde_json::to_string(&ev).unwrap();
810        assert!(s.contains("\"type\":\"log\""));
811        assert!(s.contains("\"level\":\"warn\""));
812        let back: AgentEvent = serde_json::from_str(&s).unwrap();
813        match back {
814            AgentEvent::Log { level, msg, .. } => {
815                assert_eq!(level, LogLevel::Warn);
816                assert_eq!(msg, "watch out");
817            }
818            _ => panic!("wrong variant"),
819        }
820    }
821
822    #[test]
823    fn event_budget_set_roundtrip() {
824        let ev = AgentEvent::BudgetSet {
825            run_id: run_id(),
826            time_limit_ms: Some(60_000),
827            max_rounds: Some(10),
828        };
829        let s = serde_json::to_string(&ev).unwrap();
830        assert!(s.contains("\"type\":\"budget_set\""));
831        let back: AgentEvent = serde_json::from_str(&s).unwrap();
832        match back {
833            AgentEvent::BudgetSet {
834                time_limit_ms,
835                max_rounds,
836                ..
837            } => {
838                assert_eq!(time_limit_ms, Some(60_000));
839                assert_eq!(max_rounds, Some(10));
840            }
841            _ => panic!("wrong variant"),
842        }
843    }
844
845    #[test]
846    fn event_report_emitted_roundtrip() {
847        let ev = AgentEvent::ReportEmitted {
848            run_id: run_id(),
849            phase_id: 2,
850            report: json!({"x": 1}),
851        };
852        let s = serde_json::to_string(&ev).unwrap();
853        assert!(s.contains("\"type\":\"report_emitted\""));
854        let back: AgentEvent = serde_json::from_str(&s).unwrap();
855        match back {
856            AgentEvent::ReportEmitted {
857                phase_id, report, ..
858            } => {
859                assert_eq!(phase_id, 2);
860                assert_eq!(report, json!({"x": 1}));
861            }
862            _ => panic!("wrong variant"),
863        }
864    }
865
866    #[test]
867    fn event_parallel_started_roundtrip() {
868        let ev = AgentEvent::ParallelStarted {
869            run_id: run_id(),
870            phase_id: 4,
871            span_id: 99,
872            count: 8,
873        };
874        let s = serde_json::to_string(&ev).unwrap();
875        assert!(s.contains("\"type\":\"parallel_started\""));
876        let back: AgentEvent = serde_json::from_str(&s).unwrap();
877        match back {
878            AgentEvent::ParallelStarted { span_id, count, .. } => {
879                assert_eq!(span_id, 99);
880                assert_eq!(count, 8);
881            }
882            _ => panic!("wrong variant"),
883        }
884    }
885
886    #[test]
887    fn event_parallel_done_roundtrip() {
888        let ev = AgentEvent::ParallelDone {
889            run_id: run_id(),
890            phase_id: 4,
891            span_id: 99,
892            ok: 3,
893            failed: 1,
894            results: json!([{"ok": true}, {"ok": false}]),
895            elapsed_ms: 250,
896        };
897        let s = serde_json::to_string(&ev).unwrap();
898        assert!(s.contains("\"type\":\"parallel_done\""));
899        let back: AgentEvent = serde_json::from_str(&s).unwrap();
900        match back {
901            AgentEvent::ParallelDone {
902                ok,
903                failed,
904                elapsed_ms,
905                ..
906            } => {
907                assert_eq!(ok, 3);
908                assert_eq!(failed, 1);
909                assert_eq!(elapsed_ms, 250);
910            }
911            _ => panic!("wrong variant"),
912        }
913    }
914
915    #[test]
916    fn event_workflow_started_roundtrip() {
917        let ev = AgentEvent::WorkflowStarted {
918            run_id: run_id(),
919            span_id: 1,
920            path: "/tmp/wf.lua".into(),
921            args: json!({"x": 1}),
922        };
923        let s = serde_json::to_string(&ev).unwrap();
924        assert!(s.contains("\"type\":\"workflow_started\""));
925        let back: AgentEvent = serde_json::from_str(&s).unwrap();
926        match back {
927            AgentEvent::WorkflowStarted { path, args, .. } => {
928                assert_eq!(path, "/tmp/wf.lua");
929                assert_eq!(args, json!({"x": 1}));
930            }
931            _ => panic!("wrong variant"),
932        }
933    }
934
935    #[test]
936    fn event_workflow_done_with_error_roundtrip() {
937        let ev = AgentEvent::WorkflowDone {
938            run_id: run_id(),
939            span_id: 1,
940            path: "/tmp/wf.lua".into(),
941            report: json!({"items": 7}),
942            elapsed_ms: 1000,
943            error: Some("boom".into()),
944        };
945        let s = serde_json::to_string(&ev).unwrap();
946        let back: AgentEvent = serde_json::from_str(&s).unwrap();
947        match back {
948            AgentEvent::WorkflowDone { error, .. } => assert_eq!(error.as_deref(), Some("boom")),
949            _ => panic!("wrong variant"),
950        }
951    }
952
953    #[test]
954    fn event_converge_started_roundtrip() {
955        let ev = AgentEvent::ConvergeStarted {
956            run_id: run_id(),
957            phase_id: 2,
958            span_id: 5,
959            items: 12,
960            max_rounds: 3,
961        };
962        let s = serde_json::to_string(&ev).unwrap();
963        assert!(s.contains("\"type\":\"converge_started\""));
964        let back: AgentEvent = serde_json::from_str(&s).unwrap();
965        match back {
966            AgentEvent::ConvergeStarted {
967                items, max_rounds, ..
968            } => {
969                assert_eq!(items, 12);
970                assert_eq!(max_rounds, 3);
971            }
972            _ => panic!("wrong variant"),
973        }
974    }
975
976    #[test]
977    fn event_converge_done_roundtrip() {
978        let ev = AgentEvent::ConvergeDone {
979            run_id: run_id(),
980            phase_id: 2,
981            span_id: 5,
982            rounds: 4,
983            converged: true,
984            surviving: 2,
985            result: json!({"winner": "a"}),
986            elapsed_ms: 800,
987            error: None,
988        };
989        let s = serde_json::to_string(&ev).unwrap();
990        let back: AgentEvent = serde_json::from_str(&s).unwrap();
991        match back {
992            AgentEvent::ConvergeDone {
993                rounds,
994                converged,
995                surviving,
996                error,
997                ..
998            } => {
999                assert_eq!(rounds, 4);
1000                assert!(converged);
1001                assert_eq!(surviving, 2);
1002                assert!(error.is_none());
1003            }
1004            _ => panic!("wrong variant"),
1005        }
1006    }
1007
1008    #[test]
1009    fn event_pipeline_started_roundtrip() {
1010        let ev = AgentEvent::PipelineStarted {
1011            run_id: run_id(),
1012            total_stages: 4,
1013            items: 10,
1014        };
1015        let s = serde_json::to_string(&ev).unwrap();
1016        assert!(s.contains("\"type\":\"pipeline_started\""));
1017        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1018        match back {
1019            AgentEvent::PipelineStarted {
1020                run_id: _,
1021                total_stages,
1022                items,
1023            } => {
1024                assert_eq!(total_stages, 4);
1025                assert_eq!(items, 10);
1026            }
1027            _ => panic!("wrong variant"),
1028        }
1029    }
1030
1031    #[test]
1032    fn event_pipeline_stage_started_roundtrip() {
1033        let ev = AgentEvent::PipelineStageStarted {
1034            run_id: run_id(),
1035            stage_index: 1,
1036            label: "scan".into(),
1037            agents_in_stage: 5,
1038        };
1039        let s = serde_json::to_string(&ev).unwrap();
1040        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1041        match back {
1042            AgentEvent::PipelineStageStarted { label, .. } => assert_eq!(label, "scan"),
1043            _ => panic!("wrong variant"),
1044        }
1045    }
1046
1047    #[test]
1048    fn event_pipeline_item_done_roundtrip() {
1049        let ev = AgentEvent::PipelineItemDone {
1050            run_id: run_id(),
1051            stage_index: 0,
1052            item_index: 3,
1053            status: AgentStatus::Ok,
1054            tokens: TokenUsage::default(),
1055            elapsed_ms: 50,
1056        };
1057        let s = serde_json::to_string(&ev).unwrap();
1058        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1059        match back {
1060            AgentEvent::PipelineItemDone { item_index, .. } => assert_eq!(item_index, 3),
1061            _ => panic!("wrong variant"),
1062        }
1063    }
1064
1065    #[test]
1066    fn event_pipeline_done_roundtrip() {
1067        let ev = AgentEvent::PipelineDone {
1068            run_id: run_id(),
1069            stages_completed: 4,
1070            total_ok: 12,
1071            total_failed: 1,
1072        };
1073        let s = serde_json::to_string(&ev).unwrap();
1074        assert!(s.contains("\"type\":\"pipeline_done\""));
1075        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1076        match back {
1077            AgentEvent::PipelineDone {
1078                run_id: _,
1079                stages_completed,
1080                total_ok,
1081                total_failed,
1082            } => {
1083                assert_eq!(stages_completed, 4);
1084                assert_eq!(total_ok, 12);
1085                assert_eq!(total_failed, 1);
1086            }
1087            _ => panic!("wrong variant"),
1088        }
1089    }
1090
1091    #[test]
1092    fn event_phase_span_started_roundtrip() {
1093        let ev = AgentEvent::PhaseSpanStarted {
1094            run_id: run_id(),
1095            span_id: 11,
1096            name: "investigate".into(),
1097            parent_id: Some(1),
1098            depth: 2,
1099            planned: 4,
1100        };
1101        let s = serde_json::to_string(&ev).unwrap();
1102        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1103        match back {
1104            AgentEvent::PhaseSpanStarted {
1105                span_id,
1106                depth,
1107                planned,
1108                ..
1109            } => {
1110                assert_eq!(span_id, 11);
1111                assert_eq!(depth, 2);
1112                assert_eq!(planned, 4);
1113            }
1114            _ => panic!("wrong variant"),
1115        }
1116    }
1117
1118    #[test]
1119    fn event_phase_span_done_roundtrip() {
1120        let ev = AgentEvent::PhaseSpanDone {
1121            run_id: run_id(),
1122            span_id: 11,
1123            name: "investigate".into(),
1124            parent_id: None,
1125            depth: 0,
1126            elapsed_ms: 500,
1127            status: "ok".into(),
1128        };
1129        let s = serde_json::to_string(&ev).unwrap();
1130        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1131        match back {
1132            AgentEvent::PhaseSpanDone {
1133                status, elapsed_ms, ..
1134            } => {
1135                assert_eq!(status, "ok");
1136                assert_eq!(elapsed_ms, 500);
1137            }
1138            _ => panic!("wrong variant"),
1139        }
1140    }
1141
1142    #[test]
1143    fn event_schema_retry_roundtrip() {
1144        let ev = AgentEvent::SchemaRetry {
1145            run_id: run_id(),
1146            agent_id: agent_id(),
1147            attempt: 2,
1148            max: 3,
1149        };
1150        let s = serde_json::to_string(&ev).unwrap();
1151        assert!(s.contains("\"type\":\"schema_retry\""));
1152        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1153        match back {
1154            AgentEvent::SchemaRetry { attempt, max, .. } => {
1155                assert_eq!(attempt, 2);
1156                assert_eq!(max, 3);
1157            }
1158            _ => panic!("wrong variant"),
1159        }
1160    }
1161
1162    #[test]
1163    fn event_plan_preview_roundtrip() {
1164        let ev = AgentEvent::PlanPreview {
1165            run_id: run_id(),
1166            reasoning: "because".into(),
1167            phases: vec![
1168                PlanPhase {
1169                    label: "scan".into(),
1170                    dynamic: false,
1171                    description: None,
1172                },
1173                PlanPhase {
1174                    label: "fix".into(),
1175                    dynamic: true,
1176                    description: Some("apply patches".into()),
1177                },
1178            ],
1179        };
1180        let s = serde_json::to_string(&ev).unwrap();
1181        assert!(s.contains("\"type\":\"plan_preview\""));
1182        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1183        match back {
1184            AgentEvent::PlanPreview { phases, .. } => assert_eq!(phases.len(), 2),
1185            _ => panic!("wrong variant"),
1186        }
1187    }
1188
1189    #[test]
1190    fn event_signal_received_roundtrip() {
1191        let ev = AgentEvent::SignalReceived {
1192            run_id: None,
1193            signal: "SIGINT".into(),
1194            ts: ts(),
1195        };
1196        let s = serde_json::to_string(&ev).unwrap();
1197        assert!(s.contains("\"type\":\"signal_received\""));
1198        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1199        match back {
1200            AgentEvent::SignalReceived { signal, run_id, .. } => {
1201                assert_eq!(signal, "SIGINT");
1202                assert!(run_id.is_none());
1203            }
1204            _ => panic!("wrong variant"),
1205        }
1206    }
1207
1208    #[test]
1209    fn event_clone_preserves_variant() {
1210        let ev = AgentEvent::RunStarted {
1211            run_id: run_id(),
1212            task: "t".into(),
1213            ts: ts(),
1214        };
1215        let cloned = ev.clone();
1216        let s1 = serde_json::to_string(&ev).unwrap();
1217        let s2 = serde_json::to_string(&cloned).unwrap();
1218        assert_eq!(s1, s2);
1219    }
1220
1221    #[test]
1222    fn event_debug_includes_type_tag() {
1223        let ev = AgentEvent::RunStarted {
1224            run_id: run_id(),
1225            task: "t".into(),
1226            ts: ts(),
1227        };
1228        let dbg = format!("{:?}", ev);
1229        assert!(dbg.contains("RunStarted"));
1230    }
1231
1232    #[test]
1233    fn event_unknown_type_fails() {
1234        let raw = json!({"type": "totally_made_up", "x": 1});
1235        let r: Result<AgentEvent, _> = serde_json::from_value(raw);
1236        assert!(r.is_err());
1237    }
1238
1239    // ── EventSender alias ────────────────────────────────────────
1240
1241    #[tokio::test]
1242    async fn event_sender_alias_is_broadcast_sender() {
1243        // Compile-time check: EventSender is a broadcast::Sender<AgentEvent>.
1244        let (tx, _rx) = tokio::sync::broadcast::channel::<AgentEvent>(4);
1245        let _alias: EventSender = tx;
1246    }
1247}