Observability 1a: Per-Turn-Latenz-Trace + Latenz-Karte

voice_metrics.py: neben den rollenden Stats jetzt ein echter Per-Turn-Trace
(TurnTrace + Ringpuffer der letzten 60 Turns). Balken-Stufen zeitlich disjunkt
(stt · vision · hirn · gen); hirn misst ab mark_brain_start() VOR dem Hermes-
Request, damit die Vision-Zeit nicht doppelt gezaehlt wird. STT (davor) und
Mem0-Retrieve (Rueckruf waehrend) werden per park()/_take_* best-effort dem Turn
zugeordnet (Ein-Nutzer-Geraet, kein Turn-ID noetig). Mem0 = Unter-Detail INNERHALB
hirn (nicht addieren) -> ehrlich, kein Doppelzaehlen.

voice.py: TurnTrace in voice_chat's gen() (commit garantiert 1x via finally, auch
bei Fehler/Abbruch). STT-Endpoint parkt seine Dauer. NEU: GET /api/voice/metrics
(schliesst die Luecke - selbstkritik-feed.sh curlte das, existierte nie -> 404) +
GET /api/voice/trace?limit=N.

memory.py: GET /api/memory?q=... (= Hermes' Mem0-Prefetch) misst + parkt die
Retrieve-Zeit.

Frontend: LatencyCard (gestapelte Balken je Turn, "Taeter" = groesste Stufe,
Fehler-Turns rot, Tooltip mit "davon Mem0 X s"), Query useVoiceTrace, in der
Zentrale unter "Stack & Telemetrie". Lokal gegen Seed-Server verifiziert: 30,3-s-
Haenger -> Hirn-Balken 94% + "davon Mem0 26,1 s".

Grenzen (ehrlich, als Fussnote in der Karte): Tool-Runden im Hermes-LLM-Loop haben
keinen Callback an MC2 -> stecken in "Antwort". Reine Telegram-Text-Turns laufen an
MC2 vorbei und erscheinen hier nicht.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
Hitonabi
2026-07-04 13:21:52 +02:00
parent 58021331e1
commit 78448220b2
12 changed files with 552 additions and 207 deletions
+12
View File
@@ -1,9 +1,12 @@
"""Memory-Endpoints (geteiltes Gedächtnis). LAN-only, kein Token in 2.0-Phase 3."""
import time
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from services import memory
from services.voice_metrics import park, record_stage
router = APIRouter(prefix="/api")
@@ -56,6 +59,15 @@ def dedupe(body: DedupeIn) -> dict:
@router.get("/memory")
def list_mem(q: str = "", category: str = "") -> list[dict]:
# Semantischer Retrieve (q gesetzt) = u.a. Hermes' Mem0-Prefetch VOR jedem Turn. Dauer messen
# + parken, damit der laufende Voice-Chat-Turn sie als Unter-Detail seiner Hirn-Zeit einsammelt.
if q:
_t0 = time.perf_counter()
res = memory.list_memories(q=q, category=category)
_ms = (time.perf_counter() - _t0) * 1000.0
record_stage("memory_retrieve", _ms)
park("retrieve", _ms)
return res
return memory.list_memories(q=q, category=category)
+88 -46
View File
@@ -20,7 +20,14 @@ from pydantic import BaseModel
from config import HERMES_API_KEY, HERMES_API_MODEL, HERMES_API_URL, LLAMA_SWAP_URL, VOICE_SERVICE_URL
from services import announce
from services.voice_metrics import Timer, record_stage # Per-Stage-Latenz (intern)
from services.voice_metrics import ( # Per-Stage-Latenz + Per-Turn-Trace (intern)
Timer,
TurnTrace,
get_metrics,
get_trace,
park,
record_stage,
)
# Injection-Schutz (Stufe 0): guard.py liegt im mcp/-Verzeichnis. Per Pfad laden (eigene MC2-Venv).
import sys as _sys
@@ -153,6 +160,20 @@ def voice_health() -> dict:
return out
@router.get("/voice/metrics")
def voice_metrics() -> dict:
"""Rollende Latenz-Stats je Stufe (avg/p50/p95/last, ms). Quelle u.a. für selbstkritik-feed.sh."""
return get_metrics()
@router.get("/voice/trace")
def voice_trace(limit: int = 20) -> dict:
"""Per-Turn-Trace: die letzten `limit` Chat-Turns mit Stufen-Breakdown (STT · Vision · Hirn ·
Generierung, Mem0 als Unter-Detail). Neueste zuerst. Für die Latenz-Ansicht im Cockpit —
damit man den EINEN langsamen Turn sieht, den ein Durchschnitt verschluckt."""
return {"turns": get_trace(limit)}
@router.get("/voice/voices")
def voice_voices() -> dict:
try:
@@ -172,8 +193,11 @@ async def voice_stt(audio: UploadFile = File(...), language: str = Form(default=
files = {"audio": (audio.filename or "rec.webm", data, audio.content_type or "audio/webm")}
try:
async with httpx.AsyncClient(timeout=_TIMEOUT) as client:
with Timer("stt"):
r = await client.post(f"{VOICE_SERVICE_URL}/stt", files=files, data={"language": language})
_t0 = time.perf_counter()
r = await client.post(f"{VOICE_SERVICE_URL}/stt", files=files, data={"language": language})
_ms = (time.perf_counter() - _t0) * 1000.0
record_stage("stt", _ms)
park("stt", _ms) # der folgende /voice/chat-Turn sammelt die STT-Dauer für seinen Trace ein
r.raise_for_status()
return r.json()
except httpx.HTTPError as exc:
@@ -265,51 +289,69 @@ async def voice_chat(body: ChatIn) -> StreamingResponse:
headers["X-Hermes-Session-Key"] = body.session_key
async def gen():
t0 = time.perf_counter()
# Per-Turn-Trace: sammelt STT (davor, geparkt) + Vision + Hirn-TTFT + Generierung + Mem0
# (Rückruf während) zu EINEM Datensatz -> die Latenz-Ansicht zeigt den einzelnen Hänger.
trace = TurnTrace(session_id=body.session_id, kind="voice")
first = True
first_content = True
# Bildschirm-Sicht INNERHALB des Streams (C2-Fix): so startet die SSE-Antwort sofort und
# der Client bekommt ein Progress-Event (-> Lucy kann eine Warte-Ansage sprechen), statt
# dass der Request bis zu 120 s "tot" hängt, während das Vision-Modell beschreibt.
user_text = body.text
imgs = [u for u in (body.images or []) if u]
if imgs:
yield b'event: hermes.vision.progress\ndata: {"note": "Bildschirm wird angeschaut"}\n\n'
with Timer("vision"):
desc = await _describe_images(imgs, body.text)
if desc:
safe_desc = wrap_untrusted(desc, "BILDSCHIRM")
user_text = f"[Bildschirm-Sicht — das ist gerade auf dem/den Schirm(en) zu sehen:\n{safe_desc}\n]\n\n{body.text}"
messages = []
if body.system:
messages.append({"role": "system", "content": body.system})
messages.append({"role": "user", "content": user_text})
payload = {"model": body.model or HERMES_API_MODEL, "messages": messages, "stream": True}
# Lucys Hirn (Qwen3.6) ist ein Thinking-Modell -> für die gesprochene Assistentin Thinking AUS,
# sonst generiert es tausende Reasoning-Token VOR der kurzen Antwort (gemessen: 11k Token, ~30s TTFB).
# Gleiches Muster wie die fast-Spur im Gateway (gateway_proxy.py) und die Mem0-Extraktion.
if os.environ.get("MC_VOICE_NO_THINK", "1") not in ("0", "false", "False"):
payload["chat_template_kwargs"] = {"enable_thinking": False}
committed = False
def _commit() -> None:
nonlocal committed
if not committed:
committed = True
trace.commit()
try:
async with httpx.AsyncClient(timeout=httpx.Timeout(None, connect=5.0)) as client:
async with client.stream(
"POST", f"{HERMES_API_URL}/v1/chat/completions", json=payload, headers=headers,
) as r:
if r.status_code != 200:
detail = (await r.aread()).decode("utf-8", "replace")[:500]
yield f"data: {{\"error\": \"Hermes {r.status_code}: {detail}\"}}\n\n".encode()
return
async for chunk in r.aiter_raw():
if first: # Time-To-First-Byte des Hermes-Streams (Verbindungs-Overhead)
record_stage("chat_ttfb", (time.perf_counter() - t0) * 1000.0)
first = False
# Erster CONTENT-Delta = echte Hirn-Latenz (Agent-Overhead + LLM-TTFT) —
# chat_ttfb misst nur den SSE-Start (~5 ms) und ist dafür blind.
if first_content and b'"content"' in chunk:
record_stage("chat_first_content", (time.perf_counter() - t0) * 1000.0)
first_content = False
yield chunk
except httpx.HTTPError as exc:
yield f"data: {{\"error\": \"Verbindung zu Hermes fehlgeschlagen: {exc}\"}}\n\n".encode()
# Bildschirm-Sicht INNERHALB des Streams (C2-Fix): so startet die SSE-Antwort sofort und
# der Client bekommt ein Progress-Event (-> Lucy kann eine Warte-Ansage sprechen), statt
# dass der Request bis zu 120 s "tot" hängt, während das Vision-Modell beschreibt.
user_text = body.text
imgs = [u for u in (body.images or []) if u]
trace.had_images = bool(imgs)
if imgs:
yield b'event: hermes.vision.progress\ndata: {"note": "Bildschirm wird angeschaut"}\n\n'
_tv = time.perf_counter()
desc = await _describe_images(imgs, body.text)
trace.note_vision((time.perf_counter() - _tv) * 1000.0)
if desc:
safe_desc = wrap_untrusted(desc, "BILDSCHIRM")
user_text = f"[Bildschirm-Sicht — das ist gerade auf dem/den Schirm(en) zu sehen:\n{safe_desc}\n]\n\n{body.text}"
messages = []
if body.system:
messages.append({"role": "system", "content": body.system})
messages.append({"role": "user", "content": user_text})
trace.mark_brain_start() # ab hier zählt die Hirn-Zeit (Vision ist schon abgeschlossen)
payload = {"model": body.model or HERMES_API_MODEL, "messages": messages, "stream": True}
# Lucys Hirn (Qwen3.6) ist ein Thinking-Modell -> für die gesprochene Assistentin Thinking AUS,
# sonst generiert es tausende Reasoning-Token VOR der kurzen Antwort (gemessen: 11k Token, ~30s TTFB).
# Gleiches Muster wie die fast-Spur im Gateway (gateway_proxy.py) und die Mem0-Extraktion.
if os.environ.get("MC_VOICE_NO_THINK", "1") not in ("0", "false", "False"):
payload["chat_template_kwargs"] = {"enable_thinking": False}
try:
async with httpx.AsyncClient(timeout=httpx.Timeout(None, connect=5.0)) as client:
async with client.stream(
"POST", f"{HERMES_API_URL}/v1/chat/completions", json=payload, headers=headers,
) as r:
if r.status_code != 200:
detail = (await r.aread()).decode("utf-8", "replace")[:500]
trace.error = f"Hermes {r.status_code}"
yield f"data: {{\"error\": \"Hermes {r.status_code}: {detail}\"}}\n\n".encode()
return
async for chunk in r.aiter_raw():
if first: # Time-To-First-Byte des Hermes-Streams (Verbindungs-Overhead)
trace.note_ttfb()
first = False
# Erster CONTENT-Delta = echte Hirn-Latenz (Agent-Overhead + Mem0 + LLM-TTFT) —
# chat_ttfb misst nur den SSE-Start (~5 ms) und ist dafür blind.
if first_content and b'"content"' in chunk:
trace.note_first_content()
first_content = False
yield chunk
except httpx.HTTPError as exc:
trace.error = "verbindung"
yield f"data: {{\"error\": \"Verbindung zu Hermes fehlgeschlagen: {exc}\"}}\n\n".encode()
finally:
_commit() # Turn immer verbuchen (auch bei Fehler/Abbruch)
return StreamingResponse(gen(), media_type="text/event-stream")
+125 -6
View File
@@ -1,24 +1,44 @@
"""
Per-Stage-Latenz-Metriken für die Voice/Lucy-Pipeline.
Per-Stage-Latenz-Metriken + Per-Turn-Trace für die Voice/Lucy-Pipeline.
Misst die Server-seitige Dauer jeder Stufe (STT, Vision-Beschreibung, Chat-TTFB, TTS) und hält
rollende Statistiken (avg/p50/p95/last) im Speicher. Macht aus Latenz-VERMUTUNGEN gemessene Fakten
— die eigentliche Voraussetzung, um gezielt zu optimieren (Stufe 5/C2 des Reviews). Anzeige im
Frontend-Overhaul (E) analog zur TokenPerformanceCard.
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", "chat_ttfb", "tts")
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:
@@ -32,6 +52,36 @@ def record_stage(stage: str, ms: float) -> None:
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)."""
@@ -48,6 +98,68 @@ class Timer:
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}
@@ -66,3 +178,10 @@ 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]