Files
mission-control-v2/backend/kern/messreihen.py
T
HitonabiandClaude Opus 5.5 704c21d839 monitoring: Messwerte mit Verlauf fuer KI-Box und Homelab
- kern/messreihen.py: Minutenwerte als JSON-Zeilen je Quelle und Tag unter
  <Datenordner>/mc2-messwerte/, Anhaengen unter flock, Aufraeumen nach 8 Tagen,
  Lesen fuer 1h/24h/7d (60 s, 5-min- und 30-min-Mittel), Luecken bleiben null,
  kaputte Zeilen werden uebersprungen, Zaehler-Raten ohne Spruenge
- KI-Box: Taktgeber im Steward (services/messwerte.py) schreibt jede Minute
  cpu, ram, gpu, Temperaturen, platte, Netz in Bytes/s und Tokens pro Minute;
  GET /api/messwerte (nur Rolle box)
- Homelab: Der Ausfuehrer schickt jede Minute in einem eigenen Faden Host-Werte
  aus /proc und die Gaeste aus einem pvesh-Aufruf an
  POST /api/homelab/ausfuehrer/messwerte; services/homelab/messwerte.py rechnet
  die Zaehler in Bytes/s um (Neustarts und Spruenge ergeben null),
  GET /api/homelab/messwerte liefert Host und Gaeste wie im Inventar
- Waechter (homelab): gelb, wenn der Host 10 Minuten ueber 95 % RAM oder 90 °C liegt
- Aufraeumen der Modell-Platte bietet mc2-messwerte nie zum Loeschen an

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-24 22:41:53 +02:00

302 lines
12 KiB
Python

"""Messreihen: Minutenwerte mit Verlauf für beide Bereiche (Monitoring, seit 24.09.2026).
Ablage: je Quelle und Tag eine Datei, darin eine JSON-Zeile je Minute:
<Datenordner>/mc2-messwerte/<quelle>-JJJJ-MM-TT.jsonl (Tag nach Berliner Zeit, kern/zeit.py)
{"t":1790000000,"cpu":12.5,"ram":41.2,…} (t = Unix-Sekunden der Messung, null = keine Messung)
Quellen: „box“ (die KI-Box; schreibt ihr Steward) sowie „pve“, „ct-<vmid>“ und „vm-<vmid>“ (das Homelab; schreibt der
Homelab-Teil, was der Ausführer jede Minute schickt). Dateien, deren Tag mehr als AUFBEWAHREN_TAGE zurückliegt, löscht
das erste Schreiben eines neuen Tages.
Schreiben: eine Zeile anhängen, unter einer Thread-Sperre und (Linux) flock auf der Datei. Fehlt am Dateiende der
Zeilenumbruch (Absturz mitten im Schreiben), beginnt die neue Zeile auf einer eigenen statt an die kaputte anzuwachsen.
Lesen (lesen()): 1h in 60-s-Schritten, 24h als 5-min-Mittel, 7d als 30-min-Mittel. Die Schritte liegen auf vollen
Vielfachen ihrer Länge, der letzte ist der laufende. Ein Schritt ohne Messung bleibt null — Lücken werden nicht
aufgefüllt. Kaputte Zeilen werden übersprungen. Die Summen je Tagesdatei werden für 5- und 30-min-Schritte gemerkt,
solange sich die Datei nicht ändert (vergangene Tage liest so nur die erste Anfrage).
"""
import json
import logging
import math
import os
import re
import threading
import time
from collections.abc import Iterable, Iterator
from datetime import datetime, timedelta
from pathlib import Path
from kern.einstellungen import einstellungen
from kern.zeit import LOCAL_TZ
try:
import fcntl # Linux: Sperre über Prozessgrenzen
except ImportError: # Windows (Entwicklung): nur die Thread-Sperre
fcntl = None
log = logging.getLogger(__name__)
ORDNER_NAME = "mc2-messwerte"
AUFBEWAHREN_TAGE = 8
# Zeitraum → (Länge, Schritt) in Sekunden.
ZEITRAEUME: dict[str, tuple[int, int]] = {"1h": (3600, 60), "24h": (24 * 3600, 300), "7d": (7 * 24 * 3600, 1800)}
MERKEN_AB_S = 300 # Summen je Datei merken; bei 60-s-Schritten (1440 je Tag) lohnt es den Speicher nicht
MERKEN_MAX = 512
LUECKE_MAX_S = 300 # Zählerstände, die weiter auseinanderliegen, ergeben keine Rate mehr
RATE_MAX = 5e9 # mehr als 5 GB/s ist kein Verkehr, sondern ein Zählersprung
_QUELLE = re.compile(r"[a-z0-9][a-z0-9-]{0,62}")
_DATEI = re.compile(r"(?P<quelle>[a-z0-9][a-z0-9-]*)-(?P<tag>\d{4}-\d{2}-\d{2})\.jsonl")
_lock = threading.Lock()
_aufgeraeumt: dict[str, str] = {} # Ordner → Tag des letzten Aufräumens (je Prozess)
_merk_lock = threading.Lock()
_merk: dict[tuple[str, int], tuple[tuple[int, int], dict[int, dict[str, list[float]]]]] = {}
def ordner() -> Path:
return einstellungen().daten_dir / ORDNER_NAME
def ist_zahl(wert: object) -> bool:
"""Eine endliche Zahl. bool zählt nicht, NaN und unendlich auch nicht."""
return isinstance(wert, (int, float)) and not isinstance(wert, bool) and math.isfinite(wert)
def zeitraum_pruefen(zeitraum: str) -> tuple[int, int]:
"""(Länge, Schritt) in Sekunden; ValueError bei einem unbekannten Zeitraum."""
if zeitraum not in ZEITRAEUME:
raise ValueError(f"Den Zeitraum „{zeitraum}“ gibt es nicht; möglich sind {', '.join(ZEITRAEUME)}.")
return ZEITRAEUME[zeitraum]
def _datum(t: float):
return datetime.fromtimestamp(t, LOCAL_TZ).date()
def _tage(von: float, bis: float) -> list[str]:
"""Alle Tage (Berliner Zeit), deren Dateien Messungen aus [von, bis) enthalten können."""
tag, letzter = _datum(von), _datum(max(von, bis - 1))
tage = []
while tag <= letzter:
tage.append(tag.isoformat())
tag += timedelta(days=1)
return tage
def _datei(quelle: str, tag: str) -> Path:
if not _QUELLE.fullmatch(quelle):
raise ValueError(f"Ungültiger Name einer Messquelle: {quelle!r}")
return ordner() / f"{quelle}-{tag}.jsonl"
# --- Schreiben -----------------------------------------------------------------------------------
def _wert(wert: object) -> int | float | None:
if not ist_zahl(wert):
return None
return wert if isinstance(wert, int) else round(float(wert), 3)
def schreiben(quelle: str, werte: dict[str, object], t: float | None = None) -> None:
"""Einen Messpunkt anhängen. werte: Name → Zahl oder None (nicht gemessen)."""
t = time.time() if t is None else float(t)
punkt: dict[str, object] = {"t": int(t)}
punkt.update({name: _wert(wert) for name, wert in werte.items() if name != "t"})
zeile = (json.dumps(punkt, separators=(",", ":")) + "\n").encode()
pfad = _datei(quelle, _datum(t).isoformat())
with _lock:
pfad.parent.mkdir(parents=True, exist_ok=True)
with open(pfad, "a+b") as f:
if fcntl is not None:
fcntl.flock(f, fcntl.LOCK_EX)
try:
f.seek(0, os.SEEK_END)
if f.tell() > 0:
f.seek(-1, os.SEEK_END)
if f.read(1) != b"\n":
zeile = b"\n" + zeile # die letzte Zeile brach ab: nicht an sie anhängen
f.write(zeile)
f.flush()
finally:
if fcntl is not None:
fcntl.flock(f, fcntl.LOCK_UN)
_einmal_am_tag_aufraeumen(t)
def _einmal_am_tag_aufraeumen(t: float) -> None:
schluessel, tag = str(ordner()), _datum(t).isoformat()
if _aufgeraeumt.get(schluessel) == tag:
return
_aufgeraeumt[schluessel] = tag
try:
aufraeumen(t)
except OSError:
log.warning("messreihen: Aufräumen in %s ging nicht", schluessel, exc_info=True)
def aufraeumen(jetzt: float | None = None) -> int:
"""Tagesdateien löschen, deren Tag mehr als AUFBEWAHREN_TAGE zurückliegt. Andere Dateien bleiben. Rückgabe: wie
viele gelöscht wurden."""
grenze = (_datum(time.time() if jetzt is None else jetzt) - timedelta(days=AUFBEWAHREN_TAGE)).isoformat()
try:
namen = os.listdir(ordner())
except FileNotFoundError:
return 0
geloescht = 0
for name in namen:
treffer = _DATEI.fullmatch(name)
if treffer and treffer["tag"] < grenze:
try:
(ordner() / name).unlink()
geloescht += 1
except FileNotFoundError:
pass # ein anderer Prozess war schneller
return geloescht
# --- Lesen ---------------------------------------------------------------------------------------
def _punkt(zeile: bytes) -> dict | None:
"""Eine Zeile als Messpunkt; None bei allem, was keiner ist (kaputt, leer, ohne Zeit)."""
zeile = zeile.strip()
if not zeile:
return None
try:
punkt = json.loads(zeile)
except ValueError: # auch UnicodeDecodeError
return None
return punkt if isinstance(punkt, dict) and ist_zahl(punkt.get("t")) else None
def _punkte_der_datei(pfad: Path) -> Iterator[dict]:
try:
roh = pfad.read_bytes()
except OSError:
return
for zeile in roh.split(b"\n"):
if (punkt := _punkt(zeile)) is not None:
yield punkt
def _eimer(pfad: Path, schritt: int, ab: float) -> dict[int, dict[str, list[float]]]:
"""Summe und Anzahl je Schritt und Wert für eine Tagesdatei. Gemerkt (ab 5-min-Schritten) wird die ganze Datei;
sonst zählen erst Messungen ab `ab`. Das Ergebnis nicht verändern: es kann gemerkt sein."""
try:
stand = pfad.stat()
except OSError:
return {}
merken = schritt >= MERKEN_AB_S
schluessel, kennung = (str(pfad), schritt), (stand.st_mtime_ns, stand.st_size)
if merken:
with _merk_lock:
gemerkt = _merk.get(schluessel)
if gemerkt and gemerkt[0] == kennung:
return gemerkt[1]
eimer: dict[int, dict[str, list[float]]] = {}
for punkt in _punkte_der_datei(pfad):
if not merken and punkt["t"] < ab:
continue
ziel = eimer.setdefault(int(punkt["t"]) // schritt * schritt, {})
for name, wert in punkt.items():
if name != "t" and ist_zahl(wert):
summe = ziel.setdefault(name, [0.0, 0])
summe[0] += wert
summe[1] += 1
if merken:
with _merk_lock:
if len(_merk) >= MERKEN_MAX:
_merk.clear()
_merk[schluessel] = (kennung, eimer)
return eimer
def _mittel(summe: list[float] | None) -> int | float | None:
"""Mittel eines Schritts: ab 100 ganzzahlig, darunter eine Nachkommastelle; ohne Messung None."""
if not summe or not summe[1]:
return None
wert = summe[0] / summe[1]
return round(wert) if abs(wert) >= 100 else round(wert, 1)
def lesen(quelle: str, zeitraum: str, namen: Iterable[str], jetzt: float | None = None) -> dict:
"""Die Reihen einer Quelle: {"zeitraum", "schritt_s", "reihen": {name: [[t, wert], …]}}. t = Beginn des Schritts
(Unix-Sekunden), wert = Mittel der Messungen darin oder None. Jede Reihe hat Länge/Schritt Einträge (60, 288, 336)."""
dauer, schritt = zeitraum_pruefen(zeitraum)
jetzt = time.time() if jetzt is None else jetzt
ende = (int(jetzt) // schritt + 1) * schritt # Ende des laufenden Schritts
beginn = ende - dauer
summen: dict[int, dict[str, list[float]]] = {}
for tag in _tage(beginn, ende):
for t0, werte in _eimer(_datei(quelle, tag), schritt, beginn).items():
if not beginn <= t0 < ende:
continue
ziel = summen.setdefault(t0, {})
for name, (summe, anzahl) in werte.items():
gesamt = ziel.setdefault(name, [0.0, 0])
gesamt[0] += summe
gesamt[1] += anzahl
raster = range(beginn, ende, schritt)
return {"zeitraum": zeitraum, "schritt_s": schritt,
"reihen": {name: [[t0, _mittel(summen.get(t0, {}).get(name))] for t0 in raster] for name in namen}}
def _ende_der_datei(pfad: Path, groesse: int = 16384) -> list[bytes]:
"""Die letzten Zeilen einer Datei, ohne sie ganz zu lesen."""
try:
with open(pfad, "rb") as f:
f.seek(0, os.SEEK_END)
laenge = f.tell()
f.seek(max(0, laenge - groesse))
roh = f.read()
except OSError:
return []
zeilen = roh.split(b"\n")
return zeilen[1:] if laenge > groesse else zeilen # die erste Zeile ist dann angeschnitten
def letzter_punkt(quelle: str, jetzt: float | None = None, max_alter_s: float | None = None) -> dict | None:
"""Die jüngste Messung einer Quelle (heute, sonst gestern). None, wenn es keine gibt oder sie älter ist als
max_alter_s — ein alter Wert ist kein aktueller."""
jetzt = time.time() if jetzt is None else jetzt
heute = _datum(jetzt)
for tag in (heute, heute - timedelta(days=1)):
for zeile in reversed(_ende_der_datei(_datei(quelle, tag.isoformat()))):
if (punkt := _punkt(zeile)) is not None:
if max_alter_s is not None and jetzt - punkt["t"] > max_alter_s:
return None
return punkt
return None
def punkte(quelle: str, von: float, bis: float) -> list[dict]:
"""Alle Messungen einer Quelle mit von ≤ t < bis, nach Zeit geordnet (für Prüfungen des Wächters)."""
gefunden = [p for tag in _tage(von, bis) for p in _punkte_der_datei(_datei(quelle, tag)) if von <= p["t"] < bis]
return sorted(gefunden, key=lambda p: p["t"])
def quellen(von: float, bis: float) -> list[str]:
"""Die Quellen, von denen es für [von, bis) eine Tagesdatei gibt."""
tage = set(_tage(von, bis))
try:
namen = os.listdir(ordner())
except OSError:
return []
return sorted({t["quelle"] for name in namen if (t := _DATEI.fullmatch(name)) and t["tag"] in tage})
# --- Zähler → Raten --------------------------------------------------------------------------------
def rate(alt: object, neu: object, dt: float) -> float | None:
"""Zuwachs je Sekunde zwischen zwei Ständen eines kumulativen Zählers (Bytes, Tokens). None, wenn ein Stand fehlt,
die Stände mehr als LUECKE_MAX_S auseinanderliegen, der Zähler kleiner wurde (Neustart, Überlauf) oder der Sprung
unmöglich groß ist — eine erfundene Rate wäre schlimmer als eine Lücke."""
if not (ist_zahl(alt) and ist_zahl(neu)) or not 0 < dt <= LUECKE_MAX_S:
return None
zuwachs = float(neu) - float(alt)
if zuwachs < 0:
return None
wert = zuwachs / dt
return None if wert > RATE_MAX else wert