feat: lower memory dedupe threshold for more aggressive cleaning
This commit is contained in:
+121
-121
@@ -1,121 +1,121 @@
|
||||
import os
|
||||
|
||||
from fastapi import APIRouter, Request
|
||||
from fastapi.responses import JSONResponse, StreamingResponse
|
||||
|
||||
from config import LLAMA_SWAP_URL
|
||||
from services.gateway_stream import record_stream_chunk, record_usage
|
||||
from services.router_logic import LANES, choose_for_lane
|
||||
from services.routing_policy import load_policy
|
||||
|
||||
router = APIRouter(prefix="/v1")
|
||||
|
||||
# Antwortsprache für IDE-/Lane-Traffic: die Coding-Modelle antworten sonst englisch
|
||||
# (User-Anforderung 03.07.2026). Leerer String (MC_GATEWAY_LANG_DIRECTIVE="") schaltet ab.
|
||||
_LANG_DIRECTIVE = os.environ.get(
|
||||
"MC_GATEWAY_LANG_DIRECTIVE",
|
||||
"Antworte dem Nutzer grundsätzlich auf Deutsch (Erklärungen, Pläne, Rückfragen, "
|
||||
"Zusammenfassungen) — auch wenn die Frage oder Tool-Anweisungen englisch sind. "
|
||||
"Quellcode, Bezeichner und Shell-Befehle bleiben unverändert.")
|
||||
|
||||
|
||||
def _inject_language(body: dict, alias: str) -> None:
|
||||
"""Deutsch-Direktive anhängen. An die ERSTE System-Message (viele Chat-Templates
|
||||
erwarten nur eine), sonst als neue System-Message. `hermes` ausgenommen — Lucys
|
||||
Persona (SOUL.md) regelt die Sprache selbst."""
|
||||
if not _LANG_DIRECTIVE or alias == "hermes":
|
||||
return
|
||||
msgs = body.get("messages")
|
||||
if not isinstance(msgs, list):
|
||||
return
|
||||
first_sys = next((m for m in msgs if isinstance(m, dict) and m.get("role") == "system"), None)
|
||||
if first_sys is None:
|
||||
msgs.insert(0, {"role": "system", "content": _LANG_DIRECTIVE})
|
||||
elif isinstance(first_sys.get("content"), str):
|
||||
first_sys["content"] = first_sys["content"].rstrip() + "\n\n" + _LANG_DIRECTIVE
|
||||
elif isinstance(first_sys.get("content"), list):
|
||||
first_sys["content"].append({"type": "text", "text": _LANG_DIRECTIVE})
|
||||
|
||||
# Virtuelle Lanes, die der Gateway zusätzlich zu den echten Modellen als „Modell" anbietet.
|
||||
_LANE_LABELS = {"coding": "Coding (Router → coder/heavy/fast)", "chat": "Chat (Router → fast/heavy)"}
|
||||
|
||||
|
||||
@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)
|
||||
data = r.json()
|
||||
# Lanes ganz oben einblenden, damit IDEs einfach „coding"/„chat" wählen können.
|
||||
lanes = [{"id": lane, "object": "model", "owned_by": "mc2-router",
|
||||
"description": _LANE_LABELS.get(lane, lane)} for lane in LANES]
|
||||
# Kontextlänge je Modell mitliefern (aus der llama-swap-Config geparst). Ohne sie
|
||||
# budgetieren Clients blind — Hermes-Subagents nahmen 256k an, schickten passende
|
||||
# max_tokens und rissen damit den echten Server-Kontext (Radar-Lauf 02.07.).
|
||||
# Rollen-Aliase (heavy/coder/hermes …) tauchen bei llama-swap NICHT als Einträge auf,
|
||||
# Clients fragen aber genau damit an → als eigene Einträge einblenden.
|
||||
ctx_map: dict[str, int] = {}
|
||||
alias_entries: list[dict] = []
|
||||
try:
|
||||
from services import llamaswap
|
||||
for m in llamaswap.list_models():
|
||||
ctx = m.get("ctx")
|
||||
if ctx:
|
||||
for api_id in m.get("api_ids", []):
|
||||
ctx_map[api_id] = ctx
|
||||
for alias in m.get("aliases", []):
|
||||
entry = {"id": alias, "object": "model", "owned_by": "mc2-alias",
|
||||
"description": f"Alias für {m['name']}"}
|
||||
if ctx:
|
||||
entry["context_length"] = ctx
|
||||
alias_entries.append(entry)
|
||||
except Exception:
|
||||
pass
|
||||
if isinstance(data, dict) and isinstance(data.get("data"), list):
|
||||
for entry in data["data"]:
|
||||
if (ctx := ctx_map.get(entry.get("id"))):
|
||||
entry.setdefault("context_length", ctx)
|
||||
data["data"] = lanes + alias_entries + data["data"]
|
||||
return JSONResponse(data, status_code=r.status_code)
|
||||
|
||||
|
||||
async def _proxy(path: str, request: Request):
|
||||
body = await request.json()
|
||||
requested = str(body.get("model") or "auto")
|
||||
if requested.lower() in ("auto", "chat", "coding"):
|
||||
lane = requested.lower()
|
||||
alias, reason = choose_for_lane(lane, body)
|
||||
body["model"] = alias
|
||||
routed = {"x-mc-routed-to": alias, "x-mc-route-reason": reason, "x-mc-lane": lane}
|
||||
else:
|
||||
alias = requested
|
||||
routed = {"x-mc-routed-to": requested}
|
||||
# fast-Spur: Thinking aus für flotte Antworten (sofern Client es nicht selbst setzt).
|
||||
pol = load_policy()
|
||||
if pol["fast_no_think"] and alias == pol["fast"] and "chat_template_kwargs" not in body:
|
||||
body["chat_template_kwargs"] = {"enable_thinking": False}
|
||||
_inject_language(body, alias)
|
||||
url = f"{LLAMA_SWAP_URL}{path}"
|
||||
|
||||
client = request.app.state.gw_client # geteilter Keep-Alive-Client (siehe app.py lifespan)
|
||||
if body.get("stream"):
|
||||
async def gen():
|
||||
async with client.stream("POST", url, json=body, timeout=None) as r:
|
||||
async for chunk in r.aiter_raw():
|
||||
record_stream_chunk(chunk, alias)
|
||||
yield chunk
|
||||
return StreamingResponse(gen(), media_type="text/event-stream", headers=routed)
|
||||
|
||||
r = await client.post(url, json=body, timeout=600.0)
|
||||
resp_json = r.json()
|
||||
record_usage(resp_json.get("usage") if isinstance(resp_json, dict) else None, alias)
|
||||
return JSONResponse(resp_json, status_code=r.status_code, headers=routed)
|
||||
|
||||
|
||||
@router.post("/chat/completions")
|
||||
async def chat_completions(request: Request):
|
||||
return await _proxy("/v1/chat/completions", request)
|
||||
|
||||
|
||||
@router.post("/completions")
|
||||
async def completions(request: Request):
|
||||
return await _proxy("/v1/completions", request)
|
||||
import os
|
||||
|
||||
from fastapi import APIRouter, Request
|
||||
from fastapi.responses import JSONResponse, StreamingResponse
|
||||
|
||||
from config import LLAMA_SWAP_URL
|
||||
from services.gateway_stream import record_stream_chunk, record_usage
|
||||
from services.router_logic import LANES, choose_for_lane
|
||||
from services.routing_policy import load_policy
|
||||
|
||||
router = APIRouter(prefix="/v1")
|
||||
|
||||
# Antwortsprache für IDE-/Lane-Traffic: die Coding-Modelle antworten sonst englisch
|
||||
# (User-Anforderung 03.07.2026). Leerer String (MC_GATEWAY_LANG_DIRECTIVE="") schaltet ab.
|
||||
_LANG_DIRECTIVE = os.environ.get(
|
||||
"MC_GATEWAY_LANG_DIRECTIVE",
|
||||
"Antworte dem Nutzer grundsätzlich auf Deutsch (Erklärungen, Pläne, Rückfragen, "
|
||||
"Zusammenfassungen) — auch wenn die Frage oder Tool-Anweisungen englisch sind. "
|
||||
"Quellcode, Bezeichner und Shell-Befehle bleiben unverändert.")
|
||||
|
||||
|
||||
def _inject_language(body: dict, alias: str) -> None:
|
||||
"""Deutsch-Direktive anhängen. An die ERSTE System-Message (viele Chat-Templates
|
||||
erwarten nur eine), sonst als neue System-Message. `hermes` ausgenommen — Lucys
|
||||
Persona (SOUL.md) regelt die Sprache selbst."""
|
||||
if not _LANG_DIRECTIVE or alias == "hermes":
|
||||
return
|
||||
msgs = body.get("messages")
|
||||
if not isinstance(msgs, list):
|
||||
return
|
||||
first_sys = next((m for m in msgs if isinstance(m, dict) and m.get("role") == "system"), None)
|
||||
if first_sys is None:
|
||||
msgs.insert(0, {"role": "system", "content": _LANG_DIRECTIVE})
|
||||
elif isinstance(first_sys.get("content"), str):
|
||||
first_sys["content"] = first_sys["content"].rstrip() + "\n\n" + _LANG_DIRECTIVE
|
||||
elif isinstance(first_sys.get("content"), list):
|
||||
first_sys["content"].append({"type": "text", "text": _LANG_DIRECTIVE})
|
||||
|
||||
# Virtuelle Lanes, die der Gateway zusätzlich zu den echten Modellen als „Modell" anbietet.
|
||||
_LANE_LABELS = {"coding": "Coding (Router → coder/heavy/fast)", "chat": "Chat (Router → fast/heavy)"}
|
||||
|
||||
|
||||
@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)
|
||||
data = r.json()
|
||||
# Lanes ganz oben einblenden, damit IDEs einfach „coding"/„chat" wählen können.
|
||||
lanes = [{"id": lane, "object": "model", "owned_by": "mc2-router",
|
||||
"description": _LANE_LABELS.get(lane, lane)} for lane in LANES]
|
||||
# Kontextlänge je Modell mitliefern (aus der llama-swap-Config geparst). Ohne sie
|
||||
# budgetieren Clients blind — Hermes-Subagents nahmen 256k an, schickten passende
|
||||
# max_tokens und rissen damit den echten Server-Kontext (Radar-Lauf 02.07.).
|
||||
# Rollen-Aliase (heavy/coder/hermes …) tauchen bei llama-swap NICHT als Einträge auf,
|
||||
# Clients fragen aber genau damit an → als eigene Einträge einblenden.
|
||||
ctx_map: dict[str, int] = {}
|
||||
alias_entries: list[dict] = []
|
||||
try:
|
||||
from services import llamaswap
|
||||
for m in llamaswap.list_models():
|
||||
ctx = m.get("ctx")
|
||||
if ctx:
|
||||
for api_id in m.get("api_ids", []):
|
||||
ctx_map[api_id] = ctx
|
||||
for alias in m.get("aliases", []):
|
||||
entry = {"id": alias, "object": "model", "owned_by": "mc2-alias",
|
||||
"description": f"Alias für {m['name']}"}
|
||||
if ctx:
|
||||
entry["context_length"] = ctx
|
||||
alias_entries.append(entry)
|
||||
except Exception:
|
||||
pass
|
||||
if isinstance(data, dict) and isinstance(data.get("data"), list):
|
||||
for entry in data["data"]:
|
||||
if (ctx := ctx_map.get(entry.get("id"))):
|
||||
entry.setdefault("context_length", ctx)
|
||||
data["data"] = lanes + alias_entries + data["data"]
|
||||
return JSONResponse(data, status_code=r.status_code)
|
||||
|
||||
|
||||
async def _proxy(path: str, request: Request):
|
||||
body = await request.json()
|
||||
requested = str(body.get("model") or "auto")
|
||||
if requested.lower() in ("auto", "chat", "coding"):
|
||||
lane = requested.lower()
|
||||
alias, reason = choose_for_lane(lane, body)
|
||||
body["model"] = alias
|
||||
routed = {"x-mc-routed-to": alias, "x-mc-route-reason": reason, "x-mc-lane": lane}
|
||||
else:
|
||||
alias = requested
|
||||
routed = {"x-mc-routed-to": requested}
|
||||
# fast-Spur: Thinking aus für flotte Antworten (sofern Client es nicht selbst setzt).
|
||||
pol = load_policy()
|
||||
if pol["fast_no_think"] and alias == pol["fast"] and "chat_template_kwargs" not in body:
|
||||
body["chat_template_kwargs"] = {"enable_thinking": False}
|
||||
_inject_language(body, alias)
|
||||
url = f"{LLAMA_SWAP_URL}{path}"
|
||||
|
||||
client = request.app.state.gw_client # geteilter Keep-Alive-Client (siehe app.py lifespan)
|
||||
if body.get("stream"):
|
||||
async def gen():
|
||||
async with client.stream("POST", url, json=body, timeout=None) as r:
|
||||
async for chunk in r.aiter_raw():
|
||||
record_stream_chunk(chunk, alias)
|
||||
yield chunk
|
||||
return StreamingResponse(gen(), media_type="text/event-stream", headers=routed)
|
||||
|
||||
r = await client.post(url, json=body, timeout=600.0)
|
||||
resp_json = r.json()
|
||||
record_usage(resp_json.get("usage") if isinstance(resp_json, dict) else None, alias)
|
||||
return JSONResponse(resp_json, status_code=r.status_code, headers=routed)
|
||||
|
||||
|
||||
@router.post("/chat/completions")
|
||||
async def chat_completions(request: Request):
|
||||
return await _proxy("/v1/chat/completions", request)
|
||||
|
||||
|
||||
@router.post("/completions")
|
||||
async def completions(request: Request):
|
||||
return await _proxy("/v1/completions", request)
|
||||
|
||||
+130
-130
@@ -1,130 +1,130 @@
|
||||
"""System-Endpoints: Live-Status + Wartung (Restart/Self-Update — auf der Box).
|
||||
|
||||
Wartung läuft als systemd-USER-Dienst → KEIN sudo/Passwort (Nordstern).
|
||||
Lokal (Windows) schlagen die Shell-Befehle harmlos fehl und werden als Fehler
|
||||
zurückgegeben statt zu crashen.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import os
|
||||
import subprocess
|
||||
|
||||
from fastapi import APIRouter
|
||||
from pydantic import BaseModel
|
||||
|
||||
import httpx
|
||||
|
||||
from config import GATEWAY_URL, HERMES_API_URL, LLAMA_SWAP_URL, MEM0_SERVICE_URL, VOICE_SERVICE_URL
|
||||
from services import backup as backup_svc
|
||||
from services import maintenance
|
||||
from services.agent import agent_status
|
||||
from services.gateway import gateway_reachable
|
||||
from services.llamaswap import engine_reachable, list_models
|
||||
from services.pricing import compute_savings
|
||||
from services.system import system_status
|
||||
from services.token_stats import get_stats
|
||||
|
||||
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()
|
||||
|
||||
|
||||
def _mem0_reachable() -> bool:
|
||||
try:
|
||||
return httpx.get(f"{MEM0_SERVICE_URL}/health", timeout=2).status_code == 200
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def _voice_reachable() -> bool:
|
||||
try:
|
||||
return httpx.get(f"{VOICE_SERVICE_URL}/health", timeout=2).status_code == 200
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
@router.get("/system/services")
|
||||
def services() -> dict:
|
||||
"""Aggregierte Erreichbarkeit aller Stack-Dienste (für die Health-Anzeige)."""
|
||||
a = agent_status()
|
||||
gw_url = f"{GATEWAY_URL}/v1"
|
||||
return {
|
||||
"services": [
|
||||
{"name": "Engine (llama-swap)", "unit": "llama-swap", "url": LLAMA_SWAP_URL, "ok": engine_reachable()},
|
||||
{"name": "Gateway (integriert)", "unit": "mission-control-2", "url": gw_url, "ok": gateway_reachable()},
|
||||
{"name": "Hermes-Gateway", "unit": "hermes-gateway", "url": HERMES_API_URL, "ok": a["gateway_reachable"]},
|
||||
{"name": "Hermes-Terminal", "unit": "hermes-terminal", "url": a["terminal_url"], "ok": a["terminal_reachable"]},
|
||||
{"name": "Mem0 (Gedächtnis)", "unit": "mem0-service", "url": MEM0_SERVICE_URL, "ok": _mem0_reachable()},
|
||||
{"name": "Voice (STT/TTS)", "unit": "voice-service", "url": VOICE_SERVICE_URL, "ok": _voice_reachable()},
|
||||
],
|
||||
"links": {
|
||||
"engine_ui": f"{LLAMA_SWAP_URL}/ui",
|
||||
"gateway": gw_url,
|
||||
"hermes_terminal": a["terminal_url"],
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@router.post("/system/backup")
|
||||
def backup() -> dict:
|
||||
return backup_svc.backup_now()
|
||||
|
||||
|
||||
@router.get("/system/backups")
|
||||
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: # noqa: BLE001
|
||||
return {"ok": False, "code": -1, "out": "", "err": str(exc)}
|
||||
|
||||
|
||||
class RestartReq(BaseModel):
|
||||
service: str
|
||||
|
||||
|
||||
@router.post("/system/restart")
|
||||
def restart(req: RestartReq) -> dict:
|
||||
"""Dienst neustarten. Delegiert an den Wartungs-Service (EINE Allowlist-Wahrheit):
|
||||
der kennt System-Dienste (llama-swap, via sudo -n) UND User-Dienste und wird auch von
|
||||
der UI (/api/maintenance/restart) genutzt. Vorher lag hier eine zweite, veraltete Liste
|
||||
(ohne llama-swap/hermes-terminal, mit Geist-Eintrag hermes-webui) → Engine-Neustart per
|
||||
Sprache/MCP schlug fehl."""
|
||||
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)."""
|
||||
# Rolle je Modell/Alias (lowercase) für die Tarif-Auflösung auflösen.
|
||||
role_map: dict[str, str | None] = {}
|
||||
try:
|
||||
for m in list_models():
|
||||
role_map[m["name"].lower()] = m.get("role")
|
||||
for alias in m.get("aliases", []):
|
||||
role_map[alias.lower()] = m.get("role")
|
||||
except Exception:
|
||||
log.warning("token_stats: list_models fehlgeschlagen, Tarife per Name", exc_info=True)
|
||||
return compute_savings(get_stats(), role_map)
|
||||
"""System-Endpoints: Live-Status + Wartung (Restart/Self-Update — auf der Box).
|
||||
|
||||
Wartung läuft als systemd-USER-Dienst → KEIN sudo/Passwort (Nordstern).
|
||||
Lokal (Windows) schlagen die Shell-Befehle harmlos fehl und werden als Fehler
|
||||
zurückgegeben statt zu crashen.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import os
|
||||
import subprocess
|
||||
|
||||
from fastapi import APIRouter
|
||||
from pydantic import BaseModel
|
||||
|
||||
import httpx
|
||||
|
||||
from config import GATEWAY_URL, HERMES_API_URL, LLAMA_SWAP_URL, MEM0_SERVICE_URL, VOICE_SERVICE_URL
|
||||
from services import backup as backup_svc
|
||||
from services import maintenance
|
||||
from services.agent import agent_status
|
||||
from services.gateway import gateway_reachable
|
||||
from services.llamaswap import engine_reachable, list_models
|
||||
from services.pricing import compute_savings
|
||||
from services.system import system_status
|
||||
from services.token_stats import get_stats
|
||||
|
||||
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()
|
||||
|
||||
|
||||
def _mem0_reachable() -> bool:
|
||||
try:
|
||||
return httpx.get(f"{MEM0_SERVICE_URL}/health", timeout=2).status_code == 200
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def _voice_reachable() -> bool:
|
||||
try:
|
||||
return httpx.get(f"{VOICE_SERVICE_URL}/health", timeout=2).status_code == 200
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
@router.get("/system/services")
|
||||
def services() -> dict:
|
||||
"""Aggregierte Erreichbarkeit aller Stack-Dienste (für die Health-Anzeige)."""
|
||||
a = agent_status()
|
||||
gw_url = f"{GATEWAY_URL}/v1"
|
||||
return {
|
||||
"services": [
|
||||
{"name": "Engine (llama-swap)", "unit": "llama-swap", "url": LLAMA_SWAP_URL, "ok": engine_reachable()},
|
||||
{"name": "Gateway (integriert)", "unit": "mission-control-2", "url": gw_url, "ok": gateway_reachable()},
|
||||
{"name": "Hermes-Gateway", "unit": "hermes-gateway", "url": HERMES_API_URL, "ok": a["gateway_reachable"]},
|
||||
{"name": "Hermes-Terminal", "unit": "hermes-terminal", "url": a["terminal_url"], "ok": a["terminal_reachable"]},
|
||||
{"name": "Mem0 (Gedächtnis)", "unit": "mem0-service", "url": MEM0_SERVICE_URL, "ok": _mem0_reachable()},
|
||||
{"name": "Voice (STT/TTS)", "unit": "voice-service", "url": VOICE_SERVICE_URL, "ok": _voice_reachable()},
|
||||
],
|
||||
"links": {
|
||||
"engine_ui": f"{LLAMA_SWAP_URL}/ui",
|
||||
"gateway": gw_url,
|
||||
"hermes_terminal": a["terminal_url"],
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@router.post("/system/backup")
|
||||
def backup() -> dict:
|
||||
return backup_svc.backup_now()
|
||||
|
||||
|
||||
@router.get("/system/backups")
|
||||
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: # noqa: BLE001
|
||||
return {"ok": False, "code": -1, "out": "", "err": str(exc)}
|
||||
|
||||
|
||||
class RestartReq(BaseModel):
|
||||
service: str
|
||||
|
||||
|
||||
@router.post("/system/restart")
|
||||
def restart(req: RestartReq) -> dict:
|
||||
"""Dienst neustarten. Delegiert an den Wartungs-Service (EINE Allowlist-Wahrheit):
|
||||
der kennt System-Dienste (llama-swap, via sudo -n) UND User-Dienste und wird auch von
|
||||
der UI (/api/maintenance/restart) genutzt. Vorher lag hier eine zweite, veraltete Liste
|
||||
(ohne llama-swap/hermes-terminal, mit Geist-Eintrag hermes-webui) → Engine-Neustart per
|
||||
Sprache/MCP schlug fehl."""
|
||||
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)."""
|
||||
# Rolle je Modell/Alias (lowercase) für die Tarif-Auflösung auflösen.
|
||||
role_map: dict[str, str | None] = {}
|
||||
try:
|
||||
for m in list_models():
|
||||
role_map[m["name"].lower()] = m.get("role")
|
||||
for alias in m.get("aliases", []):
|
||||
role_map[alias.lower()] = m.get("role")
|
||||
except Exception:
|
||||
log.warning("token_stats: list_models fehlgeschlagen, Tarife per Name", exc_info=True)
|
||||
return compute_savings(get_stats(), role_map)
|
||||
|
||||
+357
-357
@@ -1,357 +1,357 @@
|
||||
"""
|
||||
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")
|
||||
"""
|
||||
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")
|
||||
|
||||
Reference in New Issue
Block a user