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