TT Lab
Get started
Learn Learning paths Courses

One Slow Subscriber Froze the Livestream

Own acknowledgements, deadlines, and cleanup

Continue in TT Lab

In one line

That the send buffer is empty, that the work is confirmed done, and that the connection has been reclaimed are three different pieces of evidence.

Why this was needed

The server logged success after finishing write, but the user's screen does not change. That is because the bytes being handed to the operating system and the other program finishing the work are different things. Conversely, the other side may have processed it but the connection dropped before the ACK arrived. A real-time program must not hide this uncertainty and must express it through the message contract, the deadline, and the recovery position.

This contract is very small. The server sends EVENT 7 followed by a newline in ASCII, and the client, assuming it has finished that work, returns ACK 7 followed by a newline. The sequence number is an integer from 0 to 2147483647. It is not an authenticated business protocol or a durable message broker; it is a teaching protocol for learning about confirmation and lifetime.

How it works

StreamWriter.write puts data into the write buffer. await drain waits until transport flow control allows more. The completion of drain is not the completion of the other application, so next you wait with reader.readline for the ACK with the same sequence number. EOF, a different sequence number, or a malformed format is handled as a ConnectionError. The send function does not arbitrarily close a persistent connection that several events write to.

The time budget is granted once per event. If the total is 1 second and drain used 0.8 seconds, only about 0.2 seconds remain for the ACK. If you give a fresh 1 second to each await, the maximum wait time grows with the number of operations. At the start, set the deadline with the monotonic clock and calculate the remaining time before each wait. After the final result comes back, also check that the deadline was not exceeded. Wall-clock adjustments are not used to judge elapsed time.

The queue worker takes over one item with get and then runs send_one. Whether it succeeds or raises, finish the cleanup of the item it took over with task_done in finally. If you call task_done before the ACK, the publisher's join is released before the actual confirmation. Conversely, if you skip task_done on an exception, the task is dead while join keeps waiting. Here we apply again the point that cleanup and work success are different records.

If the worker function took over ownership of the connection, it must call writer.close even if it is canceled while waiting on an empty queue or fails while waiting for the ACK. Also put a finite deadline on wait_closed, and if the close does not finish, reclaim it by force with transport.abort. Discarding or retrying the items still left in the queue is the subscriber manager's responsibility. Do not mix one worker cleaning up its own in-flight item with emptying the whole queue.

Reconnecting needs the sequence number last processed with certainty. If 8 arrives at Cursor(last=7), it is continuous progress, and if 7 arrives again, it is a duplicate. If 10 arrives first, 8 and 9 are missing, so you raise Gap and do not move last. You must recover the omission with a log replay or a trustworthy snapshot before moving on. If you simply adopt the larger sequence number, the omission is hidden permanently.

The Cursor you build here is a small component that judges continuity in memory. When accept is True, it advances last immediately, so you must not use it as is before calling a real payment function that can fail. This course has no feature for atomically storing the application of the work and the checkpoint. Memory also disappears when the process restarts. Passing the duplicate-detection test does not guarantee exactly-once processing.

What it looks like in the field

In production, subscriber disconnects, publisher cancellations, and process shutdowns overlap at different moments. If only the normal ACK test succeeds, you miss leaks on these shutdown paths. You have to check that no sockets or tasks remain even when you disconnect repeatedly. The order also matters: close the listener first to block new connections, clean up the existing tasks and connections, and then wait for the server to shut down. In Python 3.12, Server.wait_closed waits for active connections to close.

The last real TCP test observes the partial failure and reclamation of two subscribers. The earlier steps check the overall deadline deterministically with a fake clock, and they also check the case where the ACK never arrives using a real asynchronous wait. The two approaches do not substitute for each other. The fake clock checks the boundary calculation, and the real connection checks the wiring and lifetime.

What you will do in the next lab

In the deadline test, instead of really waiting 1 second several times, you can move an injected clock to check the boundaries. But if a fake reader returns immediately, that example can pass even without any code that cuts off an infinite wait. That is why you also use a reader that never returns, separately. If you remove a failure test because it takes long, that very defect becomes an infinite wait in production. That is why you keep both the calculation boundary and the real cancellation.

In step 8 you complete the queue policies, the cursor, the overall ACK deadline, the worker, and broadcast. Put all the code in /root/realtime/delivery.py, and the checker prepares the server on an ephemeral port. Before you look at the answer, first read which layer's evidence the failure message asks for.