TT Lab
Get started
Learn Learning paths Courses

Real-Time Communication — WebSocket, gRPC Streaming and WebRTC

Operate a long-lived WebSocket hub

Continue in TT Lab

Goal

Build a publish/subscribe hub and a reconnecting client with websockets, and make them withstand, in turn, slow subscribers, idle timeouts, peers that vanish silently, and the mass disconnect caused by a deployment.

Why it matters

Most WebSocket service failures come not from the protocol but from time. When one person slows down, the hub's memory fills up, the load balancer cuts quiet connections first, when a phone goes into a tunnel the other side vanishes without a close, and a single deployment makes tens of thousands reattach at the same moment. The five rules in this lab — a per-subscriber limit, cutting with 1013, a ping shorter than the idle limit, waiting on closure together, and reconnecting with jitter and resuming by seq — are needed just the same for chat, quotes, or the partial results of a voice AI.

Steps

  1. Mix jitter into the reconnect interval — In /root/rt/wsops/client.py, create backoff(attempt, base=0.2, cap=5.0, rng=random). Return the seconds to wait before the attempt-th retry. The value is decided with a single rng.uniform(0, min(cap, base * 2 ** attempt)) (full jitter). rng is either the random module or a random.Random object.
  2. Decide whether to reconnect from the close code — In /root/rt/wsops/client.py, add should_reconnect(code). Return True for 1001 (going away), 1006 (dropped without a close frame), 1011 (internal server error), 1012 (service restart), and 1013 (try again later), and False for every other code.
  3. The hub carries one person's words to everyone — In /root/rt/wsops/hub.py, create the run(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0) coroutine. Accept three paths with websockets.asyncio.server.serve. For each text message that comes in on /pub, attach a seq that increases from 1, make it a JSON string {"seq": n, "data": }, and send it to all subscribers attached on /sub. /stats sends one JSON {"subscribers": , "dropped": , "seq": } and ends. Give each subscriber its own asyncio.Queue and a loop that drains that queue and sends.
  4. Cut only the one slow subscriber — Limit the subscriber queue's size to queue. If the queue is full and cannot take more, remove that subscriber from the list, raise dropped by 1, and close with close code 1013 and the reason "slow consumer". The grader publishes 8000 messages of 8000 bytes with one non-reading subscriber and one well-reading subscriber attached, and checks whether the well-reading one receives all of them, how much the hub's memory grew, and whether the non-reading one receives 1013.
  5. Keep quiet connections from being cut by the load balancer — Pass ping_interval and ping_timeout to serve exactly as run's arguments. The grader puts a relay that cuts connections quiet for 2.5 seconds in front of the hub, waits 6 seconds with a subscriber that does not send pings itself, and then publishes one message. That message must reach the subscriber.
  6. Clear out a subscriber that vanished silently — Fix the subscriber loop so that it waits not only on the queue but on the connection closing too. Wait on ws.wait_closed() and q.get() together with asyncio.wait(..., return_when=FIRST_COMPLETED), and if the connection closes first, remove it from the subscriber list and end the loop. Pass close_timeout=1.0 to serve too. The grader attaches a subscriber that only does the handshake and does not answer pings, and then checks whether the subscribers count in /stats returns to 0 within 6 seconds.
  7. Resume from where it broke — Give the hub a ring buffer that holds the most recent ring messages, and when a client attaches with /sub?last=N, have it first resend the messages with seq greater than N and then continue with live messages. If N+1 has already been pushed out of the buffer, first send {"reset": true, "seq": }. And in /root/rt/wsops/client.py, add the consume(url, out_path, stop_after, base=0.1, cap=1.0) coroutine. Attach to url and append the seq of each received message to out_path, one per line, and when the connection breaks, wait with should_reconnect and backoff and then reattach passing the last recorded seq as last. Finish when seq reaches stop_after. The grader cuts the connection twice in the middle of publishing and checks whether the file has 1 through stop_after exactly once each with none missing.

Notes

Mix jitter into the reconnect interval

In /root/rt/wsops/client.py, create backoff(attempt, base=0.2, cap=5.0, rng=random). Return the seconds to wait before the attempt-th retry. The value is decided with a single rng.uniform(0, min(cap, base * 2 ** attempt)) (full jitter). rng is either the random module or a random.Random object.

Without jitter, tens of thousands of clients cut at once come back at exactly the same moment and knock the freshly recovered server over again. Without a cap, the tenth retry would come several minutes later. The grader passes a seeded random.Random and computes the same value to compare.

Decide whether to reconnect from the close code

In /root/rt/wsops/client.py, add should_reconnect(code). Return True for 1001 (going away), 1006 (dropped without a close frame), 1011 (internal server error), 1012 (service restart), and 1013 (try again later), and False for every other code.

If you reattach to a connection that was cut with 1008 (policy violation) or 1002 (protocol error), it will be cut again for exactly the same reason. A retry is meaningful only when the other side's circumstances can change. 1000 means someone closed it on purpose.

The hub carries one person's words to everyone

In /root/rt/wsops/hub.py, create the run(host, port, ring=1000, queue=100, ping_interval=1.0, ping_timeout=2.0) coroutine. Accept three paths with websockets.asyncio.server.serve. For each text message that comes in on /pub, attach a seq that increases from 1, make it a JSON string {"seq": n, "data": }, and send it to all subscribers attached on /sub. /stats sends one JSON {"subscribers": , "dropped": , "seq": } and ends. Give each subscriber its own asyncio.Queue and a loop that drains that queue and sends.

If the publishing side calls await send for each subscriber in turn, then when one is slow, everyone behind them waits. A per-subscriber queue is the device that separates that waiting person by person. Publishing only puts into the queue and does not wait.

Cut only the one slow subscriber

Limit the subscriber queue's size to queue. If the queue is full and cannot take more, remove that subscriber from the list, raise dropped by 1, and close with close code 1013 and the reason "slow consumer". The grader publishes 8000 messages of 8000 bytes with one non-reading subscriber and one well-reading subscriber attached, and checks whether the well-reading one receives all of them, how much the hub's memory grew, and whether the non-reading one receives 1013.

If the queue has no upper bound, the share of messages for one slow subscriber piles up in the hub's memory without end. Whether discarding messages or cutting the connection is better is decided by the nature of the data. For a stream where order and omissions matter, it is more honest to cut and have them reconnect than to keep sending with holes.

Keep quiet connections from being cut by the load balancer

Pass ping_interval and ping_timeout to serve exactly as run's arguments. The grader puts a relay that cuts connections quiet for 2.5 seconds in front of the hub, waits 6 seconds with a subscriber that does not send pings itself, and then publishes one message. That message must reach the subscriber.

websockets' default ping interval is 20 seconds, so behind equipment whose idle limit is shorter than that, quiet connections are cut first. The default idle limit of a common AWS ALB is 60 seconds, and corporate proxies are often shorter. The interval has to be shorter than the shortest idle limit.

Clear out a subscriber that vanished silently

Fix the subscriber loop so that it waits not only on the queue but on the connection closing too. Wait on ws.wait_closed() and q.get() together with asyncio.wait(..., return_when=FIRST_COMPLETED), and if the connection closes first, remove it from the subscriber list and end the loop. Pass close_timeout=1.0 to serve too. The grader attaches a subscriber that only does the handshake and does not answer pings, and then checks whether the subscribers count in /stats returns to 0 within 6 seconds.

Even if the library closes the connection on a ping timeout, the coroutine waiting on the queue does not wake until a new message arrives. Meanwhile, the dead connection stays in the subscriber list, inflating the count and receiving and piling up messages. And the default time to wait for an answer after sending a close to an unresponsive peer is 10 seconds.

Resume from where it broke

Give the hub a ring buffer that holds the most recent ring messages, and when a client attaches with /sub?last=N, have it first resend the messages with seq greater than N and then continue with live messages. If N+1 has already been pushed out of the buffer, first send {"reset": true, "seq": }. And in /root/rt/wsops/client.py, add the consume(url, out_path, stop_after, base=0.1, cap=1.0) coroutine. Attach to url and append the seq of each received message to out_path, one per line, and when the connection breaks, wait with should_reconnect and backoff and then reattach passing the last recorded seq as last. Finish when seq reaches stop_after. The grader cuts the connection twice in the middle of publishing and checks whether the file has 1 through stop_after exactly once each with none missing.

You have to reattach based not on "what was received" but on "what was processed and recorded." If the connection breaks after receiving but before recording, that message is lost forever. If the same seq arrives twice, record it only once.