1mod buffered;
44mod channel;
45mod policy;
46
47#[cfg(feature = "tracing")]
48mod tracing_reporter;
49
50pub use buffered::BufferedReporter;
51pub use channel::ChannelReporter;
52pub use policy::{FailurePolicy, FallibleObserver, PolicyReporter};
53
54#[cfg(feature = "tracing")]
55pub use tracing_reporter::TracingReporter;
56
57use std::io::{self, Write};
58use std::time::SystemTime;
59
60use agentkit_core::{Item, ItemKind, Part, SessionId, TokenUsage, Usage};
61use agentkit_loop::{AgentEvent, LoopObserver, ObservedEvent, TurnResult};
62use serde::Serialize;
63use thiserror::Error;
64
65#[derive(Debug, Error)]
72pub enum ReportError {
73 #[error("io error: {0}")]
75 Io(#[from] io::Error),
76 #[error("serialization error: {0}")]
78 Serialize(#[from] serde_json::Error),
79 #[error("channel send failed")]
81 ChannelSend,
82}
83
84#[derive(Clone, Debug, PartialEq, Serialize)]
90pub struct EventEnvelope<'a> {
91 pub timestamp: SystemTime,
93 #[serde(skip_serializing_if = "Option::is_none")]
95 pub session_id: Option<&'a SessionId>,
96 pub event: &'a AgentEvent,
98}
99
100#[derive(Default)]
120pub struct CompositeReporter {
121 children: Vec<Box<dyn LoopObserver>>,
122}
123
124impl CompositeReporter {
125 pub fn new() -> Self {
127 Self::default()
128 }
129
130 pub fn with_observer(mut self, observer: impl LoopObserver + 'static) -> Self {
136 self.children.push(Box::new(observer));
137 self
138 }
139
140 pub fn push(&mut self, observer: impl LoopObserver + 'static) -> &mut Self {
149 self.children.push(Box::new(observer));
150 self
151 }
152}
153
154impl LoopObserver for CompositeReporter {
155 fn handle_event(&self, event: ObservedEvent) {
156 if self.children.is_empty() {
157 return;
158 }
159 let last = self.children.len() - 1;
160 for child in &self.children[..last] {
161 child.handle_event(event.clone());
162 }
163 self.children[last].handle_event(event);
164 }
165}
166
167pub struct JsonlReporter<W> {
192 writer: std::sync::Mutex<W>,
193 flush_each_event: bool,
194 include_session_id: bool,
195 errors: std::sync::Mutex<Vec<ReportError>>,
196}
197
198impl<W> JsonlReporter<W>
199where
200 W: Write,
201{
202 pub fn new(writer: W) -> Self {
204 Self {
205 writer: std::sync::Mutex::new(writer),
206 flush_each_event: true,
207 include_session_id: false,
208 errors: std::sync::Mutex::new(Vec::new()),
209 }
210 }
211
212 pub fn with_flush_each_event(mut self, flush_each_event: bool) -> Self {
214 self.flush_each_event = flush_each_event;
215 self
216 }
217
218 pub fn with_session_id(mut self, include_session_id: bool) -> Self {
223 self.include_session_id = include_session_id;
224 self
225 }
226
227 pub fn take_errors(&self) -> Vec<ReportError> {
229 std::mem::take(&mut *self.errors.lock().unwrap_or_else(|e| e.into_inner()))
230 }
231
232 fn record_result(&self, result: Result<(), ReportError>) {
233 if let Err(error) = result {
234 self.errors
235 .lock()
236 .unwrap_or_else(|e| e.into_inner())
237 .push(error);
238 }
239 }
240
241 pub fn into_inner(self) -> W {
243 self.writer.into_inner().unwrap_or_else(|e| e.into_inner())
244 }
245}
246
247impl<W> LoopObserver for JsonlReporter<W>
248where
249 W: Write + Send,
250{
251 fn handle_event(&self, event: ObservedEvent) {
252 let result = (|| -> Result<(), ReportError> {
253 let envelope = EventEnvelope {
254 timestamp: SystemTime::now(),
255 session_id: self.include_session_id.then_some(event.session_id.as_ref()),
256 event: &event.event,
257 };
258 let mut buf = serde_json::to_vec(&envelope)?;
259 buf.push(b'\n');
260 let mut writer = self.writer.lock().unwrap_or_else(|e| e.into_inner());
261 writer.write_all(&buf)?;
262 if self.flush_each_event {
263 writer.flush()?;
264 }
265 Ok(())
266 })();
267 self.record_result(result);
268 }
269}
270
271#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
273pub struct UsageTotals {
274 pub input_tokens: u64,
276 pub output_tokens: u64,
278 pub reasoning_tokens: u64,
280 pub cached_input_tokens: u64,
282 pub cache_write_input_tokens: u64,
284}
285
286#[derive(Clone, Debug, Default, PartialEq)]
288pub struct CostTotals {
289 pub amount: f64,
291 pub currency: Option<String>,
293}
294
295#[derive(Clone, Debug, Default, PartialEq)]
299pub struct UsageSummary {
300 pub events_seen: usize,
302 pub usage_events_seen: usize,
305 pub turn_results_seen: usize,
307 pub totals: UsageTotals,
309 pub cost: Option<CostTotals>,
311}
312
313#[derive(Default)]
337pub struct UsageReporter {
338 summary: std::sync::Mutex<UsageSummary>,
339}
340
341impl UsageReporter {
342 pub fn new() -> Self {
344 Self::default()
345 }
346
347 pub fn summary(&self) -> UsageSummary {
349 self.summary
350 .lock()
351 .unwrap_or_else(|e| e.into_inner())
352 .clone()
353 }
354
355 fn absorb(summary: &mut UsageSummary, usage: &Usage) {
356 summary.usage_events_seen += 1;
357 if let Some(tokens) = &usage.tokens {
358 summary.totals.input_tokens += tokens.input_tokens;
359 summary.totals.output_tokens += tokens.output_tokens;
360 summary.totals.reasoning_tokens += tokens.reasoning_tokens.unwrap_or_default();
361 summary.totals.cached_input_tokens += tokens.cached_input_tokens.unwrap_or_default();
362 summary.totals.cache_write_input_tokens +=
363 tokens.cache_write_input_tokens.unwrap_or_default();
364 }
365 if let Some(cost) = &usage.cost {
366 let totals = summary.cost.get_or_insert_with(CostTotals::default);
367 totals.amount += cost.amount;
368 if totals.currency.is_none() {
369 totals.currency = Some(cost.currency.clone());
370 }
371 }
372 }
373}
374
375impl LoopObserver for UsageReporter {
376 fn handle_event(&self, event: ObservedEvent) {
377 let event = event.event;
378 let mut summary = self.summary.lock().unwrap_or_else(|e| e.into_inner());
379 summary.events_seen += 1;
380 match event {
381 AgentEvent::UsageUpdated(usage) => Self::absorb(&mut summary, &usage),
382 AgentEvent::TurnFinished(TurnResult {
383 usage: Some(usage), ..
384 }) => {
385 summary.turn_results_seen += 1;
386 Self::absorb(&mut summary, &usage);
387 }
388 AgentEvent::TurnFinished(_) => {
389 summary.turn_results_seen += 1;
390 }
391 _ => {}
392 }
393 }
394}
395
396#[derive(Clone, Debug, Default, PartialEq)]
401pub struct TranscriptView {
402 pub items: Vec<Item>,
404}
405
406#[derive(Default)]
427pub struct TranscriptReporter {
428 transcript: std::sync::Mutex<TranscriptView>,
429}
430
431impl TranscriptReporter {
432 pub fn new() -> Self {
434 Self::default()
435 }
436
437 pub fn transcript(&self) -> TranscriptView {
439 self.transcript
440 .lock()
441 .unwrap_or_else(|e| e.into_inner())
442 .clone()
443 }
444}
445
446impl LoopObserver for TranscriptReporter {
447 fn handle_event(&self, event: ObservedEvent) {
448 let event = event.event;
449 let mut transcript = self.transcript.lock().unwrap_or_else(|e| e.into_inner());
450 match event {
451 AgentEvent::InputAccepted { items, .. } => {
452 transcript.items.extend(items);
453 }
454 AgentEvent::TurnFinished(result) => {
455 transcript.items.extend(result.items);
456 }
457 _ => {}
458 }
459 }
460}
461
462pub struct StdoutReporter<W> {
482 writer: std::sync::Mutex<W>,
483 show_usage: bool,
484 errors: std::sync::Mutex<Vec<ReportError>>,
485}
486
487impl<W> StdoutReporter<W>
488where
489 W: Write,
490{
491 pub fn new(writer: W) -> Self {
493 Self {
494 writer: std::sync::Mutex::new(writer),
495 show_usage: true,
496 errors: std::sync::Mutex::new(Vec::new()),
497 }
498 }
499
500 pub fn with_usage(mut self, show_usage: bool) -> Self {
502 self.show_usage = show_usage;
503 self
504 }
505
506 pub fn take_errors(&self) -> Vec<ReportError> {
508 std::mem::take(&mut *self.errors.lock().unwrap_or_else(|e| e.into_inner()))
509 }
510
511 fn record_result(&self, result: Result<(), ReportError>) {
512 if let Err(error) = result {
513 self.errors
514 .lock()
515 .unwrap_or_else(|e| e.into_inner())
516 .push(error);
517 }
518 }
519}
520
521impl<W> LoopObserver for StdoutReporter<W>
522where
523 W: Write + Send,
524{
525 fn handle_event(&self, event: ObservedEvent) {
526 let event = event.event;
527 let result = (|| -> Result<(), ReportError> {
528 let mut buf: Vec<u8> = Vec::new();
529 write_stdout_event(&mut buf, &event, self.show_usage)?;
530 let mut writer = self.writer.lock().unwrap_or_else(|e| e.into_inner());
531 writer.write_all(&buf)?;
532 writer.flush()?;
533 Ok(())
534 })();
535 self.record_result(result);
536 }
537}
538
539fn write_stdout_event<W>(
540 writer: &mut W,
541 event: &AgentEvent,
542 show_usage: bool,
543) -> Result<(), ReportError>
544where
545 W: Write,
546{
547 match event {
548 AgentEvent::RunStarted { session_id } => {
549 writeln!(writer, "[run] started session={session_id}")?;
550 }
551 AgentEvent::TurnStarted {
552 session_id,
553 turn_id,
554 } => {
555 writeln!(writer, "[turn] started session={session_id} turn={turn_id}")?;
556 }
557 AgentEvent::InputAccepted { items, .. } => {
558 writeln!(writer, "[input] accepted items={}", items.len())?;
559 }
560 AgentEvent::ContentDelta(delta) => {
561 writeln!(writer, "[delta] {delta:?}")?;
562 }
563 AgentEvent::ToolCallRequested(call) => {
564 writeln!(writer, "[tool] call {} {}", call.name, call.input)?;
565 }
566 AgentEvent::ToolExecutionStarted(call) => {
567 writeln!(writer, "[tool] started {} {}", call.name, call.id)?;
568 }
569 AgentEvent::ToolExecutionProgress(result) => {
570 writeln!(
571 writer,
572 "[tool] progress call_id={} is_error={}",
573 result.call_id, result.is_error
574 )?;
575 }
576 AgentEvent::ToolResultReceived(result) => {
577 writeln!(
578 writer,
579 "[tool] result call_id={} is_error={}",
580 result.call_id, result.is_error
581 )?;
582 }
583 AgentEvent::ApprovalRequired(request) => {
584 writeln!(
585 writer,
586 "[approval] {} {:?}",
587 request.summary, request.reason
588 )?;
589 }
590 AgentEvent::ApprovalResolved { approved } => {
591 writeln!(writer, "[approval] resolved approved={approved}")?;
592 }
593 AgentEvent::ToolCatalogChanged(event) => {
594 writeln!(
595 writer,
596 "[tools] catalog changed source={} added={} removed={} changed={}",
597 event.source,
598 event.added.len(),
599 event.removed.len(),
600 event.changed.len()
601 )?;
602 }
603 AgentEvent::MutationStarted {
604 turn_id,
605 mutator,
606 point,
607 ..
608 } => {
609 writeln!(
610 writer,
611 "[mutation] started turn={} mutator={mutator} point={point:?}",
612 turn_id
613 .as_ref()
614 .map(ToString::to_string)
615 .unwrap_or_else(|| "none".into()),
616 )?;
617 }
618 AgentEvent::MutationFinished {
619 turn_id,
620 mutator,
621 dirty,
622 ..
623 } => {
624 writeln!(
625 writer,
626 "[mutation] finished turn={} mutator={mutator} dirty={dirty}",
627 turn_id
628 .as_ref()
629 .map(ToString::to_string)
630 .unwrap_or_else(|| "none".into()),
631 )?;
632 }
633 AgentEvent::UsageUpdated(usage) if show_usage => {
634 writeln!(writer, "[usage] {}", format_usage(usage))?;
635 }
636 AgentEvent::UsageUpdated(_) => {}
637 AgentEvent::Warning { message } => {
638 writeln!(writer, "[warning] {message}")?;
639 }
640 AgentEvent::RunFailed { message } => {
641 writeln!(writer, "[error] {message}")?;
642 }
643 AgentEvent::TurnFinished(result) => {
644 writeln!(
645 writer,
646 "[turn] finished reason={:?} items={}",
647 result.finish_reason,
648 result.items.len()
649 )?;
650 for item in &result.items {
651 write_item_summary(writer, item)?;
652 }
653 if show_usage && let Some(usage) = &result.usage {
654 writeln!(writer, "[usage] {}", format_usage(usage))?;
655 }
656 }
657 _ => {}
658 }
659
660 writer.flush()?;
661 Ok(())
662}
663
664fn write_item_summary<W>(writer: &mut W, item: &Item) -> Result<(), ReportError>
665where
666 W: Write,
667{
668 writeln!(writer, " [{}]", item_kind_name(item.kind))?;
669 for part in &item.parts {
670 match part {
671 Part::Text(text) => writeln!(writer, " [text] {}", text.text)?,
672 Part::Reasoning(reasoning) => {
673 if let Some(summary) = &reasoning.summary {
674 writeln!(writer, " [reasoning] {summary}")?;
675 } else {
676 writeln!(writer, " [reasoning]")?;
677 }
678 }
679 Part::ToolCall(call) => {
680 writeln!(writer, " [tool-call] {} {}", call.name, call.input)?
681 }
682 Part::ToolResult(result) => writeln!(
683 writer,
684 " [tool-result] call={} error={}",
685 result.call_id, result.is_error
686 )?,
687 Part::Structured(value) => writeln!(writer, " [structured] {}", value.value)?,
688 Part::Media(media) => writeln!(
689 writer,
690 " [media] {:?} {}",
691 media.modality, media.mime_type
692 )?,
693 Part::File(file) => writeln!(
694 writer,
695 " [file] {}",
696 file.name.as_deref().unwrap_or("<unnamed>")
697 )?,
698 Part::Custom(custom) => writeln!(writer, " [custom] {}", custom.kind)?,
699 }
700 }
701 Ok(())
702}
703
704fn item_kind_name(kind: ItemKind) -> &'static str {
705 match kind {
706 ItemKind::System => "system",
707 ItemKind::Developer => "developer",
708 ItemKind::User => "user",
709 ItemKind::Assistant => "assistant",
710 ItemKind::Tool => "tool",
711 ItemKind::Context => "context",
712 ItemKind::Notification => "notification",
713 }
714}
715
716fn format_usage(usage: &Usage) -> String {
717 match &usage.tokens {
718 Some(TokenUsage {
719 input_tokens,
720 output_tokens,
721 reasoning_tokens,
722 cached_input_tokens,
723 cache_write_input_tokens,
724 }) => format!(
725 "input={} output={} reasoning={} cached_input={} cache_write_input={}",
726 input_tokens,
727 output_tokens,
728 reasoning_tokens.unwrap_or_default(),
729 cached_input_tokens.unwrap_or_default(),
730 cache_write_input_tokens.unwrap_or_default()
731 ),
732 None => "no token usage".into(),
733 }
734}
735
736#[cfg(test)]
737mod tests {
738 use super::*;
739 use agentkit_core::{FinishReason, MetadataMap, SessionId, TextPart};
740 use agentkit_loop::TurnResult;
741
742 #[test]
743 fn usage_reporter_accumulates_usage_events_and_turn_results() {
744 let reporter = UsageReporter::new();
745
746 reporter.handle_event(observed(AgentEvent::UsageUpdated(Usage {
747 tokens: Some(TokenUsage {
748 input_tokens: 10,
749 output_tokens: 5,
750 reasoning_tokens: Some(2),
751 cached_input_tokens: Some(1),
752 cache_write_input_tokens: Some(7),
753 }),
754 cost: None,
755 metadata: MetadataMap::new(),
756 })));
757
758 reporter.handle_event(observed(AgentEvent::TurnFinished(TurnResult {
759 turn_id: "turn-1".into(),
760 finish_reason: FinishReason::Completed,
761 items: Vec::new(),
762 usage: Some(Usage {
763 tokens: Some(TokenUsage {
764 input_tokens: 3,
765 output_tokens: 4,
766 reasoning_tokens: Some(1),
767 cached_input_tokens: None,
768 cache_write_input_tokens: None,
769 }),
770 cost: None,
771 metadata: MetadataMap::new(),
772 }),
773 metadata: MetadataMap::new(),
774 })));
775
776 let summary = reporter.summary();
777 assert_eq!(summary.events_seen, 2);
778 assert_eq!(summary.usage_events_seen, 2);
779 assert_eq!(summary.turn_results_seen, 1);
780 assert_eq!(summary.totals.input_tokens, 13);
781 assert_eq!(summary.totals.output_tokens, 9);
782 assert_eq!(summary.totals.reasoning_tokens, 3);
783 assert_eq!(summary.totals.cached_input_tokens, 1);
784 assert_eq!(summary.totals.cache_write_input_tokens, 7);
785 }
786
787 #[test]
788 fn transcript_reporter_tracks_inputs_and_outputs() {
789 let reporter = TranscriptReporter::new();
790
791 reporter.handle_event(observed(AgentEvent::InputAccepted {
792 session_id: SessionId::new("session-1"),
793 items: vec![Item {
794 id: None,
795 kind: ItemKind::User,
796 parts: vec![Part::Text(TextPart {
797 text: "hello".into(),
798 metadata: MetadataMap::new(),
799 })],
800 metadata: MetadataMap::new(),
801 usage: None,
802 finish_reason: None,
803 created_at: None,
804 }],
805 }));
806
807 reporter.handle_event(observed(AgentEvent::TurnFinished(TurnResult {
808 turn_id: "turn-1".into(),
809 finish_reason: FinishReason::Completed,
810 items: vec![Item {
811 id: None,
812 kind: ItemKind::Assistant,
813 parts: vec![Part::Text(TextPart {
814 text: "hi".into(),
815 metadata: MetadataMap::new(),
816 })],
817 metadata: MetadataMap::new(),
818 usage: None,
819 finish_reason: None,
820 created_at: None,
821 }],
822 usage: None,
823 metadata: MetadataMap::new(),
824 })));
825
826 assert_eq!(reporter.transcript().items.len(), 2);
827 assert_eq!(reporter.transcript().items[0].kind, ItemKind::User);
828 assert_eq!(reporter.transcript().items[1].kind, ItemKind::Assistant);
829 }
830
831 #[test]
832 fn jsonl_reporter_serializes_event_envelopes() {
833 let reporter = JsonlReporter::new(Vec::new());
834 reporter.handle_event(observed(AgentEvent::RunStarted {
835 session_id: SessionId::new("session-1"),
836 }));
837
838 let output = String::from_utf8(reporter.into_inner()).unwrap();
839 assert!(output.contains("\"RunStarted\""));
840 assert!(output.contains("session-1"));
841 assert!(!output.contains("\"session_id\":\"s1\""));
842 }
843
844 #[test]
845 fn jsonl_reporter_can_include_observed_session_id() {
846 let reporter = JsonlReporter::new(Vec::new()).with_session_id(true);
847 reporter.handle_event(observed(AgentEvent::ContentDelta(
848 agentkit_core::Delta::AppendText {
849 part_id: "part-1".into(),
850 chunk: "hello".into(),
851 },
852 )));
853
854 let output = String::from_utf8(reporter.into_inner()).unwrap();
855 assert!(output.contains("\"session_id\":\"s1\""));
856 }
857
858 fn run_started_event() -> AgentEvent {
859 AgentEvent::RunStarted {
860 session_id: SessionId::new("s1"),
861 }
862 }
863
864 fn observed(event: AgentEvent) -> ObservedEvent {
865 ObservedEvent {
866 session_id: std::sync::Arc::new(SessionId::new("s1")),
867 event,
868 }
869 }
870
871 #[test]
872 fn buffered_reporter_flushes_at_capacity() {
873 let reporter = BufferedReporter::new(UsageReporter::new(), 2);
874 reporter.handle_event(observed(run_started_event()));
875 assert_eq!(reporter.pending(), 1);
876 assert_eq!(reporter.inner().summary().events_seen, 0);
877
878 reporter.handle_event(observed(run_started_event()));
879 assert_eq!(reporter.pending(), 0);
880 assert_eq!(reporter.inner().summary().events_seen, 2);
881 }
882
883 #[test]
884 fn buffered_reporter_manual_flush() {
885 let reporter = BufferedReporter::new(UsageReporter::new(), 0);
886 reporter.handle_event(observed(run_started_event()));
887 reporter.handle_event(observed(run_started_event()));
888 assert_eq!(reporter.pending(), 2);
889
890 reporter.flush();
891 assert_eq!(reporter.pending(), 0);
892 assert_eq!(reporter.inner().summary().events_seen, 2);
893 }
894
895 #[test]
896 fn buffered_reporter_flushes_on_drop() {
897 let inner = {
898 let reporter = BufferedReporter::new(UsageReporter::new(), 100);
899 reporter.handle_event(observed(run_started_event()));
900 reporter.handle_event(observed(run_started_event()));
901 assert_eq!(reporter.inner().summary().events_seen, 0);
902 assert_eq!(reporter.pending(), 2);
905 reporter
906 };
907 assert_eq!(inner.inner().summary().events_seen, 0);
911 }
914
915 #[test]
916 fn channel_reporter_delivers_events() {
917 let (reporter, rx) = ChannelReporter::pair();
918 reporter.handle_event(observed(run_started_event()));
919 reporter.handle_event(observed(run_started_event()));
920
921 let events: Vec<_> = rx.try_iter().collect();
922 assert_eq!(events.len(), 2);
923 assert_eq!(events[0].session_id.0, "s1");
924 }
925
926 #[test]
927 fn channel_reporter_survives_dropped_receiver() {
928 let (reporter, rx) = ChannelReporter::pair();
929 drop(rx);
930 reporter.handle_event(observed(run_started_event()));
932 }
933
934 #[test]
935 fn channel_reporter_fallible_returns_error_on_dropped_receiver() {
936 let (reporter, rx) = ChannelReporter::pair();
937 drop(rx);
938
939 let result = reporter.try_handle_event(&observed(run_started_event()));
940 assert!(matches!(result, Err(ReportError::ChannelSend)));
941 }
942
943 #[test]
944 fn policy_reporter_ignore_swallows_errors() {
945 let (reporter, rx) = ChannelReporter::pair();
946 drop(rx);
947 let policy = PolicyReporter::new(reporter, FailurePolicy::Ignore);
948 policy.handle_event(observed(run_started_event()));
949 assert!(policy.take_errors().is_empty());
950 }
951
952 #[test]
953 fn policy_reporter_accumulate_collects_errors() {
954 let (reporter, rx) = ChannelReporter::pair();
955 drop(rx);
956 let policy = PolicyReporter::new(reporter, FailurePolicy::Accumulate);
957 policy.handle_event(observed(run_started_event()));
958 policy.handle_event(observed(run_started_event()));
959
960 let errors = policy.take_errors();
961 assert_eq!(errors.len(), 2);
962 assert!(matches!(errors[0], ReportError::ChannelSend));
963 }
964
965 #[test]
966 #[should_panic(expected = "reporter error: channel send failed")]
967 fn policy_reporter_fail_fast_panics() {
968 let (reporter, rx) = ChannelReporter::pair();
969 drop(rx);
970 let policy = PolicyReporter::new(reporter, FailurePolicy::FailFast);
971 policy.handle_event(observed(run_started_event()));
972 }
973}