"""SQLite time-series persistence for HyperSwap telemetry, swap events and autotune runs. Everything the arbitrator learns used to die with the process (SWITCH_HISTORY was an in-memory deque of 50). This module keeps it on disk so we can answer the question the whole app exists to answer: does a given overclock profile actually deliver more tok/s? Design notes: * WAL mode + a single writer thread -> the 1Hz sampler never blocks the event loop. * Telemetry rows are batched and flushed every FLUSH_INTERVAL_S. * Retention pruning runs opportunistically, keeping the DB bounded and small. """ import logging import os import queue import sqlite3 import threading import time from typing import Any, Dict, List, Optional logger = logging.getLogger("telemetry_store") _BASE = os.path.dirname(os.path.abspath(__file__)) DB_PATH = os.environ.get("HYPERSWAP_DB", os.path.join(_BASE, "hyperswap.db")) TELEMETRY_RETENTION_DAYS = float(os.environ.get("HYPERSWAP_TELEMETRY_RETENTION_DAYS", "14")) EVENT_RETENTION_DAYS = float(os.environ.get("HYPERSWAP_EVENT_RETENTION_DAYS", "180")) FLUSH_INTERVAL_S = 2.0 PRUNE_INTERVAL_S = 3600.0 SCHEMA = """ PRAGMA journal_mode=WAL; PRAGMA synchronous=NORMAL; CREATE TABLE IF NOT EXISTS telemetry ( ts REAL NOT NULL, profile TEXT, gpu_util_pct REAL, mem_util_pct REAL, temp_c REAL, power_w REAL, power_limit_w REAL, fan_pct REAL, clock_sm_mhz REAL, clock_mem_mhz REAL, vram_used_bytes INTEGER, ollama_bytes INTEGER, comfy_bytes INTEGER, system_bytes INTEGER, ram_used_bytes INTEGER, ram_cached_bytes INTEGER, pcie_tx_kbps INTEGER, pcie_rx_kbps INTEGER, throttle_reasons TEXT ); CREATE INDEX IF NOT EXISTS idx_telemetry_ts ON telemetry(ts); CREATE TABLE IF NOT EXISTS events ( ts REAL NOT NULL, event_type TEXT, source TEXT, target TEXT, profile TEXT, duration_ms REAL, load_duration_ms REAL, yield_confirm_ms REAL, tokens_per_sec REAL, bytes_loaded INTEGER, load_gbps REAL, cache_status TEXT, detail TEXT ); CREATE INDEX IF NOT EXISTS idx_events_ts ON events(ts); CREATE INDEX IF NOT EXISTS idx_events_type ON events(event_type); CREATE TABLE IF NOT EXISTS autotune_runs ( ts REAL NOT NULL, profile TEXT, knob TEXT, core_offset_mhz INTEGER, mem_offset_mhz INTEGER, tokens_per_sec REAL, load_gbps REAL, temp_c REAL, power_w REAL, stable INTEGER, instability TEXT, note TEXT ); CREATE INDEX IF NOT EXISTS idx_autotune_ts ON autotune_runs(ts); """ class _Writer(threading.Thread): """Single background writer: batches telemetry, commits events immediately.""" daemon = True def __init__(self) -> None: super().__init__(name="telemetry-writer") self.q: "queue.Queue[Optional[tuple]]" = queue.Queue(maxsize=10000) self._stop = threading.Event() self._last_prune = 0.0 def run(self) -> None: conn = sqlite3.connect(DB_PATH) conn.executescript(SCHEMA) conn.commit() pending: List[tuple] = [] last_flush = time.time() while not self._stop.is_set(): try: item = self.q.get(timeout=0.5) except queue.Empty: item = None if item is not None: kind, sql, params = item if kind == "telemetry": pending.append((sql, params)) else: try: conn.execute(sql, params) conn.commit() except Exception as e: logger.warning("event write failed: %s", e) now = time.time() if pending and (now - last_flush) >= FLUSH_INTERVAL_S: try: for sql, params in pending: conn.execute(sql, params) conn.commit() except Exception as e: logger.warning("telemetry flush failed: %s", e) pending.clear() last_flush = now if now - self._last_prune > PRUNE_INTERVAL_S: self._last_prune = now try: conn.execute("DELETE FROM telemetry WHERE ts < ?", (now - TELEMETRY_RETENTION_DAYS * 86400,)) conn.execute("DELETE FROM events WHERE ts < ?", (now - EVENT_RETENTION_DAYS * 86400,)) conn.commit() except Exception as e: logger.debug("prune failed: %s", e) # Drain anything still queued before closing. Without this, rows submitted but # not yet dequeued are lost on shutdown -- which is exactly when the last events # before a restart matter most. try: while True: try: item = self.q.get_nowait() except queue.Empty: break if item is None: continue _kind, sql, params = item pending.append((sql, params)) except Exception: pass try: for sql, params in pending: conn.execute(sql, params) conn.commit() conn.close() except Exception: pass def stop(self) -> None: self._stop.set() _writer: Optional[_Writer] = None _writer_lock = threading.Lock() def start() -> None: global _writer with _writer_lock: if _writer is None or not _writer.is_alive(): _writer = _Writer() _writer.start() logger.info("telemetry store started at %s", DB_PATH) def stop() -> None: global _writer with _writer_lock: if _writer is not None: _writer.stop() _writer.join(timeout=3.0) _writer = None def _submit(kind: str, sql: str, params: tuple) -> None: w = _writer if w is None: return try: w.q.put_nowait((kind, sql, params)) except queue.Full: logger.debug("telemetry queue full, dropping sample") _TELEMETRY_SQL = """ INSERT INTO telemetry (ts, profile, gpu_util_pct, mem_util_pct, temp_c, power_w, power_limit_w, fan_pct, clock_sm_mhz, clock_mem_mhz, vram_used_bytes, ollama_bytes, comfy_bytes, system_bytes, ram_used_bytes, ram_cached_bytes, pcie_tx_kbps, pcie_rx_kbps, throttle_reasons) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?) """ def record_telemetry(gpu: Dict[str, Any], ram: Dict[str, Any], profile: Optional[str] = None, throttle_reasons: Optional[str] = None) -> None: """Persist a single 1Hz telemetry sample. Never raises.""" try: if not gpu.get("available"): return bd = gpu.get("breakdown", {}) _submit("telemetry", _TELEMETRY_SQL, ( time.time(), profile, gpu.get("gpu_util_pct"), gpu.get("mem_util_pct"), gpu.get("temperature_c"), gpu.get("power_w"), gpu.get("power_limit_w"), gpu.get("fan_pct"), gpu.get("clock_graphics_mhz"), gpu.get("clock_mem_mhz"), gpu.get("vram_used_bytes"), int(bd.get("ollama_gb", 0) * (1024 ** 3)), int(bd.get("comfyui_gb", 0) * (1024 ** 3)), int(bd.get("system_gb", 0) * (1024 ** 3)), ram.get("used_bytes"), ram.get("cached_bytes"), gpu.get("pcie_tx_kbps"), gpu.get("pcie_rx_kbps"), throttle_reasons, )) except Exception as e: logger.debug("record_telemetry failed: %s", e) _EVENT_SQL = """ INSERT INTO events (ts, event_type, source, target, profile, duration_ms, load_duration_ms, yield_confirm_ms, tokens_per_sec, bytes_loaded, load_gbps, cache_status, detail) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?) """ def record_event(event: Dict[str, Any], profile: Optional[str] = None) -> None: """Persist a swap/yield/purge event. Never raises.""" try: _submit("event", _EVENT_SQL, ( event.get("ts", time.time()), event.get("event_type"), event.get("source"), event.get("target"), profile or event.get("profile"), event.get("duration_ms"), event.get("load_duration_ms"), event.get("yield_confirm_ms"), event.get("tokens_per_sec"), event.get("bytes_loaded"), event.get("load_gbps"), event.get("cache_status"), event.get("detail"), )) except Exception as e: logger.debug("record_event failed: %s", e) _AUTOTUNE_SQL = """ INSERT INTO autotune_runs (ts, profile, knob, core_offset_mhz, mem_offset_mhz, tokens_per_sec, load_gbps, temp_c, power_w, stable, instability, note) VALUES (?,?,?,?,?,?,?,?,?,?,?,?) """ def record_autotune(row: Dict[str, Any]) -> None: try: _submit("event", _AUTOTUNE_SQL, ( row.get("ts", time.time()), row.get("profile"), row.get("knob"), row.get("core_offset_mhz"), row.get("mem_offset_mhz"), row.get("tokens_per_sec"), row.get("load_gbps"), row.get("temp_c"), row.get("power_w"), 1 if row.get("stable") else 0, row.get("instability"), row.get("note"), )) except Exception as e: logger.debug("record_autotune failed: %s", e) # ------------------------------------------------------------------ queries def _read_conn() -> sqlite3.Connection: conn = sqlite3.connect(f"file:{DB_PATH}?mode=ro", uri=True, timeout=5.0) conn.row_factory = sqlite3.Row return conn def _rows(sql: str, params: tuple = ()) -> List[Dict[str, Any]]: if not os.path.exists(DB_PATH): return [] try: with _read_conn() as conn: return [dict(r) for r in conn.execute(sql, params).fetchall()] except Exception as e: logger.debug("query failed: %s", e) return [] def profile_comparison(days: float = 7.0) -> List[Dict[str, Any]]: """The headline question: which overclock profile actually produces more tok/s? Joins decode throughput from switch events against thermals sampled while that profile was active. """ since = time.time() - days * 86400 perf = _rows(""" SELECT profile, COUNT(*) AS swaps, AVG(tokens_per_sec) AS avg_tok_s, MAX(tokens_per_sec) AS max_tok_s, AVG(load_gbps) AS avg_load_gbps, AVG(load_duration_ms) AS avg_load_ms FROM events WHERE ts > ? AND event_type = 'LLM Model Switch' AND tokens_per_sec > 0 GROUP BY profile """, (since,)) thermals = {r["profile"]: r for r in _rows(""" SELECT profile, AVG(temp_c) AS avg_temp_c, MAX(temp_c) AS max_temp_c, AVG(power_w) AS avg_power_w, AVG(clock_sm_mhz) AS avg_clock_sm, AVG(clock_mem_mhz) AS avg_clock_mem, COUNT(*) AS samples FROM telemetry WHERE ts > ? AND gpu_util_pct > 5 GROUP BY profile """, (since,))} out = [] for row in perf: merged = dict(row) merged.update(thermals.get(row["profile"], {})) for k, v in list(merged.items()): if isinstance(v, float): merged[k] = round(v, 2) out.append(merged) out.sort(key=lambda r: r.get("avg_tok_s") or 0, reverse=True) return out def swap_stats(days: float = 7.0) -> Dict[str, Any]: since = time.time() - days * 86400 by_type = _rows(""" SELECT event_type, COUNT(*) AS n, AVG(duration_ms) AS avg_ms, MIN(duration_ms) AS min_ms, MAX(duration_ms) AS max_ms, AVG(yield_confirm_ms) AS avg_confirm_ms FROM events WHERE ts > ? GROUP BY event_type ORDER BY n DESC """, (since,)) cache = _rows(""" SELECT cache_status, COUNT(*) AS n, AVG(load_gbps) AS avg_gbps FROM events WHERE ts > ? AND cache_status IS NOT NULL GROUP BY cache_status """, (since,)) models = _rows(""" SELECT target AS model, COUNT(*) AS loads, AVG(tokens_per_sec) AS avg_tok_s, AVG(load_gbps) AS avg_gbps, MAX(ts) AS last_used FROM events WHERE ts > ? AND event_type = 'LLM Model Switch' GROUP BY target ORDER BY loads DESC LIMIT 25 """, (since,)) return {"by_type": by_type, "by_cache_status": cache, "by_model": models, "window_days": days} def model_usage_ranking(days: float = 30.0) -> List[Dict[str, Any]]: """Recency+frequency score per model, used to prioritise the RAM warm budget.""" since = time.time() - days * 86400 now = time.time() rows = _rows(""" SELECT target AS model, COUNT(*) AS loads, MAX(ts) AS last_used FROM events WHERE ts > ? AND target IS NOT NULL AND event_type IN ('LLM Model Switch','Model Warm') GROUP BY target """, (since,)) for r in rows: age_h = max((now - (r["last_used"] or since)) / 3600.0, 0.01) # frequency, decayed by recency (half-life ~24h) r["score"] = round(r["loads"] * (0.5 ** (age_h / 24.0)) + 1.0 / age_h, 4) r["age_hours"] = round(age_h, 2) rows.sort(key=lambda r: r["score"], reverse=True) return rows def timeseries(hours: float = 6.0, buckets: int = 240) -> List[Dict[str, Any]]: """Downsampled history for long-range dashboard charts.""" since = time.time() - hours * 3600 width = max((hours * 3600) / max(buckets, 1), 1.0) return _rows(""" SELECT CAST(ts / ? AS INTEGER) * ? AS bucket_ts, AVG(gpu_util_pct) AS gpu_util_pct, AVG(temp_c) AS temp_c, AVG(power_w) AS power_w, AVG(fan_pct) AS fan_pct, AVG(vram_used_bytes) AS vram_used_bytes, AVG(ollama_bytes) AS ollama_bytes, AVG(comfy_bytes) AS comfy_bytes, AVG(ram_cached_bytes) AS ram_cached_bytes, AVG(clock_sm_mhz) AS clock_sm_mhz, AVG(clock_mem_mhz) AS clock_mem_mhz FROM telemetry WHERE ts > ? GROUP BY bucket_ts ORDER BY bucket_ts """, (width, width, since)) def recent_events(limit: int = 50) -> List[Dict[str, Any]]: return _rows("SELECT * FROM events ORDER BY ts DESC LIMIT ?", (limit,)) def autotune_history(limit: int = 200) -> List[Dict[str, Any]]: return _rows("SELECT * FROM autotune_runs ORDER BY ts DESC LIMIT ?", (limit,)) def db_info() -> Dict[str, Any]: info = {"path": DB_PATH, "exists": os.path.exists(DB_PATH)} if info["exists"]: info["size_mb"] = round(os.path.getsize(DB_PATH) / (1024 ** 2), 2) for tbl in ("telemetry", "events", "autotune_runs"): r = _rows(f"SELECT COUNT(*) AS n FROM {tbl}") info[f"{tbl}_rows"] = r[0]["n"] if r else 0 r = _rows("SELECT MIN(ts) AS a, MAX(ts) AS b FROM telemetry") if r and r[0]["a"]: info["coverage_hours"] = round((r[0]["b"] - r[0]["a"]) / 3600.0, 2) return info