use harness_core::{Event, Hook, HookOutcome, World};
use std::collections::HashMap;
use std::sync::Mutex;
use std::time::Instant;
pub struct TelemetryHook {
run: Mutex<Option<tracing::Span>>,
tool_starts: Mutex<HashMap<String, Instant>>,
model_start: Mutex<Option<Instant>>,
awaiting_first_token: Mutex<bool>,
totals: Mutex<RunTotals>,
}
#[derive(Default)]
struct RunTotals {
started: Option<Instant>,
input_tokens: u64,
output_tokens: u64,
cached_input_tokens: u64,
cache_write_input_tokens: u64,
model_calls: u64,
tool_calls: u64,
tool_failures: u64,
compactions: u64,
tokens_saved: u64,
model_ms: u64,
tool_ms: u64,
repeats: u64,
first_token_ms: u64,
}
impl TelemetryHook {
pub fn new() -> Self {
Self {
run: Mutex::new(None),
tool_starts: Mutex::new(HashMap::new()),
model_start: Mutex::new(None),
awaiting_first_token: Mutex::new(false),
totals: Mutex::new(RunTotals::default()),
}
}
fn in_run<F: FnOnce()>(&self, f: F) {
let guard = self.run.lock().unwrap();
match &*guard {
Some(span) => span.in_scope(f),
None => f(),
}
}
}
impl Default for TelemetryHook {
fn default() -> Self {
Self::new()
}
}
impl Hook for TelemetryHook {
fn name(&self) -> &str {
"telemetry"
}
fn matches(&self, _ev: &Event<'_>) -> bool {
true
}
fn fire(&self, ev: &Event<'_>, _world: &mut World) -> HookOutcome {
match ev {
Event::SessionStart { source } => {
let span = tracing::info_span!(
target: "harness.telemetry",
"agent_run",
"gen_ai.operation.name" = "invoke_agent",
source = format!("{source:?}")
);
span.in_scope(|| {
tracing::info!(target: "harness.telemetry", event = "run.start");
});
*self.run.lock().unwrap() = Some(span);
*self.totals.lock().unwrap() = RunTotals {
started: Some(Instant::now()),
..Default::default()
};
}
Event::PreModel { .. } => {
*self.model_start.lock().unwrap() = Some(Instant::now());
*self.awaiting_first_token.lock().unwrap() = true;
}
Event::ModelTokenDelta { .. } => {
let mut awaiting = self.awaiting_first_token.lock().unwrap();
if !*awaiting {
return HookOutcome::Allow;
}
*awaiting = false;
drop(awaiting);
let ttft = self
.model_start
.lock()
.unwrap()
.map(|s| s.elapsed().as_millis() as u64)
.unwrap_or(0);
{
let mut t = self.totals.lock().unwrap();
if t.first_token_ms == 0 {
t.first_token_ms = ttft;
}
}
self.in_run(|| {
tracing::info!(
target: "harness.telemetry",
event = "model.first_token",
ttft_ms = ttft,
);
});
}
Event::Heartbeat { iter } => self.in_run(|| {
tracing::info!(target: "harness.telemetry", event = "iter", iter = *iter);
}),
Event::PostModel { out } => self.in_run(|| {
let waited = self
.model_start
.lock()
.unwrap()
.take()
.map(|s| s.elapsed().as_millis() as u64)
.unwrap_or(0);
{
let mut t = self.totals.lock().unwrap();
t.model_calls += 1;
t.model_ms += waited;
t.input_tokens += out.usage.input_tokens as u64;
t.output_tokens += out.usage.output_tokens as u64;
t.cached_input_tokens += out.usage.cached_input_tokens as u64;
t.cache_write_input_tokens += out.usage.cache_write_input_tokens as u64;
}
let stop = format!("{:?}", out.stop_reason);
tracing::info!(
target: "harness.telemetry",
event = "model.complete",
"gen_ai.operation.name" = "chat",
"gen_ai.usage.input_tokens" = out.usage.input_tokens,
"gen_ai.usage.output_tokens" = out.usage.output_tokens,
"gen_ai.usage.cached_input_tokens" = out.usage.cached_input_tokens,
"gen_ai.usage.cache_write_input_tokens" = out.usage.cache_write_input_tokens,
"gen_ai.response.finish_reasons" = %stop,
input_tokens = out.usage.input_tokens,
output_tokens = out.usage.output_tokens,
cached_input_tokens = out.usage.cached_input_tokens,
cache_write_input_tokens = out.usage.cache_write_input_tokens,
tool_calls = out.tool_calls.len(),
stop = %stop,
duration_ms = waited,
);
}),
Event::PreToolUse { action } => {
self.tool_starts
.lock()
.unwrap()
.insert(action.call_id.clone(), Instant::now());
}
Event::PostToolUse { action, result } => {
let duration_ms = self
.tool_starts
.lock()
.unwrap()
.remove(&action.call_id)
.map(|s| s.elapsed().as_millis() as u64)
.unwrap_or(0);
let repeat = result
.content
.get("repeat_of_earlier_call")
.and_then(|v| v.as_bool())
.unwrap_or(false);
{
let mut t = self.totals.lock().unwrap();
t.tool_calls += 1;
t.tool_ms += duration_ms;
if repeat {
t.repeats += 1;
}
if !result.ok {
t.tool_failures += 1;
}
}
self.in_run(|| {
tracing::info!(
target: "harness.telemetry",
event = "tool.call",
"gen_ai.operation.name" = "execute_tool",
"gen_ai.tool.name" = %action.tool,
ok = result.ok,
duration_ms,
tool = %action.tool, );
});
}
Event::PostSensor { sensor, signals } => self.in_run(|| {
tracing::debug!(
target: "harness.telemetry",
event = "sensor",
sensor = %sensor,
signals = signals.len(),
);
}),
Event::PostCompact {
stage,
before,
after,
} => self.in_run(|| {
let saved = before.saturating_sub(*after);
{
let mut t = self.totals.lock().unwrap();
t.compactions += 1;
t.tokens_saved += saved as u64;
}
tracing::info!(
target: "harness.telemetry",
event = "compact",
stage = format!("{stage:?}"),
tokens_before = *before,
tokens_after = *after,
tokens_saved = saved,
);
}),
Event::BudgetWarning { ratio } => self.in_run(|| {
tracing::warn!(
target: "harness.telemetry",
event = "budget.warning",
ratio = *ratio,
);
}),
Event::SessionEnd => {
let t = std::mem::take(&mut *self.totals.lock().unwrap());
self.in_run(|| {
tracing::info!(
target: "harness.telemetry",
event = "run.end",
"gen_ai.usage.input_tokens" = t.input_tokens,
"gen_ai.usage.output_tokens" = t.output_tokens,
"gen_ai.usage.cached_input_tokens" = t.cached_input_tokens,
"gen_ai.usage.cache_write_input_tokens" = t.cache_write_input_tokens,
total_tokens = t.input_tokens + t.output_tokens,
model_calls = t.model_calls,
tool_calls = t.tool_calls,
tool_failures = t.tool_failures,
compactions = t.compactions,
tokens_saved = t.tokens_saved,
model_ms = t.model_ms,
tool_ms = t.tool_ms,
repeat_calls = t.repeats,
first_token_ms = t.first_token_ms,
duration_ms = t.started.map(|s| s.elapsed().as_millis() as u64).unwrap_or(0),
);
});
*self.run.lock().unwrap() = None;
}
_ => {}
}
HookOutcome::Allow
}
}