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>
This commit is contained in:
400
telemetry_store.py
Normal file
400
telemetry_store.py
Normal file
@@ -0,0 +1,400 @@
|
||||
"""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
|
||||
Reference in New Issue
Block a user