my_voice_assistant_v2/app/api/admin.py

285 lines
12 KiB
Python
Raw Normal View History

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.auth import is_admin_user, require_admin, require_admin_or_user
from app.config import settings
from app.dependencies import get_store
from app.schemas import MemoryCreate, MemoryOut, UserCreate, UserCreated, UserUpdate
router = APIRouter()
_CONFIG_DIR = Path(__file__).resolve().parents[2] / "config"
_ALLOWED_SECTIONS = {"abbreviations", "units", "terms"}
_ALLOWED_LANGS = {"de", "en"}
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}", 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,
}
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"'},
)
# ── 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()
_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}
# ── 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()
proc = await asyncio.create_subprocess_exec(
"journalctl", "--user", "-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