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!(
306            serde_json::to_string(&RunStatus::Completed).unwrap(),
307            "\"completed\""
308        );
309        assert_eq!(
310            serde_json::to_string(&RunStatus::Failed).unwrap(),
311            "\"failed\""
312        );
313        assert_eq!(
314            serde_json::to_string(&RunStatus::Cancelled).unwrap(),
315            "\"cancelled\""
316        );
317        assert_eq!(
318            serde_json::to_string(&RunStatus::Partial).unwrap(),
319            "\"partial\""
320        );
321    }
322
323    #[test]
324    fn run_status_deserializes_from_snake_case() {
325        assert_eq!(
326            serde_json::from_str::<RunStatus>("\"completed\"").unwrap(),
327            RunStatus::Completed
328        );
329        assert_eq!(
330            serde_json::from_str::<RunStatus>("\"failed\"").unwrap(),
331            RunStatus::Failed
332        );
333        assert_eq!(
334            serde_json::from_str::<RunStatus>("\"cancelled\"").unwrap(),
335            RunStatus::Cancelled
336        );
337        assert_eq!(
338            serde_json::from_str::<RunStatus>("\"partial\"").unwrap(),
339            RunStatus::Partial
340        );
341    }
342
343    #[test]
344    fn run_status_equality_and_copy() {
345        let a = RunStatus::Completed;
346        let b = a; // Copy semantics
347        let c = a;
348        assert_eq!(a, b);
349        assert_eq!(a, c);
350    }
351
352    #[test]
353    fn run_status_unknown_variant_fails() {
354        let r: Result<RunStatus, _> = serde_json::from_str("\"unknown\"");
355        assert!(r.is_err());
356    }
357
358    // ── LogLevel ─────────────────────────────────────────────────
359
360    #[test]
361    fn log_level_serializes_as_lowercase() {
362        assert_eq!(
363            serde_json::to_string(&LogLevel::Trace).unwrap(),
364            "\"trace\""
365        );
366        assert_eq!(
367            serde_json::to_string(&LogLevel::Debug).unwrap(),
368            "\"debug\""
369        );
370        assert_eq!(serde_json::to_string(&LogLevel::Info).unwrap(), "\"info\"");
371        assert_eq!(serde_json::to_string(&LogLevel::Warn).unwrap(), "\"warn\"");
372        assert_eq!(
373            serde_json::to_string(&LogLevel::Error).unwrap(),
374            "\"error\""
375        );
376    }
377
378    #[test]
379    fn log_level_deserializes_from_lowercase() {
380        assert_eq!(
381            serde_json::from_str::<LogLevel>("\"warn\"").unwrap(),
382            LogLevel::Warn
383        );
384        assert_eq!(
385            serde_json::from_str::<LogLevel>("\"error\"").unwrap(),
386            LogLevel::Error
387        );
388    }
389
390    #[test]
391    fn log_level_equality_and_copy() {
392        let a = LogLevel::Info;
393        let b = a;
394        assert_eq!(a, b);
395    }
396
397    // ── ProgressDelta ────────────────────────────────────────────
398
399    #[test]
400    fn progress_delta_message_roundtrip() {
401        let d = ProgressDelta::Message {
402            text: "hello".into(),
403        };
404        let s = serde_json::to_string(&d).unwrap();
405        assert!(s.contains("\"kind\":\"message\""));
406        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
407        match back {
408            ProgressDelta::Message { text } => assert_eq!(text, "hello"),
409            _ => panic!("wrong variant"),
410        }
411    }
412
413    #[test]
414    fn progress_delta_tool_call_roundtrip() {
415        let d = ProgressDelta::ToolCall {
416            name: "bash".into(),
417            summary: "run ls".into(),
418        };
419        let s = serde_json::to_string(&d).unwrap();
420        assert!(s.contains("\"kind\":\"tool_call\""));
421        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
422        match back {
423            ProgressDelta::ToolCall { name, summary } => {
424                assert_eq!(name, "bash");
425                assert_eq!(summary, "run ls");
426            }
427            _ => panic!("wrong variant"),
428        }
429    }
430
431    #[test]
432    fn progress_delta_file_edit_roundtrip() {
433        let d = ProgressDelta::FileEdit {
434            path: PathBuf::from("src/main.rs"),
435        };
436        let s = serde_json::to_string(&d).unwrap();
437        assert!(s.contains("\"kind\":\"file_edit\""));
438        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
439        match back {
440            ProgressDelta::FileEdit { path } => assert_eq!(path, PathBuf::from("src/main.rs")),
441            _ => panic!("wrong variant"),
442        }
443    }
444
445    #[test]
446    fn progress_delta_tokens_roundtrip() {
447        let d = ProgressDelta::Tokens {
448            usage: TokenUsage {
449                input: 10,
450                output: 20,
451                cache_read: 1,
452                cache_write: 2,
453            },
454        };
455        let s = serde_json::to_string(&d).unwrap();
456        assert!(s.contains("\"kind\":\"tokens\""));
457        let back: ProgressDelta = serde_json::from_str(&s).unwrap();
458        match back {
459            ProgressDelta::Tokens { usage } => {
460                assert_eq!(usage.input, 10);
461                assert_eq!(usage.output, 20);
462            }
463            _ => panic!("wrong variant"),
464        }
465    }
466
467    // ── PlanPhase ────────────────────────────────────────────────
468
469    #[test]
470    fn plan_phase_minimal_roundtrip() {
471        let raw = json!({"label": "review"});
472        let p: PlanPhase = serde_json::from_value(raw.clone()).unwrap();
473        assert_eq!(p.label, "review");
474        assert!(!p.dynamic);
475        assert!(p.description.is_none());
476        let back = serde_json::to_value(&p).unwrap();
477        assert_eq!(
478            back,
479            json!({"label": "review", "dynamic": false, "description": null})
480        );
481    }
482
483    #[test]
484    fn plan_phase_full_roundtrip() {
485        let p = PlanPhase {
486            label: "implement".into(),
487            dynamic: true,
488            description: Some("build the thing".into()),
489        };
490        let s = serde_json::to_string(&p).unwrap();
491        let back: PlanPhase = serde_json::from_str(&s).unwrap();
492        assert_eq!(back.label, "implement");
493        assert!(back.dynamic);
494        assert_eq!(back.description.as_deref(), Some("build the thing"));
495    }
496
497    // ── AgentEvent variants ──────────────────────────────────────
498
499    #[test]
500    fn event_run_started_roundtrip() {
501        let ev = AgentEvent::RunStarted {
502            run_id: run_id(),
503            task: "investigate".into(),
504            ts: ts(),
505        };
506        let s = serde_json::to_string(&ev).unwrap();
507        assert!(s.contains("\"type\":\"run_started\""));
508        let back: AgentEvent = serde_json::from_str(&s).unwrap();
509        match back {
510            AgentEvent::RunStarted { task, .. } => assert_eq!(task, "investigate"),
511            _ => panic!("wrong variant"),
512        }
513    }
514
515    #[test]
516    fn event_phase_started_minimal_roundtrip() {
517        let ev = AgentEvent::PhaseStarted {
518            run_id: run_id(),
519            phase_id: 3,
520            label: "scan".into(),
521            planned: 5,
522            parent_span_id: None,
523            description: None,
524            role: None,
525            ts: DateTime::<Utc>::from_timestamp(0, 0).unwrap(),
526        };
527        let s = serde_json::to_string(&ev).unwrap();
528        assert!(s.contains("\"type\":\"phase_started\""));
529        let back: AgentEvent = serde_json::from_str(&s).unwrap();
530        match back {
531            AgentEvent::PhaseStarted {
532                phase_id,
533                label,
534                planned,
535                ..
536            } => {
537                assert_eq!(phase_id, 3);
538                assert_eq!(label, "scan");
539                assert_eq!(planned, 5);
540            }
541            _ => panic!("wrong variant"),
542        }
543    }
544
545    #[test]
546    fn event_phase_started_omits_optional_fields() {
547        let ev = AgentEvent::PhaseStarted {
548            run_id: run_id(),
549            phase_id: 0,
550            label: "x".into(),
551            planned: 0,
552            parent_span_id: None,
553            description: None,
554            role: None,
555            ts: ts(),
556        };
557        let s = serde_json::to_string(&ev).unwrap();
558        // Optional fields with #[serde(default)] deserialize when omitted but
559        // serialize as JSON null when None. Verify they can be omitted on the
560        // input side by deserializing a payload that drops them.
561        let back: AgentEvent = serde_json::from_str(&s).unwrap();
562        match back {
563            AgentEvent::PhaseStarted {
564                parent_span_id,
565                description,
566                role,
567                ..
568            } => {
569                assert!(parent_span_id.is_none());
570                assert!(description.is_none());
571                assert!(role.is_none());
572            }
573            _ => panic!("wrong variant"),
574        }
575        assert!(s.contains("\"parent_span_id\":null"));
576        assert!(s.contains("\"description\":null"));
577        assert!(s.contains("\"role\":null"));
578    }
579
580    #[test]
581    fn event_agent_started_roundtrip() {
582        let ev = AgentEvent::AgentStarted {
583            run_id: run_id(),
584            phase_id: 1,
585            agent_id: agent_id(),
586            prompt_preview: "find bugs".into(),
587            model: Some("gpt-4".into()),
588            description: Some("audit".into()),
589            role: Some("reviewer".into()),
590            name: Some("agent-1".into()),
591            agent_seq: 7,
592        };
593        let s = serde_json::to_string(&ev).unwrap();
594        assert!(s.contains("\"type\":\"agent_started\""));
595        let back: AgentEvent = serde_json::from_str(&s).unwrap();
596        match back {
597            AgentEvent::AgentStarted {
598                agent_seq,
599                model,
600                name,
601                ..
602            } => {
603                assert_eq!(agent_seq, 7);
604                assert_eq!(model.as_deref(), Some("gpt-4"));
605                assert_eq!(name.as_deref(), Some("agent-1"));
606            }
607            _ => panic!("wrong variant"),
608        }
609    }
610
611    #[test]
612    fn event_agent_progress_roundtrip() {
613        let ev = AgentEvent::AgentProgress {
614            run_id: run_id(),
615            agent_id: agent_id(),
616            delta: ProgressDelta::Message { text: "ok".into() },
617        };
618        let s = serde_json::to_string(&ev).unwrap();
619        assert!(s.contains("\"type\":\"agent_progress\""));
620        let back: AgentEvent = serde_json::from_str(&s).unwrap();
621        match back {
622            AgentEvent::AgentProgress { delta, .. } => match delta {
623                ProgressDelta::Message { text } => assert_eq!(text, "ok"),
624                _ => panic!("wrong delta"),
625            },
626            _ => panic!("wrong variant"),
627        }
628    }
629
630    #[test]
631    fn event_acp_raw_roundtrip() {
632        let ev = AgentEvent::AcpRaw {
633            run_id: run_id(),
634            agent_id: agent_id(),
635            kind: "agent_message_chunk".into(),
636            raw: json!({"chunk": "hi"}),
637        };
638        let s = serde_json::to_string(&ev).unwrap();
639        assert!(s.contains("\"type\":\"acp_raw\""));
640        assert!(s.contains("\"kind\":\"agent_message_chunk\""));
641        let back: AgentEvent = serde_json::from_str(&s).unwrap();
642        match back {
643            AgentEvent::AcpRaw { kind, raw, .. } => {
644                assert_eq!(kind, "agent_message_chunk");
645                assert_eq!(raw, json!({"chunk": "hi"}));
646            }
647            _ => panic!("wrong variant"),
648        }
649    }
650
651    #[test]
652    fn event_agent_done_minimal_roundtrip() {
653        // Required fields only; optional fields should default.
654        let raw = json!({
655            "type": "agent_done",
656            "run_id": run_id(),
657            "agent_id": agent_id(),
658            "status": "Ok",
659            "tokens": {"input": 0, "output": 0, "cache_read": 0, "cache_write": 0},
660            "elapsed_ms": 12,
661        });
662        let ev: AgentEvent = serde_json::from_value(raw).unwrap();
663        match ev {
664            AgentEvent::AgentDone {
665                status,
666                elapsed_ms,
667                name,
668                agent_seq,
669                output,
670                findings,
671                prompt,
672                retry_count,
673                ..
674            } => {
675                assert_eq!(status, AgentStatus::Ok);
676                assert_eq!(elapsed_ms, 12);
677                assert!(name.is_none());
678                assert_eq!(agent_seq, 0);
679                assert_eq!(output, serde_json::Value::Null);
680                assert!(findings.is_empty());
681                assert_eq!(prompt, "");
682                assert_eq!(retry_count, 0);
683            }
684            _ => panic!("wrong variant"),
685        }
686    }
687
688    #[test]
689    fn event_agent_done_full_roundtrip() {
690        let finding = Finding {
691            kind: "x".into(),
692            severity: Severity::Low,
693            title: "t".into(),
694            detail: "d".into(),
695            location: Some(Location {
696                file: PathBuf::from("f.rs"),
697                line: Some(1),
698            }),
699            evidence: vec!["e1".into()],
700            data: json!({"k": 1}),
701        };
702        let ev = AgentEvent::AgentDone {
703            run_id: run_id(),
704            agent_id: agent_id(),
705            status: AgentStatus::Ok,
706            tokens: TokenUsage {
707                input: 1,
708                output: 2,
709                cache_read: 0,
710                cache_write: 0,
711            },
712            elapsed_ms: 42,
713            name: Some("agent-A".into()),
714            agent_seq: 3,
715            output: json!({"answer": "yes"}),
716            findings: vec![finding],
717            prompt: "do it".into(),
718            retry_count: 1,
719        };
720        let s = serde_json::to_string(&ev).unwrap();
721        assert!(s.contains("\"type\":\"agent_done\""));
722        let back: AgentEvent = serde_json::from_str(&s).unwrap();
723        match back {
724            AgentEvent::AgentDone {
725                status,
726                elapsed_ms,
727                name,
728                agent_seq,
729                output,
730                findings,
731                prompt,
732                retry_count,
733                ..
734            } => {
735                assert_eq!(status, AgentStatus::Ok);
736                assert_eq!(elapsed_ms, 42);
737                assert_eq!(name.as_deref(), Some("agent-A"));
738                assert_eq!(agent_seq, 3);
739                assert_eq!(output, json!({"answer": "yes"}));
740                assert_eq!(findings.len(), 1);
741                assert_eq!(prompt, "do it");
742                assert_eq!(retry_count, 1);
743            }
744            _ => panic!("wrong variant"),
745        }
746    }
747
748    #[test]
749    fn event_phase_done_minimal_roundtrip() {
750        let raw = json!({
751            "type": "phase_done",
752            "run_id": run_id(),
753            "phase_id": 0,
754            "ok": 4,
755            "failed": 1,
756        });
757        let ev: AgentEvent = serde_json::from_value(raw).unwrap();
758        match ev {
759            AgentEvent::PhaseDone { ok, failed, ts, .. } => {
760                assert_eq!(ok, 4);
761                assert_eq!(failed, 1);
762                // Default timestamp on the optional field
763                let _ = ts;
764            }
765            _ => panic!("wrong variant"),
766        }
767    }
768
769    #[test]
770    fn event_run_done_roundtrip() {
771        let ev = AgentEvent::RunDone {
772            run_id: run_id(),
773            status: RunStatus::Completed,
774            total_tokens: TokenUsage {
775                input: 100,
776                output: 50,
777                cache_read: 0,
778                cache_write: 0,
779            },
780            report: json!({"summary": "ok"}),
781            ts: ts(),
782        };
783        let s = serde_json::to_string(&ev).unwrap();
784        assert!(s.contains("\"type\":\"run_done\""));
785        let back: AgentEvent = serde_json::from_str(&s).unwrap();
786        match back {
787            AgentEvent::RunDone { status, report, .. } => {
788                assert_eq!(status, RunStatus::Completed);
789                assert_eq!(report, json!({"summary": "ok"}));
790            }
791            _ => panic!("wrong variant"),
792        }
793    }
794
795    #[test]
796    fn event_log_roundtrip() {
797        let ev = AgentEvent::Log {
798            run_id: run_id(),
799            agent_id: None,
800            level: LogLevel::Warn,
801            msg: "watch out".into(),
802        };
803        let s = serde_json::to_string(&ev).unwrap();
804        assert!(s.contains("\"type\":\"log\""));
805        assert!(s.contains("\"level\":\"warn\""));
806        let back: AgentEvent = serde_json::from_str(&s).unwrap();
807        match back {
808            AgentEvent::Log { level, msg, .. } => {
809                assert_eq!(level, LogLevel::Warn);
810                assert_eq!(msg, "watch out");
811            }
812            _ => panic!("wrong variant"),
813        }
814    }
815
816    #[test]
817    fn event_budget_set_roundtrip() {
818        let ev = AgentEvent::BudgetSet {
819            run_id: run_id(),
820            time_limit_ms: Some(60_000),
821            max_rounds: Some(10),
822        };
823        let s = serde_json::to_string(&ev).unwrap();
824        assert!(s.contains("\"type\":\"budget_set\""));
825        let back: AgentEvent = serde_json::from_str(&s).unwrap();
826        match back {
827            AgentEvent::BudgetSet {
828                time_limit_ms,
829                max_rounds,
830                ..
831            } => {
832                assert_eq!(time_limit_ms, Some(60_000));
833                assert_eq!(max_rounds, Some(10));
834            }
835            _ => panic!("wrong variant"),
836        }
837    }
838
839    #[test]
840    fn event_report_emitted_roundtrip() {
841        let ev = AgentEvent::ReportEmitted {
842            run_id: run_id(),
843            phase_id: 2,
844            report: json!({"x": 1}),
845        };
846        let s = serde_json::to_string(&ev).unwrap();
847        assert!(s.contains("\"type\":\"report_emitted\""));
848        let back: AgentEvent = serde_json::from_str(&s).unwrap();
849        match back {
850            AgentEvent::ReportEmitted {
851                phase_id, report, ..
852            } => {
853                assert_eq!(phase_id, 2);
854                assert_eq!(report, json!({"x": 1}));
855            }
856            _ => panic!("wrong variant"),
857        }
858    }
859
860    #[test]
861    fn event_parallel_started_roundtrip() {
862        let ev = AgentEvent::ParallelStarted {
863            run_id: run_id(),
864            phase_id: 4,
865            span_id: 99,
866            count: 8,
867        };
868        let s = serde_json::to_string(&ev).unwrap();
869        assert!(s.contains("\"type\":\"parallel_started\""));
870        let back: AgentEvent = serde_json::from_str(&s).unwrap();
871        match back {
872            AgentEvent::ParallelStarted { span_id, count, .. } => {
873                assert_eq!(span_id, 99);
874                assert_eq!(count, 8);
875            }
876            _ => panic!("wrong variant"),
877        }
878    }
879
880    #[test]
881    fn event_parallel_done_roundtrip() {
882        let ev = AgentEvent::ParallelDone {
883            run_id: run_id(),
884            phase_id: 4,
885            span_id: 99,
886            ok: 3,
887            failed: 1,
888            results: json!([{"ok": true}, {"ok": false}]),
889            elapsed_ms: 250,
890        };
891        let s = serde_json::to_string(&ev).unwrap();
892        assert!(s.contains("\"type\":\"parallel_done\""));
893        let back: AgentEvent = serde_json::from_str(&s).unwrap();
894        match back {
895            AgentEvent::ParallelDone {
896                ok,
897                failed,
898                elapsed_ms,
899                ..
900            } => {
901                assert_eq!(ok, 3);
902                assert_eq!(failed, 1);
903                assert_eq!(elapsed_ms, 250);
904            }
905            _ => panic!("wrong variant"),
906        }
907    }
908
909    #[test]
910    fn event_workflow_started_roundtrip() {
911        let ev = AgentEvent::WorkflowStarted {
912            run_id: run_id(),
913            span_id: 1,
914            path: "/tmp/wf.lua".into(),
915            args: json!({"x": 1}),
916        };
917        let s = serde_json::to_string(&ev).unwrap();
918        assert!(s.contains("\"type\":\"workflow_started\""));
919        let back: AgentEvent = serde_json::from_str(&s).unwrap();
920        match back {
921            AgentEvent::WorkflowStarted { path, args, .. } => {
922                assert_eq!(path, "/tmp/wf.lua");
923                assert_eq!(args, json!({"x": 1}));
924            }
925            _ => panic!("wrong variant"),
926        }
927    }
928
929    #[test]
930    fn event_workflow_done_with_error_roundtrip() {
931        let ev = AgentEvent::WorkflowDone {
932            run_id: run_id(),
933            span_id: 1,
934            path: "/tmp/wf.lua".into(),
935            report: json!({"items": 7}),
936            elapsed_ms: 1000,
937            error: Some("boom".into()),
938        };
939        let s = serde_json::to_string(&ev).unwrap();
940        let back: AgentEvent = serde_json::from_str(&s).unwrap();
941        match back {
942            AgentEvent::WorkflowDone { error, .. } => assert_eq!(error.as_deref(), Some("boom")),
943            _ => panic!("wrong variant"),
944        }
945    }
946
947    #[test]
948    fn event_converge_started_roundtrip() {
949        let ev = AgentEvent::ConvergeStarted {
950            run_id: run_id(),
951            phase_id: 2,
952            span_id: 5,
953            items: 12,
954            max_rounds: 3,
955        };
956        let s = serde_json::to_string(&ev).unwrap();
957        assert!(s.contains("\"type\":\"converge_started\""));
958        let back: AgentEvent = serde_json::from_str(&s).unwrap();
959        match back {
960            AgentEvent::ConvergeStarted {
961                items, max_rounds, ..
962            } => {
963                assert_eq!(items, 12);
964                assert_eq!(max_rounds, 3);
965            }
966            _ => panic!("wrong variant"),
967        }
968    }
969
970    #[test]
971    fn event_converge_done_roundtrip() {
972        let ev = AgentEvent::ConvergeDone {
973            run_id: run_id(),
974            phase_id: 2,
975            span_id: 5,
976            rounds: 4,
977            converged: true,
978            surviving: 2,
979            result: json!({"winner": "a"}),
980            elapsed_ms: 800,
981            error: None,
982        };
983        let s = serde_json::to_string(&ev).unwrap();
984        let back: AgentEvent = serde_json::from_str(&s).unwrap();
985        match back {
986            AgentEvent::ConvergeDone {
987                rounds,
988                converged,
989                surviving,
990                error,
991                ..
992            } => {
993                assert_eq!(rounds, 4);
994                assert!(converged);
995                assert_eq!(surviving, 2);
996                assert!(error.is_none());
997            }
998            _ => panic!("wrong variant"),
999        }
1000    }
1001
1002    #[test]
1003    fn event_pipeline_started_roundtrip() {
1004        let ev = AgentEvent::PipelineStarted {
1005            run_id: run_id(),
1006            total_stages: 4,
1007            items: 10,
1008        };
1009        let s = serde_json::to_string(&ev).unwrap();
1010        assert!(s.contains("\"type\":\"pipeline_started\""));
1011        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1012        match back {
1013            AgentEvent::PipelineStarted {
1014                run_id: _,
1015                total_stages,
1016                items,
1017            } => {
1018                assert_eq!(total_stages, 4);
1019                assert_eq!(items, 10);
1020            }
1021            _ => panic!("wrong variant"),
1022        }
1023    }
1024
1025    #[test]
1026    fn event_pipeline_stage_started_roundtrip() {
1027        let ev = AgentEvent::PipelineStageStarted {
1028            run_id: run_id(),
1029            stage_index: 1,
1030            label: "scan".into(),
1031            agents_in_stage: 5,
1032        };
1033        let s = serde_json::to_string(&ev).unwrap();
1034        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1035        match back {
1036            AgentEvent::PipelineStageStarted { label, .. } => assert_eq!(label, "scan"),
1037            _ => panic!("wrong variant"),
1038        }
1039    }
1040
1041    #[test]
1042    fn event_pipeline_item_done_roundtrip() {
1043        let ev = AgentEvent::PipelineItemDone {
1044            run_id: run_id(),
1045            stage_index: 0,
1046            item_index: 3,
1047            status: AgentStatus::Ok,
1048            tokens: TokenUsage::default(),
1049            elapsed_ms: 50,
1050        };
1051        let s = serde_json::to_string(&ev).unwrap();
1052        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1053        match back {
1054            AgentEvent::PipelineItemDone { item_index, .. } => assert_eq!(item_index, 3),
1055            _ => panic!("wrong variant"),
1056        }
1057    }
1058
1059    #[test]
1060    fn event_pipeline_done_roundtrip() {
1061        let ev = AgentEvent::PipelineDone {
1062            run_id: run_id(),
1063            stages_completed: 4,
1064            total_ok: 12,
1065            total_failed: 1,
1066        };
1067        let s = serde_json::to_string(&ev).unwrap();
1068        assert!(s.contains("\"type\":\"pipeline_done\""));
1069        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1070        match back {
1071            AgentEvent::PipelineDone {
1072                run_id: _,
1073                stages_completed,
1074                total_ok,
1075                total_failed,
1076            } => {
1077                assert_eq!(stages_completed, 4);
1078                assert_eq!(total_ok, 12);
1079                assert_eq!(total_failed, 1);
1080            }
1081            _ => panic!("wrong variant"),
1082        }
1083    }
1084
1085    #[test]
1086    fn event_phase_span_started_roundtrip() {
1087        let ev = AgentEvent::PhaseSpanStarted {
1088            run_id: run_id(),
1089            span_id: 11,
1090            name: "investigate".into(),
1091            parent_id: Some(1),
1092            depth: 2,
1093            planned: 4,
1094        };
1095        let s = serde_json::to_string(&ev).unwrap();
1096        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1097        match back {
1098            AgentEvent::PhaseSpanStarted {
1099                span_id,
1100                depth,
1101                planned,
1102                ..
1103            } => {
1104                assert_eq!(span_id, 11);
1105                assert_eq!(depth, 2);
1106                assert_eq!(planned, 4);
1107            }
1108            _ => panic!("wrong variant"),
1109        }
1110    }
1111
1112    #[test]
1113    fn event_phase_span_done_roundtrip() {
1114        let ev = AgentEvent::PhaseSpanDone {
1115            run_id: run_id(),
1116            span_id: 11,
1117            name: "investigate".into(),
1118            parent_id: None,
1119            depth: 0,
1120            elapsed_ms: 500,
1121            status: "ok".into(),
1122        };
1123        let s = serde_json::to_string(&ev).unwrap();
1124        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1125        match back {
1126            AgentEvent::PhaseSpanDone {
1127                status, elapsed_ms, ..
1128            } => {
1129                assert_eq!(status, "ok");
1130                assert_eq!(elapsed_ms, 500);
1131            }
1132            _ => panic!("wrong variant"),
1133        }
1134    }
1135
1136    #[test]
1137    fn event_schema_retry_roundtrip() {
1138        let ev = AgentEvent::SchemaRetry {
1139            run_id: run_id(),
1140            agent_id: agent_id(),
1141            attempt: 2,
1142            max: 3,
1143        };
1144        let s = serde_json::to_string(&ev).unwrap();
1145        assert!(s.contains("\"type\":\"schema_retry\""));
1146        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1147        match back {
1148            AgentEvent::SchemaRetry { attempt, max, .. } => {
1149                assert_eq!(attempt, 2);
1150                assert_eq!(max, 3);
1151            }
1152            _ => panic!("wrong variant"),
1153        }
1154    }
1155
1156    #[test]
1157    fn event_plan_preview_roundtrip() {
1158        let ev = AgentEvent::PlanPreview {
1159            run_id: run_id(),
1160            reasoning: "because".into(),
1161            phases: vec![
1162                PlanPhase {
1163                    label: "scan".into(),
1164                    dynamic: false,
1165                    description: None,
1166                },
1167                PlanPhase {
1168                    label: "fix".into(),
1169                    dynamic: true,
1170                    description: Some("apply patches".into()),
1171                },
1172            ],
1173        };
1174        let s = serde_json::to_string(&ev).unwrap();
1175        assert!(s.contains("\"type\":\"plan_preview\""));
1176        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1177        match back {
1178            AgentEvent::PlanPreview { phases, .. } => assert_eq!(phases.len(), 2),
1179            _ => panic!("wrong variant"),
1180        }
1181    }
1182
1183    #[test]
1184    fn event_signal_received_roundtrip() {
1185        let ev = AgentEvent::SignalReceived {
1186            run_id: None,
1187            signal: "SIGINT".into(),
1188            ts: ts(),
1189        };
1190        let s = serde_json::to_string(&ev).unwrap();
1191        assert!(s.contains("\"type\":\"signal_received\""));
1192        let back: AgentEvent = serde_json::from_str(&s).unwrap();
1193        match back {
1194            AgentEvent::SignalReceived { signal, run_id, .. } => {
1195                assert_eq!(signal, "SIGINT");
1196                assert!(run_id.is_none());
1197            }
1198            _ => panic!("wrong variant"),
1199        }
1200    }
1201
1202    #[test]
1203    fn event_clone_preserves_variant() {
1204        let ev = AgentEvent::RunStarted {
1205            run_id: run_id(),
1206            task: "t".into(),
1207            ts: ts(),
1208        };
1209        let cloned = ev.clone();
1210        let s1 = serde_json::to_string(&ev).unwrap();
1211        let s2 = serde_json::to_string(&cloned).unwrap();
1212        assert_eq!(s1, s2);
1213    }
1214
1215    #[test]
1216    fn event_debug_includes_type_tag() {
1217        let ev = AgentEvent::RunStarted {
1218            run_id: run_id(),
1219            task: "t".into(),
1220            ts: ts(),
1221        };
1222        let dbg = format!("{:?}", ev);
1223        assert!(dbg.contains("RunStarted"));
1224    }
1225
1226    #[test]
1227    fn event_unknown_type_fails() {
1228        let raw = json!({"type": "totally_made_up", "x": 1});
1229        let r: Result<AgentEvent, _> = serde_json::from_value(raw);
1230        assert!(r.is_err());
1231    }
1232
1233    // ── EventSender alias ────────────────────────────────────────
1234
1235    #[tokio::test]
1236    async fn event_sender_alias_is_broadcast_sender() {
1237        // Compile-time check: EventSender is a broadcast::Sender<AgentEvent>.
1238        let (tx, _rx) = tokio::sync::broadcast::channel::<AgentEvent>(4);
1239        let _alias: EventSender = tx;
1240    }
1241}