//! Per-turn analytics sink. One sqlite row per claude turn captures: //! identity (`model`, `wake_from`, `result_kind`), timing (`started_at`, //! `ended_at`, `duration_ms`), cost (token counts), behaviour (tool-call //! count + per-tool breakdown), and post-turn snapshot metrics //! (`open_threads_count`, `open_reminders_count`). //! //! Lives next to `hyperhive-events.sqlite` in the agent's state dir //! so the host-side state vacuum sweep can reach both. Schema is //! intentionally append-only — every column has a default so future //! additions don't break old readers; new columns land via //! `ALTER TABLE ... ADD COLUMN ... DEFAULT ...` in the migration //! block. //! //! Writes are best-effort: a failed insert logs a warning and lets //! the turn loop continue. The next turn either succeeds or the //! operator sees the journal trail. use std::path::{Path, PathBuf}; use std::sync::Mutex; use anyhow::{Context, Result}; use rusqlite::{Connection, params}; /// SQL bootstrap. CREATE TABLE IF NOT EXISTS so first-boot agents /// and existing ones converge on the same shape. The base table is /// fresh-install only; additive migrations land via `MIGRATIONS` /// below as try-and-ignore ALTERs so existing dbs catch up. const SCHEMA: &str = " CREATE TABLE IF NOT EXISTS turn_stats ( id INTEGER PRIMARY KEY AUTOINCREMENT, started_at INTEGER NOT NULL, ended_at INTEGER NOT NULL, duration_ms INTEGER NOT NULL, model TEXT NOT NULL, wake_from TEXT NOT NULL, input_tokens INTEGER NOT NULL DEFAULT 0, output_tokens INTEGER NOT NULL DEFAULT 0, cache_read_input_tokens INTEGER NOT NULL DEFAULT 0, cache_creation_input_tokens INTEGER NOT NULL DEFAULT 0, last_input_tokens INTEGER NOT NULL DEFAULT 0, last_output_tokens INTEGER NOT NULL DEFAULT 0, last_cache_read_input_tokens INTEGER NOT NULL DEFAULT 0, last_cache_creation_input_tokens INTEGER NOT NULL DEFAULT 0, tool_call_count INTEGER NOT NULL DEFAULT 0, tool_call_breakdown_json TEXT, open_threads_count INTEGER, open_reminders_count INTEGER, result_kind TEXT NOT NULL, note TEXT ); CREATE INDEX IF NOT EXISTS idx_turn_stats_started ON turn_stats (started_at DESC); "; /// Additive column migrations. Each runs unconditionally and ignores /// `duplicate column name` errors — sqlite < 3.35 lacks /// `ADD COLUMN IF NOT EXISTS`, so try-and-ignore is the portable path. /// New columns MUST carry a default so existing rows decode. const MIGRATIONS: &[&str] = &[ "ALTER TABLE turn_stats ADD COLUMN last_input_tokens INTEGER NOT NULL DEFAULT 0", "ALTER TABLE turn_stats ADD COLUMN last_output_tokens INTEGER NOT NULL DEFAULT 0", "ALTER TABLE turn_stats ADD COLUMN last_cache_read_input_tokens INTEGER NOT NULL DEFAULT 0", "ALTER TABLE turn_stats ADD COLUMN last_cache_creation_input_tokens INTEGER NOT NULL DEFAULT 0", ]; /// One row to be inserted. `Option`-wrapped fields default to NULL /// when the harness couldn't gather them (e.g. socket roundtrip for /// `open_threads` failed) so a partial row beats no row. #[derive(Debug, Clone)] pub struct TurnStatRow { pub started_at: i64, pub ended_at: i64, pub duration_ms: i64, pub model: String, pub wake_from: String, /// Cumulative across every inference in the turn (cost signal). pub input_tokens: u64, pub output_tokens: u64, pub cache_read_input_tokens: u64, pub cache_creation_input_tokens: u64, /// Last inference's usage — the actual context size at turn end. pub last_input_tokens: u64, pub last_output_tokens: u64, pub last_cache_read_input_tokens: u64, pub last_cache_creation_input_tokens: u64, pub tool_call_count: u64, /// Per-tool breakdown as JSON: `{"Read":12,"Bash":3,...}`. None /// when no tools were called (saves a sqlite write of `"{}"`). pub tool_call_breakdown_json: Option, pub open_threads_count: Option, pub open_reminders_count: Option, /// `"ok" | "failed" | "prompt_too_long"`. pub result_kind: &'static str, pub note: Option, } /// Thin sqlite wrapper. Cloning is cheap (Arc-shared connection). #[derive(Clone)] pub struct TurnStats { inner: std::sync::Arc>, } impl TurnStats { /// Open the per-agent stats db, creating the file + schema if /// missing. Returns `None` when the db can't be opened (read-only /// fs in tests, missing state dir) — the harness logs and /// continues without a sink rather than failing the turn loop. #[must_use] pub fn open_default() -> Option { let path = default_path(); match Self::open(&path) { Ok(s) => Some(s), Err(e) => { tracing::warn!( error = ?e, path = %path.display(), "turn_stats: open failed; per-turn analytics disabled" ); None } } } fn open(path: &Path) -> Result { if let Some(parent) = path.parent() { let _ = std::fs::create_dir_all(parent); } let conn = Connection::open(path) .with_context(|| format!("open turn_stats db {}", path.display()))?; conn.execute_batch(SCHEMA) .context("apply turn_stats schema")?; for stmt in MIGRATIONS { // Ignore "duplicate column name" — the migration already ran. // Any other error is logged but doesn't fail open() because the // base schema works and we'd rather keep the harness alive than // crash on an upgrade hiccup. if let Err(e) = conn.execute(stmt, []) { let msg = e.to_string(); if !msg.contains("duplicate column name") { tracing::warn!(error = %msg, stmt, "turn_stats migration failed"); } } } Ok(Self { inner: std::sync::Arc::new(Mutex::new(conn)), }) } /// Insert a row. Best-effort — logs + swallows errors so a sqlite /// hiccup (locked db, full disk) doesn't crash the harness. /// /// # Panics /// /// Panics if the internal lock is poisoned. pub fn record(&self, row: &TurnStatRow) { let conn = self.inner.lock().unwrap(); let res = conn.execute( "INSERT INTO turn_stats ( started_at, ended_at, duration_ms, model, wake_from, input_tokens, output_tokens, cache_read_input_tokens, cache_creation_input_tokens, last_input_tokens, last_output_tokens, last_cache_read_input_tokens, last_cache_creation_input_tokens, tool_call_count, tool_call_breakdown_json, open_threads_count, open_reminders_count, result_kind, note ) VALUES ( ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19 )", params![ row.started_at, row.ended_at, row.duration_ms, row.model, row.wake_from, i64::try_from(row.input_tokens).unwrap_or(i64::MAX), i64::try_from(row.output_tokens).unwrap_or(i64::MAX), i64::try_from(row.cache_read_input_tokens).unwrap_or(i64::MAX), i64::try_from(row.cache_creation_input_tokens).unwrap_or(i64::MAX), i64::try_from(row.last_input_tokens).unwrap_or(i64::MAX), i64::try_from(row.last_output_tokens).unwrap_or(i64::MAX), i64::try_from(row.last_cache_read_input_tokens).unwrap_or(i64::MAX), i64::try_from(row.last_cache_creation_input_tokens).unwrap_or(i64::MAX), i64::try_from(row.tool_call_count).unwrap_or(i64::MAX), row.tool_call_breakdown_json, row.open_threads_count .map(|n| i64::try_from(n).unwrap_or(i64::MAX)), row.open_reminders_count .map(|n| i64::try_from(n).unwrap_or(i64::MAX)), row.result_kind, row.note, ], ); if let Err(e) = res { tracing::warn!(error = ?e, "turn_stats: insert failed"); } } /// Token counts from the most recently inserted row, if any. /// Returns `(ctx, cost)` — both backfill `Bus` on startup so the /// per-agent web UI's ctx + cost badges paint with real numbers on /// cold load instead of waiting for the next `TokenUsageChanged` /// SSE event. Best-effort: any sqlite error returns `(None, None)`. /// /// Pre-migration rows (before the `last_*_tokens` columns existed) /// have last-inference zeros — those rows yield `ctx = None` so the /// badge stays empty until the next real turn rather than showing a /// misleading 0. /// # Panics /// /// Panics if the internal lock is poisoned. #[must_use] pub fn last_usage( &self, ) -> ( Option, Option, ) { let conn = self.inner.lock().unwrap(); conn.query_row( "SELECT input_tokens, output_tokens, cache_read_input_tokens, cache_creation_input_tokens, last_input_tokens, last_output_tokens, last_cache_read_input_tokens, last_cache_creation_input_tokens FROM turn_stats ORDER BY started_at DESC LIMIT 1", [], |row| { let g = |i: usize| -> rusqlite::Result { Ok(u64::try_from(row.get::<_, i64>(i)?).unwrap_or(0)) }; let cost = crate::events::TokenUsage { input_tokens: g(0)?, output_tokens: g(1)?, cache_read_input_tokens: g(2)?, cache_creation_input_tokens: g(3)?, }; let last = crate::events::TokenUsage { input_tokens: g(4)?, output_tokens: g(5)?, cache_read_input_tokens: g(6)?, cache_creation_input_tokens: g(7)?, }; let ctx = if last == crate::events::TokenUsage::default() { None } else { Some(last) }; Ok((ctx, Some(cost))) }, ) .unwrap_or((None, None)) } } fn default_path() -> PathBuf { crate::paths::state_dir().join("hyperhive-turn-stats.sqlite") }