47f7a85510
Die Ampel-Nachruestung deckte 317/337 vorbestehende ruff-Verstoesse im ganzen Repo auf. Aufgeraeumt: - ruff.toml: intentionale Muster als Projekt-Politik ausgenommen (BLE001 blind-except, S110/S112 try-except-pass/continue, PLW1510 subprocess-best-effort, B008 FastAPI- Depends/File-Idiom, EXE001 Shebang, + wenige Stil-Regeln). __init__.py-Re-Exports geschuetzt (F401). - ruff --fix: 128 mechanische (Import-Sortierung, PEP585/604-Annotationen, tote Imports, ueberfluessige noqa) auto-behoben. - 12 echte Reste von Hand: PERF402/102, PLC3002 (Lambda->walrus), ISC004 (String-Concat geklammert), F841/RUF059 (ungenutzte Vars), PIE810 (startswith-Tuple), UP031 (f-string), UP035 (veraltete typing-Imports). Ergebnis: 'ruff check .' = 0, 'compileall' grün. Kein Verhaltenswechsel (nur Stil/Modernisierung). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
129 lines
4.3 KiB
Python
129 lines
4.3 KiB
Python
"""24-h-Metrik-Verlauf für die Cockpit-Zeitachse (User-Wunsch 16.07.: „WANN war Last?").
|
|
|
|
Ein leichter Sammler (Lifespan-Task in app.py) legt alle SAMPLE_S Sekunden einen
|
|
kompakten Punkt in einen Ringpuffer: CPU/RAM/GPU/Disk in Prozent + die Token-
|
|
GESAMTZÄHLER (das Frontend rechnet Raten aus den Deltas). Der Puffer hält exakt
|
|
24 h — Älteres fällt automatisch raus (deque maxlen), nichts wächst unbegrenzt.
|
|
Periodisch wird auf Platte persistiert (übersteht MC2-Restarts/Deploys); beim
|
|
Laden fliegt alles raus, was älter als 24 h ist.
|
|
|
|
Bewusst NICHT im Steward: die Historie ist reine UI-Ware, und ein paar Sekunden
|
|
Lücke pro Deploy sind egal.
|
|
"""
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import os
|
|
import time
|
|
from collections import deque
|
|
|
|
from config import MODELS_DIR
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
SAMPLE_S = 10 # Abtast-Takt
|
|
RETENTION_S = 24 * 3600 # 24 h — danach löschen und neu bauen (Ringpuffer)
|
|
MAXLEN = RETENTION_S // SAMPLE_S # 8640 Punkte
|
|
FLUSH_S = 300 # höchstens alle 5 min auf Platte
|
|
MAX_RETURN_POINTS = 300 # Endpoint downsampled auf ~diese Punktzahl
|
|
|
|
HISTORY_PATH = MODELS_DIR / "mc2-metrics-history.json"
|
|
|
|
# Punkt = [t, cpu, ram, gpu, disk, tp, tc] (t = Epoch-Sekunden; tp/tc = Token-Totale)
|
|
_points: deque = deque(maxlen=MAXLEN)
|
|
_loaded = False
|
|
_last_flush = 0.0
|
|
|
|
|
|
def _load() -> None:
|
|
global _loaded
|
|
if _loaded:
|
|
return
|
|
_loaded = True
|
|
try:
|
|
raw = json.loads(HISTORY_PATH.read_text(encoding="utf-8"))
|
|
cutoff = time.time() - RETENTION_S
|
|
for p in raw if isinstance(raw, list) else []:
|
|
if isinstance(p, list) and len(p) == 7 and isinstance(p[0], (int, float)) and p[0] >= cutoff:
|
|
_points.append(p)
|
|
if _points:
|
|
log.info("metrics_history: %d Punkte aus %s geladen", len(_points), HISTORY_PATH)
|
|
except FileNotFoundError:
|
|
pass
|
|
except Exception:
|
|
log.warning("metrics_history: %s nicht lesbar — starte leer", HISTORY_PATH, exc_info=True)
|
|
|
|
|
|
def _flush() -> None:
|
|
global _last_flush
|
|
_last_flush = time.time()
|
|
try:
|
|
tmp = HISTORY_PATH.with_suffix(".tmp")
|
|
tmp.write_text(json.dumps(list(_points), separators=(",", ":")), encoding="utf-8")
|
|
tmp.replace(HISTORY_PATH)
|
|
except OSError:
|
|
log.warning("metrics_history: %s nicht schreibbar", HISTORY_PATH, exc_info=True)
|
|
|
|
|
|
def _sample() -> None:
|
|
import psutil
|
|
|
|
from services.system import _gpu_sysfs
|
|
from services.token_stats import get_stats
|
|
|
|
vm = psutil.virtual_memory()
|
|
disk = None
|
|
try:
|
|
du = psutil.disk_usage(str(MODELS_DIR) if MODELS_DIR.exists() else os.getcwd())
|
|
disk = du.percent
|
|
except Exception:
|
|
pass
|
|
gpu = None
|
|
try:
|
|
g = _gpu_sysfs()
|
|
gpu = g.get("busy_percent") if g else None
|
|
except Exception:
|
|
pass
|
|
tp = tc = None
|
|
try:
|
|
ts = get_stats()
|
|
tp, tc = ts.get("prompt_tokens"), ts.get("completion_tokens")
|
|
except Exception:
|
|
pass
|
|
# interval=None: nicht-blockierende CPU-Messung seit dem letzten Aufruf (10 s her — ideal).
|
|
_points.append([int(time.time()), psutil.cpu_percent(interval=None), vm.percent, gpu, disk, tp, tc])
|
|
|
|
|
|
async def sampler_loop() -> None:
|
|
"""Dauer-Sammler fürs App-Lifespan. Fehler eines Ticks reißen die Schleife nie."""
|
|
_load()
|
|
while True:
|
|
try:
|
|
await asyncio.to_thread(_sample)
|
|
except Exception:
|
|
log.debug("metrics_history: Sample fehlgeschlagen", exc_info=True)
|
|
if time.time() - _last_flush >= FLUSH_S:
|
|
try:
|
|
await asyncio.to_thread(_flush)
|
|
except Exception:
|
|
pass
|
|
await asyncio.sleep(SAMPLE_S)
|
|
|
|
|
|
def history(minutes: int = 60) -> dict:
|
|
"""Punkte der letzten N Minuten, auf ≤ MAX_RETURN_POINTS ausgedünnt (Stride)."""
|
|
_load()
|
|
minutes = max(1, min(int(minutes), RETENTION_S // 60))
|
|
cutoff = time.time() - minutes * 60
|
|
pts = [p for p in _points if p[0] >= cutoff]
|
|
stride = max(1, len(pts) // MAX_RETURN_POINTS)
|
|
pts = pts[::stride]
|
|
return {
|
|
"sample_s": SAMPLE_S * stride,
|
|
"points": [
|
|
{"t": p[0], "cpu": p[1], "ram": p[2], "gpu": p[3], "disk": p[4], "tp": p[5], "tc": p[6]}
|
|
for p in pts
|
|
],
|
|
}
|