import asyncio import time from pathlib import Path import httpx import yaml from fastapi import APIRouter, Depends, HTTPException, Request, WebSocket, WebSocketDisconnect from fastapi.responses import FileResponse, RedirectResponse 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 CAPABILITY_COOKIE, 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.get("/admin/login") async def admin_app_login(request: Request): """Admin-Selbstanmeldung in die App — Ersatz fuer einen verlorenen Zugangslink. Liegt hinter Authelia (nginx: Location /api/admin/ bzw. /admin-login mit auth_request + error_page-Redirect). Erkennt den SSO-Admin via Remote-User, mintet ein frisches Capability-Token, setzt es als va_token-Cookie und leitet in die App (/). Damit meldet sich der Admin allein mit seinem Authelia-Passwort an, ohne persoenlichen Link. Jeder Aufruf rotiert das Token (alte Links DIESES Admin-Nutzers werden ungueltig) — fuer eine Recovery-Funktion korrekt und ein Sicherheitsplus: jede SSO-Anmeldung erzeugt eine frische Capability. """ user = _current_user_or_none(request) if not is_admin_user(user): raise HTTPException(status_code=403, detail="Admin privileges required") result = get_store().reset_token(user.id) if result is None: raise HTTPException(status_code=404, detail="Admin-Nutzer nicht gefunden.") _, raw_token = result log_admin_action(request, "admin_app_login", user_id=user.id) resp = RedirectResponse("/", status_code=303) resp.set_cookie(CAPABILITY_COOKIE, raw_token, httponly=True, secure=True, samesite="lax", max_age=31536000) return resp @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) _OR_MODELS_CACHE = {"ts": 0.0, "ids": []} @router.get("/admin/openrouter/models", dependencies=[Depends(require_admin)]) async def openrouter_models(): """Live-Liste der OpenRouter-Modell-IDs (10 Min gecacht) für die Modell-Auswahl. So stehen in der Konfiguration immer aktuelle, gültige Modellnamen — statt fest verdrahteter, teils veralteter/erfundener Namen. """ now = time.time() if _OR_MODELS_CACHE["ids"] and now - _OR_MODELS_CACHE["ts"] < 600: return {"models": _OR_MODELS_CACHE["ids"], "cached": True} key = settings.openrouter_api_key.strip() try: async with httpx.AsyncClient(timeout=8.0) as client: headers = {"Authorization": f"Bearer {key}"} if key else {} r = await client.get("https://openrouter.ai/api/v1/models", headers=headers) r.raise_for_status() data = r.json() ids = sorted(m["id"] for m in data.get("data", []) if m.get("id")) _OR_MODELS_CACHE.update(ts=now, ids=ids) return {"models": ids, "cached": False} except Exception as exc: # best effort -> alte Liste / leer return {"models": _OR_MODELS_CACHE["ids"], "error": str(exc)} @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", ""), **{k: (u.prefs or {}).get(k, "") for k in ( "full_name", "street", "postal_code", "city", "phone_landline", "phone_mobile", "birth_date", "medical_notes", "emergency_contacts", "emergency_phones", "emergency_cooldown_minutes", )}, "location_consent": bool((u.prefs or {}).get("location_consent", False)), "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