"""Dashboard telemetry history schema and queries, adapted from the old Athena dashboard.""" from __future__ import annotations import sqlite3 import threading import time from pathlib import Path from typing import Any HISTORY_INTERVAL = 15 DETAIL_RETENTION_DAYS = 21 class HistoryStore: """Small persistent telemetry store; never stores prompts or responses.""" TOKEN_KEYS = ("prompt_tokens_total", "prompt_tokens_cached_total", "tokens_predicted_total") def close(self) -> None: with self._lock: self._db.close() def __init__(self, path: Path) -> None: path.parent.mkdir(parents=True, exist_ok=True) self._db = sqlite3.connect(path, check_same_thread=False) self._db.row_factory = sqlite3.Row self._lock = threading.Lock() with self._db: self._db.executescript(""" PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL; CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT NOT NULL); CREATE TABLE IF NOT EXISTS samples_raw ( ts INTEGER NOT NULL, gpu_index INTEGER NOT NULL, gpu_util REAL, memory_used_mib REAL, temperature_c REAL, power_w REAL, profile TEXT, model TEXT ); CREATE INDEX IF NOT EXISTS idx_samples_raw_ts ON samples_raw(ts); CREATE TABLE IF NOT EXISTS samples_hourly ( hour_ts INTEGER NOT NULL, gpu_index INTEGER NOT NULL, gpu_util_sum REAL, gpu_util_max REAL, memory_used_sum REAL, memory_used_max REAL, temperature_sum REAL, temperature_max REAL, power_sum REAL, samples INTEGER NOT NULL, PRIMARY KEY(hour_ts, gpu_index) ); CREATE TABLE IF NOT EXISTS slot_samples_raw ( ts INTEGER NOT NULL, slot_id INTEGER NOT NULL, context_used REAL, context_total REAL, processing INTEGER NOT NULL DEFAULT 0, profile TEXT, model TEXT ); CREATE INDEX IF NOT EXISTS idx_slot_samples_raw_ts ON slot_samples_raw(ts); CREATE TABLE IF NOT EXISTS slot_samples_hourly ( hour_ts INTEGER NOT NULL, slot_id INTEGER NOT NULL, context_percent_sum REAL, context_percent_max REAL, context_used_max REAL, context_total_max REAL, busy_samples INTEGER NOT NULL, samples INTEGER NOT NULL, PRIMARY KEY(hour_ts, slot_id) ); CREATE TABLE IF NOT EXISTS token_totals ( id INTEGER PRIMARY KEY CHECK(id=1), prompt_tokens INTEGER NOT NULL DEFAULT 0, cached_tokens INTEGER NOT NULL DEFAULT 0, output_tokens INTEGER NOT NULL DEFAULT 0 ); INSERT OR IGNORE INTO token_totals(id) VALUES(1); CREATE TABLE IF NOT EXISTS token_hourly ( hour_ts INTEGER NOT NULL, profile TEXT NOT NULL, model TEXT NOT NULL, prompt_tokens INTEGER NOT NULL DEFAULT 0, cached_tokens INTEGER NOT NULL DEFAULT 0, output_tokens INTEGER NOT NULL DEFAULT 0, PRIMARY KEY(hour_ts, profile, model) ); CREATE TABLE IF NOT EXISTS model_events ( ts INTEGER NOT NULL, previous_profile TEXT, profile TEXT, previous_model TEXT, model TEXT ); CREATE INDEX IF NOT EXISTS idx_model_events_ts ON model_events(ts); """) def _meta(self, key: str) -> str | None: row = self._db.execute("SELECT value FROM meta WHERE key=?", (key,)).fetchone() return str(row[0]) if row else None def _set_meta(self, key: str, value: Any) -> None: self._db.execute( "INSERT INTO meta(key,value) VALUES(?,?) ON CONFLICT(key) DO UPDATE SET value=excluded.value", (key, str(value)), ) def record(self, snapshot: dict[str, Any]) -> None: ts = int(snapshot.get("timestamp") or time.time()) router = snapshot.get("router") or {} image = router.get("image") or {} image_active = image.get("phase") not in (None, "idle") profile = str("image" if image_active else (router.get("current_profile") or "unknown")) model = str((image.get("model") if image_active else (router.get("upstream") or {}).get("model")) or "unknown") metrics = ((router.get("llama_telemetry") or {}).get("metrics") or {}) runtime = snapshot.get("llama_runtime") or {} runtime_id = f"{runtime.get('pid', 'none')}:{runtime.get('model_file') or model}" hour = ts - ts % 3600 with self._lock, self._db: for gpu in snapshot.get("gpus") or []: if gpu.get("index") is None: continue self._db.execute( "INSERT INTO samples_raw VALUES(?,?,?,?,?,?,?,?)", (ts, int(gpu["index"]), gpu.get("gpu_percent"), gpu.get("memory_used_mib"), gpu.get("temperature_c"), gpu.get("power_w"), profile, model), ) slots = ((router.get("llama_telemetry") or {}).get("slots") or []) for slot in slots: if slot.get("id") is None: continue used = max(0, float(slot.get("context_used") or 0)) total = max(0, float(slot.get("n_ctx") or 0)) self._db.execute( "INSERT INTO slot_samples_raw VALUES(?,?,?,?,?,?,?)", (ts, int(slot["id"]), used, total, int(bool(slot.get("processing"))), profile, model), ) previous_runtime = self._meta("counter_runtime") deltas: list[int] = [] for key in self.TOKEN_KEYS: current = max(0, int(float(metrics.get(key) or 0))) previous = int(self._meta(f"counter_{key}") or 0) delta = current - previous if previous_runtime == runtime_id and current >= previous else current deltas.append(max(0, delta)) self._set_meta(f"counter_{key}", current) self._set_meta("counter_runtime", runtime_id) if any(deltas): self._db.execute( "UPDATE token_totals SET prompt_tokens=prompt_tokens+?, cached_tokens=cached_tokens+?, output_tokens=output_tokens+? WHERE id=1", deltas, ) self._db.execute( """INSERT INTO token_hourly VALUES(?,?,?,?,?,?) ON CONFLICT(hour_ts,profile,model) DO UPDATE SET prompt_tokens=prompt_tokens+excluded.prompt_tokens, cached_tokens=cached_tokens+excluded.cached_tokens, output_tokens=output_tokens+excluded.output_tokens""", (hour, profile, model, *deltas), ) previous_profile = self._meta("last_profile") previous_model = self._meta("last_model") if previous_profile is not None and (profile != previous_profile or model != previous_model): self._db.execute( "INSERT INTO model_events VALUES(?,?,?,?,?)", (ts, previous_profile, profile, previous_model, model), ) self._set_meta("last_profile", profile) self._set_meta("last_model", model) def compact(self) -> None: cutoff = int(time.time()) - DETAIL_RETENTION_DAYS * 86400 with self._lock, self._db: self._db.execute(""" INSERT INTO samples_hourly SELECT ts-ts%3600, gpu_index, SUM(gpu_util), MAX(gpu_util), SUM(memory_used_mib), MAX(memory_used_mib), SUM(temperature_c), MAX(temperature_c), SUM(power_w), COUNT(*) FROM samples_raw WHERE ts < ? GROUP BY ts-ts%3600, gpu_index ON CONFLICT(hour_ts,gpu_index) DO UPDATE SET gpu_util_sum=gpu_util_sum+excluded.gpu_util_sum, gpu_util_max=MAX(gpu_util_max,excluded.gpu_util_max), memory_used_sum=memory_used_sum+excluded.memory_used_sum, memory_used_max=MAX(memory_used_max,excluded.memory_used_max), temperature_sum=temperature_sum+excluded.temperature_sum, temperature_max=MAX(temperature_max,excluded.temperature_max), power_sum=power_sum+excluded.power_sum, samples=samples+excluded.samples """, (cutoff,)) self._db.execute("DELETE FROM samples_raw WHERE ts < ?", (cutoff,)) self._db.execute(""" INSERT INTO slot_samples_hourly SELECT ts-ts%3600, slot_id, SUM(CASE WHEN context_total>0 THEN 100.0*context_used/context_total ELSE 0 END), MAX(CASE WHEN context_total>0 THEN 100.0*context_used/context_total ELSE 0 END), MAX(context_used), MAX(context_total), SUM(processing), COUNT(*) FROM slot_samples_raw WHERE ts < ? GROUP BY ts-ts%3600, slot_id ON CONFLICT(hour_ts,slot_id) DO UPDATE SET context_percent_sum=context_percent_sum+excluded.context_percent_sum, context_percent_max=MAX(context_percent_max,excluded.context_percent_max), context_used_max=MAX(context_used_max,excluded.context_used_max), context_total_max=MAX(context_total_max,excluded.context_total_max), busy_samples=busy_samples+excluded.busy_samples, samples=samples+excluded.samples """, (cutoff,)) self._db.execute("DELETE FROM slot_samples_raw WHERE ts < ?", (cutoff,)) def query(self, range_name: str) -> dict[str, Any]: ranges = { "1h": (3600, 60), "24h": (86400, 300), "7d": (7 * 86400, 1800), "21d": (21 * 86400, 3600), "all": (0, 3600), } seconds, bucket = ranges.get(range_name, ranges["24h"]) now = int(time.time()) start = 0 if seconds == 0 else now - seconds detail_cutoff = now - DETAIL_RETENTION_DAYS * 86400 with self._lock: points: list[dict[str, Any]] = [] slot_points: list[dict[str, Any]] = [] if start < detail_cutoff: for row in self._db.execute(""" SELECT hour_ts ts,gpu_index,gpu_util_sum/samples gpu_util,gpu_util_max, memory_used_sum/samples memory_used_mib,memory_used_max, temperature_sum/samples temperature_c,temperature_max, power_sum/samples power_w FROM samples_hourly WHERE hour_ts>=? ORDER BY hour_ts,gpu_index """, (start,)): points.append(dict(row)) raw_start = max(start, detail_cutoff) for row in self._db.execute(f""" SELECT (ts/{bucket})*{bucket} ts,gpu_index,AVG(gpu_util) gpu_util,MAX(gpu_util) gpu_util_max, AVG(memory_used_mib) memory_used_mib,MAX(memory_used_mib) memory_used_max, AVG(temperature_c) temperature_c,MAX(temperature_c) temperature_max, AVG(power_w) power_w FROM samples_raw WHERE ts>=? GROUP BY (ts/{bucket}),gpu_index ORDER BY ts,gpu_index """, (raw_start,)): points.append(dict(row)) if start < detail_cutoff: for row in self._db.execute(""" SELECT hour_ts ts,slot_id,context_percent_sum/samples context_percent, context_percent_max,context_used_max context_used, context_total_max context_total, 1.0*busy_samples/samples busy_ratio FROM slot_samples_hourly WHERE hour_ts>=? ORDER BY hour_ts,slot_id """, (start,)): slot_points.append(dict(row)) for row in self._db.execute(f""" SELECT (ts/{bucket})*{bucket} ts,slot_id, AVG(CASE WHEN context_total>0 THEN 100.0*context_used/context_total ELSE 0 END) context_percent, MAX(CASE WHEN context_total>0 THEN 100.0*context_used/context_total ELSE 0 END) context_percent_max, MAX(context_used) context_used,MAX(context_total) context_total, AVG(processing) busy_ratio FROM slot_samples_raw WHERE ts>=? GROUP BY (ts/{bucket}),slot_id ORDER BY ts,slot_id """, (raw_start,)): slot_points.append(dict(row)) totals = dict(self._db.execute("SELECT * FROM token_totals WHERE id=1").fetchone()) range_tokens = dict(self._db.execute( "SELECT COALESCE(SUM(prompt_tokens),0) prompt_tokens, COALESCE(SUM(cached_tokens),0) cached_tokens, COALESCE(SUM(output_tokens),0) output_tokens FROM token_hourly WHERE hour_ts>=?", (start,), ).fetchone()) profile_usage = [dict(row) for row in self._db.execute(""" SELECT profile,model,SUM(prompt_tokens) prompt_tokens, SUM(cached_tokens) cached_tokens,SUM(output_tokens) output_tokens FROM token_hourly WHERE hour_ts>=? GROUP BY profile,model ORDER BY SUM(prompt_tokens+cached_tokens+output_tokens) DESC """, (start,))] events = [dict(row) for row in self._db.execute( "SELECT * FROM model_events WHERE ts>=? ORDER BY ts DESC LIMIT 50", (start,) )] raw_info = dict(self._db.execute( "SELECT COUNT(*) rows, MIN(ts) oldest, MAX(ts) newest FROM samples_raw" ).fetchone()) slot_raw_info = dict(self._db.execute( "SELECT COUNT(*) rows, MIN(ts) oldest, MAX(ts) newest FROM slot_samples_raw" ).fetchone()) points.sort(key=lambda item: (item["ts"], item["gpu_index"])) slot_points.sort(key=lambda item: (item["ts"], item["slot_id"])) return { "range": range_name if range_name in ranges else "24h", "detail_retention_days": DETAIL_RETENTION_DAYS, "sample_interval_seconds": HISTORY_INTERVAL, "points": points, "slot_points": slot_points, "token_totals": totals, "range_tokens": range_tokens, "profile_usage": profile_usage, "model_events": events, "storage": raw_info, "slot_storage": slot_raw_info, }