78448220b2
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>
358 lines
16 KiB
Python
358 lines
16 KiB
Python
"""
|
|
Voice-Endpoints für „Mit Hermes reden" (Browser-Voice + 3D-Avatar).
|
|
|
|
Dünner Layer: STT/TTS werden zum Voice-Sidecar (:8650) geproxyt; der Chat geht an den
|
|
Hermes-`api_server` (:8642, OpenAI-kompatibel) — denselben vollen Agenten mit Tools +
|
|
geteiltem Mem0 wie CLI/Telegram. Mit stabilem `X-Hermes-Session-Id` hält die Plattform den
|
|
Transcript server-seitig, daher schickt der Client je Turn nur die neue User-Nachricht.
|
|
|
|
LAN-only (kein Token in der 2.0-Phase), wie die übrigen MC2-Endpoints.
|
|
"""
|
|
|
|
import logging
|
|
import os
|
|
import time
|
|
|
|
import httpx
|
|
from fastapi import APIRouter, File, Form, HTTPException, UploadFile
|
|
from fastapi.responses import Response, StreamingResponse
|
|
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 ( # 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
|
|
_GUARD_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), "mcp")
|
|
if _GUARD_DIR not in _sys.path:
|
|
_sys.path.insert(0, _GUARD_DIR)
|
|
try:
|
|
from guard import wrap_untrusted
|
|
except Exception: # den Voice-Pfad nie wegen des Filters lahmlegen
|
|
def wrap_untrusted(text: str, label: str = "") -> str:
|
|
return text
|
|
|
|
log = logging.getLogger(__name__)
|
|
router = APIRouter(prefix="/api")
|
|
|
|
# Bildschirm-Sicht: das DEDIZIERTE Vision-Modell (Qwen3-VL-8B) beschreibt das Bild; die Beschreibung
|
|
# geht als TEXT an Hermes -> Lucy behält ihr volles Hirn/Gedächtnis UND nutzt das bessere VL-Modell
|
|
# (statt der schwächeren Vision der fast-MoE). Per Env abschaltbar/umstellbar.
|
|
VISION_MODEL = os.environ.get("MC_VISION_MODEL", "vision")
|
|
# Knappe Beschreibung = schnellere VL-Generierung UND weniger Hermes-Kontext-Bloat (B2).
|
|
VISION_MAX_TOKENS = int(os.environ.get("MC_VISION_MAX_TOKENS", "280"))
|
|
|
|
|
|
async def _describe_images(image_urls: list[str], hint: str) -> str:
|
|
"""Lässt das Vision-Modell die Screenshots (1 je Monitor) knapp beschreiben (Deutsch).
|
|
Mehrere Bilder gehen in EINER Nachricht ans VL-Modell. Leerer String bei Fehler."""
|
|
multi = len(image_urls) > 1
|
|
intro = (f"Hier sind {len(image_urls)} Screenshots (je ein Monitor). Beschreibe auf Deutsch in höchstens "
|
|
"5 kurzen Sätzen das Wesentliche (pro Monitor: App/Fenster, wichtige Inhalte, sichtbarer Text/Code). "
|
|
"Keine Einleitung, keine Wiederholung der Frage. "
|
|
if multi else
|
|
"Beschreibe auf Deutsch in höchstens 5 kurzen Sätzen das Wesentliche auf diesem Screenshot "
|
|
"(App/Fenster, wichtige Inhalte, sichtbarer Text/Code). Keine Einleitung. ")
|
|
content: list = [{"type": "text", "text": intro + "Frage des Nutzers dazu: " + hint}]
|
|
for u in image_urls:
|
|
content.append({"type": "image_url", "image_url": {"url": u}})
|
|
try:
|
|
# 45 s statt 120 s: Qwen3-VL braucht warm ~5 s; wenn es 45 s nicht schafft, ist etwas
|
|
# kaputt und Lucy soll lieber ohne Bildschirm-Kontext antworten als ewig hängen.
|
|
async with httpx.AsyncClient(timeout=httpx.Timeout(float(os.environ.get("MC_VISION_TIMEOUT", "45")), connect=5.0)) as client:
|
|
r = await client.post(f"{LLAMA_SWAP_URL}/v1/chat/completions", json={
|
|
"model": VISION_MODEL, "max_tokens": VISION_MAX_TOKENS, "stream": False,
|
|
"messages": [{"role": "user", "content": content}],
|
|
})
|
|
r.raise_for_status()
|
|
return (r.json().get("choices") or [{}])[0].get("message", {}).get("content", "").strip()
|
|
except Exception as exc:
|
|
log.warning("Vision-Beschreibung fehlgeschlagen: %s", exc)
|
|
return ""
|
|
|
|
_TIMEOUT = httpx.Timeout(120.0, connect=5.0) # Chatterbox-TTS auf CPU darf dauern
|
|
|
|
|
|
class TTSIn(BaseModel):
|
|
text: str
|
|
engine: str = "piper"
|
|
voice: str = ""
|
|
language: str = ""
|
|
ref_path: str = ""
|
|
|
|
|
|
class ChatIn(BaseModel):
|
|
text: str # die neue User-Äußerung (STT-Ergebnis)
|
|
session_id: str # stabiler Voice-Faden → server-seitiger Transcript
|
|
session_key: str = "" # optional: Langzeit-Memory-Scope
|
|
system: str = "" # optionaler ephemerer System-Prompt (z.B. „antworte knapp/gesprochen")
|
|
model: str = ""
|
|
images: list[str] = [] # optionale Bildschirm-Sicht: ein data:-URL je Monitor (Lucys „Augen")
|
|
|
|
|
|
class AnnounceIn(BaseModel):
|
|
text: str # die Meldung (wird von Lucy gesprochen)
|
|
subject: str = "" # kurze Betreffzeile (z.B. "[Update]")
|
|
source: str = "" # Absender (sentry/notify/cron …) — nur fürs Log/Panel
|
|
priority: str = "normal" # 'silent' = nur im Verlauf zeigen, nicht sprechen
|
|
|
|
|
|
class AlarmIn(BaseModel):
|
|
text: str # die Alarm-Meldung
|
|
subject: str = "[Alarm]" # Betreff (Telegram-Präfix)
|
|
source: str = "alarm" # Absender fürs Log/Panel (z.B. "lucy-watchdog")
|
|
|
|
|
|
@router.post("/alarm")
|
|
def alarm(body: AlarmIn) -> dict:
|
|
"""Lucy-UNABHÄNGIGER Alarm-Weg: schickt direkt auf Telegram (und legt die Meldung in den
|
|
Briefkasten). Für Absender, die NICHT auf die sprechende Lucy zählen können — allen voran
|
|
der PC-seitige Lucy-Watchdog, wenn die Desktop-App selbst hängt (dann nützt der Briefkasten
|
|
nichts, weil niemand ihn vorliest → Telegram ist der einzige verlässliche Kanal). LAN-only
|
|
wie alle MC2-Endpoints."""
|
|
text = (body.text or "").strip()
|
|
if not text:
|
|
raise HTTPException(400, "Leere Meldung.")
|
|
subject = (body.subject or "[Alarm]").strip()
|
|
try:
|
|
item = announce.add(text, subject, body.source or "alarm", "normal")
|
|
except ValueError as exc:
|
|
raise HTTPException(400, str(exc))
|
|
announce.notify_telegram(subject, text) # best-effort Telegram (posix/bash; Windows = No-op)
|
|
return {"ok": True, "item": item}
|
|
|
|
|
|
@router.post("/voice/announce")
|
|
def voice_announce(body: AnnounceIn) -> dict:
|
|
"""Meldung in den Briefkasten legen (Lucy-Proaktivität). Absender: Health-Wächter,
|
|
notify.sh (Updates/Radar/Telegram-Spiegel), Hermes-cron. LAN-only wie alle MC2-Endpoints."""
|
|
try:
|
|
return {"ok": True, "item": announce.add(body.text, body.subject, body.source, body.priority)}
|
|
except ValueError as exc:
|
|
raise HTTPException(400, str(exc))
|
|
|
|
|
|
@router.get("/voice/announcements")
|
|
def voice_announcements(after: int | None = None, limit: int = 20) -> dict:
|
|
"""Neue Meldungen nach Cursor `after` abholen (Lucy pollt). Ohne `after` nur den
|
|
aktuellen Cursor-Stand (latest) — Erststart plappert so keine alten Meldungen nach."""
|
|
return announce.list_after(after, limit)
|
|
|
|
|
|
@router.get("/voice/health")
|
|
def voice_health() -> dict:
|
|
"""Erreichbarkeit des Voice-Sidecars + ob der Hermes-API-Key gesetzt ist."""
|
|
out: dict = {"sidecar": False, "hermes_key": bool(HERMES_API_KEY)}
|
|
try:
|
|
r = httpx.get(f"{VOICE_SERVICE_URL}/health", timeout=httpx.Timeout(5.0))
|
|
out["sidecar"] = r.status_code == 200
|
|
out["detail"] = r.json() if r.status_code == 200 else None
|
|
except Exception as exc: # noqa: BLE001
|
|
out["error"] = str(exc)
|
|
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:
|
|
r = httpx.get(f"{VOICE_SERVICE_URL}/voices", timeout=httpx.Timeout(10.0))
|
|
r.raise_for_status()
|
|
return r.json()
|
|
except Exception as exc: # noqa: BLE001
|
|
raise HTTPException(502, f"Voice-Sidecar nicht erreichbar: {exc}")
|
|
|
|
|
|
@router.post("/voice/stt")
|
|
async def voice_stt(audio: UploadFile = File(...), language: str = Form(default="")) -> dict:
|
|
"""Mikro-Audio → Text (Proxy auf Sidecar /stt)."""
|
|
data = await audio.read()
|
|
if not data:
|
|
raise HTTPException(400, "Leeres Audio.")
|
|
files = {"audio": (audio.filename or "rec.webm", data, audio.content_type or "audio/webm")}
|
|
try:
|
|
async with httpx.AsyncClient(timeout=_TIMEOUT) as client:
|
|
_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:
|
|
raise HTTPException(502, f"STT fehlgeschlagen: {exc}")
|
|
|
|
|
|
@router.post("/voice/turn")
|
|
async def voice_turn(audio: UploadFile = File(...)) -> dict:
|
|
"""Semantische Turn-Detection (Smart Turn v3): war die Äußerung fertig? Proxy → Sidecar."""
|
|
data = await audio.read()
|
|
if not data:
|
|
raise HTTPException(400, "Leeres Audio.")
|
|
files = {"audio": (audio.filename or "rec.wav", data, audio.content_type or "audio/wav")}
|
|
try:
|
|
async with httpx.AsyncClient(timeout=httpx.Timeout(10.0, connect=3.0)) as client:
|
|
with Timer("turn"):
|
|
r = await client.post(f"{VOICE_SERVICE_URL}/turn", files=files)
|
|
r.raise_for_status()
|
|
return r.json()
|
|
except httpx.HTTPError as exc:
|
|
# Turn-Check ist eine Optimierung — bei Ausfall lieber sofort antworten als hängen.
|
|
log.warning("Turn-Check fehlgeschlagen: %s", exc)
|
|
return {"complete": True, "probability": 1.0, "engine": "fallback"}
|
|
|
|
|
|
@router.post("/voice/reference")
|
|
async def voice_set_reference(audio: UploadFile = File(...)) -> dict:
|
|
"""Klon-Referenz (z.B. ElevenLabs-Erzeugnis) hochladen → Chatterbox nutzt sie. Proxy → Sidecar."""
|
|
data = await audio.read()
|
|
if not data:
|
|
raise HTTPException(400, "Leeres Audio.")
|
|
files = {"audio": (audio.filename or "ref.wav", data, audio.content_type or "audio/mpeg")}
|
|
try:
|
|
async with httpx.AsyncClient(timeout=_TIMEOUT) as client:
|
|
r = await client.post(f"{VOICE_SERVICE_URL}/reference", files=files)
|
|
r.raise_for_status()
|
|
return r.json()
|
|
except httpx.HTTPError as exc:
|
|
raise HTTPException(502, f"Referenz-Upload fehlgeschlagen: {exc}")
|
|
|
|
|
|
@router.get("/voice/reference")
|
|
def voice_get_reference() -> dict:
|
|
try:
|
|
r = httpx.get(f"{VOICE_SERVICE_URL}/reference", timeout=httpx.Timeout(8.0))
|
|
r.raise_for_status()
|
|
return r.json()
|
|
except Exception as exc: # noqa: BLE001
|
|
return {"active": False, "error": str(exc)}
|
|
|
|
|
|
@router.delete("/voice/reference")
|
|
def voice_clear_reference() -> dict:
|
|
try:
|
|
r = httpx.delete(f"{VOICE_SERVICE_URL}/reference", timeout=httpx.Timeout(8.0))
|
|
r.raise_for_status()
|
|
return r.json()
|
|
except httpx.HTTPError as exc:
|
|
raise HTTPException(502, f"Löschen fehlgeschlagen: {exc}")
|
|
|
|
|
|
@router.post("/voice/tts")
|
|
async def voice_tts(body: TTSIn) -> Response:
|
|
"""Text → Sprache (Proxy auf Sidecar /tts), liefert WAV-Bytes."""
|
|
try:
|
|
async with httpx.AsyncClient(timeout=_TIMEOUT) as client:
|
|
with Timer("tts"):
|
|
r = await client.post(f"{VOICE_SERVICE_URL}/tts", json=body.model_dump())
|
|
r.raise_for_status()
|
|
return Response(content=r.content, media_type=r.headers.get("content-type", "audio/wav"))
|
|
except httpx.HTTPError as exc:
|
|
raise HTTPException(502, f"TTS fehlgeschlagen: {exc}")
|
|
|
|
|
|
@router.post("/voice/chat")
|
|
async def voice_chat(body: ChatIn) -> StreamingResponse:
|
|
"""Neue User-Äußerung → Hermes-Agent (api_server, streamend). SSE wird 1:1 durchgereicht.
|
|
|
|
Mit `X-Hermes-Session-Id` hält die Plattform den Verlauf — wir senden nur die neue Nachricht.
|
|
Auth per Bearer (API_SERVER_KEY); ohne Key liefert :8642 ein 401."""
|
|
if not HERMES_API_KEY:
|
|
raise HTTPException(503, "HERMES_API_KEY/API_SERVER_KEY nicht gesetzt — Agent-Auth fehlt.")
|
|
|
|
headers = {
|
|
"Authorization": f"Bearer {HERMES_API_KEY}",
|
|
"X-Hermes-Session-Id": body.session_id,
|
|
}
|
|
if body.session_key:
|
|
headers["X-Hermes-Session-Key"] = body.session_key
|
|
|
|
async def gen():
|
|
# 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
|
|
committed = False
|
|
|
|
def _commit() -> None:
|
|
nonlocal committed
|
|
if not committed:
|
|
committed = True
|
|
trace.commit()
|
|
|
|
try:
|
|
# 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")
|