import asyncio from pathlib import Path import yaml from fastapi import APIRouter, Depends, HTTPException, Request, WebSocket, WebSocketDisconnect from fastapi.responses import FileResponse from pydantic import BaseModel from app.admin_llm import ( LlmControlError, llm_status, restart_gateway_detached, switch_backend, ) from app.audit import log_admin_action from app.auth import is_admin_user, require_admin, require_admin_or_user from app.config import settings from app.dependencies import get_store from app.runtime_config import RUNTIME_SETTABLE, invalidate_cache, runtime_settings from app.schemas import AdminUserPrefsUpdate, MemoryCreate, MemoryOut, UserCreate, UserCreated, UserUpdate, UserAdminUpdate router = APIRouter() _CONFIG_DIR = Path(__file__).resolve().parents[2] / "config" _ALLOWED_SECTIONS = {"abbreviations", "units", "terms"} _ALLOWED_LANGS = {"de", "en", "fr", "es", "it", "nl", "ru", "zh"} class PronunciationEntry(BaseModel): section: str key: str value: str @router.get("/admin/request-headers") async def request_headers(request: Request, key: str | None = None): """Discovery: zeigt die eingehenden HTTP-Header + Quell-IP. Hilft, hinter dem Reverse-Proxy/SSO den richtigen Identitaets-Header (TRUSTED_AUTH_HEADER) festzustellen. Henne-Ei: Beim Einrichten kennt das Gateway den SSO-Admin noch nicht. Daher ist der Endpoint auch per Admin-Key aufrufbar - als Query (`?key=...`, im Browser durchs SSO bequem) oder X-Admin-Key-Header. Nach dem Setup wieder meiden bzw. den Key rotieren (er landet sonst in Proxy-Logs). """ expected = settings.admin_api_key.strip() provided = key or request.headers.get("x-admin-key") ok = bool(expected) and provided is not None and provided.strip() == expected if not ok and not is_admin_user(_current_user_or_none(request)): raise HTTPException(status_code=403, detail="Admin privileges required") return { "client": request.client.host if request.client else None, "headers": dict(request.headers), } def _current_user_or_none(request: Request): from app.auth import authenticate, _bearer_token client_host = request.client.host if request.client else "" token = _bearer_token(request.headers.get("authorization")) return authenticate(request.headers, client_host, token) @router.post("/admin/users", response_model=UserCreated, dependencies=[Depends(require_admin)]) async def create_user(payload: UserCreate): """Legt einen Nutzer an und gibt das Bearer-Token EINMALIG zurueck.""" user, token = get_store().create_user(payload.display_name) return UserCreated(user_id=user.id, display_name=user.display_name, token=token) @router.delete("/admin/users/{user_id}", dependencies=[Depends(require_admin)]) async def delete_user(user_id: str): """Loescht einen Nutzer und alle seine Daten (Sessions, Nachrichten, Erinnerungen, Nutzungsdaten). Der anonyme Nutzer kann nicht geloescht werden.""" try: deleted = get_store().delete_user(user_id) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) if not deleted: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") return {"deleted": user_id} @router.post("/admin/users/{user_id}/token", response_model=UserCreated, dependencies=[Depends(require_admin)]) async def reset_token(user_id: str): """Stellt einen neuen Bearer-Token aus; der alte wird sofort ungueltig. Der neue Token wird EINMALIG zurueckgegeben und danach nicht mehr angezeigt.""" result = get_store().reset_token(user_id) if result is None: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") user, token = result return UserCreated(user_id=user.id, display_name=user.display_name, token=token) @router.put("/admin/users/{user_id}/admin", dependencies=[Depends(require_admin)]) async def set_user_admin(user_id: str, payload: UserAdminUpdate): """Schaltet Admin-Rechte (persistentes DB-Flag) für einen Nutzer ein/aus.""" user = get_store().set_user_admin(user_id, payload.is_admin) if user is None: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") return {"user_id": user.id, "is_admin": user.is_admin} @router.put("/admin/users/{user_id}/prefs", dependencies=[Depends(require_admin)]) async def update_user_prefs(user_id: str, payload: AdminUserPrefsUpdate): """Setzt Nutzer-Einstellungen (z. B. erlaubte Sprachen). Merge mit bestehenden Prefs.""" store = get_store() current = store.get_user(user_id) if current is None: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") new = {k: v for k, v in payload.model_dump().items() if v is not None} updated = store.set_user_prefs(user_id, {**current.prefs, **new}) return {"user_id": updated.id, "prefs": updated.prefs} @router.put("/admin/users/{user_id}", dependencies=[Depends(require_admin)]) async def update_user(user_id: str, payload: UserUpdate): """Aktualisiert den Anzeigenamen eines Nutzers (z. B. nach erstem SSO-Login).""" user = get_store().update_display_name(user_id, payload.display_name) if user is None: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") return {"user_id": user.id, "display_name": user.display_name} @router.post("/admin/users/{user_id}/memories", response_model=MemoryOut, dependencies=[Depends(require_admin)]) async def add_user_memory(user_id: str, payload: MemoryCreate): """Legt eine Erinnerung fuer einen Nutzer an (Admin kann Kontext vorbelegen).""" store = get_store() if store.get_user(user_id) is None: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") memory = store.add_memory(user_id, payload.content) return MemoryOut(id=memory.id, content=memory.content, created_at=memory.created_at) @router.get("/admin/users", dependencies=[Depends(require_admin_or_user)]) async def list_users(): """Listet die Nutzer (ohne Secrets). Fuer Admins (SSO/ADMIN_USERS) oder ADMIN_API_KEY.""" return [ { "user_id": u.id, "display_name": u.display_name, "external_id": u.external_id, "created_at": u.created_at, "allowed_languages": (u.prefs or {}).get("allowed_languages", ""), "is_admin": is_admin_user(u), } for u in get_store().list_users() ] @router.get("/admin/users/{user_id}/memories", dependencies=[Depends(require_admin)]) async def get_user_memories(user_id: str): """Gibt alle Erinnerungen eines Nutzers zurueck.""" store = get_store() if store.get_user(user_id) is None: raise HTTPException(status_code=404, detail=f"Nutzer {user_id!r} nicht gefunden.") return [ {"id": m.id, "content": m.content, "created_at": m.created_at} for m in store.get_memories(user_id) ] @router.delete("/admin/users/{user_id}/memories/{memory_id}", dependencies=[Depends(require_admin)]) async def delete_user_memory(user_id: str, memory_id: int): """Loescht eine einzelne Erinnerung eines Nutzers.""" deleted = get_store().delete_memory(user_id, memory_id) if not deleted: raise HTTPException(status_code=404, detail="Erinnerung nicht gefunden.") return {"deleted": memory_id} @router.get("/admin/users/{user_id}/sessions", dependencies=[Depends(require_admin)]) async def list_user_sessions(user_id: str): """Listet alle Sessions eines Nutzers (neueste zuerst).""" return get_store().list_sessions_for_user(user_id) @router.get("/admin/sessions/{session_id}/messages", dependencies=[Depends(require_admin)]) async def get_session_messages(session_id: str, limit: int = 200): """Gibt alle Nachrichten einer Session zurueck (Gespraechs-Browser).""" return get_store().get_messages_for_session(session_id, limit) @router.get("/admin/emergency-events", dependencies=[Depends(require_admin)]) async def list_emergency_events(limit: int = 50): """Listet alle protokollierten Notfall-Ereignisse (neueste zuerst).""" return get_store().list_emergency_events(limit) @router.get("/admin/users/{user_id}/usage", dependencies=[Depends(require_admin)]) async def get_user_usage(user_id: str): """Nutzungsstatistik eines Nutzers (letzte 30 Tage).""" return get_store().get_usage_for_user(user_id) @router.get("/admin/usage", dependencies=[Depends(require_admin)]) async def get_all_usage(): """Aggregierte Nutzungsstatistik aller Nutzer.""" return get_store().get_all_usage() # ── Datenbank-Export ──────────────────────────────────────────────────────── @router.get("/admin/db-export", dependencies=[Depends(require_admin)]) async def export_db(): """Laed die SQLite-Datenbank als Datei herunter (Backup).""" path = Path(settings.db_path) if not path.exists(): raise HTTPException(status_code=404, detail="Datenbank nicht gefunden.") return FileResponse( path, media_type="application/octet-stream", filename="voice-assistant.db", headers={"Content-Disposition": 'attachment; filename="voice-assistant.db"'}, ) # ── LLM-/System-Status (read-only) ────────────────────────────────────────── @router.get("/admin/llm/status", dependencies=[Depends(require_admin)]) async def get_llm_status(): """Read-only Status: aktives Backend, Modell, GPU-Auslastung, Dienste.""" return await llm_status() class BackendSwitch(BaseModel): backend: str model: str | None = None @router.post("/admin/llm/backend", dependencies=[Depends(require_admin)]) async def post_llm_backend(payload: BackendSwitch, request: Request): """Wechselt das LLM-Backend (Allowlist-validiert, detached). Greift voll erst, wenn das Gateway als systemd-Dienst läuft; sonst muss es manuell neu starten.""" try: result = await switch_backend(payload.backend, payload.model) except LlmControlError as exc: log_admin_action(request, "llm_backend_switch_rejected", backend=payload.backend, model=payload.model, error=str(exc)) raise HTTPException(status_code=422, detail=str(exc)) log_admin_action(request, "llm_backend_switch", backend=payload.backend, model=payload.model) return result @router.post("/admin/gateway/restart", dependencies=[Depends(require_admin)]) async def post_gateway_restart(request: Request): """Startet das Gateway (systemd-User-Dienst) neu — losgelöst, Self-Restart-sicher.""" log_admin_action(request, "gateway_restart") return restart_gateway_detached() # ── Aussprache-Lexikon CRUD ───────────────────────────────────────────────── def _read_pronunciation(lang: str) -> dict: path = _CONFIG_DIR / f"pronunciation.{lang}.yaml" if not path.exists(): return {"abbreviations": {}, "units": {}, "terms": {}} data = yaml.safe_load(path.read_text(encoding="utf-8")) or {} return { "abbreviations": dict(data.get("abbreviations") or {}), "units": dict(data.get("units") or {}), "terms": dict(data.get("terms") or {}), } def _write_pronunciation(lang: str, data: dict) -> None: path = _CONFIG_DIR / f"pronunciation.{lang}.yaml" path.write_text( yaml.dump(data, allow_unicode=True, default_flow_style=False, sort_keys=False), encoding="utf-8", ) # LRU-Cache des Normalizers invalidieren, damit die Aenderung sofort greift. from app.pipeline.tts_normalizer import _load_lexicon _load_lexicon.cache_clear() @router.get("/admin/pronunciation/{lang}", dependencies=[Depends(require_admin)]) async def get_pronunciation(lang: str): """Gibt alle Eintraege des Aussprache-Lexikons zurueck.""" if lang not in _ALLOWED_LANGS: raise HTTPException(status_code=400, detail=f"Sprache muss eine von {_ALLOWED_LANGS} sein.") return _read_pronunciation(lang) @router.post("/admin/pronunciation/{lang}", dependencies=[Depends(require_admin)]) async def add_pronunciation(lang: str, entry: PronunciationEntry): """Fuegt einen Eintrag zum Aussprache-Lexikon hinzu oder ueberschreibt ihn.""" if lang not in _ALLOWED_LANGS: raise HTTPException(status_code=400, detail=f"Sprache muss eine von {_ALLOWED_LANGS} sein.") if entry.section not in _ALLOWED_SECTIONS: raise HTTPException(status_code=400, detail=f"Section muss eine von {_ALLOWED_SECTIONS} sein.") if not entry.key.strip() or not entry.value.strip(): raise HTTPException(status_code=422, detail="key und value duerfen nicht leer sein.") data = _read_pronunciation(lang) data[entry.section][entry.key.strip()] = entry.value.strip() # Nach jedem Einfügen: Sektion alphabetisch aufsteigend nach Schlüssel sortieren. data[entry.section] = dict( sorted(data[entry.section].items(), key=lambda kv: kv[0].lower()) ) _write_pronunciation(lang, data) return {"section": entry.section, "key": entry.key.strip(), "value": entry.value.strip()} @router.delete("/admin/pronunciation/{lang}/{section}/{key:path}", dependencies=[Depends(require_admin)]) async def delete_pronunciation(lang: str, section: str, key: str): """Loescht einen Eintrag aus dem Aussprache-Lexikon.""" if lang not in _ALLOWED_LANGS: raise HTTPException(status_code=400, detail="Unbekannte Sprache.") if section not in _ALLOWED_SECTIONS: raise HTTPException(status_code=400, detail="Unbekannte Section.") data = _read_pronunciation(lang) if key not in data[section]: raise HTTPException(status_code=404, detail=f"Eintrag '{key}' nicht gefunden.") del data[section][key] _write_pronunciation(lang, data) return {"deleted": key} # ── Laufzeit-Konfiguration ───────────────────────────────────────────────── class ConfigValue(BaseModel): value: str @router.get("/admin/config", dependencies=[Depends(require_admin)]) async def get_runtime_config(): """Gibt alle überschreibbaren Einstellungen mit aktuellem Wert zurück.""" overrides = get_store().get_config_overrides() result = [] for key, (label, type_str, hint) in RUNTIME_SETTABLE.items(): base_val = getattr(settings, key, None) effective_val = getattr(runtime_settings, key, base_val) result.append({ "key": key, "label": label, "type": type_str, "hint": hint, "base_value": str(base_val) if base_val is not None else "", "override_value": overrides.get(key), "effective_value": str(effective_val) if effective_val is not None else "", "is_overridden": key in overrides, }) return result @router.put("/admin/config/{key}", dependencies=[Depends(require_admin)]) async def set_runtime_config(key: str, body: ConfigValue, request: Request): """Setzt eine Laufzeit-Einstellung (wirkt sofort, kein Neustart nötig).""" if key not in RUNTIME_SETTABLE: raise HTTPException(status_code=400, detail=f"Nicht überschreibbar: {key!r}") value = body.value.strip() get_store().set_config_override(key, value) invalidate_cache() log_admin_action(request, "config_set", key=key, value=value) return {"key": key, "value": value} @router.delete("/admin/config/{key}", dependencies=[Depends(require_admin)]) async def delete_runtime_config(key: str, request: Request): """Entfernt eine Laufzeit-Einstellung (fällt auf .env-Wert zurück).""" if key not in RUNTIME_SETTABLE: raise HTTPException(status_code=400, detail=f"Nicht überschreibbar: {key!r}") deleted = get_store().delete_config_override(key) invalidate_cache() if not deleted: raise HTTPException(status_code=404, detail=f"Kein Override für {key!r} gesetzt.") log_admin_action(request, "config_reset", key=key) return {"deleted": key} # ── Live-Log (journalctl → WebSocket) ────────────────────────────────────── @router.websocket("/admin/log") async def admin_log_ws(websocket: WebSocket, key: str | None = None): """Streamt den systemd-Journal-Log des Voice-Assistant-Service live.""" # Auth: Admin-Key als Query-Param ODER SSO-Identitaet via Cookie/Header. from app.auth import authenticate, _bearer_token client_host = websocket.client.host if websocket.client else "" token = _bearer_token(websocket.headers.get("authorization")) or key user = authenticate(websocket.headers, client_host, token) if not is_admin_user(user): await websocket.close(code=1008) return await websocket.accept() # Hinweis, falls die laufende Instanz NICHT der systemd-Dienst ist (z. B. manueller # `uvicorn --reload`-Start). Dann hat das Journal dieser Unit keine aktuellen Zeilen, # und der Log-Tab bliebe sonst kommentarlos leer. try: check = await asyncio.create_subprocess_exec( "systemctl", "is-active", "voice-assistant.service", stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL, ) out, _ = await check.communicate() if out.decode().strip() != "active": await websocket.send_text( "⚠ Der systemd-Dienst 'voice-assistant.service' ist nicht aktiv — " "die laufende Instanz wurde vermutlich manuell gestartet (uvicorn --reload). " "Live-Logs erscheinen hier nur, wenn das Gateway als Dienst läuft " "(sudo systemctl start voice-assistant.service). Manuelle Starts loggen ins Terminal." ) except Exception: pass # stdbuf -oL: zeilenweise Pufferung erzwingen, damit neue Log-Zeilen sofort # (nicht erst blockweise) im Browser ankommen und durchscrollen. proc = await asyncio.create_subprocess_exec( "stdbuf", "-oL", "-eL", "journalctl", "-f", "-u", "voice-assistant.service", "-n", "100", "--no-pager", "-o", "short", stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.STDOUT, ) try: while True: line = await proc.stdout.readline() if not line: break await websocket.send_text(line.decode("utf-8", "replace").rstrip()) except (WebSocketDisconnect, Exception): pass finally: try: proc.terminate() except Exception: pass