Streaming to a browser
Turn the engine's event stream into Server-Sent Events, one frame at a time.
stream() yields four kinds of event, in a
fixed order: an AnswerToken per model fragment,
then one Citation per resolved marker, then a
Usage, then exactly one terminal
AnswerComplete. That order is the contract, and
it maps onto SSE frames with no buffering and no special cases — including the not-found path,
which streams as ordinary tokens and completes ungrounded with no citations.
The whole integration is one match statement.
The program
import asyncio
import json
from ragsage import (
AnswerComplete,
AnswerEvent,
AnswerToken,
Citation,
IngestionPipeline,
QueryEngine,
RawSource,
Scope,
Usage,
)
from ragsage.fakes import FakeEngineKit
DOC = b"Employees accrue 20 days of paid leave per year, and may carry over 5 days."
def to_sse(event: AnswerEvent) -> str:
"""Render one engine event as an SSE frame.
Exhaustive on purpose: the wildcard raises rather than dropping an event, so a
new event type shows up as a loud failure instead of a silently truncated stream.
"""
match event:
case AnswerToken(text=text):
name, payload = "token", {"text": text}
case Citation() as citation:
name, payload = "citation", {
"marker": citation.marker,
"document_id": citation.document_id,
"page": citation.page,
"quote": citation.quote,
}
case Usage() as usage:
name, payload = "usage", {
"sources": usage.sources,
"completion_tokens": usage.completion_tokens,
}
case AnswerComplete() as done:
name, payload = "complete", {
"outcome": done.outcome.value,
"grounded": done.grounded,
}
case _:
raise AssertionError(f"unhandled event: {event!r}")
return f"event: {name}\ndata: {json.dumps(payload)}\n\n"
async def main() -> None:
kit = FakeEngineKit()
scope = Scope(namespace="local")
pipeline = IngestionPipeline(
parser=kit.parser,
classifier=kit.classifier,
chunker=kit.chunker,
contextualizer=kit.contextualizer,
embedder=kit.embedder,
vector_store=kit.vector_store,
lexical_store=kit.lexical_store,
document_store=kit.document_store,
llm=kit.llm,
cache=kit.cache,
)
engine = QueryEngine(
embedder=kit.embedder,
vector_store=kit.vector_store,
lexical_store=kit.lexical_store,
reranker=kit.reranker,
llm=kit.llm,
)
await pipeline.ingest(RawSource(name="leave.txt", content=DOC), scope)
async for event in engine.stream("How much leave do employees accrue?", scope):
print(to_sse(event), end="")
asyncio.run(main())$ python sse.py
event: token
data: {"text": "Employees "}
event: token
data: {"text": "accrue "}
... one frame per fragment ...
event: token
data: {"text": "[1]"}
event: citation
data: {"marker": 1, "document_id": "4faee5bf77cf0e3f", "page": 1, "quote": "Employees accrue 20 days of paid leave per year, and may carry over 5 days."}
event: usage
data: {"sources": 1, "completion_tokens": 16}
event: complete
data: {"outcome": "answered", "grounded": true}Reading the stream
Markers arrive before the citations they refer to. The [1] is part of the answer text,
emitted as a token like any other; the citation frame that resolves it to a document and
page comes after the prose is complete. A UI that renders markers as it streams should treat
them as placeholders and hydrate them when the citation frames land — the marker numbers are
what pair them up.
complete is the one frame you must handle. It carries the final text and the
Outcome, and it arrives exactly once. If your client
saw no citation frames and grounded is false, the engine declined to answer from the
corpus and the UI should say so rather than presenting the prose as sourced.
Nothing needs a not-found branch. A question the corpus can't support streams the
not-found message as ordinary tokens and finishes with "outcome": "not_found", so the
consumer's happy path already covers it.
Behind a web framework
The generator above is the response body. Any ASGI framework's SSE response takes an async iterable of strings, so wiring it up adds no ragsage-specific code:
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
app = FastAPI()
@app.get("/ask")
async def ask(question: str, tenant: str) -> StreamingResponse:
async def frames():
async for event in sage.stream(question, Scope(namespace=tenant)):
yield to_sse(event)
return StreamingResponse(frames(), media_type="text/event-stream")Don't drain the stream before returning it
With the assembled façade, RagSage.stream()
holds a scoped database session open for the whole stream — retrieval happens as the
stream is consumed, not before it starts. Collecting the events into a list first would
keep the transaction open for no reason and throw away the latency the streaming bought
you. Pass the generator straight through; if the client disconnects, closing it releases
the connection.
Wanting the whole answer at once instead? query()
is a thin buffer over this same stream, so its result is byte-for-byte what the frames above
would have assembled.
Next
- Multi-turn chat —
stream()takes the samehistory=argument. - Tenants and filters — where that
Scopein the route handler comes from.