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