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