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