Real-Time Communication — WebSocket, gRPC Streaming and WebRTC
Handle gRPC deadlines, cancellation, flow control, keepalive, retries and shutdown
Goal
Build a relay server that passes the deadline and cancellation on to the back service, a streaming server that produces at the pace of a slow consumer, a channel that declares keepalive and retry policies, and a server that shuts down gracefully on SIGTERM, and confirm with server-side numbers.
Why it matters
A real-time service has the shape of several gRPC calls linked in a chain. If one place in the chain breaks the deadline or the cancellation, resources are spent on work nobody is waiting for; if you bypass flow control, one slow consumer fills the server's memory; if keepalive is mismatched, a quiet stream gets cut by an intermediate device; and if you misunderstand retries, you receive a stream's results twice. Servers whose streams break at every deployment are also common. Each of these is blocked by one option line or one callback, but if you don't know that one line, finding the cause takes days.
Steps
- Pass the received deadline to the back — First generate the code with /opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpcc --grpc_python_out=/root/rt/grpcc /opt/fixtures/rt/grpc/meter.proto. Then create serve(port, backend) in /root/rt/grpcc/front.py. Listen for the Meter service on 127.0.0.1:port, but implement only Work. Work passes the same request on to Meter.Work at the backend address, setting the remaining deadline of its own call (context.time_remaining()) as the timeout as it is. If the call has no deadline, do not set a timeout. If an error occurs in the back, call context.abort with the same status code and description.
- If the front is canceled, cancel the back too — Fix Work so that it makes the back call asynchronously with stub.Work.future(...) and uses context.add_callback to cancel that future when its own call ends or is canceled. Wait for the result with future.result(). The grader sets a 30-step job of 100ms each with a 10-second deadline, cancels it at 0.4 seconds, and looks 2 seconds later at how many steps the back service did.
- Produce slowly to match a slow consumer — Create Meter and serve(port) in /root/rt/grpcc/server.py. Feed sends request.n Chunks (seq, data), where data is request.size bytes. GetStats returns the running total of Chunks produced as Stats(feed_produced). The grader requests 3000 chunks of 64KiB, reads just one and pauses for 3 seconds, and checks how many the server produced in the meantime and how much the server's memory grew.
- Choose the interval between the idle timeout and the ping policy — Create make_channel(target) in /root/rt/grpcc/client.py. Return a grpc.insecure_channel with keepalive options set. The grader puts a relay that cuts connections quiet for 3 seconds, and behind it a server that rejects pings more frequent than 1 second with a GOAWAY (too_many_pings), and sends a call through your channel during which no bytes pass for 7 seconds.
- Declare retries through the channel configuration — Add make_retry_channel(target) to /root/rt/grpcc/client.py. Put a retryPolicy for the whole rt.Meter service as JSON into the grpc.service_config option. It is maxAttempts 4, initialBackoff "0.1s", maxBackoff "1s", backoffMultiplier 2, and retryableStatusCodes ["UNAVAILABLE"]. Turn grpc.enable_retries on with 1 too.
- A stream is not retried after the first response — Add collect(target, n, fail_at, fail_before) to /root/rt/grpcc/client.py. Open a channel with make_retry_channel, call Countdown(n=n, interval_ms=10, fail_at=fail_at, fail_before=fail_before, fail_unavailable=True), and return the list of received values and the name of the status code it ended with. Do not call it again in the application.
- Finish in-progress calls even during a deployment — Attach a SIGTERM handler to serve in /root/rt/grpcc/server.py so that on receiving the signal, it calls server.stop(3) to reject new calls and give in-progress calls 3 seconds. Also keep in Meter a Countdown that sends from n down to 1, as in the earlier lab. The grader starts a Countdown of 4 items at 300ms intervals, sends SIGTERM at 0.3 seconds, and checks whether that stream arrives to the end, whether new calls are rejected, and whether the process ends.
Notes
- The working folder is /root/rt/grpcc. Create it first with mkdir -p /root/rt/grpcc.
- The counterpart for the back service and the client steps is the reference server. You can start it yourself with /opt/rt-lab/bin/python /opt/fixtures/rt/grpc/refserver.py --port 50052 --min-ping-ms 1000.
- Start the front service like cd /root/rt/grpcc && /opt/rt-lab/bin/python -c "import front; front.serve(50061, '127.0.0.1:50052')". The grader picks a free port and starts it separately.
- Two common mistakes. Passing time_remaining() of a call with no deadline as it is, and calling a call that was cut mid-stream again from the start in the application so that you receive values twice.
- 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.
Pass the received deadline to the back
First generate the code with /opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpcc --grpc_python_out=/root/rt/grpcc /opt/fixtures/rt/grpc/meter.proto. Then create serve(port, backend) in /root/rt/grpcc/front.py. Listen for the Meter service on 127.0.0.1:port, but implement only Work. Work passes the same request on to Meter.Work at the backend address, setting the remaining deadline of its own call (context.time_remaining()) as the timeout as it is. If the call has no deadline, do not set a timeout. If an error occurs in the back, call context.abort with the same status code and description.
If you receive a 0.6-second deadline in the front but call the back with no deadline, the back works to the end even after the front's client has given up. Passing on the remaining deadline is deadline propagation. On a call with no deadline, time_remaining() becomes a very large value, so do not pass it as it is.
If the front is canceled, cancel the back too
Fix Work so that it makes the back call asynchronously with stub.Work.future(...) and uses context.add_callback to cancel that future when its own call ends or is canceled. Wait for the result with future.result(). The grader sets a 30-step job of 100ms each with a 10-second deadline, cancels it at 0.4 seconds, and looks 2 seconds later at how many steps the back service did.
A deadline is conveyed through a header, but cancellation is not. When the client cancels, only the front server's context becomes inactive, and the back call that the front server is stuck waiting on does not know. Cutting it off is the front server's job.
Produce slowly to match a slow consumer
Create Meter and serve(port) in /root/rt/grpcc/server.py. Feed sends request.n Chunks (seq, data), where data is request.size bytes. GetStats returns the running total of Chunks produced as Stats(feed_produced). The grader requests 3000 chunks of 64KiB, reads just one and pauses for 3 seconds, and checks how many the server produced in the meantime and how much the server's memory grew.
If you yield one at a time from a generator, gRPC takes out the next only when the flow control window allows. If you keep a separate production thread and fill an unbounded queue ahead of time, you bypass that window and the consumer's share piles up in server memory.
Choose the interval between the idle timeout and the ping policy
Create make_channel(target) in /root/rt/grpcc/client.py. Return a grpc.insecure_channel with keepalive options set. The grader puts a relay that cuts connections quiet for 3 seconds, and behind it a server that rejects pings more frequent than 1 second with a GOAWAY (too_many_pings), and sends a call through your channel during which no bytes pass for 7 seconds.
grpc.keepalive_time_ms is the ping interval, grpc.keepalive_timeout_ms is the time to wait for an answer, grpc.keepalive_permit_without_calls is whether to ping even when there are no calls, and grpc.http2.max_pings_without_data is the upper bound on pings sent without data (0 means unlimited). If it is too sparse, the relay cuts it, and if too frequent, the server cuts it.
Declare retries through the channel configuration
Add make_retry_channel(target) to /root/rt/grpcc/client.py. Put a retryPolicy for the whole rt.Meter service as JSON into the grpc.service_config option. It is maxAttempts 4, initialBackoff "0.1s", maxBackoff "1s", backoffMultiplier 2, and retryableStatusCodes ["UNAVAILABLE"]. Turn grpc.enable_retries on with 1 too.
If you write retries as an application loop, the backoff, the cap, and which status codes to retry get scattered across calls. Errors that would give the same result if tried again, such as INVALID_ARGUMENT, are not put in the retry list.
A stream is not retried after the first response
Add collect(target, n, fail_at, fail_before) to /root/rt/grpcc/client.py. Open a channel with make_retry_channel, call Countdown(n=n, interval_ms=10, fail_at=fail_at, fail_before=fail_before, fail_unavailable=True), and return the list of received values and the name of the status code it ended with. Do not call it again in the application.
The retry policy retries only errors that occur before the server sends response headers or the first message. After that, the call is committed, and even the same UNAVAILABLE comes up as it is. If you call again from the beginning a stream that already has received values, you receive those values twice.
Finish in-progress calls even during a deployment
Attach a SIGTERM handler to serve in /root/rt/grpcc/server.py so that on receiving the signal, it calls server.stop(3) to reject new calls and give in-progress calls 3 seconds. Also keep in Meter a Countdown that sends from n down to 1, as in the earlier lab. The grader starts a Countdown of 4 items at 300ms intervals, sends SIGTERM at 0.3 seconds, and checks whether that stream arrives to the end, whether new calls are rejected, and whether the process ends.
When Kubernetes takes down a Pod, it sends SIGTERM and then sends SIGKILL after terminationGracePeriodSeconds. Without a handler, Python dies immediately on SIGTERM, and every client that was streaming receives UNAVAILABLE. Inside the signal handler, only call stop and do not wait.