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 },
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 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 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 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 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 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 SchemaRetry {
220 run_id: RunId,
221 agent_id: AgentId,
222 attempt: u32,
223 max: u32,
224 },
225 PlanPreview {
229 run_id: RunId,
230 reasoning: String,
231 phases: Vec<PlanPhase>,
232 },
233 SignalReceived {
237 run_id: Option<RunId>,
238 signal: String,
239 ts: DateTime<Utc>,
240 },
241}
242
243#[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 #[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; 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 #[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 #[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 #[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 #[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 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 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 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 #[tokio::test]
1236 async fn event_sender_alias_is_broadcast_sender() {
1237 let (tx, _rx) = tokio::sync::broadcast::channel::<AgentEvent>(4);
1239 let _alias: EventSender = tx;
1240 }
1241}