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