Spaces:
Sleeping
Sleeping
| """SSE streaming endpoint for the Agent. | |
| Frontend ``POST /api/agent/chat`` with JSON body, gets back a Server-Sent | |
| Events stream where each event is one of: | |
| event: agent_started | data: {model, mode, tools} | |
| event: assistant_text | data: {text} | |
| event: thinking | data: {text} | |
| event: tool_use | data: {id, name, input} | |
| event: tool_result | data: {tool_use_id, content, is_error} | |
| event: artifact | data: {artifact: {...}} | |
| event: result | data: {session_id, total_cost_usd, duration_ms, ...} | |
| event: error | data: {error} | |
| event: done | data: {} | |
| The frontend can resume a session by passing back ``session_id`` from the | |
| ``result`` event in the next request. | |
| """ | |
| from __future__ import annotations | |
| import json | |
| import logging | |
| from collections.abc import AsyncIterator | |
| from fastapi import APIRouter | |
| from fastapi.responses import StreamingResponse | |
| from app.agent.direct_orchestrator import stream_agent_direct | |
| from app.agent.orchestrator import stream_agent as stream_agent_sdk | |
| from app.api.schemas import AgentChatRequest | |
| from app.core.config import get_settings | |
| logger = logging.getLogger(__name__) | |
| router = APIRouter(prefix="/api/agent", tags=["agent"]) | |
| def _sse_format(event: str, payload: dict) -> bytes: | |
| """Encode one SSE event. Multi-line data is fine — we serialise the whole | |
| payload as a single JSON string on one data: line for easy client parsing.""" | |
| data = json.dumps(payload, ensure_ascii=False, default=str) | |
| return f"event: {event}\ndata: {data}\n\n".encode() | |
| async def _event_stream(req: AgentChatRequest) -> AsyncIterator[bytes]: | |
| runtime = get_settings().agent_runtime.lower() | |
| streamer = stream_agent_sdk if runtime == "sdk" else stream_agent_direct | |
| try: | |
| async for ev in streamer( | |
| prompt=req.prompt, | |
| session_id=req.session_id, | |
| mode=req.mode, | |
| ): | |
| event_name = ev.get("type", "message") | |
| yield _sse_format(event_name, ev) | |
| except Exception as e: | |
| logger.exception("agent stream crashed") | |
| yield _sse_format("error", {"error": f"{type(e).__name__}: {e}"}) | |
| finally: | |
| yield _sse_format("done", {}) | |
| async def chat(req: AgentChatRequest) -> StreamingResponse: | |
| return StreamingResponse( | |
| _event_stream(req), | |
| media_type="text/event-stream", | |
| headers={ | |
| "Cache-Control": "no-cache, no-transform", | |
| "Connection": "keep-alive", | |
| "X-Accel-Buffering": "no", # disable nginx buffering | |
| }, | |
| ) | |