TT Lab
Get started
Learn Learning paths Courses

One Slow Subscriber Froze the Livestream

Rescue a livestream by isolating a slow subscriber

Continue in TT Lab

Goal

Implement queue policies, ACK deadlines, and connection lifetime to isolate a slow subscriber on real TCP.

Why it matters

One subscriber's delay can block the whole publish or keep consuming memory. You build a program that distinguishes how much data loss is acceptable from work confirmation and that reclaims resources even on cancellation. You need the basics of Python functions, exceptions, and async/await. You use only the standard library, and no installation or internet is needed.

Steps

  1. Put a limit on a queue that piles up endlessly — Implement make_queue(capacity) in delivery.py. Return an empty asyncio.Queue whose maxsize equals the input. Allow only a positive int; a bool, float, string, 0, or negative is a ValueError.
  2. Wait, but leave no canceled publish behind — Add async offer_wait(queue, item). Wait until space opens, insert, and return None. A cancellation while waiting propagates as CancelledError and preserves the existing items and their order. Do not insert the canceled item later.
  3. Drop the old coordinates and keep the current position — Add a synchronous function offer_latest(queue, item). If there is space, insert and return None; if full, remove the one item that has waited longest, put in the new item, and return the removed item. Also match the discarded item with one task_done.
  4. Report a stalled subscriber with a policy exception — Implement an Exception subclass SlowConsumer and a synchronous offer_disconnect(queue, item). If there is space, insert and return None. If full, raise SlowConsumer without changing the existing queue. Closing the connection is not this function's responsibility.
  5. Skip duplicates and do not hide gaps — Add an Exception subclass Gap and Cursor(last=-1). last is a public per-instance attribute and an int in the range -1–2147483647. accept(seq) accepts only an int from 0–2147483647. If seq<=last, return False and keep the state; if seq==last+1, update last and return True; for a larger sequence number, raise Gap and keep the state. A wrong type such as bool, or an out-of-range value, is a ValueError. It is an in-memory cursor and does not implement durable business processing.
  6. Give sending and work confirmation one time budget — Create async send_one(reader, writer, seq, timeout, clock=None). seq is a valid sequence number as in step 5, and timeout is a positive finite int/float excluding bool; anything else is a ValueError. The default clock is time.monotonic, and an injected clock returns seconds with no arguments. Write ASCII EVENT, the sequence number, and a newline once with writer.write, wait for drain, and then you must receive, with reader.readline, an ACK with exactly the same sequence number and a newline to return None. EOF or a wrong ACK is a ConnectionError, and exceeding the overall deadline is a TimeoutError. drain and the ACK share the same deadline, and the deadline is checked even after completion. This function does not close the writer and propagates errors and cancellation.
  7. Reclaim the connection even when empty or failing — Add async serve_queue(reader, writer, queue, timeout). A positive finite timeout is passed in. Repeatedly process the seq you got with send_one, and call task_done for the item you took over exactly once, only after that call finishes, fails, or is canceled. On every exit path, including waiting on an empty queue, call writer.close and then wait for wait_closed within timeout. On a timeout or ConnectionError while waiting for the close, reclaim with writer.transport.abort, and do not hide the original send error or cancellation. Cleaning up items still left in the queue is the caller's responsibility.
  8. Disconnect the one slow subscriber and keep the livestream going — Complete a synchronous broadcast(queues, item, disconnect). queues is a name-to-queue dictionary, and you apply offer_disconnect to each queue in an iteration snapshot. Remove from the dictionary only the target that raised SlowConsumer, call disconnect(name) once, and then keep delivering to the other subscribers. Propagate other exceptions, and a normal return is None. The callback handles, without raising, canceling that worker, reclaiming the connection, and cleaning up the remaining queue. In the real TCP check, slow holds the first ACK and fast confirms every time. fast must receive all events 0–8 and only slow must be disconnected.

Notes

Put all functions in a single file, /root/realtime/delivery.py. Create the working folder with mkdir -p /root/realtime. Leave the not-yet-implemented functions in the example as skeletons, but do not overwrite the functions you have finished. Grading also checks the contracts of earlier steps. Running one step has a 5-second limit, which is not a limit on the time you take to write. The TCP server, ephemeral port, and subscribers are prepared and shut down by the checker. This lab does not guarantee kernel buffer saturation, internet performance, durable delivery, or real reconnect recovery. Files do not persist when the lab session ends, so keep any code you need separately before it ends.

Put a limit on a queue that piles up endlessly

Implement make_queue(capacity) in delivery.py. Return an empty asyncio.Queue whose maxsize equals the input. Allow only a positive int; a bool, float, string, 0, or negative is a ValueError.

In asyncio.Queue, 0 means unlimited. Split the type check and the range check, and use the standard library.

Wait, but leave no canceled publish behind

Add async offer_wait(queue, item). Wait until space opens, insert, and return None. A cancellation while waiting propagates as CancelledError and preserves the existing items and their order. Do not insert the canceled item later.

put_nowait and await put behave differently when the queue is full. Do not turn a cancellation into a normal success.

Drop the old coordinates and keep the current position

Add a synchronous function offer_latest(queue, item). If there is space, insert and return None; if full, remove the one item that has waited longest, put in the new item, and return the removed item. Also match the discarded item with one task_done.

With 10 and 11 in a capacity 2 queue, if 12 arrives, 11 and 12 remain. An in-flight item cannot be removed from this queue.

Report a stalled subscriber with a policy exception

Implement an Exception subclass SlowConsumer and a synchronous offer_disconnect(queue, item). If there is space, insert and return None. If full, raise SlowConsumer without changing the existing queue. Closing the connection is not this function's responsibility.

One function judges the saturation policy and the connection owner reclaims. Quietly discarding or waiting is a different policy.

Skip duplicates and do not hide gaps

Add an Exception subclass Gap and Cursor(last=-1). last is a public per-instance attribute and an int in the range -1–2147483647. accept(seq) accepts only an int from 0–2147483647. If seq<=last, return False and keep the state; if seq==last+1, update last and return True; for a larger sequence number, raise Gap and keep the state. A wrong type such as bool, or an out-of-range value, is a ValueError. It is an in-memory cursor and does not implement durable business processing.

The order of the duplicate check and the gap check matters. If 9 comes after 7, there is no evidence that 8 was processed.

Give sending and work confirmation one time budget

Create async send_one(reader, writer, seq, timeout, clock=None). seq is a valid sequence number as in step 5, and timeout is a positive finite int/float excluding bool; anything else is a ValueError. The default clock is time.monotonic, and an injected clock returns seconds with no arguments. Write ASCII EVENT, the sequence number, and a newline once with writer.write, wait for drain, and then you must receive, with reader.readline, an ACK with exactly the same sequence number and a newline to return None. EOF or a wrong ACK is a ConnectionError, and exceeding the overall deadline is a TimeoutError. drain and the ACK share the same deadline, and the deadline is checked even after completion. This function does not close the writer and propagates errors and cancellation.

Give asyncio.wait_for the remaining time every time. The test injects a clock and checks whether a 0.8-second drain and a 0.3-second ACK exceed the total of 1 second.

Reclaim the connection even when empty or failing

Add async serve_queue(reader, writer, queue, timeout). A positive finite timeout is passed in. Repeatedly process the seq you got with send_one, and call task_done for the item you took over exactly once, only after that call finishes, fails, or is canceled. On every exit path, including waiting on an empty queue, call writer.close and then wait for wait_closed within timeout. On a timeout or ConnectionError while waiting for the close, reclaim with writer.transport.abort, and do not hide the original send error or cancellation. Cleaning up items still left in the queue is the caller's responsibility.

The try/finally for item ownership and the try/finally for connection ownership have different scopes. If you call task_done before the ACK, join is released early.

Disconnect the one slow subscriber and keep the livestream going

Complete a synchronous broadcast(queues, item, disconnect). queues is a name-to-queue dictionary, and you apply offer_disconnect to each queue in an iteration snapshot. Remove from the dictionary only the target that raised SlowConsumer, call disconnect(name) once, and then keep delivering to the other subscribers. Propagate other exceptions, and a normal return is None. The callback handles, without raising, canceling that worker, reclaiming the connection, and cleaning up the remaining queue. In the real TCP check, slow holds the first ACK and fast confirms every time. fast must receive all events 0–8 and only slow must be disconnected.

Do not return after the first disconnect. The checker starts a real server and two clients, so you do not need to keep a server running yourself.