harness_loop/telemetry.rs
1//! `TelemetryHook` — maps the agent's lifecycle [`Event`] stream onto
2//! structured `tracing` spans and events, so a run becomes observable in any
3//! `tracing` subscriber.
4//!
5//! Why `tracing` and not a hard OpenTelemetry dependency? Because `tracing` is
6//! the idiomatic Rust instrumentation seam: the library emits spans + events,
7//! and the *binary* chooses the exporter. Attach
8//! [`tracing-opentelemetry`](https://docs.rs/tracing-opentelemetry) with an
9//! OTLP pipeline and every span below is exported to Jaeger / Tempo / any OTLP
10//! backend with **zero changes here**; attach `tracing_subscriber::fmt().json()`
11//! and you get newline-delimited JSON for log pipelines. One instrumentation,
12//! many backends.
13//!
14//! Span/event shape (target `harness.telemetry`):
15//!
16//! ```text
17//! agent_run (span, fields: source)
18//! ├─ run.start
19//! ├─ iter (iter)
20//! ├─ model.complete (input_tokens, output_tokens, cached_input_tokens, tool_calls, stop)
21//! ├─ tool.call (tool, ok, duration_ms)
22//! ├─ sensor (sensor, signals)
23//! ├─ compact (stage)
24//! ├─ budget.warning (ratio)
25//! └─ run.end
26//! ```
27//!
28//! Wire it like any hook:
29//! ```ignore
30//! let loop_ = AgentLoop::new(model).with_hook(std::sync::Arc::new(TelemetryHook::new()));
31//! ```
32
33use harness_core::{Event, Hook, HookOutcome, World};
34use std::collections::HashMap;
35use std::sync::Mutex;
36use std::time::Instant;
37
38/// Emits a span per run and a structured event per model call, tool call,
39/// sensor, compaction, and budget warning. See the module docs for the OTLP
40/// bridge.
41pub struct TelemetryHook {
42 /// The current run's span. Events are recorded inside it so an OTLP exporter
43 /// nests them under one trace.
44 run: Mutex<Option<tracing::Span>>,
45 /// `call_id -> dispatch start`, so `tool.call` can report a duration.
46 tool_starts: Mutex<HashMap<String, Instant>>,
47}
48
49impl TelemetryHook {
50 pub fn new() -> Self {
51 Self {
52 run: Mutex::new(None),
53 tool_starts: Mutex::new(HashMap::new()),
54 }
55 }
56
57 /// Run `f` inside the current run span (if any), so its events attach to the
58 /// run's trace. Falls back to the ambient subscriber if no run is active.
59 fn in_run<F: FnOnce()>(&self, f: F) {
60 let guard = self.run.lock().unwrap();
61 match &*guard {
62 Some(span) => span.in_scope(f),
63 None => f(),
64 }
65 }
66}
67
68impl Default for TelemetryHook {
69 fn default() -> Self {
70 Self::new()
71 }
72}
73
74impl Hook for TelemetryHook {
75 fn name(&self) -> &str {
76 "telemetry"
77 }
78 fn matches(&self, _ev: &Event<'_>) -> bool {
79 true
80 }
81
82 fn fire(&self, ev: &Event<'_>, _world: &mut World) -> HookOutcome {
83 match ev {
84 Event::SessionStart { source } => {
85 let span = tracing::info_span!(
86 target: "harness.telemetry",
87 "agent_run",
88 source = format!("{source:?}")
89 );
90 span.in_scope(|| {
91 tracing::info!(target: "harness.telemetry", event = "run.start");
92 });
93 *self.run.lock().unwrap() = Some(span);
94 }
95 Event::Heartbeat { iter } => self.in_run(|| {
96 tracing::info!(target: "harness.telemetry", event = "iter", iter = *iter);
97 }),
98 Event::PostModel { out } => self.in_run(|| {
99 tracing::info!(
100 target: "harness.telemetry",
101 event = "model.complete",
102 input_tokens = out.usage.input_tokens,
103 output_tokens = out.usage.output_tokens,
104 cached_input_tokens = out.usage.cached_input_tokens,
105 tool_calls = out.tool_calls.len(),
106 stop = format!("{:?}", out.stop_reason),
107 );
108 }),
109 Event::PreToolUse { action } => {
110 self.tool_starts
111 .lock()
112 .unwrap()
113 .insert(action.call_id.clone(), Instant::now());
114 }
115 Event::PostToolUse { action, result } => {
116 let duration_ms = self
117 .tool_starts
118 .lock()
119 .unwrap()
120 .remove(&action.call_id)
121 .map(|s| s.elapsed().as_millis() as u64)
122 .unwrap_or(0);
123 self.in_run(|| {
124 tracing::info!(
125 target: "harness.telemetry",
126 event = "tool.call",
127 tool = %action.tool,
128 ok = result.ok,
129 duration_ms,
130 );
131 });
132 }
133 Event::PostSensor { sensor, signals } => self.in_run(|| {
134 tracing::debug!(
135 target: "harness.telemetry",
136 event = "sensor",
137 sensor = %sensor,
138 signals = signals.len(),
139 );
140 }),
141 Event::PostCompact { stage } => self.in_run(|| {
142 tracing::debug!(
143 target: "harness.telemetry",
144 event = "compact",
145 stage = format!("{stage:?}"),
146 );
147 }),
148 Event::BudgetWarning { ratio } => self.in_run(|| {
149 tracing::warn!(
150 target: "harness.telemetry",
151 event = "budget.warning",
152 ratio = *ratio,
153 );
154 }),
155 Event::SessionEnd => {
156 self.in_run(|| {
157 tracing::info!(target: "harness.telemetry", event = "run.end");
158 });
159 *self.run.lock().unwrap() = None;
160 }
161 _ => {}
162 }
163 HookOutcome::Allow
164 }
165}