phase1b: Backend entruempelt und robuster (35 tote Routen raus, Sperren, ehrliche Update-Pruefung)
Ampel / ampel (push) Successful in 26s

Ballast raus:
- 35 Routen ohne Nutzer entfernt (agent/*, fit, roles, ctx, drafts, groups, routing/policy,
  system/history, system/self-update, maintenance/reboot, zeitmaschine/inhalt, zeitplan,
  voice/health|metrics|trace|voices|reference|tts). Von 95 auf 60.
- Tote Module geloescht: agent-Router, roles, agent_aktivitaet, metrics_history (samt
  10-s-Sampler), voice_metrics, migrate_config, parse_mc2_timeout, scripts/.
- Unbenutzte Funktionen und Konstanten entfernt (Modell-Upgrade-Empfehlung, Draft-/Kontext-
  Setzer, Konsole, PC-Ausfuehrer-Probe, Routing-Policy-Editor ...).

Robuster:
- Jobs in eigener Prozessgruppe (Abbrechen beendet wirklich alles), Zeitlimit je Job-Art,
  start_job_exklusiv: zwei Klicks starten kein doppeltes Update mehr; alte Jobs raeumen sich auf.
- Update-Pruefung meldet Fehler (pruef_fehler, Lampe "Pruefung unklar") statt "aktuell".
- Nach jedem Update sofort neu pruefen (update_stand) statt 10 Minuten alten Stand zeigen.
- llama-swap-Config: Sperre (RLock + flock) fuer UI, Radar, Aufraeumen und Hirn-Umstellung.
- Hermes-Config: bei Lesefehler nichts schreiben, atomar, mit Sicherung.
- Live-Strom und Gateway-Warnung blockieren den Event-Loop nicht mehr (Lucy, OpenChamber).
- Gateway antwortet bei Engine-Ausfall im OpenAI-Fehlerformat (502) statt nacktem 500.
- Abgestuerzte Waechter-Pruefung wird ein gelber Hinweis statt still zu verschwinden.
- Download laedt nur den gewuenschten Quant (vorher bei Fehlen alle Teile aller Varianten),
  Download-Jobs in Gruppe "download"; HF-Suche kodiert den Suchbegriff.
- Herkunftspruefung: schreibende /api-Aufrufe fremder Webseiten werden abgelehnt (keine
  Anmeldung, User-Entscheid); Skripte, Desktop-Lucy und /v1 unveraendert.
- Modellpfade: Eintragen und Loeschen nur innerhalb von MODELS_DIR.
- Dienste-Liste fragt keine abgebauten Dienste mehr ab (PC-Ausfuehrer haette 3 s gekostet).
- SSE-Fehlerzeilen von /api/voice/chat als gueltiges JSON.
- mission-control-2.service: --timeout-graceful-shutdown 3 (Neustart ohne 10-s-Haenger).

Tests: 92 gruen (neu: Herkunft, Quant-Auswahl, abgestuerzte Pruefung).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Hitonabi
2026-09-24 15:08:44 +02:00
co-authored by Claude Opus 5.5
parent c8718fcde0
commit e9f488b56c
42 changed files with 799 additions and 2197 deletions
-52
View File
@@ -1,52 +0,0 @@
"""Agent-Endpoint: Hermes-Status + WebUI-Link (MC verlinkt nur, betreibt nicht)."""
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from services import agent_aktivitaet
from services.agent import agent_status, hermes_brain_info, set_agent_brain, update_brain_model
router = APIRouter(prefix="/api")
class BrainReq(BaseModel):
model: str
class SetBrainReq(BaseModel):
model_id: str
@router.get("/agent/status")
def status() -> dict:
return agent_status()
@router.get("/agent/brain")
def brain_info() -> dict:
"""Aktuelles Agent-Hirn (hermes) + bestes NousResearch-Hermes-Update."""
return hermes_brain_info()
@router.post("/agent/brain/set")
def set_brain(body: SetBrainReq) -> dict:
"""Setzt ein installiertes Modell als Agent-Hirn (Alias + warm-Gruppe + Config)."""
res = set_agent_brain(body.model_id)
if not res.get("ok"):
raise HTTPException(400, res.get("reason", "Fehler beim Setzen des Agent-Hirns"))
return res
@router.post("/agent/brain")
def set_brain_model(body: BrainReq) -> dict:
ok = update_brain_model(body.model)
return {"ok": ok}
@router.get("/agent/aktivitaet")
def aktivitaet(limit: int = 60) -> dict:
"""Werkzeug-Verlauf des Agenten (v3-Umbau P6).
Gelesen aus Hermes' eigenem Log — MC2 patcht dort nichts, es schaut nur zu. Was
NICHT drin steht (Denkstrom, Werkzeug-Argumente), verspricht die Ansicht auch nicht.
"""
return agent_aktivitaet.uebersicht(max(1, min(limit, 300)))
+23 -10
View File
@@ -5,7 +5,6 @@ Box-Wart-Schnittstellen (Umbau 09/2026): alles, was der Cockpit-Startbildschirm
GET /api/hinweise Hinweise + Verlauf des Wächters
POST /api/hinweise/{id}/aktion/{aktion} Knopf eines Hinweises ausführen
GET /api/modelle/nutzung Wer nutzt die Modelle (7 Tage, 24 h je Stunde)
GET /api/zeitplan Heute gelaufen / demnächst geplant
GET /api/updates/verlauf Was die letzten Update-Läufe wirklich gebracht haben
POST /api/updates/festgehalten/{b}/freigeben Festgehaltenen Baustein wieder freigeben
GET /api/modelle/aufraeumen Was auf der Modell-Platte ungenutzt Platz belegt
@@ -40,11 +39,12 @@ router = APIRouter(prefix="/api", tags=["boxwart"])
LOCAL_TZ = ZoneInfo(os.environ.get("MC_LOCAL_TZ", "Europe/Berlin"))
_updates_lock = threading.Lock()
_updates_cache: dict = {"ts": 0.0, "daten": None, "laeuft": False}
_updates_cache: dict = {"ts": 0.0, "daten": None, "laeuft": False, "stand": -1}
UPDATES_CACHE_S = 600
def _updates_holen() -> None:
stand = maintenance.update_stand
try:
daten = maintenance.updates()
except Exception:
@@ -53,6 +53,7 @@ def _updates_holen() -> None:
if daten is not None or _updates_cache["daten"] is None:
_updates_cache["daten"] = daten or {}
_updates_cache["ts"] = time.time()
_updates_cache["stand"] = stand
_updates_cache["laeuft"] = False
@@ -61,7 +62,10 @@ def _updates_gecacht() -> dict:
gern zehn Sekunden und mehr. Der Start wartet darauf nie: Er bekommt den letzten Stand (oder
zunächst nichts), frisch geholt wird im Hintergrund."""
with _updates_lock:
veraltet = _updates_cache["daten"] is None or time.time() - _updates_cache["ts"] > UPDATES_CACHE_S
# Nach jedem abgeschlossenen Update zählt maintenance.update_stand hoch — dann sofort neu holen,
# statt bis zu 10 Minuten „Aktualisieren“ für schon eingespielte Updates zu zeigen (24.09.2026).
veraltet = (_updates_cache["daten"] is None or time.time() - _updates_cache["ts"] > UPDATES_CACHE_S
or _updates_cache["stand"] != maintenance.update_stand)
if veraltet and not _updates_cache["laeuft"]:
_updates_cache["laeuft"] = True
threading.Thread(target=_updates_holen, name="updates-holen", daemon=True).start()
@@ -83,6 +87,15 @@ def _bausteine(u: dict) -> list[dict]:
return liste
NAMEN_UNKLAR = {"os": "Betriebssystem", "engine": "Motor", "swap": "llama-swap", "hermes": "Hermes"}
def _unklar(u: dict) -> list[dict]:
"""Bausteine, deren Update-Prüfung scheiterte — ehrlich „unbekannt“ statt „aktuell“."""
return [{"id": k, "name": NAMEN_UNKLAR.get(k, k), "grund": grund}
for k, grund in (u.get("pruef_fehler") or {}).items() if grund]
def _zeit_kurz(ts: float | None) -> str:
if not ts:
return ""
@@ -143,12 +156,17 @@ def _lampen(stand: dict, u: dict) -> list[dict]:
n = len(_bausteine(u))
fest = len(update_verlauf.festgehalten())
unklar = len(_unklar(u))
if fest:
lampen.append(lampe("updates", "Updates", "warn", f"{fest} festgehalten"))
elif _updates_cache["daten"] is None:
lampen.append(lampe("updates", "Updates", "aus", "wird geprüft"))
elif n:
lampen.append(lampe("updates", "Updates", "info", f"{n} bereit"))
elif unklar:
lampen.append(lampe("updates", "Updates", "warn", "Prüfung unklar"))
else:
lampen.append(lampe("updates", "Updates", "info" if n else "ok", f"{n} bereit" if n else "aktuell"))
lampen.append(lampe("updates", "Updates", "ok", "aktuell"))
return lampen
@@ -183,7 +201,7 @@ def start() -> dict:
"verlauf": stand["verlauf"][:20],
"lampen": _lampen(stand, u),
"box": _box(),
"updates": {"bausteine": _bausteine(u), "modelle": u.get("model_list") or [],
"updates": {"bausteine": _bausteine(u), "unklar": _unklar(u),
"festgehalten": update_verlauf.festgehalten(),
"gelesen": _updates_cache["daten"] is not None},
"flugplan": zeitplan.flugplan(),
@@ -208,11 +226,6 @@ def nutzung() -> dict:
return modell_nutzung.nutzung()
@router.get("/zeitplan")
def zeitplan_route() -> dict:
return zeitplan.flugplan()
@router.get("/updates/verlauf")
def updates_verlauf(anzahl: int = 8) -> dict:
"""Die letzten Update-Läufe mit dem Ergebnis je Baustein (aus dem Meldeprotokoll der Box)."""
+6 -6
View File
@@ -28,7 +28,6 @@ import json
import logging
from pathlib import Path
from config import MODELS_DIR
from fastapi import APIRouter, Request
from fastapi.responses import StreamingResponse
@@ -69,8 +68,6 @@ def _fingerprints() -> dict[str, object]:
)
except Exception:
pass
# Erinnerungen: Datei-mtime (das Auftragsbuch ist seit 04.09.2026 ausgebaut)
fp["reminders"] = _mtime(MODELS_DIR / "mc2-reminders.json")
return fp
@@ -80,7 +77,9 @@ async def _strom(request: Request, mit_metrik: bool):
Die Basislinie entsteht JE VERBINDUNG (der Client hat beim Verbinden frisch geladen) —
ein globaler Snapshot würde bei mehreren Clients Events verschlucken.
"""
alt = _fingerprints()
# Die Abdrücke fragen u. a. llama-swap ab (bis 2 s). Im Event-Loop blockierte das jede andere
# Anfrage dieses Prozesses, auch Lucys /v1-Weiterleitung — seit 24.09.2026 im Thread.
alt = await asyncio.to_thread(_fingerprints)
yield ": verbunden\n\n"
seit_abdruck = 0.0
@@ -97,7 +96,8 @@ async def _strom(request: Request, mit_metrik: bool):
if mit_metrik:
try:
from services.system import metrik_punkt
yield f"event: metrik\ndata: {json.dumps(metrik_punkt())}\n\n"
punkt = await asyncio.to_thread(metrik_punkt)
yield f"event: metrik\ndata: {json.dumps(punkt)}\n\n"
seit_ping = 0.0
except Exception:
# Ein kaputter Messpunkt darf den Strom nicht reißen — die Ansicht fällt
@@ -106,7 +106,7 @@ async def _strom(request: Request, mit_metrik: bool):
if seit_abdruck >= ABDRUCK_S:
seit_abdruck = 0.0
neu = _fingerprints()
neu = await asyncio.to_thread(_fingerprints)
keys = [k for k, v in neu.items() if k in alt and v != alt[k]]
alt.update(neu)
if keys:
+24 -4
View File
@@ -161,10 +161,20 @@ def _inject_language(body: dict, alias: str) -> None:
elif isinstance(first_sys.get("content"), list):
first_sys["content"].append({"type": "text", "text": _LANG_DIRECTIVE})
def _upstream_fehler(text: str, status: int = 502, headers: dict | None = None) -> JSONResponse:
"""Fehler im OpenAI-Format statt eines nackten 500 (24.09.2026): Hermes, OpenChamber und Lucy
können damit umgehen und zeigen den Grund an."""
return JSONResponse({"error": {"message": text, "type": "upstream_error"}},
status_code=status, headers=headers)
@router.get("/models")
async def models(request: Request):
client = request.app.state.gw_client # geteilter Keep-Alive-Client (siehe app.py lifespan)
r = await client.get(f"{LLAMA_SWAP_URL}/v1/models", timeout=10.0)
try:
r = await client.get(f"{LLAMA_SWAP_URL}/v1/models", timeout=10.0)
except httpx.HTTPError as exc:
return _upstream_fehler(f"Engine nicht erreichbar: {exc.__class__.__name__}")
try:
data = r.json()
except ValueError:
@@ -275,7 +285,10 @@ async def _proxy(path: str, request: Request):
url = f"{LLAMA_SWAP_URL}{path}"
if body.get("stream"):
req = client.build_request("POST", url, json=body, timeout=None)
r = await client.send(req, stream=True)
try:
r = await client.send(req, stream=True)
except httpx.HTTPError as exc:
return _upstream_fehler(f"Engine nicht erreichbar: {exc.__class__.__name__}", headers=routed)
if r.status_code != 200:
await r.aread()
try:
@@ -293,8 +306,15 @@ async def _proxy(path: str, request: Request):
await r.aclose()
return StreamingResponse(gen(), media_type="text/event-stream", headers=routed)
r = await client.post(url, json=body, timeout=600.0)
resp_json = r.json()
try:
r = await client.post(url, json=body, timeout=600.0)
except httpx.HTTPError as exc:
return _upstream_fehler(f"Engine nicht erreichbar: {exc.__class__.__name__}", headers=routed)
try:
resp_json = r.json()
except ValueError:
log.warning("%s: Engine-Antwort ist kein JSON (HTTP %s): %.200s", path, r.status_code, r.text)
return _upstream_fehler(f"Engine lieferte keine gültige Antwort (HTTP {r.status_code}).", headers=routed)
if isinstance(resp_json, dict):
record_usage(resp_json.get("usage"), alias)
if any(isinstance(c, dict) and c.get("finish_reason") == "length"
-5
View File
@@ -78,11 +78,6 @@ def update_all() -> dict:
return maintenance.update_all_job()
@router.post("/maintenance/reboot")
def reboot() -> dict:
return maintenance.reboot()
@router.post("/maintenance/restart")
def restart(body: RestartReq) -> dict:
return maintenance.restart_service(body.service)
+10 -102
View File
@@ -1,11 +1,12 @@
"""Modelle-Endpoints: Liste (mit Caps), Discover, Fit, Register, Groups."""
"""Modelle-Endpoints: Liste, Discover, Suche und Installation bei Hugging Face, Jobs, Rollen,
Laden/Entladen/Löschen. Fit-, Kontext-, Draft- und Gruppen-Routen sind am 24.09.2026 entfallen
(ohne Nutzer seit dem Box-Wart-Umbau)."""
import psutil
from config import HF_DOWNLOAD_ENV, MODELS_DIR
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from services import budget, discover, geheimnisse, hf, jobengine, llamaswap
from services.fit import evaluate_fit, max_ctx_for
router = APIRouter(prefix="/api")
@@ -29,27 +30,6 @@ def discover_models(force: bool = False) -> dict:
return {**data, "sys_ram_gb": round(ram, 1)}
@router.get("/fit")
def fit(params_b: float = 0, quant: str = "Q4_K_M", ctx: int = 8192,
name: str = "", role: str = "") -> dict:
"""Hardware-Fit-Vorschau. params_b<=0 → aus KATALOG (echte Metadaten, MoE-bewusst)
oder sonst aus dem Namen geschätzt. assigned_ctx = der ctx, der TATSÄCHLICH vergeben
würde: SETUP-BEWUSST (neben Hirn/warmem Set), nicht nur gegen den Gesamt-RAM.
So sieht die 'Erweiterte Ansicht' vor dem Download Ampel + echten ctx."""
ram = _ram_gb()
pb = params_b if params_b > 0 else budget.params_b_for(name)
saw = budget.setup_aware_ctx(pb, quant, role=role or None)
return {
"params_b": round(pb, 1),
"fit": evaluate_fit(pb, quant, ctx, ram, name=name),
"optimal_ctx": max_ctx_for(pb, quant, ram), # Roh-Obergrenze (Modell allein)
"assigned_ctx": saw["ctx"], # setup-bewusst vergeben
"budget": {"gtt_gb": saw["gtt_gb"], "reserved_gb": saw["reserved_gb"],
"budget_gb": saw["budget_gb"], "mode": saw["mode"]},
"sys_ram_gb": round(ram, 1),
}
class RegisterReq(BaseModel):
model_path: str
role: str | None = None
@@ -61,6 +41,10 @@ class RegisterReq(BaseModel):
@router.post("/models/register")
def register(req: RegisterReq) -> dict:
# Nur Dateien aus dem Modellordner (24.09.2026): der Pfad landet im Startbefehl der Engine.
for pfad in (req.model_path, req.mmproj_path):
if pfad and not llamaswap.im_modellordner(pfad):
raise HTTPException(400, f"Pfad liegt nicht im Modellordner ({MODELS_DIR}): {pfad}")
try:
model_id = llamaswap.register_model(
req.model_path, role=req.role, ctx=req.ctx, ttl=req.ttl,
@@ -148,8 +132,9 @@ def install(req: InstallReq) -> dict:
env = dict(HF_DOWNLOAD_ENV)
if token := geheimnisse.hf_token():
env["HF_TOKEN"] = token
job_id = jobengine.start_job(args, f"download {req.repo}", env=env,
on_done=_apply_role if role else None)
# Gruppe „download“: so taucht der Download in der Oberfläche mit Fortschritt und Abbrechen auf.
job_id = jobengine.start_job(args, f"Download {repo}", env=env, group="download",
on_done=_apply_role if role else None, zeitlimit_s=6 * 3600)
jobengine.attach_download_progress(job_id, str(target), info["total_bytes"])
return {"ok": True, "job_id": job_id, "model_id": model_id, "model_path": model_path,
"total_bytes": info["total_bytes"], "files": len(info["files"])}
@@ -169,14 +154,6 @@ class RoleReq(BaseModel):
role: str | None = None
@router.get("/roles/{role}/recommend")
def recommend_role(role: str) -> dict:
"""Welches installierte Modell passt am besten auf diese Rolle? (Capability + setup-
bewusster Fit). Basis für 'Empfohlen'-Hinweis + Auto-Pick im Rollen-Zuweisungs-Modal."""
from services import roles
return roles.recommend_for_role(role)
@router.post("/models/{model_id}/role")
def set_model_role(model_id: str, body: RoleReq) -> dict:
new_role = (body.role or "").strip().lower()
@@ -201,54 +178,6 @@ def set_model_role(model_id: str, body: RoleReq) -> dict:
return {"ok": True}
class CtxReq(BaseModel):
ctx: int
@router.get("/models/{model_id}/ctx/auto")
def auto_ctx(model_id: str) -> dict:
"""Setup-bewusster Optimal-ctx für ein bestehendes Modell (Rolle/Params/Quant +
aktuelles Setup). Basis für den 'Auto'-Button an der Modellkarte."""
m = next((x for x in llamaswap.list_models() if x["name"] == model_id), None)
if not m:
raise HTTPException(404, "Modell nicht gefunden")
saw = budget.setup_aware_ctx_for_model(m)
return {"model_id": model_id, "current_ctx": m.get("ctx"),
"params_b": round(budget.params_of_model(m), 1), "quant": m.get("quant"),
"role": m.get("role"), **saw}
@router.post("/models/{model_id}/ctx")
def set_model_ctx(model_id: str, body: CtxReq) -> dict:
if not llamaswap.set_ctx(model_id, body.ctx):
raise HTTPException(404, "Modell nicht gefunden")
return {"ok": True}
@router.get("/models/drafts")
def list_drafts(target: str = "") -> dict:
"""Verfügbare Draft-Modelle + ihre Vocab-Kompatibilität zum Ziel-Modell
(target = GGUF-Pfad). Basis für die idiotensichere Spec-Draft-Auswahl im UI."""
return llamaswap.drafts_for(target)
class DraftReq(BaseModel):
draft_path: str | None = None
@router.post("/models/{model_id}/draft")
def set_model_draft(model_id: str, body: DraftReq) -> dict:
"""Setzt/entfernt den Speculative-Decoding-Draft eines Modells. Inkompatible
(oder nicht prüfbare) Drafts werden serverseitig abgelehnt."""
try:
res = llamaswap.set_spec_draft(model_id, body.draft_path)
except PermissionError as exc:
raise HTTPException(500, str(exc))
if not res["ok"]:
raise HTTPException(400 if "kompatib" in res["reason"].lower() else 404, res["reason"])
return res
@router.post("/models/unload")
def unload_all_models() -> dict:
import httpx
@@ -297,24 +226,3 @@ def delete(model_id: str) -> dict:
if not llamaswap.delete_model(model_id):
raise HTTPException(404, "Modell nicht gefunden")
return {"ok": True}
@router.get("/groups")
def groups() -> dict:
return {"groups": llamaswap.list_groups()}
class GroupReq(BaseModel):
group: str
members: list[str]
swap: bool = False
persist: bool = False
@router.put("/groups")
def set_group(req: GroupReq) -> dict:
try:
llamaswap.set_group(req.group, req.members, swap=req.swap, persist=req.persist)
except PermissionError as exc:
raise HTTPException(500, str(exc))
return {"ok": True}
+3 -33
View File
@@ -1,9 +1,8 @@
"""Routing-Endpoints: Lane-Summary (chat/coding) + UI-editierbare Policy (hot-reload)."""
"""Routing-Übersicht (Lanes, Gateway erreichbar). Der Policy-Editor ist seit dem Box-Wart-Umbau
aus der Oberfläche raus; die Policy selbst liest der Gateway weiter aus mc2-routing.json."""
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from fastapi import APIRouter
from services import gateway
from services.routing_policy import policy_meta, save_policy
router = APIRouter(prefix="/api")
@@ -11,32 +10,3 @@ router = APIRouter(prefix="/api")
@router.get("/routing")
def routing() -> dict:
return {**gateway.routing_summary(), "gateway_reachable": gateway.gateway_reachable()}
@router.get("/routing/policy")
def get_policy() -> dict:
"""Aktuelle Policy + Defaults (für „Zurücksetzen“) + Feld-Spezifikation für den Editor."""
return policy_meta()
class PolicyPatch(BaseModel):
fast: str | None = None
heavy: str | None = None
coder: str | None = None
coder_lite: str | None = None
heavy_chars: int | None = None
coding_escalate_chars: int | None = None
fast_no_think: bool | None = None
@router.put("/routing/policy")
def put_policy(patch: PolicyPatch) -> dict:
"""Teil-Update der Routing-Policy. Validiert, persistiert atomar, sofort wirksam (hot-reload)."""
fields = {k: v for k, v in patch.model_dump().items() if v is not None}
if not fields:
raise HTTPException(status_code=400, detail="Keine Felder zum Aktualisieren.")
try:
new_policy = save_policy(fields)
except (ValueError, TypeError) as e:
raise HTTPException(status_code=400, detail=f"Ungültige Policy: {e}")
return {"policy": new_policy}
+1 -30
View File
@@ -1,4 +1,4 @@
"""System-Endpoints: Live-Status + Wartung (Restart/Self-Update — auf der Box).
"""System-Endpoints: Live-Status, Dienste-Liste, Sicherung, Neustart, Token-Verbrauch.
Wartung läuft als systemd-USER-Dienst → KEIN sudo/Passwort (Nordstern).
Lokal (Windows) schlagen die Shell-Befehle harmlos fehl und werden als Fehler
@@ -26,22 +26,11 @@ log = logging.getLogger(__name__)
router = APIRouter(prefix="/api")
# Quelle für Self-Update (auf der Box ~/mission-control-v2).
SOURCE_DIR = os.path.expanduser(os.environ.get("MC2_SOURCE_DIR", "~/mission-control-v2"))
@router.get("/system/status")
def status() -> dict:
return system_status()
@router.get("/system/history")
def history(minutes: int = 60) -> dict:
"""Metrik-Verlauf (max. 24 h) für die Cockpit-Zeitachse — s. services/metrics_history."""
from services import metrics_history
return metrics_history.history(minutes)
def _voice_reachable() -> bool:
try:
return httpx.get(f"{VOICE_SERVICE_URL}/health", timeout=2).status_code == 200
@@ -124,15 +113,6 @@ def backups() -> dict:
return {"backups": backup_svc.list_backups()}
def _run(cmd: list[str], cwd: str | None = None) -> dict:
try:
p = subprocess.run(cmd, cwd=cwd, capture_output=True, text=True, timeout=180)
return {"ok": p.returncode == 0, "code": p.returncode,
"out": (p.stdout or "")[-2000:], "err": (p.stderr or "")[-2000:]}
except Exception as exc:
return {"ok": False, "code": -1, "out": "", "err": str(exc)}
class RestartReq(BaseModel):
service: str
@@ -147,15 +127,6 @@ def restart(req: RestartReq) -> dict:
return maintenance.restart_service(req.service)
@router.post("/system/self-update")
def self_update() -> dict:
"""git pull (Source) → venv-Deps → Dienst-Restart. Auf der Box; lokal Fehler."""
pull = _run(["git", "fetch", "--all"], cwd=SOURCE_DIR)
reset = _run(["git", "reset", "--hard", "origin/main"], cwd=SOURCE_DIR)
restart_res = _run(["systemctl", "--user", "restart", "mission-control-2"])
return {"pull": pull, "reset": reset, "restart": restart_res}
@router.get("/system/token-stats")
def token_stats() -> dict:
"""Token-Verbrauch + Cloud-Ersparnis. Logik im pricing-Service (SSoT)."""
+289 -425
View File
@@ -1,425 +1,289 @@
"""
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 +
eigenem Gedächtnis 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
# Injection-Schutz (Stufe 0): guard.py liegt im mcp/-Verzeichnis. Per Pfad laden (eigene MC2-Venv).
import sys as _sys
import time
import httpx
from config import HERMES_API_KEY, HERMES_API_MODEL, HERMES_API_URL, LLAMA_SWAP_URL, LUCY_STIMME_URL, VOICE_SERVICE_URL
from fastapi import APIRouter, File, Form, HTTPException, UploadFile
from fastapi.responses import Response, StreamingResponse
from pydantic import BaseModel
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,
)
_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:
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). 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:
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:
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 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).
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 + Gedächtnis + 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")
# ---------------------------------------------------------------------------------------------
# Lucys Stimme ins LAN reichen (Raphael-Umbau 04.09.2026). lucy-stimme.service (:8021, pocket-tts
# german_24l) bindet nur Loopback; die Desktop-Lucy am PC spricht seit dem Umbau nicht mehr mit
# einem eigenen pocket_server, sondern mit DIESEM — dieselbe Stimme wie die Telegram-Sprachnachrichten.
# Dünner Proxy, API 1:1 (pocket_server: /health, /tts -> WAV, /tts/stream -> PCM16 + X-Sample-Rate).
# LAN-only wie alle MC2-Endpoints.
class LucyTtsIn(BaseModel):
text: str
emo: str | None = None # Stimmungs-Profil (pocket_server EMO_PROFILES); Raphael-Lucy setzt keins
@router.get("/lucy/stimme/health")
def lucy_stimme_health() -> dict:
"""Bereitschaft von Lucys Stimme (pocket_server /health: status ok|loading)."""
try:
r = httpx.get(f"{LUCY_STIMME_URL}/health", timeout=httpx.Timeout(5.0))
r.raise_for_status()
return r.json()
except Exception as exc:
raise HTTPException(502, f"Lucys Stimme (:8021) nicht erreichbar: {exc}")
@router.post("/lucy/stimme/tts")
async def lucy_stimme_tts(body: LucyTtsIn) -> Response:
"""Text -> WAV (ganzer Text). Warm-up der Desktop-Lucy + Jobs, die eine Datei brauchen."""
try:
async with httpx.AsyncClient(timeout=httpx.Timeout(600.0, connect=5.0)) as client:
with Timer("lucy_tts"):
r = await client.post(f"{LUCY_STIMME_URL}/tts", json=body.model_dump(exclude_none=True))
r.raise_for_status()
return Response(content=r.content, media_type=r.headers.get("content-type", "audio/wav"),
headers={k: v for k, v in r.headers.items() if k.lower().startswith("x-")})
except httpx.HTTPStatusError as exc:
raise HTTPException(exc.response.status_code, f"Lucys Stimme: {exc.response.text[:200]}")
except httpx.HTTPError as exc:
raise HTTPException(502, f"Lucys Stimme nicht erreichbar: {exc}")
@router.post("/lucy/stimme/tts/stream")
async def lucy_stimme_tts_stream(body: LucyTtsIn) -> StreamingResponse:
"""Text -> rohes PCM16-mono, satzweise gestreamt (Samplerate im Header X-Sample-Rate).
Der Live-Pfad der Desktop-Lucy: erstes Audio nach dem ersten Satz. Der Upstream-Stream bleibt
offen, solange der Client liest — bricht der Client ab (Barge-in), schließt httpx den Upstream."""
client = httpx.AsyncClient(timeout=httpx.Timeout(None, connect=5.0))
try:
req = client.build_request("POST", f"{LUCY_STIMME_URL}/tts/stream", json=body.model_dump(exclude_none=True))
upstream = await client.send(req, stream=True)
except httpx.HTTPError as exc:
await client.aclose()
raise HTTPException(502, f"Lucys Stimme nicht erreichbar: {exc}")
if upstream.status_code != 200:
detail = (await upstream.aread()).decode("utf-8", "replace")[:200]
await upstream.aclose(); await client.aclose()
raise HTTPException(upstream.status_code, f"Lucys Stimme: {detail}")
async def gen():
try:
async for chunk in upstream.aiter_raw():
yield chunk
finally:
await upstream.aclose()
await client.aclose()
return StreamingResponse(gen(), media_type="application/octet-stream",
headers={"X-Sample-Rate": upstream.headers.get("x-sample-rate", "24000")})
"""
Sprach- und Melde-Endpunkte für Lucy (Desktop-App, Telegram-Spiegel, Wächter).
Dünner Layer: STT und die Turn-Erkennung gehen an den Voice-Sidecar (:8650, schläft seit
24.09.2026 bis zur Android-App), der Chat an den Hermes-`api_server` (:8642, OpenAI-kompatibel) —
derselbe volle Agent mit Werkzeugen und Gedächtnis wie Telegram. Mit stabilem
`X-Hermes-Session-Id` hält die Plattform den Verlauf, der Client schickt je Turn nur die neue
Nachricht. Lucys Stimme (pocket-tts, :8021) wird ins LAN gereicht.
Abgebaut am 24.09.2026 (ohne Nutzer): Voice-Health, Latenz-Metriken und -Trace, Stimmenliste,
Klon-Referenz und das alte Piper/Chatterbox-TTS.
"""
import json
import logging
import os
# Injection-Schutz (Stufe 0): guard.py liegt im mcp/-Verzeichnis. Per Pfad laden (eigene MC2-Venv).
import sys as _sys
import httpx
from config import HERMES_API_KEY, HERMES_API_MODEL, HERMES_API_URL, LLAMA_SWAP_URL, LUCY_STIMME_URL, VOICE_SERVICE_URL
from fastapi import APIRouter, File, Form, HTTPException, UploadFile
from fastapi.responses import Response, StreamingResponse
from pydantic import BaseModel
from services import announce
_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 Bild-Modell beschreibt die Screenshots; die Beschreibung geht als TEXT an
# Hermes -> Lucy behält ihr volles Hirn und Gedächtnis. Per Env umstellbar.
VISION_MODEL = os.environ.get("MC_VISION_MODEL", "vision")
# Knappe Beschreibung = schnellere Generierung UND weniger Kontext für Hermes.
VISION_MAX_TOKENS = int(os.environ.get("MC_VISION_MAX_TOKENS", "280"))
_STT_TIMEOUT = httpx.Timeout(120.0, connect=5.0)
def _sse_fehler(text: str) -> bytes:
"""SSE-Fehlerzeile als gültiges JSON (Anführungszeichen im Text brachen früher den Lucy-Client)."""
return f"data: {json.dumps({'error': text}, ensure_ascii=False)}\n\n".encode()
async def _describe_images(image_urls: list[str], hint: str) -> str:
"""Lässt das Bild-Modell die Screenshots (1 je Monitor) knapp beschreiben (Deutsch).
Mehrere Bilder gehen in EINER Nachricht ans 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: Wenn das Bild-Modell so lange braucht, 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 ""
class ChatIn(BaseModel):
text: str # die neue User-Äußerung (STT-Ergebnis)
session_id: str # stabiler Gesprächsfaden → server-seitiger Verlauf
session_key: str = "" # optional an Hermes durchgereicht (X-Hermes-Session-Key)
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. Der Betreff „[Alarm]“ geht
auch nachts sofort raus (notify.sh)."""
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: Wächter, notify.sh
(Updates/Radar/Telegram-Spiegel), Hermes-Cron."""
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.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=_STT_TIMEOUT) as client:
r = await client.post(f"{VOICE_SERVICE_URL}/stt", files=files, data={"language": language})
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:
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/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():
# Bildschirm-Sicht INNERHALB des Streams: so startet die SSE-Antwort sofort und der Client
# bekommt ein Progress-Event (-> Lucy kann eine Warte-Ansage sprechen), statt dass der
# Request hängt, während das Bild-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'
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 ist ein Thinking-Modell -> für die gesprochene Assistentin Thinking AUS, sonst
# generiert es tausende Reasoning-Token VOR der kurzen Antwort (gemessen: 11k Token, ~30 s).
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]
yield _sse_fehler(f"Hermes {r.status_code}: {detail}")
return
async for chunk in r.aiter_raw():
yield chunk
except httpx.HTTPError as exc:
yield _sse_fehler(f"Verbindung zu Hermes fehlgeschlagen: {exc}")
return StreamingResponse(gen(), media_type="text/event-stream")
# ---------------------------------------------------------------------------------------------
# Lucys Stimme ins LAN reichen (Raphael-Umbau 04.09.2026). lucy-stimme.service (:8021, pocket-tts
# german_24l) bindet nur Loopback; die Desktop-Lucy am PC spricht mit DIESEM — dieselbe Stimme wie
# die Telegram-Sprachnachrichten. Dünner Proxy, API 1:1 (pocket_server: /health, /tts -> WAV,
# /tts/stream -> PCM16 + X-Sample-Rate).
class LucyTtsIn(BaseModel):
text: str
emo: str | None = None # Stimmungs-Profil (pocket_server EMO_PROFILES); Raphael-Lucy setzt keins
@router.get("/lucy/stimme/health")
def lucy_stimme_health() -> dict:
"""Bereitschaft von Lucys Stimme (pocket_server /health: status ok|loading)."""
try:
r = httpx.get(f"{LUCY_STIMME_URL}/health", timeout=httpx.Timeout(5.0))
r.raise_for_status()
return r.json()
except Exception as exc:
raise HTTPException(502, f"Lucys Stimme (:8021) nicht erreichbar: {exc}")
@router.post("/lucy/stimme/tts")
async def lucy_stimme_tts(body: LucyTtsIn) -> Response:
"""Text -> WAV (ganzer Text). Warm-up der Desktop-Lucy + Jobs, die eine Datei brauchen."""
try:
async with httpx.AsyncClient(timeout=httpx.Timeout(600.0, connect=5.0)) as client:
r = await client.post(f"{LUCY_STIMME_URL}/tts", json=body.model_dump(exclude_none=True))
r.raise_for_status()
return Response(content=r.content, media_type=r.headers.get("content-type", "audio/wav"),
headers={k: v for k, v in r.headers.items() if k.lower().startswith("x-")})
except httpx.HTTPStatusError as exc:
raise HTTPException(exc.response.status_code, f"Lucys Stimme: {exc.response.text[:200]}")
except httpx.HTTPError as exc:
raise HTTPException(502, f"Lucys Stimme nicht erreichbar: {exc}")
@router.post("/lucy/stimme/tts/stream")
async def lucy_stimme_tts_stream(body: LucyTtsIn) -> StreamingResponse:
"""Text -> rohes PCM16-mono, satzweise gestreamt (Samplerate im Header X-Sample-Rate).
Der Live-Pfad der Desktop-Lucy: erstes Audio nach dem ersten Satz. Der Upstream-Stream bleibt
offen, solange der Client liest — bricht der Client ab (Barge-in), schließt httpx den Upstream."""
client = httpx.AsyncClient(timeout=httpx.Timeout(None, connect=5.0))
try:
req = client.build_request("POST", f"{LUCY_STIMME_URL}/tts/stream", json=body.model_dump(exclude_none=True))
upstream = await client.send(req, stream=True)
except httpx.HTTPError as exc:
await client.aclose()
raise HTTPException(502, f"Lucys Stimme nicht erreichbar: {exc}")
if upstream.status_code != 200:
detail = (await upstream.aread()).decode("utf-8", "replace")[:200]
await upstream.aclose(); await client.aclose()
raise HTTPException(upstream.status_code, f"Lucys Stimme: {detail}")
async def gen():
try:
async for chunk in upstream.aiter_raw():
yield chunk
finally:
await upstream.aclose()
await client.aclose()
return StreamingResponse(gen(), media_type="application/octet-stream",
headers={"X-Sample-Rate": upstream.headers.get("x-sample-rate", "24000")})
-11
View File
@@ -31,17 +31,6 @@ def list_snapshots() -> dict:
return {"available": os.name == "posix", "backups": backup_svc.list_backups()}
@router.get("/zeitmaschine/inhalt")
def snapshot_contents(file: str) -> dict:
"""Top-Level-Komponenten eines Snapshots (hermes, llama-swap, MANIFEST …)."""
if not _FILE_RX.match(file):
raise HTTPException(400, "Ungültiger Snapshot-Name.")
p = backup_svc.BACKUP_DIR / file
if not p.is_file():
raise HTTPException(404, "Snapshot nicht gefunden.")
return {"file": file, "components": backup_svc.snapshot_components(p)}
@router.post("/zeitmaschine/restore")
def restore(body: RestoreIn) -> dict:
"""Wiederherstellung starten (detached). restore.sh macht VORHER selbst ein