Skip to content

Relaying LLM Token Streams Through FastAPI

Most LLM features in web products run through a backend: the browser must not hold the API key, the server adds context and enforces quotas, and it relays tokens back as they arrive. In FastAPI that relay is a StreamingResponse fed by the SDK's stream, and its correctness hinges on one question — when the browser goes away, does the upstream stream stop? Measured on Python 3.14 with FastAPI 0.142 on uvicorn, the anthropic SDK 1.11.0 and a local mock of the Messages API streaming a 2,000-token answer at 5 ms per token: with the upstream stream opened inside the response generator, a browser that left after 50 events at 0.31 s caused the mock to record a disconnect at 50 tokens, with nothing in flight two seconds later. With a background task reading the upstream into a queue that the response drained, the same disconnect left the upstream running — 441 tokens sent two seconds later, still streaming. The first event reached the browser after 56–64 ms in both designs, and Starlette's GZipMiddleware did not buffer the event stream (first event at 64 ms with gzip accepted). This guide builds a relay that stops when its reader does.

Prerequisites

1. Own the upstream stream inside the response generator

Open the SDK stream inside the async generator that produces the response body. When the client disconnects, Starlette cancels the generator, the async with exits, and the upstream connection closes:

import json
from contextlib import asynccontextmanager

from anthropic import AsyncAnthropic
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse


@asynccontextmanager
async def lifespan(app: FastAPI):
    app.state.llm = AsyncAnthropic()
    yield
    await app.state.llm.close()

app = FastAPI(lifespan=lifespan)


def sse(data: dict) -> str:
    return f"data: {json.dumps(data)}\n\n"


@app.post("/chat")
async def chat(request: Request, body: ChatRequest):
    llm = request.app.state.llm

    async def events():
        async with llm.messages.stream(model="claude-sonnet-5-5", max_tokens=2000,
                                       messages=body.messages) as stream:
            async for text in stream.text_stream:
                yield sse({"text": text})
        yield sse({"done": True})

    return StreamingResponse(events(), media_type="text/event-stream",
                             headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})

Measured: the browser closed its connection after 50 events; the mock saw one disconnect and had sent exactly 50 tokens, and nothing was in flight two seconds later. The generator is the only owner of the upstream stream, so its lifetime is the response's lifetime. The client is created once in the lifespan and shared, as in managing startup and shutdown with ASGI lifespan.

Verify: close the browser mid-answer; the provider's usage for that request shows output tokens close to what was delivered.

A browser leaves mid-answer A sequence of 7 messages between 4 participants. A browser leaves mid-answer browser FastAPI response SDK stream model server tokens text_stream yields data: {...} events connection closed (after 50 events) generator cancelled async with exits: close() upstream closed at 50 tokens One owner for the upstream stream: the generator that writes the response.

2. Do not move the upstream into a background task

The decoupled design — a task reads the LLM into a queue, the response reads the queue — appears in code that wants to cache answers, fan out to several listeners, or keep generating "in case the user comes back":

@app.get("/chat-queued")
async def chat_queued(request: Request):
    q: asyncio.Queue = asyncio.Queue()

    async def producer():
        async with llm.messages.stream(**req) as stream:
            async for text in stream.text_stream:
                await q.put(text)
        await q.put(None)

    asyncio.create_task(producer())                  # nobody cancels this

    async def events():
        while (t := await q.get()) is not None:
            yield sse({"text": t})

    return StreamingResponse(events(), media_type="text/event-stream")

Measured: when the browser left at 0.31 s, the response generator was cancelled but the producer was not. Two seconds later the mock had sent 441 tokens, recorded no disconnect, and was still streaming; left alone it would have generated all 2,000 into a queue that no one would read. If decoupling is genuinely needed, tie the producer's lifetime to the response: create it inside the generator and cancel it in a finally, or run both in a TaskGroup inside the generator.

Verify: for every create_task in a streaming endpoint, identify the code that cancels it when the client disconnects.

What the upstream did after the browser left at 0.31 s A grid of 3 rows by 4 columns. What the upstream did after the browser left at 0.31 s relay design first event upstream tokens 2 s later upstream state generator owns the stream 64 ms 50 closed (1 disconnect) same, behind GZipMiddleware 64 ms 50 closed (1 disconnect) background producer + queue 56 ms 441 still streaming FastAPI 0.142, Starlette 1.7, uvicorn; local mock at 5 ms per token, 2,000-token answers.

3. Keep proxies and middleware from buffering

Server-sent events only feel live if every hop forwards bytes as they arrive. Measured in this setup, Starlette's GZipMiddleware passed text/event-stream responses through without delaying them — the first event arrived at 64 ms with Accept-Encoding: gzip, the same as without the middleware. Other hops are less predictable, so set the headers that switch buffering off explicitly:

SSE_HEADERS = {
    "Cache-Control": "no-cache",        # no intermediate caching of a live stream
    "X-Accel-Buffering": "no",          # nginx: do not buffer this response
    "Connection": "keep-alive",
}


def sse_response(events) -> StreamingResponse:
    return StreamingResponse(events, media_type="text/event-stream", headers=SSE_HEADERS)

Test through the full production path — CDN, load balancer, ingress — not just against uvicorn directly: a single buffering proxy turns a 60 ms first token into an answer that appears all at once at the end. A periodic SSE comment (: keep-alive\n\n) every 15–30 seconds during long pauses also keeps idle-timeout proxies from cutting the connection, as discussed in streaming server-sent events from asyncio.

Verify: through the production ingress, the first event's arrival time is within a few milliseconds of uvicorn's.

4. Report errors inside the stream

After the first byte, the HTTP status is fixed at 200. Upstream failures that happen later — an overload event, a timeout — must be sent as events the browser understands, and the stream must end cleanly:

async def events():
    try:
        async with llm.messages.stream(**req) as stream:
            async for text in stream.text_stream:
                yield sse({"text": text})
            final = await stream.get_final_message()
        yield sse({"done": True, "output_tokens": final.usage.output_tokens})
    except anthropic.APIStatusError as exc:
        log.warning("upstream error mid-stream: %s %s", exc.status_code, exc.message)
        yield sse({"error": "The model is busy. Please try again.", "retryable": True})
    except (anthropic.APIConnectionError, httpx2.TransportError):
        yield sse({"error": "Connection to the model was lost.", "retryable": True})

An error event mid-stream raised APIStatusError with status code 200 in testing; catching it lets the browser show a message instead of a truncated answer with no explanation. Do not catch BaseException here: CancelledError must propagate so the disconnect path in step 1 still closes the upstream. Retrying automatically on the server is possible but awkward once tokens have reached the user — the retry starts the answer over — so a client-side "try again" with the partial answer cleared is usually clearer. Retry mechanics are in retrying LLM API calls on 429 and overload errors.

Verify: an injected mid-stream upstream error produces an error event and a closed response, not a hung connection.

5. Bound streams per user and per process

Each relayed stream holds an upstream connection, a response, and a share of the event loop's CPU — about 52–97 µs per token measured in streaming LLM tokens asynchronously. Cap them:

from collections import defaultdict

PER_USER = defaultdict(lambda: asyncio.Semaphore(2))
PROCESS = asyncio.Semaphore(200)


@app.post("/chat")
async def chat(request: Request, body: ChatRequest, user: User = Depends(current_user)):
    user_slot = PER_USER[user.id]
    if user_slot.locked() or PROCESS.locked():
        raise HTTPException(429, "Too many concurrent answers")

    async def events():
        async with user_slot, PROCESS:
            async with llm.messages.stream(**request_for(body)) as stream:
                async for text in stream.text_stream:
                    yield sse({"text": text})

    return StreamingResponse(events(), media_type="text/event-stream")

The slots are acquired inside the generator, so they are released when the stream ends for any reason — completion, error or disconnect. Checking locked() before starting rejects excess requests immediately with a status code, instead of accepting them and making the browser wait in silence. The per-process cap should come from the CPU budget and the provider's concurrency limit divided by the number of processes, as in fanning out LLM calls with bounded concurrency.

Verify: a user opening a third concurrent answer gets an immediate 429; the slot frees as soon as one of their streams ends or disconnects.

How should this relay be structured? A decision on Who reads the answer with 4 outcomes. How should this relay be structured? Who reads the answer? one browser, live stream inside the response generator stopped at 50 tokens must finish even if the user leaves a job that stores the result not an untied task several listeners shared producer, cancelled at zero listeners refcount any design per-user and per-process caps slots released in the generator Measured: an untied producer kept streaming after its reader left.

Verification

An LLM relay is correct when:

  • The upstream stream is opened inside the response generator, or its producer is cancelled with the response.
  • A browser disconnect closes the upstream within a few tokens, visible in provider usage.
  • Mid-stream upstream errors become SSE error events, and CancelledError is never swallowed.
  • No hop buffers the stream, and concurrent streams are capped per user and per process.

Diagnostic Hook: log, per relayed answer, the tokens delivered to the browser and the output tokens reported by the provider. If reported tokens regularly exceed delivered ones by more than a few, disconnects are not reaching the upstream — look for background tasks or swallowed cancellations in the streaming endpoint.

Pitfalls & edge cases

  • Background producer tasks. Measured: 441 tokens and still streaming after the browser left.
  • Catching BaseException in the generator. It swallows the cancellation that closes the upstream.
  • Buffering proxies. The answer arrives all at once at the end.
  • Unbounded concurrent streams. Each costs CPU per token and a provider concurrency slot.

Frequently Asked Questions

How do I stream LLM responses to the browser with FastAPI?

Return a StreamingResponse with media_type text/event-stream whose async generator opens the SDK stream with async with and yields an SSE data line per text chunk. The first event arrived in about 60 ms in testing.

Does FastAPI stop the LLM request when the client disconnects?

It cancels the response generator. If the generator owns the upstream stream, the SDK closes it — the mock stopped at 50 tokens. If a separate task reads the upstream, it keeps running unless you cancel it.

Does GZipMiddleware break server-sent events?

Not in the tested Starlette 1.7: an event stream behind GZipMiddleware delivered its first event in 64 ms, the same as without it. Other proxies may buffer; send Cache-Control: no-cache and X-Accel-Buffering: no and test the full path.

How do I report an LLM error after streaming has started?

The HTTP status is already 200, so send an SSE event such as {error, retryable} and end the stream; catch APIStatusError, which the SDK raises with status code 200 for mid-stream errors.