""" Per-Stage-Latenz-Metriken + Per-Turn-Trace für die Voice/Lucy-Pipeline. Zwei Sichten auf dieselben Messungen: 1. **Rollende Stats** (`record_stage`/`Timer`/`get_metrics`) — avg/p50/p95/last je Stufe über die letzten N Messungen. Gut für Trends („Wie schnell ist STT im Schnitt?"). 2. **Per-Turn-Trace** (`TurnTrace`/`get_trace`) — jeder einzelne Chat-Turn als eigener Datensatz mit seinem Stufen-Breakdown. Nur so sieht man den EINEN langsamen Turn (der 30-s-Hänger), den ein Durchschnitt verschluckt. Das war der eigentliche Anlass (Telegram-/Voice-Hänger diagnostizieren). **Was MC2 messen kann — und was nicht:** MC2 proxyt den Chat nur an Hermes (:8642). Die Stufen STT, Vision, Hirn-TTFT (Zeit bis zum ersten Inhalts-Token) und Generierung sind hier direkt messbar. Der **Mem0-Retrieve** läuft zwar in Hermes, ruft aber MC2s `/api/memory` per HTTP zurück → messbar und als Unter-Detail INNERHALB der Hirn-Zeit ausgewiesen (kein Doppelzählen). Die **Tool-Runden** dagegen laufen im Hermes-LLM-Loop ohne Callback an MC2 → für MC2 unsichtbar, sie stecken im „Generierung"-Bucket. STT (davor) und Mem0-Retrieve (währenddessen) sind separate HTTP-Requests ohne Turn-ID. Auf einem EIN-Nutzer-Gerät genügt eine schlanke Best-Effort-Korrelation: die zuletzt gemessene STT-Dauer bzw. der letzte Retrieve werden global „geparkt" und vom nächsten Chat-Turn eingesammelt (mit Frist-/Reihenfolge- Check). Kein Turn-ID-Durchreichen durch den Lucy-Client nötig. In-Memory + thread-safe (keine Datei-I/O — Latenz-Telemetrie ist transient, Restart = Reset). """ import threading import time import uuid from collections import deque _LOCK = threading.Lock() _MAX = 200 _STAGES: dict[str, deque] = {} _TURNS: deque = deque(maxlen=60) # letzte N vollständige Chat-Turns (Per-Turn-Trace) # Bekannte Stufen (für stabile UI-Reihenfolge); unbekannte werden trotzdem erfasst. STAGES = ("stt", "vision", "memory_retrieve", "chat_ttfb", "chat_first_content", "tts") # Best-effort-Korrelation (Ein-Nutzer-Gerät): zuletzt gemessene STT-Dauer / Mem0-Retrieve, je # (ms, perf_counter-Zeitstempel). Der nächste passende Chat-Turn sammelt sie ein und leert sie. _PARKED: dict[str, tuple[float, float] | None] = {"stt": None, "retrieve": None} def record_stage(stage: str, ms: float) -> None: """Eine gemessene Stage-Dauer (ms) verbuchen. No-op bei negativen Werten.""" if ms is None or ms < 0: return with _LOCK: dq = _STAGES.get(stage) if dq is None: dq = _STAGES[stage] = deque(maxlen=_MAX) dq.append(float(ms)) def park(kind: str, ms: float) -> None: """Eine Messung, die NICHT im Chat-Request selbst passiert (STT davor, Mem0-Retrieve als Rückruf während), global parken, damit der nächste Chat-Turn sie einsammeln kann.""" if ms is None or ms < 0 or kind not in _PARKED: return with _LOCK: _PARKED[kind] = (float(ms), time.perf_counter()) def _take_stt(max_age: float = 20.0) -> float | None: """Geparkte STT-Dauer einsammeln, wenn frisch (STT liegt VOR dem Turn-Start).""" with _LOCK: v = _PARKED.get("stt") if v and (time.perf_counter() - v[1]) <= max_age: _PARKED["stt"] = None return round(v[0], 1) return None def _take_retrieve(since_perf: float, max_age: float = 90.0) -> float | None: """Geparkten Mem0-Retrieve einsammeln, wenn er NACH dem Turn-Start kam (Rückruf während des Turns) und frisch ist.""" with _LOCK: v = _PARKED.get("retrieve") if v and v[1] >= since_perf and (time.perf_counter() - v[1]) <= max_age: _PARKED["retrieve"] = None return round(v[0], 1) return None class Timer: """Context-Manager: misst die verstrichene Zeit und verbucht sie auf `stage`. Funktioniert um `await`-Aufrufe herum (enter → await → exit).""" def __init__(self, stage: str) -> None: self.stage = stage self._t0 = 0.0 def __enter__(self) -> "Timer": self._t0 = time.perf_counter() return self def __exit__(self, *exc) -> None: record_stage(self.stage, (time.perf_counter() - self._t0) * 1000.0) class TurnTrace: """Ein Per-Turn-Trace für den Voice/Lucy-Chatpfad. In `voice.py` über die Dauer eines Chat-Turns gehalten; `commit()` schreibt den Datensatz in den Ringpuffer UND speist die rollenden Stats. Balken-Stufen (zeitlich DISJUNKT, ergeben zusammen den Turn): stt · vision · hirn · gen. Unter-Detail: mem0 (Teil VON hirn, wird separat ausgewiesen, aber NICHT zum Balken addiert).""" def __init__(self, session_id: str = "", kind: str = "voice") -> None: self.id = uuid.uuid4().hex[:8] self.ts = time.time() self.perf0 = time.perf_counter() self._brain0 = self.perf0 # Referenz für die Hirn-Zeit (nach Vision neu gesetzt) self.session_id = (session_id or "")[:24] self.kind = kind self.vision_ms: float | None = None # Bildschirm-Beschreibung (falls Bilder) self.hirn_ms: float | None = None # Zeit bis zum ersten Inhalts-Token (Agent + Mem0 + TTFT) self.had_images = False self.error: str | None = None def note_vision(self, ms: float) -> None: self.vision_ms = round(ms, 1) record_stage("vision", ms) def mark_brain_start(self) -> None: """Startpunkt der Hirn-Zeit — direkt VOR dem Hermes-Request, damit die Vision-Zeit NICHT in die Hirn-Zeit gezählt wird (sonst Doppelzählung mit dem Vision-Segment).""" self._brain0 = time.perf_counter() def note_ttfb(self) -> None: record_stage("chat_ttfb", (time.perf_counter() - self._brain0) * 1000.0) # SSE-Start (~5 ms), nur rollend def note_first_content(self) -> None: """Erster Inhalts-Delta = echte Hirn-Latenz (Agent-Overhead + Mem0 + LLM-TTFT).""" ms = (time.perf_counter() - self._brain0) * 1000.0 self.hirn_ms = round(ms, 1) record_stage("chat_first_content", ms) def commit(self) -> dict: total = (time.perf_counter() - self.perf0) * 1000.0 stt = _take_stt() # rollend bereits in /voice/stt erfasst mem0 = _take_retrieve(self.perf0) # rollend bereits in /api/memory erfasst # Generierung = alles nach dem ersten Inhalts-Token bis Stream-Ende. gen = round(total - self.hirn_ms, 1) if self.hirn_ms is not None else None rec = { "id": self.id, "ts": round(self.ts, 3), "session": self.session_id, "kind": self.kind, "had_images": self.had_images, "stt_ms": stt, "vision_ms": self.vision_ms, "hirn_ms": self.hirn_ms, "gen_ms": gen if (gen is None or gen >= 0) else 0.0, "mem0_ms": mem0, "total_ms": round(total, 1), "error": self.error, } with _LOCK: _TURNS.append(rec) return rec def _summary(vals: list[float]) -> dict: if not vals: return {"count": 0} s = sorted(vals) n = len(s) return { "count": n, "avg_ms": round(sum(s) / n, 1), "p50_ms": round(s[n // 2], 1), "p95_ms": round(s[min(n - 1, int(n * 0.95))], 1), "last_ms": round(vals[-1], 1), } def get_metrics() -> dict: """Rollende Zusammenfassung je Stufe.""" with _LOCK: return {stage: _summary(list(dq)) for stage, dq in _STAGES.items()} def get_trace(limit: int = 20) -> list[dict]: """Die letzten `limit` Chat-Turns (neueste zuerst) mit Stufen-Breakdown.""" limit = max(1, min(int(limit or 20), _TURNS.maxlen or 60)) with _LOCK: return list(_TURNS)[-limit:][::-1]