Tests. First automated coverage for the project: 164 tests, 2.7s, no GPU or network. An autouse fixture stubs overclock_manager._sh -- the single choke point for every nvidia-smi/nvidia-settings write -- so no test can mutate the card. They deliberately pin the empirically measured constants that would otherwise rot silently: the cold and warm load figures behind the cache-hit thresholds, the warm_confident residency rule, and the busy/stalled yield split. One test asserts RAM_HIT_GBPS stays at or below the measured 2.63 GB/s warm load, so the old physically unreachable 5.0 GB/s bar cannot come back. Three bugs the suite surfaced, now fixed: - autotune._subsample(values, 1) divided by zero; the early return only covered len(values) <= max_steps. - telemetry_store.stop() flushed its local pending list but never drained the queue, silently losing rows submitted just before a shutdown -- exactly when the last events matter. - ram_optimizer.page_residency's zero-byte short-circuit omitted keys every other return path provides, so a 0-byte file was planned for warming. Reclaim. The README has claimed bidirectional arbitration from the start, but only one direction was ever automatic. Establishing what actually happens took a controlled test with the service stopped: with ComfyUI holding 6.83 GB, Ollama does not spill to the CPU on this box -- it aborts with "cudaMalloc failed: out of memory", because n_gpu_layers is pinned to 99 and it will not reduce the layer count. So both failure modes are handled: _check_ollama_starved watches size_vram < size for the default configuration where Ollama does spill, and switch_ollama_model catches the hard OOM, reclaims VRAM from an idle ComfyUI and retries once. The request that returned HTTP 500 from Ollama directly now succeeds through HyperSwap, loading at 3.85 GB/s after reclaiming 6.83 GB. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
416 lines
15 KiB
Python
416 lines
15 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)
|
|
# 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
|