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