Handbooks / Python asyncio / Chapter 6

Async WebSockets & Streaming

41 pages · ~79 min✓ Reviewed

Builds on Async Databases.

Part 1 · Async WebSockets & Streaming: Real-Time Connections in Python

Async WebSockets & Streaming: Real-Time Connections in Python

Plain HTTP is a question-and-answer protocol: the client asks, the server replies, and the conversation is over. That is fine for loading a page, but it falls apart for chat messages, live prices, multiplayer games or a model typing out its answer word by word. In those cases the server has something to say before anyone asks, and both sides may talk at the same time. A WebSocket keeps one long-lived connection open so either side can send whenever it likes, and Server-Sent Events (SSE) give you a simpler one-way stream over ordinary HTTP.

asyncio is a natural fit for this work. Each connection becomes one cheap task rather than one operating-system thread, so a single process can hold thousands of mostly idle sockets. The catch is that all of those tasks share one event loop: a single blocking call such as time.sleep() or a synchronous database driver freezes every connected client at once. Much of this chapter is about keeping that loop free and about the failure modes of connections that live for hours instead of milliseconds.

We start with what actually happens on the wire: the handshake, the frames and the close codes. Then you build servers and clients with the websockets library, aiohttp and FastAPI/Starlette. After that come the production concerns: ping/pong heartbeats, backpressure with bounded queues, broadcasting to many clients, and reconnecting with jittered backoff. The chapter ends with SSE and async generators, plus a cheatsheet. Afterwards you should be able to choose between WebSocket, SSE and polling, write an endpoint that survives slow or dead peers, and build a client that reconnects and resumes without losing events.

Before you start

You need Python 3.11 or newer and a working grasp of async/await, tasks and asyncio.Queue from the earlier asyncio chapters. Install the libraries with pip install websockets aiohttp "fastapi" "uvicorn[standard]"; the [standard] extra matters, because without a WebSocket backend Uvicorn cannot complete the upgrade. A Redis server is only needed for the multi-worker broadcasting example, and you can skip it on a first read.

Part 2 · WebSocket Fundamentals: Handshake, Frames, Full Duplex

One connection, both directions

A WebSocket (defined in RFC 6455) is a single, long-lived TCP connection that stays open after it is set up. It is full duplex: the client and the server can each send a message at any moment, without waiting for the other side to ask. Plain HTTP can't do this, because there the client asks and the server answers.

There are two URL schemes. ws:// runs on port 80 with no encryption. wss:// wraps the same protocol in TLS and runs on port 443. In production, always use wss://.

Why wss:// passes more reliably

TLS on port 443 goes through proxies untouched. Plain ws:// on port 80 is often mangled by middleboxes that don't understand the upgrade.

A WebSocket doesn't begin as its own protocol. It starts life as an ordinary HTTP request, and both sides then agree to switch. That agreement is the opening handshake.

Upgrade step by step
  1. 1Client sends GETUpgrade: websocket
  2. 2Server checks versionmust be 13
  3. 3Server hashes key + GUIDbuilds Sec-WebSocket-Accept
  4. 4Server replies 101client verifies the accept value
  5. 5Frames onlyno more HTTP

The request is a normal HTTP/1.1 GET. Four headers turn it into an upgrade request: Upgrade: websocket, Connection: Upgrade, a random Sec-WebSocket-Key, and Sec-WebSocket-Version: 13.

http
GET /chat HTTP/1.1
Host: example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Sec-WebSocket-Version: 13

The client's side of the handshake

If the server agrees, it answers 101 Switching Protocols. The reply carries a Sec-WebSocket-Accept header, which proves the server really speaks WebSocket and isn't a plain HTTP server echoing headers back. The value is computed as base64(SHA1(key + GUID)), where the GUID is a fixed string defined by the RFC.

python
import base64, hashlib

GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"

def accept(key: str) -> str:
    raw = (key + GUID).encode()
    digest = hashlib.sha1(raw).digest()
    return base64.b64encode(digest).decode()

print(accept("dGhlIHNhbXBsZSBub25jZQ=="))

The key above is the sample from the RFC

output
s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
http
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=

The server's reply

After the 101, the socket stops speaking HTTP. From that point it carries only WebSocket frames, in both directions, until one side closes.

Frames on the wire

Everything after the handshake travels as frames. Each frame starts with a small header whose most important field is the 4-bit opcode. It tells the receiver what kind of frame this is.

OpcodeFrameKind
0x0continuationdata
0x1text (UTF-8)data
0x2binarydata
0x8closecontrol
0x9pingcontrol
0xApongcontrol

The last three are control frames. They manage the connection itself rather than carry your data. Close ends the session, and ping and pong check that the peer is still alive.

FieldBitsMeaning
FIN1this is the last fragment
opcode4text / binary / control
MASK1set to 1 on client frames
payload len7 / 16 / 64size of the data
masking key32client to server only

Masking: client frames only

Every frame a client sends must be masked. The client picks a random 4-byte key and XORs each payload byte with the key, cycling through its four bytes. Frames from the server to the client are unmasked. A server that receives an unmasked client frame closes the connection with code 1002 (protocol error). Masking exists to stop crafted bytes from poisoning proxy caches.

Masking is not encryption

The key travels in the frame right next to the data, so anyone watching can undo it. For secrecy you need wss://.

The sketch below builds the masked client frame for the text Hi, then reads the header fields back and unmasks the payload.

python
key = bytes([1, 2, 3, 4])
payload = "Hi".encode()
masked = bytes(b ^ key[i % 4] for i, b in enumerate(payload))
frame = bytes([0x81, 0x80 | len(payload)]) + key + masked
print(frame.hex(" "))

b0, b1 = frame[0], frame[1]
print("FIN", b0 >> 7, "opcode", hex(b0 & 0x0F))
print("MASK", b1 >> 7, "len", b1 & 0x7F)

unmasked = bytes(b ^ key[i % 4] for i, b in enumerate(frame[6:]))
print(unmasked.decode())
output
81 82 01 02 03 04 49 6b
FIN 1 opcode 0x1
MASK 1 len 2
Hi

Fragmentation

A single message may be split across several frames. The first fragment carries the real opcode (text or binary) with FIN=0. The following fragments use opcode 0x0, continuation, and the last one sets FIN=1. Control frames such as a ping can arrive in between the fragments. Libraries reassemble the fragments, so your code sees whole messages.

python
# (fin, opcode, data)
frames = [(0, 0x1, b"Hel"), (1, 0x9, b""), (0, 0x0, b"lo, "), (1, 0x0, b"world")]

buf = b""
for fin, op, data in frames:
    if op == 0x9:
        print("ping -> pong")
        continue
    buf += data
    if fin:
        print("message:", buf.decode())

What a library does for you behind the scenes

output
ping -> pong
message: Hello, world

Choosing a transport

WebSocket is not the only way to get fresh data to a browser or client. The right choice depends on who needs to talk and how often.

TechniqueDirectionCost
Short pollingclient asksa new request every time
Long pollingserver to clientone held request per event
SSEserver to client only, over HTTPone open HTTP stream
WebSockettwo-waylow per-message overhead

The per-message cost is the real difference. An HTTP request repeats about 0.5 to 1 KB of headers every time. A server-sent event adds only a data: prefix and a blank line. A WebSocket server frame adds 2 to 10 bytes of header, and a client frame adds 4 more for the mask.

Which transport?

Use a WebSocket when both sides send often: chat, live collaboration, games and trading ticks. When data only flows from server to client, prefer SSE or plain HTTP. SSE reconnects by itself, while a WebSocket needs your own backoff logic. If updates are rare, minutes apart, simple polling is enough.

One task per connection, and never block

In asyncio, each WebSocket connection gets its own coroutine, run as its own task. An idle connection is just a suspended coroutine waiting for data, so thousands of them cost only memory, not threads.

ModelPer idle connection10k connections
Thread per connectionan OS thread plus its stackheavy, with many context switches
asyncio taska few KBfits in one process

There is one catch. All those tasks share a single event loop thread, and a task only gives up control at an await. Any blocking call freezes every connection: time.sleep, a synchronous database driver, requests.get, or a long CPU loop. Pings then miss their timeouts, and peers see the connection closed with code 1011.

Blocking callAsync fix
time.sleep(1)await asyncio.sleep(1)
sync DB driverasync driver, or asyncio.to_thread
requests.get()aiohttp or httpx
CPU-heavy loopto_thread or a process pool

Prefer natively async libraries. When you have no choice, asyncio.to_thread runs the blocking function in a worker thread while the loop keeps serving other connections. For heavy CPU work, a process pool also avoids the GIL. The demo below runs a ticker task next to a 0.3 second job and counts how many ticks it gets.

python
import asyncio, time

async def ticker(counter):
    while True:
        await asyncio.sleep(0.05)
        counter[0] += 1

async def run(work):
    counter = [0]
    t = asyncio.create_task(ticker(counter))
    await asyncio.sleep(0)
    await work()
    t.cancel()
    return counter[0]

async def blocking():
    time.sleep(0.3)

async def offloaded():
    await asyncio.to_thread(time.sleep, 0.3)

async def main():
    print("blocking ticks:", await run(blocking))
    print("offloaded ticks >= 3:", await run(offloaded) >= 3)

asyncio.run(main())
output
blocking ticks: 0
offloaded ticks >= 3: True
One sleep stalls every socket

A single time.sleep(2) inside one handler stops all other connections for two seconds. Replace it with await asyncio.sleep(2), or push the work into asyncio.to_thread.

Key takeaways

The handshake is an HTTP GET answered by a 101, and after that there are only frames. Client frames are masked and server frames are not. Control frames are close 0x8, ping 0x9 and pong 0xA. One event loop can serve thousands of idle sockets, but only if no handler ever blocks it.

Part 3 · The websockets Library: Server and Client

A Server in a Few Lines

The websockets package is a library that does one job: WebSockets over asyncio. Its current API lives in the websockets.asyncio namespace, so you import the server side with from websockets.asyncio.server import serve. Older tutorials import from websockets.server or websockets.legacy. Those are the older interfaces, so avoid them in new code.

serve is an async context manager. You give it a handler coroutine, a host and a port. An empty string for the host means every network interface. Inside the async with block the server is already listening, and await server.serve_forever() keeps it running until the program is cancelled.

Now the handler. It takes exactly one argument, the connection object, which this section calls ws. The library starts one task per connection and calls your handler inside it. When the handler returns, the connection is closed for you. The URL path the client asked for is on the handshake request, as ws.request.path. Earlier versions passed the path as a second argument, handler(ws, path). That signature was removed, so a two-argument handler fails as soon as a client connects.

python
import asyncio
from websockets.asyncio.server import serve

async def handler(ws):
    print("client asked for", ws.request.path)
    await ws.send("welcome")

async def main():
    async with serve(handler, "", 8765) as server:
        await server.serve_forever()

asyncio.run(main())

Needs the third-party websockets package, so it is shown as a server skeleton rather than a runnable check.

The five-line echo server

Once you know that the handler takes one argument and that the connection can be iterated, an echo server is almost nothing. The handler reads each incoming message and sends the same value back. The loop is two lines, and the rest is the boilerplate that starts the server.

python
import asyncio
from websockets.asyncio.server import serve

async def echo(ws):
    async for m in ws: await ws.send(m)

async def main():
    async with serve(echo, "", 8765) as s: await s.serve_forever()

asyncio.run(main())

A complete echo server. The handler is the two lines `async def echo(ws)` and `async for m in ws: await ws.send(m)`.

Common mistake: the old signature

Writing async def handler(ws, path) is the most common failure when you copy older code. The path is no longer a parameter, so read ws.request.path instead.

Sending and Receiving Messages

There are two ways to read from a connection. The first is the receive loop, async for message in ws:. Each pass of the loop gives you one complete message, even if the peer split it into several frames on the wire. A text message arrives as a str and a binary message as bytes.

What happens at the end of the loop tells you how the connection ended. If the peer closes normally, with code 1000 or 1001, the loop simply finishes and the code after it runs. If the connection ends in any other way, such as a dropped network or a protocol error, the loop raises ConnectionClosedError instead. So a handler that must clean up after a crash has to catch that exception.

python
import logging
from websockets.exceptions import ConnectionClosedError

log = logging.getLogger("ws")

async def handler(ws):
    try:
        async for m in ws:
            await ws.send(m)
        # reached after a normal close (1000 / 1001)
    except ConnectionClosedError:
        log.warning("dropped, code %s", ws.close_code)

The handler only catches the abnormal case. A normal close just falls out of the loop.

The second way is await ws.recv(), which reads exactly one message and returns it. Use it when the conversation has a fixed shape, for example a login message that must arrive first. Unlike the loop, recv() has no quiet way to stop. When the connection is closed it always raises: ConnectionClosedOK for a normal close (1000 or 1001) and ConnectionClosedError for anything else. Both inherit from ConnectionClosed, so catching the base class covers either case.

ExceptionWhen it is raisedClose codes
ConnectionClosedOKThe peer closed normally1000 / 1001
ConnectionClosedErrorThe connection ended abnormallyAny other code
ConnectionClosedBase class of bothCatch it to handle either

Sending is one call, await ws.send(x), and the type of x decides the frame type. A str goes out as a text frame and bytes as a binary frame. You can also pass an iterable of chunks, such as a list or a generator of strings or bytes. The library then sends one fragmented message, one frame per chunk, so a large payload never has to sit in memory as one block. The peer's recv() or receive loop still sees a single whole message.

You pass to send()The peer receives
a strOne text frame
bytesOne binary frame
an iterable of str or bytes chunksOne message, sent as several fragments
python
from websockets.asyncio.server import serve
from websockets.exceptions import ConnectionClosed

async def handler(ws):
    first = await ws.recv()            # exactly one message
    await ws.send("hello " + first)    # text frame
    await ws.send(b"\x01\x02")          # binary frame
    await ws.send(["part 1, ", "part 2"])  # one fragmented message

async def run(handler):
    try:
        async with serve(handler, "", 8765) as s:
            await s.serve_forever()
    except ConnectionClosed:
        pass

recv() reads one message. send() picks the frame type from the value you pass.

Common mistake: leaving the exception uncaught

A handler that loops with async for and never catches ConnectionClosedError logs a traceback every time a client's network drops. Wrap the loop in try and catch ConnectionClosedError. Do the same around a bare recv(), since it raises on every close.

The Client, Two-Way Traffic and Size Limits

The client mirrors the server. Import connect from websockets.asyncio.client and use it as an async context manager: async with connect("ws://host:8765") as ws:. The ws you get has the same send, recv and async for as the server side. Leaving the async with block closes the connection with a normal close, so you do not call close() yourself in the usual case.

python
import asyncio
from websockets.asyncio.client import connect

async def main():
    async with connect("ws://localhost:8765") as ws:
        await ws.send("hello")
        print(await ws.recv())
    # leaving the block closed the connection

asyncio.run(main())

A client that sends one message and prints the reply.

A client like that can only alternate: send, then wait for a reply. Real applications often need both directions at once, for example printing server pushes while the user types. A single loop that mixes send() and recv() cannot do that, because while it waits on one it ignores the other. The fix is to run a reader task and a writer task side by side under asyncio.TaskGroup. If either task fails, the group cancels the other and re-raises the error. When both finish, the block ends and the connection closes.

python
import asyncio
from websockets.asyncio.client import connect

async def reader(ws):
    async for m in ws:
        print("got", m)

async def writer(ws):
    for i in range(3):
        await ws.send(f"tick {i}")
        await asyncio.sleep(1)

async def main():
    async with connect("ws://localhost:8765") as ws:
        async with asyncio.TaskGroup() as tg:
            tg.create_task(reader(ws))
            tg.create_task(writer(ws))

asyncio.run(main())

Reading and writing happen at the same time. The group waits for both tasks.

The last setting in this section protects the receiver. By default a message may be at most 1 MiB (max_size, 1,048,576 bytes). If a peer sends a bigger one, the library closes the connection with code 1009 (message too big). Raise the limit with serve(handler, "", 8765, max_size=N), or pass max_size to connect() on the client. A client that suddenly starts getting 1009 closes is usually uploading something larger than the server expected.

SettingDefaultWhat happens when it is exceeded
max_size1 MiBThe connection is closed with 1009 (message too big)
Common mistake: one loop for both directions

Putting send() and recv() in the same loop makes each wait on the other. Use a reader task and a writer task in a TaskGroup whenever traffic flows both ways.

Rejecting Clients Before the Upgrade

Checking credentials inside the handler works, but it is late. By then the handshake has finished and the server has already spent a task on the client. The library offers an earlier point. serve takes a process_request(connection, request) hook that runs on the HTTP upgrade request, before the 101 response goes out. The hook inspects the request, for example its headers. If it returns None, the upgrade continues as normal. If it returns an HTTP response, the client receives that response, 401 or 403, and no WebSocket is ever opened.

python
import asyncio
from websockets.asyncio.server import serve

SECRET = "s3cret"

def auth(connection, request):
    if request.headers.get("X-Token") != SECRET:
        return connection.respond(401, "bad token\n")
    # returning None lets the upgrade continue

async def handler(ws):
    async for m in ws:
        await ws.send(m)

async def main():
    async with serve(handler, "", 8765, process_request=auth) as s:
        await s.serve_forever()

asyncio.run(main())

A missing or wrong X-Token header gets a plain 401 and the handler never starts.

The hook returnsWhat the client sees
NoneThe upgrade goes ahead (101)
connection.respond(401, ...)HTTP 401: the token is missing or wrong
connection.respond(403, ...)HTTP 403: valid user, but no access

Browsers cannot set custom headers on a WebSocket, so a browser client usually puts its token in the URL, as in new WebSocket("/ws?t=..."). The hook then reads it from request.path with the standard library's urlsplit and parse_qs. Keep the hook quick, because a slow check delays every incoming upgrade. Also use wss://, because query strings end up in access logs, and use short-lived tokens.

python
from urllib.parse import urlsplit, parse_qs
from websockets.asyncio.server import serve

def bad(token):
    return token != "s3cret"

def auth(conn, req):
    query = urlsplit(req.path).query
    tok = parse_qs(query).get("t")
    if not tok or bad(tok[0]):
        return conn.respond(403, "no\n")

# serve(handler, "", 8765, process_request=auth)

Token read from the query string. Only the call site is shown for serve.

Where you rejectClient seesServer cost
process_requestHTTP 401 / 403No task is created
Handler closes with 1008A close frameHandshake plus a task
Never checkedAn open socketA security hole
What to remember

Import serve and connect from websockets.asyncio. The handler takes one argument and finds the path at ws.request.path. Use async for for a stream of messages and recv() for one. Run a reader and a writer under a TaskGroup when traffic flows both ways, and do authentication in process_request.

Part 4 · aiohttp WebSocket Server and Client

The aiohttp server: prepare, loop, return

aiohttp is an asyncio HTTP server and client with WebSocket support built in. A WebSocket route is just a normal GET route. Inside the handler you turn that request into a WebSocket by creating a web.WebSocketResponse and calling prepare() on it. That call performs the upgrade handshake you saw earlier (the 101 Switching Protocols reply).

The heartbeat=30 argument in the next example makes aiohttp ping the peer every 30 seconds. Treat it as a setting for now. The page on options explains exactly how it works.

python
from aiohttp import web, WSMsgType

async def ws_handler(request):
    ws = web.WebSocketResponse(heartbeat=30)
    await ws.prepare(request)

    async for msg in ws:
        if msg.type == WSMsgType.TEXT:
            await ws.send_str(f"echo: {msg.data}")
        elif msg.type == WSMsgType.BINARY:
            await ws.send_bytes(msg.data)
        elif msg.type == WSMsgType.ERROR:
            print("socket error:", ws.exception())

    return ws

Fragment: a complete echo handler. It needs aiohttp and a running app to execute.

Two details matter here. First, the handler must finally return ws. aiohttp expects a response object from every handler, and the WebSocket is that response. Second, async for msg in ws yields a WSMessage for each incoming message. Each one has a .type and a .data, so you branch on the type. The loop also stops by itself when a CLOSE message arrives, so you never handle that case.

msg.typeWhat msg.data holdsWhat to do
WSMsgType.TEXTa strprocess it, usually reply
WSMsgType.BINARYbytesprocess it, usually reply
WSMsgType.ERRORthe error is in ws.exception()log it
CLOSEthe close codenothing: the loop ends on its own

Unlike the websockets library, aiohttp hands you a message object and makes you look at its type. In return you get one place that sees every kind of event, including errors.

Common mistake: forgetting return ws

If the handler ends without return ws, aiohttp has no response object to finish the request with, and you get a confusing error after the socket has already closed. Always make return ws the last line.

Sending and receiving come in matching helpers. Pick the one that fits your payload and aiohttp does the encoding for you.

TextBytesJSON
Sendsend_strsend_bytessend_json
Receivereceive_strreceive_bytesreceive_json

send_json serialises a Python object to a text frame. receive_json reads one message and parses it. The receive_* helpers fit a handler that expects a fixed conversation, such as a login message first. async for fits a handler that serves messages of any kind until the peer leaves.

The client and its options

On the client side you open the connection with session.ws_connect(url). A ClientSession owns the connection pool, so you create one session and reuse it for every connection. Creating a new session per connection throws away the pool and leaks resources.

python
import aiohttp

async def talk(url):
    async with aiohttp.ClientSession() as s:
        async with s.ws_connect(url, heartbeat=30) as ws:
            await ws.send_json({"op": 1})
            reply = await ws.receive_json()
            print(reply)

# Reuse one session for many sockets:
async def talk_many(urls):
    async with aiohttp.ClientSession() as s:
        for url in urls:
            async with s.ws_connect(url, heartbeat=30) as ws:
                await ws.send_str("hello")

Fragment: the first function opens one socket, the second reuses the same session for several.

Leaving the inner async with closes the WebSocket cleanly. Leaving the outer one closes the session. The same loop async for msg in ws works on the client and gives you the same WSMessage objects.

Several options control how a socket behaves. Most of them can be set on both WebSocketResponse and ws_connect.

OptionEffect
heartbeat=Nsends a ping every N seconds; if no pong arrives in time, the connection is closed
autoping=Truethe default: answers incoming pings with pongs for you
receive_timeoutcaps how long each receive() may wait
max_msg_sizelargest message accepted; 4 MiB by default, 0 means unlimited

Heartbeat is the dead-peer detector. Every N seconds aiohttp sends a ping, and if the pong does not come back in time it closes the connection. That is how you find out about a peer that vanished without a goodbye. Autoping is the other direction: when the peer pings you, aiohttp answers automatically, so your code never sees ping or pong messages. Set the same heartbeat on server and client so both ends notice a dead link.

receive_timeout bounds each receive() call. When it runs out, a timeout error is raised instead of waiting forever. max_msg_size protects your memory: a message larger than the limit closes the connection.

Common mistake: max_msg_size=0

Setting max_msg_size=0 turns the limit off, so a single client can send a huge message and exhaust your memory. Keep the limit close to the largest message you really expect, especially on a public endpoint.

Safe sends and graceful shutdown

A socket can close at any moment, including between your last receive and your next send. Check ws.closed before sending. If you send after the peer has gone, aiohttp raises ConnectionResetError. After the receive loop ends, ws.close_code tells you how the connection ended.

python
import asyncio

async def push(ws, text):
    if ws.closed:
        return False
    try:
        await ws.send_str(text)
    except ConnectionResetError:
        return False  # the peer closed between the check and the send
    return True

async def wait_for_hello(ws):
    try:
        return await ws.receive(timeout=10)
    except asyncio.TimeoutError:
        await ws.close()
        return None

Fragment: the check avoids most errors, the except covers the small race that remains.

Shutdown needs a list of every open socket, and that list must not keep dead sockets alive. Keep the sockets in a weakref.WeakSet stored on the app. Add each socket after prepare(). A socket whose handler has finished drops out of the set on its own. The example below uses a plain class in place of a real socket to show that behaviour.

python
import gc
import weakref

class FakeSocket:
    pass

sockets = weakref.WeakSet()
a, b = FakeSocket(), FakeSocket()
sockets.add(a)
sockets.add(b)
print("open:", len(sockets))

del b  # the handler for b finished
gc.collect()
print("open:", len(sockets))
output
open: 2
open: 1

With the set in place, the shutdown hook walks over a copy of it and closes each socket with the going away code, 1001. Clients that see this code know the server is restarting and can reconnect. Register the hook in app.on_shutdown.

python
import weakref
from aiohttp import web, WSCloseCode

async def on_shutdown(app):
    for ws in set(app["sockets"]):
        await ws.close(code=WSCloseCode.GOING_AWAY, message=b"restarting")

def make_app():
    app = web.Application()
    app["sockets"] = weakref.WeakSet()
    app.add_routes([web.get("/ws", ws_handler)])
    app.on_shutdown.append(on_shutdown)
    return app

if __name__ == "__main__":
    web.run_app(make_app(), port=8080)

Fragment: ws_handler is the handler from the first page, with request.app["sockets"].add(ws) after prepare().

Common mistakes with closing

Sending without checking ws.closed, so the send raises ConnectionResetError. Passing a str as the close message, which must be bytes. Iterating the live set while sockets remove themselves, so copy it with set(...) first.

aiohttp or websockets?

Both libraries speak the same protocol, so the choice depends on what else your program does. aiohttp is a bundle: an HTTP server, an HTTP client and WebSocket support in one package. websockets does only WebSockets and follows the specification strictly.

Aspectaiohttpwebsockets
ScopeHTTP server and client plus WebSocketsWebSocket only
Stylea bundled frameworkfocused and spec-strict
Receiveyou branch on msg.typeyou get plain str or bytes
Closethe loop ends quietly on CLOSEasync for ends quietly on a normal close; recv() and send() raise ConnectionClosedOK or ConnectionClosedError
Which one?

Pick aiohttp when the same app also serves HTTP routes or makes HTTP calls. Pick websockets for a standalone WebSocket service where you want the leanest API.

Key takeaways

Call await ws.prepare(request), then always return ws. Set heartbeat=30 on both ends. Reuse one ClientSession. Check ws.closed before sending. Close sockets with GOING_AWAY in on_shutdown, and never set max_msg_size=0 on a public endpoint.

Part 5 · FastAPI / Starlette WebSocket Endpoints

The endpoint, accept() and the receive/send methods

FastAPI sits on top of Starlette, and Starlette gives you a WebSocket route with the same decorator style as an HTTP route. You register it with @app.websocket("/ws") and take a WebSocket parameter. The handler is a coroutine that lives as long as the connection does, so one connected client means one running call of your function.

The connection starts as an HTTP upgrade request that is still pending. Nothing is agreed until you call await websocket.accept(), which completes the handshake. Until then you can receive nothing and send nothing, so accept comes first, unless you are rejecting the client (covered in the next page).

python
from fastapi import FastAPI, WebSocket

app = FastAPI()

@app.websocket("/ws")
async def ep(websocket: WebSocket):
    await websocket.accept()          # must come before any send
    async for m in websocket.iter_text():
        await websocket.send_text(m)  # echo each message back

The smallest useful endpoint: accept, then echo

The socket offers one receive method and one send method per payload type. Text frames carry a str, binary frames carry bytes, and the JSON helpers encode or decode text frames for you. For a plain loop you can use the iter_* helpers, which are async generators that yield one message per turn.

TextBytesJSON
Receive onereceive_text()receive_bytes()receive_json()
Send onesend_text(s)send_bytes(b)send_json(obj)
Loop over alliter_text()iter_bytes()iter_json()
Common mistake: sending before accept()

Calling send_text() before accept() raises a RuntimeError, because the handshake has not finished. The same happens when you forget accept() entirely. Make await websocket.accept() the first line of every endpoint that lets the client in.

Disconnects, parameters and authentication

When the client leaves

A client can close the tab or lose its network at any moment. When that happens, the next receive_* call raises WebSocketDisconnect, and its .code attribute holds the close code the peer sent (1000 for a normal close, 1001 when a browser tab goes away). The iter_* loops swallow this exception and simply end, so use them when you do not need the code. Use the receive_* methods in a try block when you want to log it.

Whatever you registered about a client, such as a set entry or a queue, has to be removed on every exit path. Put the cleanup in finally, so a crash in your own code does not leave a dead socket in the registry.

python
from fastapi import WebSocketDisconnect

clients: set[WebSocket] = set()

@app.websocket("/ws")
async def ep(websocket: WebSocket):
    await websocket.accept()
    clients.add(websocket)
    try:
        while True:
            text = await websocket.receive_text()
            await websocket.send_text(text)
    except WebSocketDisconnect as e:
        print("client left with code", e.code)
    finally:
        clients.discard(websocket)   # never leak a dead socket

Add in try, discard in finally

Path, query and Depends

Parameters work the same way as in HTTP routes. A name in the path such as /ws/{room} becomes an argument, a plain argument with a default becomes a query parameter, and Depends() runs a dependency before your handler body. One difference matters in practice: a browser's WebSocket constructor cannot set request headers, so an Authorization header never arrives from browser JavaScript. Pass the token in the query string (new WebSocket("/ws?token=...")) or send it as the first message after connecting.

python
from fastapi import Depends, Query, WebSocketException, status

async def current_user(websocket: WebSocket, token: str = Query("")):
    user = await verify(token)       # your own check
    if user is None:
        raise WebSocketException(code=status.WS_1008_POLICY_VIOLATION)
    return user

@app.websocket("/ws/{room}")
async def ep(websocket: WebSocket, room: str, user=Depends(current_user)):
    await websocket.accept()
    await websocket.send_json({"room": room, "hello": user})

The same Depends() you use on HTTP routes, with a token from the query string

To refuse a client, do it before accept(). You can call await websocket.close(code=status.WS_1008_POLICY_VIOLATION), or raise WebSocketException with that code, as the dependency above does. The upgrade is then denied and the client sees the code. HTTPException does not apply here, because there is no HTTP response to attach it to once the route is a WebSocket.

Before accept(), you doWhat happens
close(code=1008)Handshake denied with that code
raise WebSocketException(1008)Closed with its code
raise HTTPException(...)Not handled as a rejection
Common mistake: expecting headers from the browser

Reading Authorization from a WebSocket request works for server-side clients and tools, but never for browser JavaScript. If the token is in the URL, remember that it ends up in access logs, so keep it short-lived and serve the endpoint over wss://.

Receiving and pushing at the same time

A handler that only loops on receive_text() can answer messages but cannot push anything on its own. If the server wants to send an alert while the client is quiet, the push waits until the next inbound message wakes the loop. The fix is to split the work into two tasks per client: one that reads from the socket, and one that sends whatever appears in that client's own queue.

Both tasks run inside an asyncio.TaskGroup. The group waits for both of them to finish. If either one raises an exception, the group cancels the other. If one simply returns, the other keeps running, so the reader has to tell the pump to stop. The usual way is to put a None sentinel on the queue when the reader ends. The queue is bounded, so a slow client cannot make the server hold unlimited memory.

The demo below uses a stand-in socket so it runs without a server. It receives two messages, queues an echo for each, and then the stand-in raises a disconnect, which makes recv_loop put the sentinel.

python
import asyncio

class FakeSocket:
    def __init__(self, incoming):
        self.incoming = list(incoming)
        self.sent = []

    async def receive_text(self):
        await asyncio.sleep(0.01)
        if not self.incoming:
            raise ConnectionError("client left")
        return self.incoming.pop(0)

    async def send_text(self, text):
        self.sent.append(text)
        print("sent:", text)

async def recv_loop(ws, q):
    try:
        while True:
            text = await ws.receive_text()
            await q.put(f"echo {text}")
    except ConnectionError:
        await q.put(None)          # tell pump to stop

async def pump(ws, q):
    while (item := await q.get()) is not None:
        await ws.send_text(item)

async def main():
    ws = FakeSocket(["hi", "there"])
    q = asyncio.Queue(maxsize=100)
    async with asyncio.TaskGroup() as tg:
        tg.create_task(recv_loop(ws, q))
        tg.create_task(pump(ws, q))
    print("both tasks finished, sent", len(ws.sent))

asyncio.run(main())

A recv task and a send task sharing one queue

output
sent: echo hi
sent: echo there
both tasks finished, sent 2

In a real endpoint the same shape applies. The reader catches WebSocketDisconnect instead of ConnectionError, the pump calls await websocket.send_json(item), and the queue is the one other parts of your app put events on. Remove the client from your registry in a finally around the whole TaskGroup block.

One reader, one writer

Never let one loop both wait for input and push output. Two tasks and one bounded queue per client keep pushes timely and keep slow clients from eating memory.

Running it, scaling it and testing it

Uvicorn needs a WebSocket backend

Uvicorn does not speak WebSocket by itself. It needs the websockets or wsproto library, and the easiest way to get one is pip install "uvicorn[standard]". With a plain uvicorn install, the upgrade request fails: the client gets a 404 and the server logs a warning that no WebSocket library is installed. The keepalive options --ws-ping-interval and --ws-ping-timeout also belong to this backend.

SituationEffectFix
Plain uvicorn installUpgrade fails with a 404 and a warningInstall uvicorn[standard]
No websockets or wsprotoNo backend to handle framesInstall one of them
--workers 4Four separate sets of connectionsUse Redis or another broker
Client closes the tabWebSocketDisconnect with .codeCatch it and clean up

Several workers do not share connections

With uvicorn app:app --workers 4, each worker is its own process with its own memory. A clients set in one process knows nothing about sockets held by the other three, so a broadcast from worker 1 reaches only a quarter of your users. To cross processes, publish each message to Redis pub/sub (or NATS or another broker). Every worker subscribes, and each one forwards the message to its own local clients.

python
import redis.asyncio as redis

r = redis.Redis()

async def fan_out_forever():
    ps = r.pubsub()
    await ps.subscribe("chat")
    async for m in ps.listen():
        if m["type"] == "message":     # skip subscribe confirmations
            for ws in list(clients):   # snapshot of local sockets
                await ws.send_text(m["data"].decode())

Runs once per worker; publish side is await r.publish("chat", data)

Testing without a server

Starlette's TestClient can open a WebSocket against your app directly, with no running server and no network. websocket_connect is a synchronous context manager, so the test is an ordinary function. The with block closes the connection when it ends.

python
from fastapi.testclient import TestClient

def test_echo():
    with TestClient(app).websocket_connect("/ws") as ws:
        ws.send_text("hi")
        assert ws.receive_text() == "hi"

Synchronous test, no server needed

Common mistake: in-memory state with many workers

A module-level clients set works perfectly in a single-process test and then silently delivers to only part of your users in production. Decide early whether you will run more than one worker.

Checklist for a FastAPI WebSocket endpoint

Call accept() first, then send. Reject bad clients before accepting, with close code 1008 or WebSocketException. Pass tokens in the query string or the first message. Deregister clients in finally.

  • Use receive_* when you need WebSocketDisconnect and .code; the iter_* loops end quietly.
  • Run a recv task and a pump task over a bounded per-client queue when the server must push.
  • Install uvicorn[standard] so the upgrade works.
  • With more than one worker, broadcast through Redis pub/sub.
  • Test with TestClient(app).websocket_connect(...).

Part 6 · Ping/Pong Heartbeats and Keepalive

Why a quiet connection needs a heartbeat

A WebSocket is one long-lived TCP connection, and a connection that carries no traffic is easy to lose without noticing. NATs, load balancers and proxies keep a table entry for every connection they forward. When an entry sits idle for too long, the box simply deletes it. It usually sends no FIN and no RST, so neither endpoint is told. Both sides still believe the link is open, and the next send either vanishes or fails much later.

The same silence hides a dead peer. If a phone drops off Wi-Fi, a laptop lid closes or a cable is pulled, the other end gets no goodbye. Without any probing, a server can hold that half-open socket for minutes and keep queuing messages for a client that is gone.

A heartbeat fixes both problems. One side sends a small ping on a timer and expects an answer within a deadline. Steady traffic keeps middlebox timers fresh. A missing answer proves the link is dead, so you can close it and free the resources.

What a heartbeat does
  1. 1Timer firesevery ping_interval
  2. 2Send pingtiny control frame
  3. 3Wait for pongup to ping_timeout
  4. 4Pong arriveslink is alive, reset timer
  5. 5No pongclose the connection
Silence is not health

An idle TCP connection tells you nothing about whether the other end is still there. Only traffic that must be answered proves it.

Ping and pong frames

The protocol has this built in as two control frames. A ping has opcode 0x9. The receiver must reply with a pong, opcode 0xA, that carries the same payload. Because the payload is echoed, the sender can match each pong to the ping that caused it. Libraries do this reply for you, so your handler code normally never sees these frames.

FrameOpcodeRule
ping0x9The peer must answer it
pong0xAEchoes the payload of the ping it answers
unsolicited pong0xAAnswers nothing; works as a one-way keepalive

The third row is a handy detail. A pong that nobody asked for is legal and is ignored by the receiver, but it still crosses every NAT and proxy on the path and refreshes their idle timers. You can send one when you only want to keep the route open and do not need proof of life in return.

The next example is a small stand-in for this exchange, built with only asyncio queues so it runs anywhere. The heartbeat task sends a ping, waits for the echoed payload with a deadline, and gives up when none arrives. The two runs use a live peer and a dead peer. The intervals are shrunk to fractions of a second so the example finishes quickly.

python
import asyncio

async def echo_peer(ping_q, pong_q, alive):
    while True:
        payload = await ping_q.get()
        if alive:
            await pong_q.put(payload)

async def heartbeat(ping_q, pong_q, interval, timeout, rounds):
    for n in range(1, rounds + 1):
        await asyncio.sleep(interval)
        payload = f"ping-{n}".encode()
        await ping_q.put(payload)
        try:
            got = await asyncio.wait_for(pong_q.get(), timeout)
        except asyncio.TimeoutError:
            return "closed 1011 keepalive ping timeout"
        print(f"round {n}: pong echoed payload = {got == payload}")
    return "still open"

async def run(alive):
    ping_q, pong_q = asyncio.Queue(), asyncio.Queue()
    peer = asyncio.create_task(echo_peer(ping_q, pong_q, alive))
    result = await heartbeat(ping_q, pong_q, 0.05, 0.1, 3)
    peer.cancel()
    print("alive" if alive else "dead", "->", result)

async def main():
    await run(True)
    await run(False)

asyncio.run(main())

A toy heartbeat: ping, wait for the echo, give up on silence

output
round 1: pong echoed payload = True
round 2: pong echoed payload = True
round 3: pong echoed payload = True
alive -> still open
dead -> closed 1011 keepalive ping timeout

Real libraries do the same thing with real frames. Control frames are also allowed to arrive in the middle of a fragmented message, so a ping never has to wait for a large upload to finish.

Configuring heartbeats in each stack

The websockets library turns heartbeats on by default. It sends a ping every ping_interval=20 seconds and waits ping_timeout=20 seconds for the pong. If none comes, it closes the connection with code 1011 and the reason keepalive ping timeout. You only need to set these values to change them, and passing None for either disables that part.

python
from websockets.asyncio.server import serve

async def main(handler):
    async with serve(handler, "", 8765,
                     ping_interval=20,
                     ping_timeout=20) as server:
        await server.serve_forever()

These are the defaults, written out

You can also send a ping yourself. ws.ping() returns a future that completes when the matching pong arrives, so awaiting it gives you the round-trip time in seconds. This is a cheap way to log connection quality or to expose a latency metric.

python
async def measure(ws):
    pong = await ws.ping()      # sends the frame, returns a future
    latency = await pong        # resolves when the pong arrives
    print(f"round trip: {latency * 1000:.1f} ms")

Two awaits: one to send, one to wait for the answer

Other stacks use different names for the same idea. aiohttp calls it heartbeat, and the value is the number of seconds between pings. Uvicorn exposes the same knobs as command-line flags for FastAPI and Starlette apps.

StackSettingValue
websocketsping_interval / ping_timeout20s / 20s
aiohttp serverWebSocketResponse(heartbeat=30)30s
aiohttp clientws_connect(url, heartbeat=30)30s
Uvicorn--ws-ping-interval 20 --ws-ping-timeout 2020s / 20s
python
from aiohttp import web, ClientSession

async def handler(request):
    ws = web.WebSocketResponse(heartbeat=30)
    await ws.prepare(request)
    async for msg in ws:
        pass
    return ws

async def client(url):
    async with ClientSession() as session:
        async with session.ws_connect(url, heartbeat=30) as ws:
            await ws.send_str("hello")

aiohttp: set heartbeat on both ends

bash
uvicorn app:app --ws-ping-interval 20 --ws-ping-timeout 20

Uvicorn serving a FastAPI or Starlette app

Proxies, browsers, TCP and the interval trade-off

Pings only help if they arrive before the middlebox gives up. A reverse proxy has its own idle timer, and for nginx the relevant setting is proxy_read_timeout, which defaults to 60 seconds. If your ping interval were 90 seconds, nginx would cut the connection after 60 seconds of silence, long before the next ping. Always keep the proxy timeout comfortably longer than the ping interval.

nginx
location /ws {
    proxy_pass http://app;
    proxy_http_version 1.1;
    proxy_set_header Upgrade $http_upgrade;
    proxy_set_header Connection "upgrade";
    proxy_read_timeout 120s;   # longer than a 20-30s ping interval
}

Idle timeout set above the ping interval

Browsers add a wrinkle. JavaScript's WebSocket object cannot send protocol-level ping frames. The browser answers the server's pings automatically, so server-side detection works fine. If the client needs to notice a dead link on its own, you have to add an application-level message such as {"type":"ping"}. The server replies with {"type":"pong"}, and the page closes and reconnects when no reply shows up in time.

Who needs to detect the dead link?

You may wonder why TCP itself does not solve this. TCP keepalive exists, but the operating system default waits about two hours before the first probe. That is far too slow for a chat or a live feed. WebSocket-level pings run in seconds, travel through the same path as your data, and work even when the OS settings are out of your control.

The last decision is the interval itself, and it is a trade-off. A shorter interval detects failures faster. It also wakes the event loop more often and sends more tiny packets, and across 10,000 or more connections that adds up quickly. A longer interval is cheaper but leaves dead peers undetected for longer and risks a proxy timeout.

Short interval (5s)Long interval (60s)
Failure detectionWithin secondsUp to a minute or more
Wakeups and trafficMany, noticeable at 10k+ connectionsFew
Proxy safetyVery safeMay exceed a 60s proxy timeout
A sensible default

Ping every 20 to 30 seconds with a timeout of about 20 seconds, and set every proxy idle timeout above that.

Common mistake

Setting the ping interval to 60 seconds or more behind nginx with its default proxy_read_timeout of 60s. The proxy drops the idle connection first, and clients see random disconnects.

Common mistake

Relying on TCP keepalive to find dead clients. With the OS default of roughly two hours, a vanished peer holds its socket and memory for a very long time.

Part 7 · Backpressure and Bounded Queues

Why a slow reader must slow the writer

Imagine a server pushing price ticks to a phone on a bad connection. The server can produce thousands of messages per second, but the phone drains them slowly. Every message the phone has not yet read has to sit somewhere: in the socket buffer, in a library queue, or in a Python list you built yourself. If nothing ever tells the producer to wait, that pile only grows.

Backpressure is the mechanism that fixes this. It means a slow consumer slows down the producer. Without it, buffers grow without limit until the process runs out of memory, and then every client on the server goes down, not just the slow one.

How a slow client pushes back along the pipe
  1. 1Client reads slowlyits TCP receive window fills
  2. 2TCP stops the senderthe kernel's send buffer fills too
  3. 3ws.send() waitsthe write buffer passed its high-water mark
  4. 4Your coroutine pausesthe producer cannot run ahead

The websockets library already builds this chain for you, as long as you cooperate. On the sending side, await ws.send(message) returns immediately while the write buffer is small. Once the buffer passes the high-water mark (write_limit, 32 KiB by default), the call waits until the buffer drains. That wait is the backpressure, so always await the send and never throw it away.

On the receiving side the library keeps incoming messages in a queue of at most max_queue messages (16 by default). When that queue is full, the library stops reading from the socket. The unread bytes then pile up in the kernel, TCP flow control shrinks the window, and the remote sender is slowed down in turn.

SideKnobDefaultWhat happens when full
Sendwrite_limit32 KiBawait ws.send() waits
Receivemax_queue16 messagesthe library stops reading the socket
NetworkTCP windowset by the kernelthe remote sender blocks
The wait is the feature

When await ws.send() pauses, nothing is wrong. That pause is the library telling your code that the client cannot keep up yet.

Bounded queues and drop policies

The library's own buffers only cover one connection's socket. The moment your code puts a queue between a producer and a client, you own the backpressure decision. The tool for this is asyncio.Queue(maxsize=N). When the queue holds N items, await q.put(x) suspends the producer until the consumer takes something out.

The example below uses a queue of size 2 with a slow consumer. Watch how the producer is held back: it can only run two items ahead of the consumer.

python
import asyncio

async def producer(q):
    for i in range(5):
        await q.put(i)
        print('put', i)

async def consumer(q):
    for _ in range(5):
        await asyncio.sleep(0.05)
        item = await q.get()
        print('got', item)

async def main():
    q = asyncio.Queue(maxsize=2)
    await asyncio.gather(producer(q), consumer(q))

asyncio.run(main())
output
put 0
put 1
got 0
put 2
got 1
put 3
got 2
put 4
got 3
got 4

Blocking is not always what you want. Sometimes the producer is shared by many clients and must not wait for one of them. Then you use put_nowait, which raises asyncio.QueueFull instead of waiting. Catching that exception is where you choose a drop policy: throw away the new item, throw away the oldest item to make room, or give up on the client.

python
import asyncio

def push_drop_newest(q, x):
    try:
        q.put_nowait(x)
    except asyncio.QueueFull:
        print('dropped newest', x)

def push_drop_oldest(q, x):
    try:
        q.put_nowait(x)
    except asyncio.QueueFull:
        old = q.get_nowait()
        q.put_nowait(x)
        print('dropped oldest', old)

def drain(q):
    items = []
    while not q.empty():
        items.append(q.get_nowait())
    return items

a = asyncio.Queue(maxsize=2)
b = asyncio.Queue(maxsize=2)
for price in (101, 102, 103, 104):
    push_drop_newest(a, price)
    push_drop_oldest(b, price)
print(drain(a), drain(b))

Two ways to react to QueueFull

output
dropped newest 103
dropped oldest 101
dropped newest 104
dropped oldest 102
[101, 102] [103, 104]

Which policy is right depends on what the data means. Ask whether an old message is still worth delivering after a newer one exists.

Choosing a policy for a full queue
PolicyHowFits
Drop newestcatch QueueFull and skip xticks where a stale value is acceptable
Drop oldestget_nowait, then putprices, cursors
Keep one slotQueue(maxsize=1) with drop oldesta single latest value
Blockawait q.put(x)chat, event logs
Disconnectclose with 1008 or 1013chat, event logs
Pick by data type

Latest-value data such as prices and cursors should drop old items or keep a single slot. Event logs and chat must not lose messages silently, so block the producer or disconnect the client with 1008 or 1013.

One writer task per client

The standard shape for pushing data to many websocket clients is a small pipeline. A producer puts messages into a bounded queue that belongs to one client. A single writer task for that client pulls from the queue and sends. Because each client has its own queue and its own writer, one slow client can only fill its own queue and never holds up the others.

The core of the writer is the single line await ws.send(await q.get()). There is one more risk, though: a client that has stopped reading entirely makes send wait forever. Wrap it in asyncio.wait_for(ws.send(m), timeout=5), and when the timeout fires, close that client so it cannot hold resources forever. The example uses a stand-in socket so it runs without a network, with a short timeout so it finishes quickly.

python
import asyncio

class FakeWS:
    def __init__(self, delay):
        self.delay = delay

    async def send(self, m):
        await asyncio.sleep(self.delay)
        print('sent', m)

    async def close(self, code, reason=''):
        print('closed', code)

async def writer(ws, q):
    try:
        while True:
            m = await q.get()
            await asyncio.wait_for(ws.send(m), timeout=0.1)
    except asyncio.TimeoutError:
        await ws.close(1013, 'too slow')

async def run(ws, messages):
    q = asyncio.Queue(maxsize=10)
    for m in messages:
        q.put_nowait(m)
    task = asyncio.create_task(writer(ws, q))
    await asyncio.sleep(0.3)
    task.cancel()

async def main():
    await run(FakeWS(0.01), ['a', 'b'])
    await run(FakeWS(1), ['c', 'd'])

asyncio.run(main())

A fast client keeps receiving; a stalled one is closed with 1013

output
sent a
sent b
closed 1013

The producer side of the same pipeline stays simple: for each connected client it calls put_nowait on that client's queue, and on QueueFull it applies the policy from the previous page. Fanning one message out to many clients is covered in the next section.

Fire-and-forget sends

Writing asyncio.create_task(ws.send(m)) for every message hides errors, because nobody awaits the task to see a ConnectionClosed exception. It also discards backpressure, because the producer never waits, so thousands of pending tasks pile up in memory.

Unbounded queues and unguarded sends

A plain asyncio.Queue() on a hot path has no limit, so a slow client turns it into an out-of-memory crash. Calling send() without await, or awaiting it with no timeout, lets one stalled peer hang its writer forever.

Part 8 · Broadcasting to Many Clients

The registry and the naive loop

Broadcasting means taking one message and delivering it to every connected client. Before you can do that you need to know who is connected, so the first piece is a registry: a plain set of live connections. Each handler adds its connection to the set when it starts and removes it when it ends.

The important detail is where those two calls go. Call clients.add(ws) just before the try, and put clients.discard(ws) in the finally. A finally block runs on a normal close, on an exception and on cancellation, so a disconnect can never leave a dead entry behind. Use discard rather than remove, because it does not raise if the entry is already gone.

python
import asyncio

clients = set()

async def handler(name, steps):
    clients.add(name)
    try:
        for _ in range(steps):
            await asyncio.sleep(0)
        raise ConnectionError('peer vanished')
    finally:
        clients.discard(name)

async def main():
    results = await asyncio.gather(
        handler('a', 2), handler('b', 1), return_exceptions=True)
    print(len(clients), [type(r).__name__ for r in results])

asyncio.run(main())

Both handlers crash, yet the registry ends up empty.

output
0 ['ConnectionError', 'ConnectionError']

With a registry in hand, the obvious broadcast is a loop: for ws in clients: await ws.send(m). It works in a demo, but it sends to one client at a time. If the third client is on a bad network and its send takes half a second, every client after it waits too, and the next broadcast cannot start until the loop finishes.

The first improvement is asyncio.gather(*(ws.send(m) for ws in clients), return_exceptions=True). It starts all the sends at once, so a slow client no longer delays the ones before or after it in the loop. The catch is that gather only returns when every send has finished, so the slowest client still sets the finish time of the whole broadcast. The example below uses stand-in connections with fixed delays to make that visible.

python
import asyncio, time

class FakeWS:
    def __init__(self, delay):
        self.delay = delay
    async def send(self, m):
        await asyncio.sleep(self.delay)

clients = [FakeWS(0.1), FakeWS(0.1), FakeWS(0.5)]

async def naive(m):
    for ws in clients:
        await ws.send(m)

async def concurrent(m):
    await asyncio.gather(
        *(ws.send(m) for ws in clients), return_exceptions=True)

async def timed(fn):
    start = time.perf_counter()
    await fn('hi')
    return round(time.perf_counter() - start, 1)

async def main():
    print('naive ', await timed(naive))
    print('gather', await timed(concurrent))

asyncio.run(main())

Delays of 0.1 s, 0.1 s and 0.5 s: the loop adds them up, gather waits for the largest.

output
naive  0.7
gather 0.5
Common mistake: a send loop with no escape hatch

With for ws in clients: await ws.send(m) one slow client stalls everyone behind it. Switching to gather helps, but with return_exceptions=True the errors come back as values, so scan the results and discard the sockets that failed.

Library broadcast and per-client queues

The websockets library ships a helper for this job: websockets.asyncio.server.broadcast(connections, message). It is an ordinary synchronous function, so you call it without await. It writes the message to each connection's transport and returns immediately. It skips connections that are already closed, but it applies no backpressure: it does not wait for a slow client to catch up, and it does not skip one whose buffer is full. A slow client's buffer keeps growing until the ping timeout finally closes that connection.

That makes broadcast() a good fit for latest-value data where losing a client now and then does not matter. When you need control over what happens to a slow client, the robust approach is to give every client its own bounded queue and its own writer task. The broadcaster only does put_nowait onto each queue, which never blocks. Each writer then does await ws.send(await q.get()) at its own pace and applies its own policy when its queue fills up: drop the oldest message, drop the newest, or disconnect the client.

Queue and writer fan-out
  1. 1Publisherone message arrives
  2. 2json.dumps onceone string for everyone
  3. 3put_nowait per clientbounded queue each
  4. 4Writer taskown drop or disconnect policy
  5. 5ws.sendpaced by that client

Notice that serialization happens once, before the loop. If you call json.dumps inside the loop you repeat identical work for every client. Fan-out costs O(N) writes per message no matter what, since each client has to receive its own copy, but the encoding should be paid for once and the same string or bytes object reused.

python
import asyncio, json

class Client:
    def __init__(self, maxsize):
        self.queue = asyncio.Queue(maxsize=maxsize)
        self.dropped = 0

clients = set()

def fan_out(message):
    data = json.dumps(message)  # once, not per client
    for c in list(clients):
        try:
            c.queue.put_nowait(data)
        except asyncio.QueueFull:
            c.queue.get_nowait()  # policy: drop the oldest
            c.queue.put_nowait(data)
            c.dropped += 1

async def main():
    fast, slow = Client(3), Client(3)
    clients.update({fast, slow})
    fast_seen = 0
    for i in range(5):
        fan_out({'n': i})
        while not fast.queue.empty():  # the fast writer keeps up
            fast.queue.get_nowait()
            fast_seen += 1
    left = []
    while not slow.queue.empty():
        left.append(json.loads(slow.queue.get_nowait())['n'])
    print('fast', fast.dropped, fast_seen)
    print('slow', slow.dropped, left)

asyncio.run(main())

The slow client never reads, so only its own queue overflows and only its own oldest messages are lost.

output
fast 0 5
slow 2 [2, 3, 4]

The fast client received all five messages and never lost any. The slow client lost its two oldest messages, and nothing it did affected the other client. Here is how the four approaches compare.

MethodWaits for clients?What a slow client does
for ws: await ws.send(m)Yes, one after anotherStalls every client behind it
gather(..., return_exceptions=True)Yes, all at onceThe slowest client sets the finish time
broadcast(conns, m)No, it is synchronousNot skipped; its buffer grows until the ping timeout closes it
Bounded queue plus writer taskNo, only put_nowaitGets its own drop or disconnect policy

Snapshots, many workers and rooms

Whenever a broadcast loop contains an await, other tasks get to run in the middle of it, and one of them may be a handler that adds or removes a client. If you are iterating the live set at that moment, Python raises RuntimeError: Set changed size during iteration. The fix is to iterate over a snapshot with list(clients), which is a cheap copy of the references.

python
clients = {'a', 'b', 'c'}

try:
    for c in clients:
        clients.discard('b')  # a handler leaving mid-broadcast
except RuntimeError as e:
    print(e)

for c in list(clients):
    clients.discard(c)
print(len(clients))

The first loop breaks; the second iterates a copy, so changing the set is safe.

output
Set changed size during iteration
0

Everything so far assumes a single process. Once you run several worker processes, or several hosts, each worker's registry only contains the clients connected to that worker, so a broadcast in one process misses everyone else. The standard answer is a message broker. The publisher sends the message to Redis pub/sub with await redis.publish('chat', data), using the redis.asyncio client rather than the blocking one. Every worker subscribes to the channel, and each one fans the message out to its own local clients. NATS plays the same role. Skip the non-message events that a subscription also yields, such as the subscribe confirmation.

Often clients care about only some of the traffic, such as one chat room or one stock symbol. For that, keep a dict of rooms: dict[str, set[ws]] maps each topic to its subscribers. Joining uses rooms.setdefault(topic, set()).add(ws). Leaving discards the connection and then deletes the room if its set is now empty. Without that last step, every room name anyone ever used stays in the dict forever.

python
rooms = {}

def join(topic, ws):
    rooms.setdefault(topic, set()).add(ws)

def leave(topic, ws):
    subs = rooms.get(topic)
    if subs is None:
        return
    subs.discard(ws)
    if not subs:
        del rooms[topic]  # empty room = leaked key

join('news', 'a'); join('news', 'b'); join('sport', 'a')
print(sorted(rooms))
leave('news', 'a'); leave('news', 'b')
print(sorted(rooms))
leave('sport', 'a'); leave('sport', 'a')
print(rooms)

Leaving twice is harmless, and the last leave removes the final room.

output
['news', 'sport']
['sport']
{}

To broadcast to a room, look up rooms.get(topic, set()), so that a missing room simply means no subscribers, then serialize once and fan out over list(subs). Which fan-out method you choose depends on how much loss you can tolerate.

SituationUse
Few clients, all fastgather with return_exceptions=True
Latest-value ticks where a closed client is acceptablebroadcast()
Events that must not be lost, slow clients expectedBounded queue plus writer task
Many workers or hostsRedis pub/sub or NATS, then local fan-out
Broadcasting checklist

Keep these rules together and a broadcast stays safe as the client count grows.

  • Add to the registry in the try and discard in the finally.
  • Serialize once, before the loop.
  • Iterate over a snapshot with list(clients).
  • Give each client a bounded queue and a writer task with its own policy.
  • Use pub/sub when you run more than one process or host.
  • Delete empty rooms when the last subscriber leaves.
Common mistake: leaking dead sockets

If discard is not in a finally, a crashed connection stays in the registry. Every later broadcast then tries to write to it, and the set grows without bound.

Part 9 · Graceful Close Codes and Reconnect with Backoff

How a WebSocket Ends: the Close Handshake and Standard Codes

A WebSocket does not just vanish when one side is done. It ends with a small conversation of its own, called the close handshake. One side sends a close frame (opcode 0x8) that carries a numeric status code and an optional human-readable reason. The reason is limited to 123 bytes, because a control frame payload may not exceed 125 bytes and two of them are taken by the code.

The other side answers by echoing a close frame back, normally with the same code. Once both frames have been exchanged, the server closes the underlying TCP connection. That order matters: the server owning the TCP close keeps the client from being left with a half-dead socket in TIME_WAIT.

The code is how the two programs tell each other why the connection ended, and, just as important, whether it is worth coming back. The first group of codes covers ordinary endings and protocol-level mistakes.

CodeNameWhen you see it
1000Normal closureThe work is done and both sides agree to stop
1001Going awayThe server is shutting down, or the user navigated away from the page
1002Protocol errorThe peer sent a broken frame, for example an unmasked client frame
1003Unsupported dataThe endpoint received a data type it cannot accept, such as binary on a text-only service

A second group describes problems on the server's side of the relationship: who is allowed in, how much data is acceptable, and whether the server itself is healthy. These are the codes that most affect what a client should do next, so we will come back to them when we write the retry policy.

CodeNameMeaningRetry?
1008Policy violationAuthentication or authorization failed, or a rule was brokenNo, fix credentials first
1009Message too bigThe message exceeded the receiver's size limitNo, send less data
1011Internal errorThe server hit an unexpected problemYes, with backoff
1012Service restartThe server is restarting, for example during a deployYes, with backoff
1013Try again laterThe server is overloaded or temporarily refusing clientsYes, with backoff
Common mistake: oversized reason

A close reason longer than 123 bytes is invalid and can turn a clean shutdown into a protocol error. Keep reasons short, and remember that non-ASCII characters take more than one byte each.

Codes Nobody Sends, Codes You Invent, and Who Gets Retried

One code in every list deserves special care: 1006 abnormal closure. It is never put on the wire. No peer ever writes 1006 into a close frame. It is a local report from your WebSocket library saying that the TCP connection died without any close frame arriving, for instance because a cable was pulled, a proxy dropped the link, or the process was killed. If you see 1006, the peer had no chance to explain itself.

The range 4000 to 4999 is reserved for your own application. The protocol gives these numbers no meaning, so you decide them. A common choice is 4001 for an expired token, which lets the client tell that case apart from a generic policy failure. Sending one is a single call, and the reason is a short string.

text
if token_expired(token):
    await ws.close(code=4001, reason='expired')
    # the client sees code 4001, refreshes its token, then reconnects

Closing with an application code (the reason stays under 123 bytes)

The server picks the code, and the code tells the client what to do next. A server that is restarting for a deploy should send a code that says come back soon, while one that has rejected a bad token should send a code that says do not come back until something changes.

Server situationCode to sendClient should
Graceful shutdown1001Reconnect with backoff
Deploy or restart1012Back off, then return
Overloaded, or a client that is too slow1013Back off, then return
Bad or missing token1008Stop and re-authenticate
Message over the size limit1009Fix the payload, do not resend it
Token expired (your own rule)4001Refresh credentials, then reconnect

This gives the client a simple reconnect policy. Network-ish and server-side failures are temporary, so retry on 1001, 1006, 1011, 1012 and 1013. Authentication failures are different: 1008 and 4xxx auth codes mean your credentials are the problem, and repeating the same request with the same token only hammers the server. Fix the credentials first, then connect again.

What to do after a close

The same rule as a tiny function makes it concrete. It treats a normal close as the end of the story, retries the temporary codes, and refuses everything else until a human or a token refresh intervenes.

python
RETRY = {1001, 1006, 1011, 1012, 1013}

def decide(code):
    if code == 1000:
        return 'stop'
    if code in RETRY:
        return 'retry'
    return 'fix first'

for code in (1000, 1001, 1006, 1008, 1013, 4001):
    print(code, decide(code))
output
1000 stop
1001 retry
1006 retry
1008 fix first
1013 retry
4001 fix first
Common mistakes

Never send 1006 yourself, since it is reserved for local reporting. And do not retry 1008 or 4001 in a loop with the same bad token: you will only generate load and log noise.

Exponential Backoff with Full Jitter

When a connection drops, reconnecting immediately is tempting, but it is a trap. If a server goes down, every client loses its connection at the same moment, and if they all retry on the same schedule they hit the recovering server in synchronized waves. That is the thundering herd. The fix has two parts: wait longer after each failure (exponential backoff), and randomize the wait (jitter) so clients spread out.

The version called full jitter picks the delay uniformly between zero and a ceiling that doubles with every attempt, up to a cap. The formula is delay = random.uniform(0, min(cap, base * 2**attempt)). Notice the delay is a random value from 0 up to the ceiling, not a fixed wait.

AttemptCeiling (base 1s, cap 30s)Actual delay
01sanywhere in 0 to 1s
12sanywhere in 0 to 2s
38sanywhere in 0 to 8s
532s, cut to 30sanywhere in 0 to 30s

Here is the formula as code. Because the delay is random, the program checks the bounds instead of printing random numbers. It prints the ceilings, then confirms that a thousand samples all fall inside the range and really do use most of it.

python
import random

BASE, CAP = 1, 30

def ceiling(attempt):
    return min(CAP, BASE * 2**attempt)

def delay(attempt):
    return random.uniform(0, ceiling(attempt))

for attempt in (0, 1, 3, 5):
    print(attempt, ceiling(attempt))

random.seed(7)
samples = [delay(5) for _ in range(1000)]
print(all(0 <= d <= 30 for d in samples), max(samples) > 25, min(samples) < 5)
output
0 1
1 2
3 8
5 30
True True True

There is one more detail: when to start counting from zero again. Resetting the attempt counter straight after the handshake succeeds is a mistake, because a flapping link can complete the handshake and drop a second later, over and over, so the counter would never grow and the client would hammer the server. Instead, record when each connection started and reset only once it has stayed up for a stable period, such as 60 seconds.

python
STABLE = 60
attempt = 0
for lived in (2, 3, 90, 5):
    if lived >= STABLE:
        attempt = 0
    print(f'lived {lived}s -> next ceiling {ceiling(attempt)}s')
    attempt += 1

Each value is how many seconds a connection survived before closing. Reuses ceiling() from above.

output
lived 2s -> next ceiling 1s
lived 3s -> next ceiling 2s
lived 90s -> next ceiling 1s
lived 5s -> next ceiling 2s

The short-lived connections keep the ceiling growing, while the 90 second one proves the link was healthy and earns a fresh start. The cap stays at 30 seconds, and jitter still applies on every attempt.

Common mistakes

A fixed delay with no jitter brings the herd back in lockstep. Resetting the attempt counter right after the handshake lets a flapping link retry at the fastest rate forever.

Auto-Reconnect with websockets and Resuming Without Losing Events

The websockets library has the retry loop built in. Writing async for ws in connect(uri): gives you a fresh connection on each pass, and the library applies its own backoff between failed attempts. Inside the loop you do your work, catch ConnectionClosed when the socket drops, and continue so the loop produces the next connection. A break stops reconnecting for good.

text
async for ws in connect(URI):
    try:
        await ws.send(auth(token))
        await ws.send(subscribe(last_id))
        async for m in ws:
            handle(m)
    except ConnectionClosed:
        continue

The real loop (needs the websockets package). Inspect the close code in the handler to stop on auth failures.

Getting a new socket is only half the job. A brand-new connection knows nothing about the old one, so after every reconnect you must do three things in order, before you treat the stream as live.

After every reconnect
  1. 1Re-authenticatesend the token again, the server forgot you
  2. 2Re-subscribeask again for your channels or rooms
  3. 3Resume from last idrequest everything after the last message you handled

The third step is what keeps you from losing events. Each message carries an id, and your handler stores the id of the last one it processed. After a reconnect you send that id, and the server replays everything newer before going live. To see it work without a network, the next example fakes the connection source with an async generator that yields a drop reason per connection. The first two sessions die with 1006 and 1012, and the third finishes cleanly.

python
import asyncio

EVENTS = [{'id': i, 'data': f'tick{i}'} for i in range(1, 7)]
state = {'last_id': 0}

class Closed(Exception):
    def __init__(self, code):
        self.code = code

async def fake_connect():
    for drop in (1006, 1012, None):
        yield drop

async def session(drop):
    sent = 0
    for e in EVENTS:
        if e['id'] <= state['last_id']:
            continue
        if drop and sent == 2:
            raise Closed(drop)
        state['last_id'] = e['id']
        sent += 1
        print('got', e['id'])

async def main():
    async for drop in fake_connect():
        print('connected, resume after', state['last_id'])
        try:
            await session(drop)
        except Closed as exc:
            print('closed', exc.code)
            continue
        break

asyncio.run(main())
output
connected, resume after 0
got 1
got 2
closed 1006
connected, resume after 2
got 3
got 4
closed 1012
connected, resume after 4
got 5
got 6

Every event from 1 to 6 arrives exactly once, even though the link died twice. Without the last_id cursor, the second connection would have started over at event 1, or skipped the events sent during the gap.

ExceptionTypical codesWhat to do
ConnectionClosedOK1000, 1001Stop, or reconnect if it was only a shutdown
ConnectionClosedError1006, 1011, and othersBack off and retry
Close received with 1008 or 4xxxauth failureFix credentials first, do not loop
Reconnect checklist

Retry only on 1001, 1006, 1011, 1012 and 1013. Use full jitter with a cap. Reset the counter only after a stable period. Re-authenticate, re-subscribe, and replay from the last seen id before going live.

Common mistake: reconnecting without resume

Reconnecting and simply listening again loses every event published while you were away. Always send your last seen message id on the new socket.

Part 10 · Server-Sent Events and Async Generator Streaming

One open response, many events

Not every real-time feature needs a two-way socket. When only the server has something to say, such as tokens from a language model, a progress bar or a notification, you can use Server-Sent Events (SSE). An SSE stream is a single ordinary HTTP response with Content-Type: text/event-stream that the server never finishes. It keeps the response open and writes small text events into it whenever it has news.

The wire format is deliberately tiny. The body is UTF-8 text made of lines. Each line is a field name, a colon and a value, and a blank line ends the event. The browser buffers lines until it sees that blank line and only then delivers the event to your code.

LineMeaning
data: helloThe payload. Several data: lines in one event are joined with newlines.
event: tokenOptional event type name. Without it the browser fires the generic message handler.
id: 42Resume cursor. The browser remembers the last one it saw.
retry: 3000How many milliseconds the browser waits before reconnecting.
: keepaliveA line starting with a colon is a comment. Browsers ignore it, so it works as a heartbeat.

Writing these strings by hand is easy to get wrong, so it helps to put the rules in one small function. This one is plain Python and builds a full event. Notice that a payload containing a newline has to be split so that every line gets its own data: prefix.

python
def sse(data, event=None, id=None, retry=None):
    lines = []
    if retry is not None:
        lines.append(f"retry: {retry}")
    if id is not None:
        lines.append(f"id: {id}")
    if event:
        lines.append(f"event: {event}")
    for part in str(data).split("\n"):
        lines.append(f"data: {part}")
    return "\n".join(lines) + "\n\n"

print(repr(sse("hello", event="token", id=42, retry=3000)))
print(repr(sse("two\nlines")))
output
'retry: 3000\nid: 42\nevent: token\ndata: hello\n\n'
'data: two\ndata: lines\n\n'
Common mistake

Forgetting the blank line. data: hi\n with a single newline is never delivered, because the browser is still waiting for the event to end. Every event must finish with \n\n.

Common mistake

Putting a raw newline inside a token. The second line has no data: prefix, so the browser treats it as a malformed field and drops it. Split on newlines, as the function above does.

Streaming from an async generator

In FastAPI and Starlette, an SSE endpoint is an async generator wrapped in a StreamingResponse. Each yield becomes a chunk written to the open connection, and between yields the event loop is free to serve other clients. The generator must yield strings like data: ...\n\n, and you must pass media_type="text/event-stream" so the browser treats the response as an event stream.

text
@app.get("/stream")
async def stream(request: Request):
    return StreamingResponse(gen(request), media_type="text/event-stream")

Framework fragment: only the call site, it needs FastAPI installed

The generator is where the interesting part lives. A client can vanish at any moment, and a loop that ignores that keeps producing events and holding database cursors for nobody. Inside the loop, check await request.is_disconnected(). And because cancellation can also arrive in the middle of a yield or an await, put cleanup in finally. When the client leaves, the server cancels the task with CancelledError, or closes the generator with GeneratorExit, and finally runs in both cases.

The example below is plain standard-library code. A tiny fake request stands in for the real one, reporting that the client is gone after three checks, so you can see the loop stop and the cleanup run.

python
import asyncio

class FakeRequest:
    def __init__(self, gone_after):
        self.checks = 0
        self.gone_after = gone_after

    async def is_disconnected(self):
        self.checks += 1
        return self.checks > self.gone_after

async def gen(req):
    n = 0
    try:
        while not await req.is_disconnected():
            n += 1
            yield f"data: tick {n}\n\n"
            await asyncio.sleep(0.01)
    finally:
        print("cleanup: released resources")

async def main():
    async for chunk in gen(FakeRequest(gone_after=3)):
        print(repr(chunk))

asyncio.run(main())
output
'data: tick 1\n\n'
'data: tick 2\n\n'
'data: tick 3\n\n'
cleanup: released resources

The polite check is not the only exit. If the server stops iterating, the generator is closed from the outside and receives GeneratorExit at the paused yield. This second example closes an endless generator after one event, and finally still runs.

python
async def endless():
    n = 0
    try:
        while True:
            n += 1
            yield f"data: {n}\n\n"
    finally:
        print("finally ran")

async def main2():
    g = endless()
    print(repr(await anext(g)))
    await g.aclose()

asyncio.run(main2())
output
'data: 1\n\n'
finally ran
Release in finally

Close database cursors, cancel upstream calls and deregister subscribers in finally. That one block covers a polite disconnect, a cancelled task and a generator closed by the framework.

Common mistake

Making the endpoint a normal function that builds a list, or a generator that is not async. StreamingResponse needs an async generator for the event loop to stay free between events.

Reconnects, resume and the library shortcuts

The browser's EventSource handles the awkward part of streaming for you. When the connection drops, it waits for the retry: delay, which defaults to a few seconds, and opens a new request on its own. That request carries a Last-Event-ID header holding the last id: it received. If you gave every event an id: line, the server can use the header to replay what the client missed and then switch to live events.

Resuming after a dropped connection
  1. 1Stream runsevery event has an id: line
  2. 2Connection dropsnetwork blip, deploy, proxy timeout
  3. 3Browser waitsthe retry: delay
  4. 4Reconnectssends Last-Event-ID: 2
  5. 5Server replaysevents after id 2, then goes live

The replay part is just a generator that reads the header. Here headers are a plain dictionary, and the client reconnected after seeing event 2, so only events 3 and 4 are sent again.

python
EVENTS = [(1, "a"), (2, "b"), (3, "c"), (4, "d")]

def replay(headers):
    last = int(headers.get("last-event-id", "0"))
    for id_, data in EVENTS:
        if id_ > last:
            yield f"id: {id_}\ndata: {data}\n\n"

for chunk in replay({"last-event-id": "2"}):
    print(repr(chunk))
output
'id: 3\ndata: c\n\n'
'id: 4\ndata: d\n\n'

You do not have to write all the plumbing yourself. The sse-starlette package provides EventSourceResponse, which sends periodic ping comments so idle streams survive, notices disconnects and stops the generator, and formats the fields for you when you yield dictionaries with keys such as event, data and id.

text
from sse_starlette.sse import EventSourceResponse

async def gen():
    yield {"event": "token", "data": "hello", "id": "1"}

return EventSourceResponse(gen())

Fragment: needs sse-starlette installed

If you use aiohttp instead, you manage the response yourself. Create a web.StreamResponse, set the content type header, call prepare to send the headers, and then write bytes for each event.

text
resp = web.StreamResponse(headers={"Content-Type": "text/event-stream"})
await resp.prepare(req)
await resp.write(b"data: x\n\n")
return resp

aiohttp fragment inside a handler; set headers before prepare()

Heartbeat quiet streams

If nothing happens for a while, write a comment such as : ping followed by a blank line every 15 seconds or so. Browsers ignore it, but proxies see traffic and do not cut the connection. EventSourceResponse does this for you.

Proxies, connection limits and choosing SSE

A stream that works on your laptop can stall in production because something between the server and the browser is helping. Reverse proxies and compression buffer responses by default so they can send fewer, larger packets. For SSE that means the events arrive in clumps or not at all, so you must tell every layer not to wait.

LayerWhat to doWhy
Proxy such as nginxSend X-Accel-Buffering: noStops the proxy from holding the body until it is full
Browsers and cachesSend Cache-Control: no-cacheKeeps the stream from being cached or replayed
Gzip middlewareExclude the stream routeCompression batches output until a block is full

There is also a browser limit. Over HTTP/1.1 a browser opens at most 6 connections per domain, and every open SSE tab uses one for as long as it lives. The seventh tab to the same site simply hangs. Over HTTP/2, many streams share one connection, so the limit effectively disappears. Serve SSE behind HTTP/2 whenever users might keep several tabs open.

Common mistake

Testing only with curl on localhost. There is no proxy or compression there, so the stream looks perfect. Check it through the real proxy, with gzip on, before you trust it.

Finally, how do you choose between SSE and a WebSocket? Look at who talks. SSE goes only from the server to the client, and only as text, but because it is plain HTTP it passes through ordinary infrastructure, uses your normal cookies and auth, and reconnects by itself. A WebSocket is two-way and can carry binary, but you build the reconnect and backoff logic yourself.

SSEWebSocket
DirectionServer to client onlyBoth ways
DataText onlyText and binary
Infrastructure and authPlain HTTP, cookies and headersUpgrade handshake, token often in the query string
ReconnectBuilt into EventSourceYou write the backoff
Best forLLM token streams, notifications, progress, logsChat, games, collaboration
Which transport?
Rule of thumb

For an LLM answer that appears word by word, or a notification feed, SSE is the right tool. Yield each token as a data: event, end with a marker the client understands, and cancel the upstream call in finally when the tab closes.

Part 11 · Cheatsheet: Async WebSockets & Streaming

Servers, loops and keepalive

Each library has one short way to start a WebSocket server and one way to read messages. Learn these first, because everything else in the chapter hangs off them.

LibraryServer one-linerReceive loopEnd of connection
websocketsserve(handler, host, port)async for m in wsthe loop ends on 1000/1001; ConnectionClosed is raised on an abnormal close
aiohttpws = WebSocketResponse() then await ws.prepare(req)async for m in wsthe loop just stops; check ws.close_code; always return ws
FastAPI / Starlette@app.websocket("/ws") then await websocket.accept()websocket.iter_text()the loop just ends; receive_text() raises WebSocketDisconnect

The Starlette iterators iter_text(), iter_bytes() and iter_json() catch the disconnect internally and stop quietly, so wrapping them in except WebSocketDisconnect never fires. If you need the close code, read with an explicit receive_text() loop. Either way, put cleanup in finally. The example below is a fragment: it assumes app, clients and log already exist, so it is shown without being run.

text
@app.websocket("/ws")
async def ep(ws: WebSocket):
    await ws.accept()              # first, before any send
    clients.add(ws)
    try:
        while True:
            m = await ws.receive_text()
            await ws.send_text(m)
    except WebSocketDisconnect as e:
        log.info("left %s", e.code)
    finally:
        clients.discard(ws)        # always deregister

Fragment: FastAPI endpoint with an explicit loop, so the close code is available

With websockets, catch ConnectionClosed (or its subclasses ConnectionClosedOK and ConnectionClosedError) around async for. With aiohttp there is no exception to catch: the loop simply stops when the client closes, so read ws.close_code after it.

Keepalive

NATs, load balancers and proxies drop idle connections without telling either side. A heartbeat is the only way to notice. Send a ping every 20 to 30 seconds and expect the pong within about 20 seconds. Your proxy's idle timeout must be longer than the ping interval, otherwise the proxy cuts the link between pings.

StackSettingValue
websocketsping_interval / ping_timeout20 s / 20 s
aiohttp server and clientheartbeat=3030 s
Uvicorn--ws-ping-interval 20 --ws-ping-timeout 2020 s / 20 s
nginxproxy_read_timeout (default 60 s)longer than the ping interval
Common mistake: no heartbeat behind a proxy

Without pings, an idle link is silently dropped by the proxy and neither end finds out. Browser JavaScript cannot send ping frames, so use an app-level {"type":"ping"} message there.

Backpressure and broadcast

Backpressure means a slow consumer slows the producer instead of letting buffers grow until the process runs out of memory. Four habits cover most of it.

  • Always await ws.send(...): the wait is the backpressure. Never fire-and-forget with create_task(ws.send(m)).
  • Put a bounded asyncio.Queue(maxsize=N) between the producer and each client, never an unbounded one on a hot path.
  • Give each consumer a drop or disconnect policy for when its queue is full.
  • Wrap sends to slow peers in asyncio.wait_for(ws.send(m), 5) and close the client on timeout.
PolicyHowFits
Drop newestskip the new item on QueueFullstale-tolerant ticks
Drop oldestget_nowait() then put_nowait()prices, cursors
Keep one slotmaxsize=1latest value only
Blockawait q.put(x)chat, event logs
Disconnectclose with 1008 or 1013clients that cannot keep up

Broadcast

Keep a registry set of live connections. Add the socket inside try and discard it in finally, so a dead socket never stays registered. For the fan-out, give every client its own bounded queue and a writer task that applies the client's own policy, and call json.dumps once before the loop rather than once per client. Iterate over a snapshot, list(clients), because the set changes while you loop.

MethodWaits?Slow client
for ws: await ws.send(m)serially, one after anotherstalls everyone
gather(..., return_exceptions=True)concurrentlythe slowest sets the finish time
broadcast(conns, data) (websockets)no, it is synchronousskips only closed connections; applies no backpressure, so a slow client's write buffer keeps growing until the ping timeout closes it
per-client queue + writer tasknoits own drop or disconnect policy

A registry only sees connections in its own process. With several workers or hosts, publish to Redis pub/sub (use redis.asyncio, never the blocking client). Each worker subscribes and fans the message out to its local clients.

Common mistake: sequential broadcast

A loop of await ws.send(m) over all clients means one slow client delays everyone behind it. Queue and writer task per client, or at least gather.

Close codes and reconnecting

The close code tells the other side whether it should come back. Learn the table, and remember that 1006 is only ever a local report that the TCP link died with no close frame. Never send it yourself.

CodeMeaningClient should
1000normal closestop
1001going away (shutdown, navigation)reconnect
1006dropped, local only, never on the wireback off and retry
1008policy violation (auth)do not retry; fix credentials
1009message too bigdo not retry as is
1011internal error (also a missed pong)back off and retry
1012service restartback off and retry
4000-4999defined by your app, for example 4001 token expiredas your app decides

Reconnect with full jitter

After a drop, wait a random time between zero and a ceiling that doubles each attempt up to a cap. Randomising the whole delay is called full jitter, and it stops thousands of clients from reconnecting at the same instant. The example prints the ceiling for each attempt and checks that a drawn delay never exceeds it.

python
import random

BASE, CAP = 1, 30  # seconds

def ceiling(n):
    return min(CAP, BASE * 2**n)

def delay(n):
    return random.uniform(0, ceiling(n))

for n in range(7):
    print(n, ceiling(n), 0 <= delay(n) <= ceiling(n))
output
0 1 True
1 2 True
2 4 True
3 8 True
4 16 True
5 30 True
6 30 True
What a client does after every new socket
  1. 1Waitrandom 0..ceiling
  2. 2Connectre-authenticate
  3. 3Re-subscribeto your channels
  4. 4Resumesend the last seen id
  5. 5Go livereset the attempt only after about 60 s stable
Common mistake: retrying auth failures

Reconnecting on 1008 or a 4xxx auth code repeats the same bad token forever. Fix the credentials first. Also avoid a fixed delay with no jitter, and resetting the attempt counter right after the handshake.

Server-Sent Events and picking a transport

Server-Sent Events (SSE) is one long HTTP response with the content type text/event-stream. The server keeps writing small text events and the browser's EventSource reconnects on its own. Each event is lines of field: value ended by a blank line. A token may contain newlines, so every line of it needs its own data: prefix. The helper below builds one event and prints it with repr so the newlines are visible.

python
def sse(data, id=None, event=None):
    out = []
    if id is not None:
        out.append(f"id: {id}")
    if event:
        out.append(f"event: {event}")
    for line in data.split("\n"):
        out.append(f"data: {line}")
    return "\n".join(out) + "\n\n"

print(repr(sse("hello", id=42)))
print(repr(sse("a\nb", event="token")))
output
'id: 42\ndata: hello\n\n'
'event: token\ndata: a\ndata: b\n\n'
NeedSSE detail
Resumegive each event an id:; the browser sends Last-Event-ID on reconnect, so replay what came after it, then go live
Heartbeata comment line such as : ping followed by a blank line, about every 15 s; browsers ignore it
Bufferingsend X-Accel-Buffering: no and Cache-Control: no-cache, and keep gzip off the stream
Cleanupput upstream cleanup in the generator's finally, which runs when the client leaves

Pick one

NeedPick
two-way, low latencyWebSocket
server push over plain HTTPSSE
rare updatesHTTP polling
Top mistakes

Blocking calls such as time.sleep or requests.get in a handler freeze every connection; use asyncio.sleep, an async client or asyncio.to_thread. Other classic slips are unbounded queues, sequential broadcasts, sending in FastAPI before accept(), and no heartbeat behind a proxy.

Remember

Await every send, bound every queue, ping every 20 to 30 seconds, close with the right code, reconnect with capped full jitter, and resume from the last id.

Part 12 · Check yourself

Quiz

Each question shows a short piece of code or a situation. Decide what happens, or what is wrong, before you open the answer.

This websockets handler is running with 500 clients connected. One client sends a message that triggers the time.sleep(2) line. What do the other 499 clients experience, and how do you fix it?
  • Every connection freezes for 2 seconds. time.sleep blocks the one event loop that all 500 tasks share, so no other handler can run, and pings and sends are delayed too.
  • Enough of these stalls can make pings miss their timeout, and peers then see a 1011 keepalive close.
  • Fix: use await asyncio.sleep(2) if you only need to wait. If the slow part is a sync DB driver or a CPU loop, use await asyncio.to_thread(fn, m), or a process pool for heavy CPU work.
async def handler(ws):
    async for m in ws:
        time.sleep(2)   # pretend this is a slow lookup
        await ws.send(m)
This FastAPI endpoint leaks entries. After clients leave, clients still holds their sockets, and the log line never prints. Why, and what is the fix?
  • The iter_text() loop catches the disconnect internally and simply stops. The except WebSocketDisconnect never fires, so the cleanup inside it never runs.
  • Do the cleanup in try/finally, so it runs however the loop ends: finally: clients.discard(websocket).
  • If you need to catch WebSocketDisconnect and read .code, write an explicit receive_text() loop instead. receive_text() is what raises the exception. iter_text() ends silently.
clients = set()

@app.websocket("/ws")
async def ep(websocket: WebSocket):
    await websocket.accept()
    clients.add(websocket)
    try:
        async for m in websocket.iter_text():
            await websocket.send_text(m)
    except WebSocketDisconnect as e:
        clients.discard(websocket)
        print('left', e.code)
A client reconnects after each drop with the code below. After a server restart, 10,000 clients all fail at once. What goes wrong, and which two changes help most?
  • With no randomness, every client waits exactly the same time. They all retry in the same instant, which is a thundering herd that can knock the server over again.
  • Use full jitter: random.uniform(0, min(cap, base * 2**attempt)). The cap stops the delay growing forever.
  • The loop also retries every close, including 1008 or 4001 auth failures. Those should stop the loop until credentials are fixed.
  • Reset attempt only after the link has stayed up for a while, not right after the handshake.
attempt = 0
while True:
    try:
        await run_client()
    except ConnectionClosed:
        await asyncio.sleep(2 ** attempt)
        attempt += 1
This broadcast works in testing with three clients. In production it sometimes raises an error, and it gets slower as clients are added. Name every problem.
  • clients is iterated live. Another handler can add or discard a socket during an await, which raises set changed size during iteration. Loop over list(clients) instead.
  • The sends run one after another, so one slow client stalls everyone behind it. Better: give each client a bounded queue and a writer task, so each client applies its own drop or disconnect policy.
  • json.dumps runs once per client. Serialize once before the loop.
  • One failing send raises and aborts the rest of the loop, so later clients miss the message. A per-client queue avoids this too.
async def broadcast_msg(msg):
    for ws in clients:
        await ws.send(json.dumps(msg))
A browser opens your SSE endpoint, but no onmessage event ever fires, even though the server logs show it yielding. What is the most likely bug in the generator?
  • A blank line ends each SSE event. Each yield ends with a single \n, so the browser keeps waiting for the event to finish and never dispatches it.
  • Fix: yield f"data: tick {i}\n\n".
  • If the events still arrive in one batch, check proxy buffering: send X-Accel-Buffering: no and Cache-Control: no-cache, and turn off gzip for the stream.
async def gen():
    for i in range(3):
        yield f"data: tick {i}\n"
        await asyncio.sleep(1)

Summary

  • A WebSocket starts as an HTTP GET with an Upgrade header, gets a 101 reply, and then carries only frames. Use wss:// in production.
  • Never block the event loop. Swap time.sleep for asyncio.sleep and push sync or CPU work into asyncio.to_thread.
  • Ending a receive loop differs by library. websockets raises ConnectionClosed on an abnormal close, aiohttp just ends the loop, and FastAPI raises WebSocketDisconnect only from the receive_*() methods, not from iter_*(). Put cleanup in finally.
  • Heartbeat every 20 to 30 seconds, and keep proxy idle timeouts longer than the ping interval.
  • Always await sends, use bounded queues, and give each client its own drop or disconnect policy so one slow client cannot stall the rest.
  • Broadcast by iterating a snapshot, serializing once, and using per-client queues with writer tasks. Use Redis pub/sub across workers.
  • Reconnect with capped exponential backoff and full jitter. Do not retry auth closes such as 1008 or 4xxx, and resume from the last seen id.
  • For server-only push such as notifications or LLM tokens, use SSE: text/event-stream, a blank line after every event, and Last-Event-ID to resume.