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)