//! One-shot backfill of historical per-agent turn stats into the OTEL //! collector as `claude_code.token.usage` metrics, so pre-OTEL history //! shows up alongside the live telemetry stream. Driven by //! `hivectl migrate-stats`. //! //! Reads each agent's `hyperhive-turn-stats.sqlite` (host path via //! `Coordinator::agent_harness_dir`) and turns the four per-turn token //! columns into **cumulative** `claude_code.token.usage` counter samples: //! one series per `(type, model)`, value = the running total at each //! turn's end time. Cumulative — not delta — because Prometheus/Mimir //! family backends silently drop delta sums (the same reason the live //! export forces cumulative temporality in harness-base.nix). The samples //! are sent as OTLP/HTTP JSON to `/v1/metrics`. Every //! datapoint carries `hyperhive-migration="true"` so the backfill is //! distinguishable from the live stream. //! //! Timestamps: `turn_stats.{started_at,ended_at}` are epoch *seconds* //! (see `hive-ag3nt` turn.rs), scaled to nanoseconds for OTLP. use anyhow::{Context, Result}; use rusqlite::Connection; use serde_json::{Value, json}; use std::collections::HashMap; use std::path::Path; use crate::coordinator::Coordinator; const METRIC: &str = "claude_code.token.usage"; /// Cap datapoints per POST so a long-lived agent's history is chunked /// into reasonably-sized OTLP requests rather than one huge body. const MAX_DATAPOINTS_PER_POST: usize = 1000; /// `claude_code.token.usage` `type` attribute value paired with the /// `turn_stats` column it maps to (column index in the SELECT below). const TOKEN_TYPES: &[(&str, usize)] = &[ ("input", 3), ("output", 4), ("cacheRead", 5), ("cacheCreation", 6), ]; /// Run the migration. `endpoint` falls back to `OTEL_EXPORTER_OTLP_ENDPOINT` /// then `HYPERHIVE_OTEL_ENDPOINT`. `only_agent` limits to one agent. /// `dry_run` reports counts without sending. /// /// # Errors /// Returns an error when no endpoint is resolvable, a stats db can't be /// read, or the collector rejects a POST. pub async fn run( endpoint: Option, only_agent: Option, dry_run: bool, ) -> Result<()> { let endpoint = endpoint .or_else(|| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok()) .or_else(|| std::env::var("HYPERHIVE_OTEL_ENDPOINT").ok()) .filter(|s| !s.is_empty()) .context("no OTLP endpoint: pass --endpoint or set OTEL_EXPORTER_OTLP_ENDPOINT")?; let url = format!("{}/v1/metrics", endpoint.trim_end_matches('/')); let hive = std::env::var("HYPERHIVE_HIVE_NAME").unwrap_or_else(|_| "unknown".to_owned()); let swarm = std::env::var("HYPERHIVE_SWARM_NAME").unwrap_or_else(|_| "unknown".to_owned()); let agents = match only_agent { Some(a) => vec![a], None => Coordinator::kept_state_names(), }; let client = reqwest::Client::new(); let mut total_points = 0usize; let mut hit_agents = 0usize; for agent in &agents { let db = Coordinator::agent_harness_dir(agent).join("hyperhive-turn-stats.sqlite"); if !db.exists() { continue; } let points = collect_agent_points(&db).with_context(|| format!("read turn_stats for {agent}"))?; if points.is_empty() { println!("{agent}: no token rows, skipping"); continue; } hit_agents += 1; total_points += points.len(); println!("{agent}: {} datapoints", points.len()); if dry_run { continue; } for chunk in points.chunks(MAX_DATAPOINTS_PER_POST) { let payload = build_payload(agent, &hive, &swarm, chunk); post(&client, &url, &payload) .await .with_context(|| format!("POST metrics for {agent}"))?; } } println!( "migrate-stats: {hit_agents} agent(s) with stats, {total_points} datapoints {}", if dry_run { "(dry-run, nothing sent)" } else { "sent" } ); Ok(()) } /// Build cumulative `claude_code.token.usage` datapoints for one agent's /// stats db. One running-total series per `(type, model)`; a datapoint is /// emitted whenever that series' value changes (token delta > 0). fn collect_agent_points(db: &Path) -> Result> { let conn = Connection::open_with_flags(db, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY) .context("open stats db read-only")?; let mut stmt = conn.prepare( "SELECT started_at, ended_at, model, input_tokens, output_tokens, \ cache_read_input_tokens, cache_creation_input_tokens \ FROM turn_stats ORDER BY ended_at ASC", )?; let rows = stmt.query_map([], |r| { Ok(( r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, String>(2)?, r.get::<_, i64>(3)?, r.get::<_, i64>(4)?, r.get::<_, i64>(5)?, r.get::<_, i64>(6)?, )) })?; // (type, model) -> (running total, series start time in nanos string) let mut series: HashMap<(&'static str, String), (u64, String)> = HashMap::new(); let mut points: Vec = Vec::new(); for row in rows { let (started, ended, model, inp, out, cache_read, cache_create) = row?; let cols = [inp, out, cache_read, cache_create]; let start_nanos = u64::try_from(started).unwrap_or(0) * 1_000_000_000; let time_nanos = (u64::try_from(ended).unwrap_or(0) * 1_000_000_000).to_string(); for &(type_, idx) in TOKEN_TYPES { let val = cols[idx - 3]; if val <= 0 { continue; } let entry = series .entry((type_, model.clone())) .or_insert_with(|| (0u64, start_nanos.to_string())); entry.0 += u64::try_from(val).unwrap_or(0); points.push(json!({ "asInt": entry.0.to_string(), "startTimeUnixNano": entry.1.clone(), "timeUnixNano": time_nanos, "attributes": [ {"key": "type", "value": {"stringValue": type_}}, {"key": "model", "value": {"stringValue": model.clone()}}, {"key": "hyperhive-migration", "value": {"stringValue": "true"}}, ], })); } } Ok(points) } /// Wrap a chunk of datapoints in an OTLP/HTTP JSON `ExportMetricsServiceRequest`. /// Resource attributes mirror the live export (`hive-serve-otel`): a fixed /// `service.name` plus this agent's name and the hive/swarm names. fn build_payload(agent: &str, hive: &str, swarm: &str, points: &[Value]) -> Value { json!({ "resourceMetrics": [{ "resource": {"attributes": [ {"key": "service.name", "value": {"stringValue": "hyperhive-agent"}}, {"key": "agent", "value": {"stringValue": agent}}, {"key": "hive", "value": {"stringValue": hive}}, {"key": "swarm", "value": {"stringValue": swarm}}, ]}, "scopeMetrics": [{ "scope": {"name": "hyperhive-migration"}, "metrics": [{ "name": METRIC, "unit": "tokens", "sum": { // 2 = AGGREGATION_TEMPORALITY_CUMULATIVE "aggregationTemporality": 2, "isMonotonic": true, "dataPoints": points, }, }], }], }], }) } /// POST one OTLP/HTTP JSON payload to `/v1/metrics`. async fn post(client: &reqwest::Client, url: &str, payload: &Value) -> Result<()> { let body = serde_json::to_vec(payload)?; let resp = client .post(url) .header("content-type", "application/json") .body(body) .send() .await .context("send to collector")?; let status = resp.status(); if !status.is_success() { let txt = resp.text().await.unwrap_or_default(); anyhow::bail!("collector returned {status}: {txt}"); } Ok(()) }