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.
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://.
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.
- 1Client sends GETUpgrade: websocket
- 2Server checks versionmust be 13
- 3Server hashes key + GUIDbuilds Sec-WebSocket-Accept
- 4Server replies 101client verifies the accept value
- 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.
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.
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
s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
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.
| Opcode | Frame | Kind |
|---|---|---|
| 0x0 | continuation | data |
| 0x1 | text (UTF-8) | data |
| 0x2 | binary | data |
| 0x8 | close | control |
| 0x9 | ping | control |
| 0xA | pong | control |
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.
| Field | Bits | Meaning |
|---|---|---|
| FIN | 1 | this is the last fragment |
| opcode | 4 | text / binary / control |
| MASK | 1 | set to 1 on client frames |
| payload len | 7 / 16 / 64 | size of the data |
| masking key | 32 | client 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.
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.
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())
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.
# (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
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.
| Technique | Direction | Cost |
|---|---|---|
| Short polling | client asks | a new request every time |
| Long polling | server to client | one held request per event |
| SSE | server to client only, over HTTP | one open HTTP stream |
| WebSocket | two-way | low 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.
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.
| Model | Per idle connection | 10k connections |
|---|---|---|
| Thread per connection | an OS thread plus its stack | heavy, with many context switches |
| asyncio task | a few KB | fits 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 call | Async fix |
|---|---|
time.sleep(1) | await asyncio.sleep(1) |
| sync DB driver | async driver, or asyncio.to_thread |
requests.get() | aiohttp or httpx |
| CPU-heavy loop | to_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.
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())
blocking ticks: 0 offloaded ticks >= 3: True
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.
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.
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.
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)`.
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.
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.
| Exception | When it is raised | Close codes |
|---|---|---|
| ConnectionClosedOK | The peer closed normally | 1000 / 1001 |
| ConnectionClosedError | The connection ended abnormally | Any other code |
| ConnectionClosed | Base class of both | Catch 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 str | One text frame |
bytes | One binary frame |
an iterable of str or bytes chunks | One message, sent as several fragments |
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.
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.
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.
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.
| Setting | Default | What happens when it is exceeded |
|---|---|---|
max_size | 1 MiB | The connection is closed with 1009 (message too big) |
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.
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 returns | What the client sees |
|---|---|
None | The 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.
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 reject | Client sees | Server cost |
|---|---|---|
process_request | HTTP 401 / 403 | No task is created |
| Handler closes with 1008 | A close frame | Handshake plus a task |
| Never checked | An open socket | A security hole |
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.
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.type | What msg.data holds | What to do |
|---|---|---|
| WSMsgType.TEXT | a str | process it, usually reply |
| WSMsgType.BINARY | bytes | process it, usually reply |
| WSMsgType.ERROR | the error is in ws.exception() | log it |
| CLOSE | the close code | nothing: 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.
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.
| Text | Bytes | JSON | |
|---|---|---|---|
| Send | send_str | send_bytes | send_json |
| Receive | receive_str | receive_bytes | receive_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.
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.
| Option | Effect |
|---|---|
heartbeat=N | sends a ping every N seconds; if no pong arrives in time, the connection is closed |
autoping=True | the default: answers incoming pings with pongs for you |
receive_timeout | caps how long each receive() may wait |
max_msg_size | largest 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.
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.
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.
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))
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.
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().
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.
| Aspect | aiohttp | websockets |
|---|---|---|
| Scope | HTTP server and client plus WebSockets | WebSocket only |
| Style | a bundled framework | focused and spec-strict |
| Receive | you branch on msg.type | you get plain str or bytes |
| Close | the loop ends quietly on CLOSE | async for ends quietly on a normal close; recv() and send() raise ConnectionClosedOK or ConnectionClosedError |
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.
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).
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.
| Text | Bytes | JSON | |
|---|---|---|---|
| Receive one | receive_text() | receive_bytes() | receive_json() |
| Send one | send_text(s) | send_bytes(b) | send_json(obj) |
| Loop over all | iter_text() | iter_bytes() | iter_json() |
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.
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.
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 do | What happens |
|---|---|
close(code=1008) | Handshake denied with that code |
raise WebSocketException(1008) | Closed with its code |
raise HTTPException(...) | Not handled as a rejection |
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.
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
sent: echo hi
sent: echo there
both tasks finished, sent 2In 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.
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.
| Situation | Effect | Fix |
|---|---|---|
Plain uvicorn install | Upgrade fails with a 404 and a warning | Install uvicorn[standard] |
No websockets or wsproto | No backend to handle frames | Install one of them |
--workers 4 | Four separate sets of connections | Use Redis or another broker |
| Client closes the tab | WebSocketDisconnect with .code | Catch 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.
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.
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
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.
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 needWebSocketDisconnectand.code; theiter_*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.
- 1Timer firesevery ping_interval
- 2Send pingtiny control frame
- 3Wait for pongup to ping_timeout
- 4Pong arriveslink is alive, reset timer
- 5No pongclose the connection
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.
| Frame | Opcode | Rule |
|---|---|---|
| ping | 0x9 | The peer must answer it |
| pong | 0xA | Echoes the payload of the ping it answers |
| unsolicited pong | 0xA | Answers 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.
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
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.
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.
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.
| Stack | Setting | Value |
|---|---|---|
| websockets | ping_interval / ping_timeout | 20s / 20s |
| aiohttp server | WebSocketResponse(heartbeat=30) | 30s |
| aiohttp client | ws_connect(url, heartbeat=30) | 30s |
| Uvicorn | --ws-ping-interval 20 --ws-ping-timeout 20 | 20s / 20s |
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
uvicorn app:app --ws-ping-interval 20 --ws-ping-timeout 20Uvicorn 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.
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.
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 detection | Within seconds | Up to a minute or more |
| Wakeups and traffic | Many, noticeable at 10k+ connections | Few |
| Proxy safety | Very safe | May exceed a 60s proxy timeout |
Ping every 20 to 30 seconds with a timeout of about 20 seconds, and set every proxy idle timeout above that.
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.
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.
- 1Client reads slowlyits TCP receive window fills
- 2TCP stops the senderthe kernel's send buffer fills too
- 3ws.send() waitsthe write buffer passed its high-water mark
- 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.
| Side | Knob | Default | What happens when full |
|---|---|---|---|
| Send | write_limit | 32 KiB | await ws.send() waits |
| Receive | max_queue | 16 messages | the library stops reading the socket |
| Network | TCP window | set by the kernel | the remote sender blocks |
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.
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())
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.
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
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.
| Policy | How | Fits |
|---|---|---|
| Drop newest | catch QueueFull and skip x | ticks where a stale value is acceptable |
| Drop oldest | get_nowait, then put | prices, cursors |
| Keep one slot | Queue(maxsize=1) with drop oldest | a single latest value |
| Block | await q.put(x) | chat, event logs |
| Disconnect | close with 1008 or 1013 | chat, event logs |
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.
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
sent a
sent b
closed 1013The 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.
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.
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.
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.
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.
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.
naive 0.7 gather 0.5
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.
- 1Publisherone message arrives
- 2json.dumps onceone string for everyone
- 3put_nowait per clientbounded queue each
- 4Writer taskown drop or disconnect policy
- 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.
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.
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.
| Method | Waits for clients? | What a slow client does |
|---|---|---|
for ws: await ws.send(m) | Yes, one after another | Stalls every client behind it |
gather(..., return_exceptions=True) | Yes, all at once | The slowest client sets the finish time |
broadcast(conns, m) | No, it is synchronous | Not skipped; its buffer grows until the ping timeout closes it |
| Bounded queue plus writer task | No, only put_nowait | Gets 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.
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.
Set changed size during iteration
0Everything 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.
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.
['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.
| Situation | Use |
|---|---|
| Few clients, all fast | gather with return_exceptions=True |
| Latest-value ticks where a closed client is acceptable | broadcast() |
| Events that must not be lost, slow clients expected | Bounded queue plus writer task |
| Many workers or hosts | Redis pub/sub or NATS, then local fan-out |
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.
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.
| Code | Name | When you see it |
|---|---|---|
| 1000 | Normal closure | The work is done and both sides agree to stop |
| 1001 | Going away | The server is shutting down, or the user navigated away from the page |
| 1002 | Protocol error | The peer sent a broken frame, for example an unmasked client frame |
| 1003 | Unsupported data | The 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.
| Code | Name | Meaning | Retry? |
|---|---|---|---|
| 1008 | Policy violation | Authentication or authorization failed, or a rule was broken | No, fix credentials first |
| 1009 | Message too big | The message exceeded the receiver's size limit | No, send less data |
| 1011 | Internal error | The server hit an unexpected problem | Yes, with backoff |
| 1012 | Service restart | The server is restarting, for example during a deploy | Yes, with backoff |
| 1013 | Try again later | The server is overloaded or temporarily refusing clients | Yes, with backoff |
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.
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 situation | Code to send | Client should |
|---|---|---|
| Graceful shutdown | 1001 | Reconnect with backoff |
| Deploy or restart | 1012 | Back off, then return |
| Overloaded, or a client that is too slow | 1013 | Back off, then return |
| Bad or missing token | 1008 | Stop and re-authenticate |
| Message over the size limit | 1009 | Fix the payload, do not resend it |
| Token expired (your own rule) | 4001 | Refresh 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.
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.
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))1000 stop 1001 retry 1006 retry 1008 fix first 1013 retry 4001 fix first
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.
| Attempt | Ceiling (base 1s, cap 30s) | Actual delay |
|---|---|---|
| 0 | 1s | anywhere in 0 to 1s |
| 1 | 2s | anywhere in 0 to 2s |
| 3 | 8s | anywhere in 0 to 8s |
| 5 | 32s, cut to 30s | anywhere 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.
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)
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.
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.
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.
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.
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.
- 1Re-authenticatesend the token again, the server forgot you
- 2Re-subscribeask again for your channels or rooms
- 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.
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())
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.
| Exception | Typical codes | What to do |
|---|---|---|
| ConnectionClosedOK | 1000, 1001 | Stop, or reconnect if it was only a shutdown |
| ConnectionClosedError | 1006, 1011, and others | Back off and retry |
| Close received with 1008 or 4xxx | auth failure | Fix credentials first, do not loop |
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.
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.
| Line | Meaning |
|---|---|
data: hello | The payload. Several data: lines in one event are joined with newlines. |
event: token | Optional event type name. Without it the browser fires the generic message handler. |
id: 42 | Resume cursor. The browser remembers the last one it saw. |
retry: 3000 | How many milliseconds the browser waits before reconnecting. |
: keepalive | A 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.
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")))
'retry: 3000\nid: 42\nevent: token\ndata: hello\n\n' 'data: two\ndata: lines\n\n'
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.
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.
@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.
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())
'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.
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())
'data: 1\n\n' finally ran
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.
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.
- 1Stream runsevery event has an id: line
- 2Connection dropsnetwork blip, deploy, proxy timeout
- 3Browser waitsthe retry: delay
- 4Reconnectssends Last-Event-ID: 2
- 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.
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))
'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.
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.
resp = web.StreamResponse(headers={"Content-Type": "text/event-stream"})
await resp.prepare(req)
await resp.write(b"data: x\n\n")
return respaiohttp fragment inside a handler; set headers before prepare()
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.
| Layer | What to do | Why |
|---|---|---|
| Proxy such as nginx | Send X-Accel-Buffering: no | Stops the proxy from holding the body until it is full |
| Browsers and caches | Send Cache-Control: no-cache | Keeps the stream from being cached or replayed |
| Gzip middleware | Exclude the stream route | Compression 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.
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.
| SSE | WebSocket | |
|---|---|---|
| Direction | Server to client only | Both ways |
| Data | Text only | Text and binary |
| Infrastructure and auth | Plain HTTP, cookies and headers | Upgrade handshake, token often in the query string |
| Reconnect | Built into EventSource | You write the backoff |
| Best for | LLM token streams, notifications, progress, logs | Chat, games, collaboration |
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.
| Library | Server one-liner | Receive loop | End of connection |
|---|---|---|---|
| websockets | serve(handler, host, port) | async for m in ws | the loop ends on 1000/1001; ConnectionClosed is raised on an abnormal close |
| aiohttp | ws = WebSocketResponse() then await ws.prepare(req) | async for m in ws | the 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.
@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.
| Stack | Setting | Value |
|---|---|---|
| websockets | ping_interval / ping_timeout | 20 s / 20 s |
| aiohttp server and client | heartbeat=30 | 30 s |
| Uvicorn | --ws-ping-interval 20 --ws-ping-timeout 20 | 20 s / 20 s |
| nginx | proxy_read_timeout (default 60 s) | longer than the ping interval |
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 withcreate_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.
| Policy | How | Fits |
|---|---|---|
| Drop newest | skip the new item on QueueFull | stale-tolerant ticks |
| Drop oldest | get_nowait() then put_nowait() | prices, cursors |
| Keep one slot | maxsize=1 | latest value only |
| Block | await q.put(x) | chat, event logs |
| Disconnect | close with 1008 or 1013 | clients 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.
| Method | Waits? | Slow client |
|---|---|---|
for ws: await ws.send(m) | serially, one after another | stalls everyone |
gather(..., return_exceptions=True) | concurrently | the slowest sets the finish time |
broadcast(conns, data) (websockets) | no, it is synchronous | skips 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 task | no | its 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.
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.
| Code | Meaning | Client should |
|---|---|---|
| 1000 | normal close | stop |
| 1001 | going away (shutdown, navigation) | reconnect |
| 1006 | dropped, local only, never on the wire | back off and retry |
| 1008 | policy violation (auth) | do not retry; fix credentials |
| 1009 | message too big | do not retry as is |
| 1011 | internal error (also a missed pong) | back off and retry |
| 1012 | service restart | back off and retry |
| 4000-4999 | defined by your app, for example 4001 token expired | as 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.
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))
0 1 True 1 2 True 2 4 True 3 8 True 4 16 True 5 30 True 6 30 True
- 1Waitrandom 0..ceiling
- 2Connectre-authenticate
- 3Re-subscribeto your channels
- 4Resumesend the last seen id
- 5Go livereset the attempt only after about 60 s stable
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.
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")))
'id: 42\ndata: hello\n\n' 'event: token\ndata: a\ndata: b\n\n'
| Need | SSE detail |
|---|---|
| Resume | give each event an id:; the browser sends Last-Event-ID on reconnect, so replay what came after it, then go live |
| Heartbeat | a comment line such as : ping followed by a blank line, about every 15 s; browsers ignore it |
| Buffering | send X-Accel-Buffering: no and Cache-Control: no-cache, and keep gzip off the stream |
| Cleanup | put upstream cleanup in the generator's finally, which runs when the client leaves |
Pick one
| Need | Pick |
|---|---|
| two-way, low latency | WebSocket |
| server push over plain HTTP | SSE |
| rare updates | HTTP polling |
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.
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.sleepblocks 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, useawait 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. Theexcept WebSocketDisconnectnever 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
WebSocketDisconnectand read.code, write an explicitreceive_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
attemptonly 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.
clientsis iterated live. Another handler can add or discard a socket during anawait, which raisesset changed size during iteration. Loop overlist(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.dumpsruns once per client. Serialize once before the loop.- One failing
sendraises 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: noandCache-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.sleepforasyncio.sleepand push sync or CPU work intoasyncio.to_thread. - Ending a receive loop differs by library. websockets raises
ConnectionClosedon an abnormal close, aiohttp just ends the loop, and FastAPI raisesWebSocketDisconnectonly from thereceive_*()methods, not fromiter_*(). Put cleanup infinally. - 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, andLast-Event-IDto resume.