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>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
bf1153cb66
commit
704c21d839
@@ -0,0 +1,301 @@
|
||||
"""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
|
||||
Reference in New Issue
Block a user