TT Lab
Get started
Learn Learning paths Courses

One Slow Connection Froze Every Other One

Serve hundreds of connections from one loop

Continue in TT Lab

Goal

Measure the queueing of a sequential server yourself, build a multiplexing server that takes the same load in a single loop, and compare the two measurements.

Why it matters

If one connection starts talking late, a sequential server can do nothing in the meantime. CPU is idle and the log is quiet, so from the metrics alone the fact that it is blocked is itself invisible. If you gather the place to wait into one spot instead of giving each connection its own, this problem goes away, but in return the program has to manage partial reads, partial writes, interest events, and idle cleanup itself. This lab builds those management costs one by one, measures before and after the change with the same ruler, and leaves what actually improved as numbers. It uses only the standard library and needs no internet access or installation.

Steps

  1. A listening socket that doesn't block — In /root/mux/server.py, create make_listener(host, port, backlog=64). On an AF_INET, SOCK_STREAM socket, turn on SO_REUSEADDR, bind to host and port, listen with backlog, and return the socket switched with setblocking(False). Port 0 is allowed, so that the operating system picks a free port. It is a ValueError if host is not a non-empty string, if port is not an int from 0 to 65535, or if backlog is not an int of 1 or more. A bool is not accepted as an int.
  2. Record the queueing as numbers — Run python3 /opt/fixtures/mux/probe.py --target sequential. Write the two values blocked and fast at the end of the output to /root/mux/report.txt as two lines, seq_blocked= and seq_fast=. The grader runs the same measurement again on the spot and compares the two values.
  3. Empty the queue on a single wake-up — Add accept_all(listener) to server.py. Repeat accept until BlockingIOError occurs, and return a list of the accepted sockets. If no connection is waiting, it is an empty list, and the accepted sockets are also switched with setblocking(False). Skip a connection that raised ConnectionAbortedError or ConnectionResetError, and keep accepting the rest of the queue.
  4. Read only what arrived and keep the rest — Add a Conn class to server.py. Conn(sock, now=0.0) sets inbox and outbox to empty bytes, closed to False, and last_active to float(now). on_readable() calls sock.recv only once, appends the received bytes to inbox, and then returns as a list only the lines that ended with a newline, with the newline stripped. The incomplete fragment stays in inbox. If recv gives b"", set closed to True and return an empty list; for BlockingIOError, change nothing and return an empty list; for ConnectionResetError, set closed to True and return an empty list. A line of length 0 is also a line.
  5. Delete only what was sent — Add queue(data), the wants_write attribute, and on_writable() to Conn. queue accepts only bytes and appends them after outbox, and it is a ValueError if it is not bytes. wants_write is whether outbox is non-empty. on_writable returns 0 without calling send if outbox is empty; otherwise it calls sock.send once, cuts off from the front of outbox only as many bytes as the returned count, and returns that count. For BlockingIOError it is 0 and leaves outbox as it is. For BrokenPipeError and ConnectionResetError, it sets closed to True and returns 0.
  6. Watch for writes only when there is something to send — Add interest(conn) to server.py. If conn's outbox is empty, it returns just selectors.EVENT_READ; if something remains, it returns selectors.EVENT_READ and selectors.EVENT_WRITE combined with OR.
  7. Choose the quiet connections — Add touch(now) to Conn and idle_keys(conns, now, idle_timeout) to server.py. touch changes last_active to float(now). idle_keys takes a dictionary from keys to Conn and returns, as a sorted list, only the keys for which now - last_active is at least idle_timeout, and it does not change the dictionary it received. It is a ValueError if idle_timeout is not a positive finite int or float, and a bool is not accepted.
  8. Take everything in a single loop and measure again to compare — Add serve(host, port, idle_timeout, ready=None, clock=None) to server.py. After it starts listening with make_listener, if ready is given, call it once with the port number actually bound, and loop while watching the listening socket and all the connections together with selectors. When it receives a line with a value after REQ, it puts a line with OK followed by the same value into that connection's outbox, updates the interest events with interest, and removes from the interest list and closes the connections that are closed and the connections chosen by idle_keys. Put an upper bound on the wait time of the multiplexing call so that the idle check runs even without events. Then run python3 /opt/fixtures/mux/probe.py --target /root/mux/server.py and python3 /opt/fixtures/mux/fdcount.py /root/mux/server.py, and add three lines to /root/mux/report.txt, mux_blocked=, mux_fast=, and mux_fds=.

Notes

A listening socket that doesn't block

In /root/mux/server.py, create make_listener(host, port, backlog=64). On an AF_INET, SOCK_STREAM socket, turn on SO_REUSEADDR, bind to host and port, listen with backlog, and return the socket switched with setblocking(False). Port 0 is allowed, so that the operating system picks a free port. It is a ValueError if host is not a non-empty string, if port is not an int from 0 to 65535, or if backlog is not an int of 1 or more. A bool is not accepted as an int.

bind and listen only change the kernel's state and do not wait for the other side. What stops the process is accept, and the one line that prevents it is the socket's mode. In the validation, you have to handle separately the fact that bool is a subtype of int.

Record the queueing as numbers

Run python3 /opt/fixtures/mux/probe.py --target sequential. Write the two values blocked and fast at the end of the output to /root/mux/report.txt as two lines, seq_blocked= and seq_fast=. The grader runs the same measurement again on the spot and compares the two values.

The measuring tool attaches one slow client first and lines up several fast clients behind it. blocked is the number of clients that took more than 1 second, and fast is the number of clients that finished within 0.2 seconds. If you write a guess, it will disagree with the re-measured value. If you keep the output in a file, it is easy to compare in the last step.

Empty the queue on a single wake-up

Add accept_all(listener) to server.py. Repeat accept until BlockingIOError occurs, and return a list of the accepted sockets. If no connection is waiting, it is an empty list, and the accepted sockets are also switched with setblocking(False). Skip a connection that raised ConnectionAbortedError or ConnectionResetError, and keep accepting the rest of the queue.

A read-readiness notification does not tell you that exactly one connection is waiting. There is no value that tells you how many, either, so repeating until the signal that there are no more comes is the only way. That signal is not an error but a normal termination condition.

Read only what arrived and keep the rest

Add a Conn class to server.py. Conn(sock, now=0.0) sets inbox and outbox to empty bytes, closed to False, and last_active to float(now). on_readable() calls sock.recv only once, appends the received bytes to inbox, and then returns as a list only the lines that ended with a newline, with the newline stripped. The incomplete fragment stays in inbox. If recv gives b"", set closed to True and return an empty list; for BlockingIOError, change nothing and return an empty list; for ConnectionResetError, set closed to True and return an empty list. A line of length 0 is also a line.

There is no guarantee that a single recv gives you one message. You can find the newline only after joining it with the earlier fragment, and the absence of a newline does not mean the input is wrong but that it hasn't fully arrived yet. b"" and an empty line may look similar as values, but their meanings are exactly opposite.

Delete only what was sent

Add queue(data), the wants_write attribute, and on_writable() to Conn. queue accepts only bytes and appends them after outbox, and it is a ValueError if it is not bytes. wants_write is whether outbox is non-empty. on_writable returns 0 without calling send if outbox is empty; otherwise it calls sock.send once, cuts off from the front of outbox only as many bytes as the returned count, and returns that count. For BlockingIOError it is 0 and leaves outbox as it is. For BrokenPipeError and ConnectionResetError, it sets closed to True and returns 0.

The return value of send is the number of bytes this call handled. If you test with short responses, they usually go out in one go, so even a wrong implementation passes. First decide what should remain in outbox when you put in 20 bytes and only 3 went out.

Watch for writes only when there is something to send

Add interest(conn) to server.py. If conn's outbox is empty, it returns just selectors.EVENT_READ; if something remains, it returns selectors.EVENT_READ and selectors.EVENT_WRITE combined with OR.

Read readiness and write readiness become true with completely different frequencies. The send buffer is empty most of the time, so write readiness is almost always true. That fact decides this function's condition.

Choose the quiet connections

Add touch(now) to Conn and idle_keys(conns, now, idle_timeout) to server.py. touch changes last_active to float(now). idle_keys takes a dictionary from keys to Conn and returns, as a sorted list, only the keys for which now - last_active is at least idle_timeout, and it does not change the dictionary it received. It is a ValueError if idle_timeout is not a positive finite int or float, and a bool is not accepted.

If you mix choosing and cutting in one function, the data structure changes in the middle of the traversal. If you keep the choosing side a pure function, you can test boundary conditions just by feeding in a clock. First decide which way to go when the elapsed time is exactly equal to the limit.

Take everything in a single loop and measure again to compare

Add serve(host, port, idle_timeout, ready=None, clock=None) to server.py. After it starts listening with make_listener, if ready is given, call it once with the port number actually bound, and loop while watching the listening socket and all the connections together with selectors. When it receives a line with a value after REQ, it puts a line with OK followed by the same value into that connection's outbox, updates the interest events with interest, and removes from the interest list and closes the connections that are closed and the connections chosen by idle_keys. Put an upper bound on the wait time of the multiplexing call so that the idle check runs even without events. Then run python3 /opt/fixtures/mux/probe.py --target /root/mux/server.py and python3 /opt/fixtures/mux/fdcount.py /root/mux/server.py, and add three lines to /root/mux/report.txt, mux_blocked=, mux_fast=, and mux_fds=.

You can just chain together the functions you built in the earlier steps. The only thing the loop newly decides is the order. Accept, read, write, recompute the interest events, and cut what needs to be cut. The fds value is a number tied to the implementation, so it may differ from someone else's answer, and you should be able to explain why it is not 200.