Real-Time Communication — WebSocket, gRPC Streaming and WebRTC
Operate a long-lived WebSocket hub
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
- 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.
- 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.
- 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.
- 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.
- 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.
- 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.
- 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
- The working folder is /root/rt/wsops. Create it first with mkdir -p /root/rt/wsops.
- To start the hub yourself, use cd /root/rt/wsops && /opt/rt-lab/bin/python -c "import asyncio, hub; asyncio.run(hub.run('127.0.0.1', 9002))". The grader picks a free port and starts it separately.
- Idle timeouts and connection cuts are made with IdleProxy and Relay.cut() in /opt/fixtures/rt/rtnet.py. You may read how they cut.
- Two common mistakes. Calling await send for each subscriber in the publishing loop so that one person stops everyone, and losing a message that was cut off before you moved last forward once you had recorded it on receipt.
- Always run Python with /opt/rt-lab/bin/python. This lab's libraries are only in that virtual environment, and if you run it with plain python3, you get a ModuleNotFoundError. It is convenient to shorten it with something like alias rpy=/opt/rt-lab/bin/python.
- The lab Pod blocks outbound connections. All communication happens on 127.0.0.1 inside the same Pod, and no installation or download is needed.
- The grader loads your code in a separate process and makes real connections. The example file is only a function skeleton, so it does not pass if left as is. Do not delete the functions you finished in earlier steps.
- When the lab session ends, the files in /root do not remain. Keep the code you need separately before you finish.
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.