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