my_voice_assistant_v3_jamulix/app/api/chat.py
Dieter Schlüter 1eb79c1f09 feat: Tageskontingent und Notfall-Eskalation (#6)
- Quota (app/quota.py): Anfragen pro Nutzer/Tag (usage-Tabelle), DAILY_REQUEST_LIMIT,
  pro Nutzer via prefs uebersteuerbar; 429 (REST) bzw. error-Event (WS); Metrik
- Notfall (app/safety/emergency.py): heuristische Erkennung (de/en); Log im Store
  (emergency_events) + optionaler Webhook (best-effort) + X-Emergency/emergency-Event;
  Metrik emergency_total; Notfaelle umgehen das Kontingent
- verdrahtet in chat/speak/transcribe + WS-Turns
- Tests: 64 gruen (+6); Doku aktualisiert (README, BEDIENUNGSANLEITUNG, Architektur, .env.example)

Hinweis: Notfall-Erkennung ist eine Heuristik (kein Lebensretter); erkannte Texte
sind sensibel -> DSGVO beachten.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 05:29:30 +02:00

138 lines
4.6 KiB
Python

from io import BytesIO
from fastapi import APIRouter, Depends, HTTPException, Query
from fastapi.responses import JSONResponse, StreamingResponse
from app.config import settings
from app.errors import RoutingError
from app.auth import require_user
from app.store import User, SessionOwnershipError
from app.dependencies import (
resolve_route,
build_orchestrator,
resolve_output_endpoint,
get_store,
)
from app.quota import enforce_quota, record_usage, QuotaExceededError
from app.safety.emergency import handle_emergency
from app.schemas import ChatRequest
router = APIRouter()
def _route_headers(route) -> dict:
return {
"X-Input-Endpoint": route.input_endpoint,
"X-Output-Endpoint": route.output_endpoint,
"X-STT-Provider": route.stt_provider,
"X-LLM-Provider": route.llm_provider,
"X-TTS-Provider": route.tts_provider,
}
@router.post("/chat")
async def chat(
payload: ChatRequest,
debug: bool = Query(
default=False,
description="Return JSON trace instead of audio response",
),
session_id: str | None = Query(
default=None,
description="Optional session id to apply a stored route",
),
user: User = Depends(require_user),
):
overrides = {
"input_endpoint": payload.input_endpoint,
"output_endpoint": payload.output_endpoint,
"language": payload.language,
"stt_provider": payload.stt_provider,
"llm_provider": payload.llm_provider,
"tts_provider": payload.tts_provider,
}
voice = payload.voice or settings.openrouter_tts_voice
store = get_store()
try:
route = resolve_route(user, session_id, overrides)
orchestrator = build_orchestrator(route)
output = await resolve_output_endpoint(route)
# Gespraechsverlauf laden (nur bei gesetzter session_id -> sonst zustandslos).
conversation = (
store.get_recent_messages(session_id, settings.history_max_messages)
if session_id
else []
)
# Langzeit-Erinnerungen sind nutzerbezogen und gelten auch ohne Session.
memories = store.get_memories(user.id)
except SessionOwnershipError as exc:
raise HTTPException(status_code=403, detail=str(exc))
except RoutingError as exc:
raise HTTPException(status_code=422, detail=str(exc))
llm_context = list(conversation)
if memories:
memory_text = "Was du ueber den Nutzer weisst:\n" + "\n".join(
f"- {m.content}" for m in memories
)
llm_context = [{"role": "system", "content": memory_text}] + llm_context
# Notfall-Erkennung zuerst (immer eskalieren, auch bei Quota-Limit).
emergency = handle_emergency(user, payload.text, store)
if emergency is None:
try:
enforce_quota(user, store)
except QuotaExceededError as exc:
raise HTTPException(status_code=429, detail=str(exc))
try:
trace, audio = await orchestrator.chat_text(
payload.text,
language=route.language,
voice=voice,
output=output,
history=llm_context,
)
except Exception as exc:
raise HTTPException(status_code=502, detail=str(exc))
record_usage(user, store, len(payload.text) + len(trace.semantic_response or ""))
# Turn persistieren (User-Eingabe + semantische Antwort) fuer das Gedaechtnis.
if session_id:
store.append_message(session_id, user.id, "user", payload.text)
store.append_message(session_id, user.id, "assistant", trace.semantic_response)
if debug:
return JSONResponse(
content={
"ok": True,
"voice": voice,
"route": route.as_dict(),
"history_len": len(conversation),
"memories_len": len(memories),
"emergency": emergency,
"trace": {
"raw_transcript": trace.raw_transcript,
"cleaned_transcript": trace.cleaned_transcript,
"semantic_response": trace.semantic_response,
"spoken_response": trace.spoken_response,
"tts_ready_text": trace.tts_ready_text,
},
}
)
headers = {
"Content-Language": route.language,
"X-Audio-Format": "pcm",
"X-Audio-Sample-Rate": "24000",
"X-Audio-Channels": "1",
"X-Audio-Sample-Width": "16",
**_route_headers(route),
}
if emergency:
headers["X-Emergency"] = emergency["category"]
return StreamingResponse(BytesIO(audio), media_type="audio/pcm", headers=headers)