Files
gpu-program-swapper/telemetry_store.py
drjones 5431144b2e Add barrier-confirmed yielding, measured residency, persistence and closed-loop tuning
Nine changes, in rough order of how much they affect real behaviour:

1. VRAM yield is now a barrier. Posting keep_alive:0 only asks Ollama to unload;
   measured here, the HTTP call returns in 63ms while the driver takes a further
   77ms to release 14.9GB. Returning inside that window is how ComfyUI ends up
   allocating into VRAM that is still occupied. instant_free_ollama_vram() polls
   NVML until the allocation is actually gone and reports request/confirm split.

2. ComfyUI VRAM is no longer purged 1.5s after every prompt, which forced a full
   checkpoint reload on each workflow iteration. It is held for 30s of genuinely
   empty queue, with an immediate purge when Ollama actually asks for the memory.

3. Cache-hit classification uses achieved bandwidth (size / load duration) rather
   than a fixed `load_duration < 2500ms`. That constant called a 12.9GB model read
   at 2.9GB/s a cold load, and a 0.5GB model read from NVMe a cache hit.

4. Page-cache residency is measured, not assumed. mincore(2) reported 128GB
   resident on a box with 46GB of page cache: the kernel only permits page-cache
   introspection on files you own, and the Ollama blobs are owned by uid ollama,
   for which mincore answers "all resident" instead of failing. Uses cachestat(2)
   where permitted and a randomised read-rate probe elsewhere, labelling which was
   used. Fixed-offset probing was self-fulfilling, so windows are random and cold
   ones are returned with FADV_DONTNEED.

5. Warming is budgeted and ranked by recency/frequency instead of reading every
   file top-to-bottom, which on 64GB of RAM just evicts whatever was warmed first.

6. Telemetry and events persist to SQLite (~0.38 MB/hour) instead of living in a
   50-entry in-memory deque, so /api/analytics/profiles can finally answer whether
   an overclock profile actually delivers more tok/s.

7. Thermal governor walks the overclock back on sustained heat or hardware
   throttling, with hysteresis, fed from the existing sampler.

8. Autotune sweeps a clock offset, benchmarks decode at each step, watches for Xid
   errors and degenerate output, and restores the profile in a finally block.

9. Stock clocks/power/fans are restored on shutdown and via systemd ExecStopPost.
   Nothing previously undid a locked clock or a manually pinned fan.

Also: one shared 1Hz telemetry sampler fanned out to SSE subscribers rather than
every client re-running the whole snapshot; wall-clock timestamps in place of the
event loop's monotonic clock; cached nvidia-smi shell-outs; quieter httpx logging.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-28 08:57:35 -07:00

401 lines
14 KiB
Python

"""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)
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