ragsage
Examples

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

On this page