Skip to main content

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}