Stream the support-chat answer to the widget over SSE
Python · LLM apps · intermediate · greenfield
Streams the support-chat answer to the widget: `POST /chat/stream` relays model deltas as SSE frames and saves the completed answer to the per-user history store. Provider hiccups are logged and the stream is closed cleanly so the browser is never left hanging on a dead connection. Text-less chunks are skipped, and each turn sends the system prompt plus at most the last 20 stored messages. Exercised against a live key — frames arrive incrementally and the answer lands in history.
A customer-support chat widget. The browser renders each `delta` frame as the assistant's reply and treats the `done` frame as "this answer is complete"; the stored history is what the next turn sends to the model as context.
Requirements
- `POST /chat/stream` accepts JSON `{user_id, message}` (both non-empty strings; `message` at most 4000 characters) and responds with `text/event-stream`. Because SSE headers are on the wire before the first frame, the HTTP status is always 200 and every outcome — success and failure — is signalled inside the stream.
- Every frame is `data: <JSON object>` followed by a blank line, and the JSON carries a `type`: `delta` (one non-empty `text` piece of the answer), `done` (the model finished normally; always the final frame of a successful stream), or `error` (a human-readable `message`; always the final frame of a failed stream).
- The answer is streamed from the openai 1.x SDK (`AsyncOpenAI`, `stream=True`) against the pinned snapshot `gpt-4o-mini-2024-07-18`. Provider failures surface as `openai.APIError` (or a subclass) from the initial call, and during chunk iteration as either `openai.APIError` or an `httpx.TransportError` (a read timeout or connection reset on the response body); the upstream response stream is closed whenever the generator exits (`async with`); a client disconnect mid-stream is out of scope for this PR — the framework finalises the generator, which releases the upstream response the same way.
- Completion contract: an answer is complete only when the stream delivers the model's terminal chunk with `finish_reason == "stop"`. A stream that ends (EOF) without that chunk, or whose terminal chunk carries any other `finish_reason` (`length`, `content_filter`), is a failure. Failure contract: when the stream fails at any point — provider exception, transport error, or an incomplete ending — the client receives exactly one terminal `error` frame, never a `done` frame, and the partial text is NOT persisted to the conversation history as the assistant's answer; only a completed answer is persisted.
- Streamed chunks may carry `delta.content` of `None` (role-only deltas) and may carry an empty `choices` list; neither produces a frame nor crashes the stream.
- History is per-user and in-memory: the user's message is appended when the turn starts; the assistant's answer is appended only after a completed answer (a completed answer with zero text is not appended, but the stream still ends with `done`). The widget sends at most one in-flight request per user, so one user's turns are processed sequentially and no turn-level locking is required. Every model call sends the constant system prompt followed by at most the 20 most recent stored messages (oldest dropped first); the system prompt is never stored and never counted in that limit.
Files touched
- app/services/chat_stream.py
- app/services/history.py
- app/routers/chat.py
--- app/services/chat_stream.py +"""Streaming chat answer: OpenAI deltas -> SSE frames -> history store.""" + +import json +import logging +import os +from collections.abc import AsyncIterator + +import httpx +import openai +from openai import AsyncOpenAI + +from app.services import history + +logger = logging.getLogger(__name__) + +MODEL = "gpt-4o-mini-2024-07-18" +MAX_HISTORY_MESSAGES = 20 + +client = AsyncOpenAI(api_key=os.environ["OPENAI_API_KEY"]) + +SYSTEM_PROMPT = ( + "You are the support assistant for Northwind Outfitters. " + "Answer the customer's question using the conversation so far. " + "Keep replies short and plain-text." +) + + +class StreamInterrupted(Exception): + """The stream ended without the model's normal `stop` completion chunk."""