- 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>
302 lines
12 KiB
Python
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
|