Testing FastAPI SSE Endpoints with httpx Permalink to this section
Part of Python & FastAPI SSE Implementation Guide, under Backend Stream Generation & Connection Management.
The standard way to test a FastAPI route — TestClient(app).get(...) — waits for the response to finish. A Server-Sent Events endpoint never finishes, so the test hangs, and in-process transports that assemble the whole body before returning make the problem worse. The reliable approach is to run the app under a real uvicorn server for the test session and read the stream incrementally with httpx. This guide builds the fixtures and a set of tests that cover what matters: headers, event format, Last-Event-ID replay, and cleanup on disconnect.
Symptom & Developer Intent Permalink to this section
- SSE tests hang until the CI job times out.
- A test passes locally and times out in CI, or the reverse.
- Tests assert on the first event and miss that the replay seam duplicates events.
- The test suite leaves background tasks running, producing warnings about pending tasks at shutdown.
- Disconnect handling is untested, and leaks reach production.
The intent is a pytest suite where every streaming test finishes in well under a second, fails loudly with a timeout instead of hanging, and covers replay and cleanup.
Root Cause Analysis Permalink to this section
Two separate mechanisms cause hangs. Test clients built for request/response semantics read the body to the end before returning control, which never happens for an infinite stream. And even with a streaming client, a test that reads “until the stream ends” waits forever.
A real server in the test process has a second benefit: disconnect behaviour is real. When the test closes its stream, uvicorn sends http.disconnect to the app exactly as in production, so the endpoint’s cleanup code runs and can be asserted on.
Step-by-Step Resolution Permalink to this section
Step 1 — Start uvicorn once per session on an ephemeral port Permalink to this section
# conftest.py
import asyncio, socket, threading, time
import pytest, uvicorn
from myapp.main import app
@pytest.fixture(scope="session")
def live_url():
sock = socket.socket(); sock.bind(("127.0.0.1", 0)); port = sock.getsockname()[1]; sock.close()
config = uvicorn.Config(app, host="127.0.0.1", port=port, log_level="warning", lifespan="on")
server = uvicorn.Server(config)
thread = threading.Thread(target=server.run, daemon=True)
thread.start()
while not server.started:
time.sleep(0.01)
yield f"http://127.0.0.1:{port}"
server.should_exit = True
thread.join(timeout=5)
The server runs in its own thread with its own event loop, so it behaves like a separate process while sharing the app’s in-memory state — which lets tests publish into the app’s hub directly.
Step 2 — Read events incrementally with a timeout Permalink to this section
# sse_testing.py
import asyncio, httpx
async def read_events(url, n, headers=None, timeout=3.0):
"""Return the first n events; raise TimeoutError instead of hanging."""
events, cur = [], {"data": []}
async with httpx.AsyncClient(timeout=None) as client:
async with client.stream("GET", url, headers=headers or {}) as res:
res.raise_for_status()
assert res.headers["content-type"].startswith("text/event-stream")
async def consume():
async for line in res.aiter_lines():
if line == "":
if cur["data"]:
events.append({**cur, "data": "\n".join(cur["data"])})
cur.clear(); cur["data"] = []
if len(events) >= n:
return
elif line.startswith(":"):
continue
else:
field, _, value = line.partition(":")
value = value[1:] if value.startswith(" ") else value
if field == "data": cur["data"].append(value)
else: cur[field] = value
await asyncio.wait_for(consume(), timeout)
return events
Leaving the async with block closes the connection, which is how each test “disconnects”.
Step 3 — Test format and headers Permalink to this section
import pytest
pytestmark = pytest.mark.anyio
async def test_stream_headers_and_first_event(live_url, hub):
hub.publish({"id": "1", "event": "order", "data": {"n": 1}})
events = await read_events(f"{live_url}/events", 1, headers={"Last-Event-ID": "0"})
assert events[0]["event"] == "order"
assert events[0]["id"] == "1"
Step 4 — Test the replay seam Permalink to this section
The most valuable test: connect with a cursor while a new event is published, and assert the client receives exactly the missed events followed by the new one.
async def test_replay_then_live_without_gap_or_duplicate(live_url, hub):
for i in (1, 2, 3):
hub.publish({"id": str(i), "event": "order", "data": {"n": i}})
task = asyncio.create_task(read_events(f"{live_url}/events", 3, headers={"Last-Event-ID": "1"}))
await asyncio.sleep(0.05) # let the stream subscribe
hub.publish({"id": "4", "event": "order", "data": {"n": 4}})
events = await task
assert [e["id"] for e in events] == ["2", "3", "4"]
Step 5 — Test cleanup on disconnect Permalink to this section
async def test_disconnect_releases_subscription(live_url, hub):
before = hub.subscriber_count()
await read_events(f"{live_url}/events", 1, headers={"Last-Event-ID": "0"}) # connects, then closes
for _ in range(50): # give the server a moment to unwind
if hub.subscriber_count() == before:
break
await asyncio.sleep(0.02)
assert hub.subscriber_count() == before
This test fails for handlers that never observe the disconnect — the leak described in handling client disconnects in Node.js SSE, in its Python form.
Step 6 — Make heartbeats testable Permalink to this section
Inject the heartbeat interval through settings so tests can set it to milliseconds:
# app: interval from settings, overridable in tests
@app.get("/events")
async def events(request: Request, settings: Settings = Depends(get_settings)):
...
await asyncio.wait_for(queue.get(), timeout=settings.heartbeat_s)
app.dependency_overrides[get_settings] = lambda: Settings(heartbeat_s=0.05)
Step 7 — Isolate state between tests Permalink to this section
A session-scoped server shares the app’s hub across tests, so a test that leaves a subscriber behind or publishes an unconsumed event can affect the next one. Reset the hub in an autouse fixture and give each test its own ids or topics:
@pytest.fixture(autouse=True)
def fresh_hub():
hub.reset() # drop buffered history and assert no subscribers remain
yield hub
assert hub.subscriber_count() == 0, "a test left a stream open"
The assertion in teardown doubles as a leak detector across the whole suite: any test whose stream was not closed fails at the point where it happened, instead of causing a confusing failure several tests later.
Validation & Monitoring Permalink to this section
Run the suite with pytest -x --timeout=30 (pytest-timeout) as a final guard. Add -W error::RuntimeWarning to catch “coroutine was never awaited” and pending-task warnings, which usually indicate a stream task not cleaned up.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Can TestClient test SSE endpoints?
Only for endpoints that eventually end. For infinite streams, in-process test clients that collect the body before returning will hang; use a real server and a streaming httpx client.
Why run uvicorn in a thread rather than a subprocess?
A thread shares the app's memory, so tests can publish into the hub and inspect subscriber counts directly. A subprocess is closer to production isolation but needs a control endpoint for both.
Do I need the anyio pytest plugin?
You need an async test runner: the anyio plugin that ships with anyio, or pytest-asyncio. Either works; the examples use the anyio marker.
How do I test sse-starlette's EventSourceResponse?
The same way. It is an ordinary streaming response under a real server; read it incrementally and assert on the parsed events. Its ping interval can be set in tests just like a hand-written heartbeat.