#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LogKind {
Output,
Runtime,
}
#[derive(Debug, Clone, PartialEq)]
pub enum TelemetryEvent {
RunStarted {
run_id: String,
agent_name: String,
model: Option<String>,
parent_run_id: Option<String>,
recovered: bool,
at_ms: i64,
},
StageEntered {
run_id: String,
stage_index: usize,
stage_name: String,
at_ms: i64,
},
StageExited {
run_id: String,
stage_index: usize,
stage_name: String,
prompt_tokens: usize,
completion_tokens: usize,
at_ms: i64,
},
InferenceCompleted {
run_id: String,
stage_name: String,
provider: String,
model: String,
latency_ms: u64,
prompt_tokens: usize,
completion_tokens: usize,
cached_tokens: usize,
success: bool,
},
ToolCallCompleted {
run_id: String,
stage_name: String,
tool_name: String,
batch_latency_ms: u64,
success: bool,
},
CompactionCompleted {
run_id: String,
stage_name: String,
success: bool,
},
RunCompleted {
run_id: String,
status: String,
prompt_tokens: usize,
completion_tokens: usize,
tool_calls: usize,
at_ms: i64,
},
Log {
run_id: String,
stage_index: usize,
kind: LogKind,
line: String,
},
}
impl TelemetryEvent {
#[must_use]
pub fn kind(&self) -> &'static str {
match self {
Self::RunStarted { .. } => "run_started",
Self::StageEntered { .. } => "stage_entered",
Self::StageExited { .. } => "stage_exited",
Self::InferenceCompleted { .. } => "inference_completed",
Self::ToolCallCompleted { .. } => "tool_call_completed",
Self::CompactionCompleted { .. } => "compaction_completed",
Self::RunCompleted { .. } => "run_completed",
Self::Log { .. } => "log",
}
}
pub fn run_id(&self) -> &str {
match self {
Self::RunStarted { run_id, .. }
| Self::StageEntered { run_id, .. }
| Self::StageExited { run_id, .. }
| Self::InferenceCompleted { run_id, .. }
| Self::ToolCallCompleted { run_id, .. }
| Self::CompactionCompleted { run_id, .. }
| Self::RunCompleted { run_id, .. }
| Self::Log { run_id, .. } => run_id,
}
}
}
pub trait TelemetrySink: Send + Sync {
fn emit(&self, event: TelemetryEvent);
fn force_flush(&self) {}
}
pub struct NoopSink;
impl TelemetrySink for NoopSink {
fn emit(&self, _event: TelemetryEvent) {}
}
#[derive(Default)]
pub struct MemorySink {
events: std::sync::Mutex<Vec<TelemetryEvent>>,
flushes: std::sync::atomic::AtomicUsize,
}
impl MemorySink {
pub fn events(&self) -> Vec<TelemetryEvent> {
self.events.lock().expect("telemetry event lock").clone()
}
pub fn flush_count(&self) -> usize {
self.flushes.load(std::sync::atomic::Ordering::SeqCst)
}
}
impl TelemetrySink for MemorySink {
fn emit(&self, event: TelemetryEvent) {
self.events
.lock()
.expect("telemetry event lock")
.push(event);
}
fn force_flush(&self) {
self.flushes
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn run_started(run_id: &str) -> TelemetryEvent {
TelemetryEvent::RunStarted {
run_id: run_id.to_string(),
agent_name: "coder".to_string(),
model: Some("claude-sonnet-5".to_string()),
parent_run_id: None,
recovered: false,
at_ms: 1_000,
}
}
#[test]
fn memory_sink_records_events_in_order() {
let sink = MemorySink::default();
sink.emit(run_started("r1"));
sink.emit(TelemetryEvent::StageEntered {
run_id: "r1".to_string(),
stage_index: 0,
stage_name: "plan".to_string(),
at_ms: 1_001,
});
let events = sink.events();
assert_eq!(events.len(), 2);
assert_eq!(events[0].kind(), "run_started");
assert_eq!(events[1].kind(), "stage_entered");
}
#[test]
fn memory_sink_counts_flushes() {
let sink = MemorySink::default();
assert_eq!(sink.flush_count(), 0);
sink.force_flush();
sink.force_flush();
assert_eq!(sink.flush_count(), 2);
}
#[test]
fn noop_sink_accepts_events_and_default_flush() {
let sink = NoopSink;
sink.emit(run_started("r1"));
let boxed: Box<dyn TelemetrySink> = Box::new(NoopSink);
boxed.force_flush();
}
#[test]
fn run_id_reaches_every_variant() {
let events = [
run_started("r1"),
TelemetryEvent::StageEntered {
run_id: "r1".to_string(),
stage_index: 0,
stage_name: "plan".to_string(),
at_ms: 0,
},
TelemetryEvent::StageExited {
run_id: "r1".to_string(),
stage_index: 0,
stage_name: "plan".to_string(),
prompt_tokens: 10,
completion_tokens: 5,
at_ms: 0,
},
TelemetryEvent::InferenceCompleted {
run_id: "r1".to_string(),
stage_name: "plan".to_string(),
provider: "anthropic".to_string(),
model: "claude-sonnet-5".to_string(),
latency_ms: 120,
prompt_tokens: 10,
completion_tokens: 5,
cached_tokens: 0,
success: true,
},
TelemetryEvent::ToolCallCompleted {
run_id: "r1".to_string(),
stage_name: "build".to_string(),
tool_name: "read_file".to_string(),
batch_latency_ms: 8,
success: true,
},
TelemetryEvent::CompactionCompleted {
run_id: "r1".to_string(),
stage_name: "build".to_string(),
success: true,
},
TelemetryEvent::RunCompleted {
run_id: "r1".to_string(),
status: "complete".to_string(),
prompt_tokens: 10,
completion_tokens: 5,
tool_calls: 1,
at_ms: 0,
},
TelemetryEvent::Log {
run_id: "r1".to_string(),
stage_index: 0,
kind: LogKind::Runtime,
line: "[Tokens: 10 in, 5 out]".to_string(),
},
];
let kinds: Vec<&str> = events.iter().map(TelemetryEvent::kind).collect();
assert_eq!(
kinds,
[
"run_started",
"stage_entered",
"stage_exited",
"inference_completed",
"tool_call_completed",
"compaction_completed",
"run_completed",
"log",
]
);
for event in &events {
assert_eq!(event.run_id(), "r1");
}
}
#[test]
fn event_clone_debug_and_eq() {
let event = run_started("r1");
let cloned = event.clone();
assert_eq!(event, cloned);
assert!(format!("{event:?}").contains("RunStarted"));
assert_ne!(LogKind::Output, LogKind::Runtime);
assert!(format!("{:?}", LogKind::Output).contains("Output"));
}
}