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 ANONYMOUS_USER_ID, User, SessionOwnershipError from app.dependencies import ( resolve_route, build_orchestrator, resolve_output_endpoint, get_store, piper_voice_for_language, voice_for_route, ) from app.core.memory_extractor import maybe_schedule_extraction from app.quota import enforce_quota, record_usage, QuotaExceededError from app.safety.emergency import record_manual_emergency from app.schemas import ChatRequest, EmergencyRequest 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, } 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) user_context_parts = [] if user.id != ANONYMOUS_USER_ID: user_context_parts.append(f"Du sprichst mit {user.display_name}.") if memories: user_context_parts.append( "Was du ueber den Nutzer weisst:\n" + "\n".join(f"- {m.content}" for m in memories) ) if user_context_parts: llm_context = [{"role": "system", "content": "\n".join(user_context_parts)}] + llm_context try: enforce_quota(user, store) except QuotaExceededError as exc: raise HTTPException(status_code=429, detail=str(exc)) # Explizit angeforderte Stimme gewinnt; sonst folgt sie der Routensprache. voice = payload.voice or voice_for_route(route.tts_provider, route.language, route.voice_gender) try: trace, audio = await orchestrator.chat_text( payload.text, language=route.language, voice=voice, output=output, history=llm_context, text_only=bool(payload.text_only), ) 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) maybe_schedule_extraction(store, user.id, session_id) if debug: return JSONResponse( content={ "ok": True, "voice": voice, "route": route.as_dict(), "history_len": len(conversation), "memories_len": len(memories), "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) @router.post("/emergency") async def trigger_emergency(payload: EmergencyRequest, user: User = Depends(require_user)): """Vom Nutzer bestätigter Notruf. Protokolliert + liefert den Hinweis (in Nutzersprache). Die eigentliche Benachrichtigung von Angehörigen folgt später. """ store = get_store() return await record_manual_emergency(user, store, payload.language)