TT Lab
Get started
Learn Learning paths Courses

Real-Time Communication — WebSocket, gRPC Streaming and WebRTC

Check the four shapes of gRPC streaming and deadlines by hand

Continue in TT Lab

Goal

Generate code from meter.proto, implement the server, the client, and bidirectional streaming, and then check how a mid-stream error and a deadline overrun look on the client and on the server respectively.

Why it matters

gRPC streaming carries several messages on top of one HTTP/2 stream. So the response status comes at the very end, the deadline is conveyed to the server through a header, and if both sides wait on each other, it stops. If you don't know these three, things like "we switched to streaming but it comes all at once," "when an error occurs, the results we received all vanish," and "the client timed out but server CPU keeps running" happen. A voice AI's partial recognition results and token streaming flow in exactly this shape.

Steps

  1. Generate code from the contract — Generate Python code from /opt/fixtures/rt/grpc/meter.proto and place it in /root/rt/grpc. With the single line /opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpc --grpc_python_out=/root/rt/grpc /opt/fixtures/rt/grpc/meter.proto, meter_pb2.py (the messages) and meter_pb2_grpc.py (the service skeleton and stubs) are created. Open the two files and check what shape of call each of Countdown, Sum, and Chat was generated as.
  2. Server streaming — send as things are produced — In /root/rt/grpc/server.py, create Meter, which inherits meter_pb2_grpc.MeterServicer, and serve(port). serve attaches Meter to a grpc.server with 16 threads, listens on 127.0.0.1:port, and waits until it ends. Countdown sends Ticks one at a time from request.n down to 1, and rests request.interval_ms milliseconds after each one. The grader measures when the first Tick arrives.
  3. Client streaming — receive to the end and answer once — Add Sum to Meter. Read the incoming Num stream to the end and return the count and sum once as Total(count, sum). If the stream ends with nothing having arrived, it is Total(count=0, sum=0).
  4. Bidirectional streaming — answer before receiving everything — Add Chat to Meter. For each incoming Line, immediately return a Line with its text converted to uppercase. The grader sends one line, and sends the next line only after receiving the answer. And it sets a 3-second deadline.
  5. A mid-stream error does not erase what was already received — Fix Countdown so that if request.fail_at is not 0, it ends with context.abort(grpc.StatusCode.ABORTED, "fail_at") after sending fail_at Ticks. And in /root/rt/grpc/client.py, create read_countdown(target, n, interval_ms, fail_at, timeout=10). Open a channel to target, call Countdown, and return a tuple of the list of received values and the name of the status code it ended with ("OK" if normal). The grader calls it against the reference server.
  6. Put the deadline on the call — In /root/rt/grpc/client.py, add countdown_deadline(target, n, interval_ms, timeout). Call Countdown with a deadline of timeout seconds, and return a tuple of the list of received values, the status code name, and the seconds the call took. The deadline must be set through the timeout argument of the stub call.
  7. When the deadline passes, the server's work stops too — Add Work and GetStats to Meter. Work works request.steps times for request.step_ms milliseconds each, but stops if context.is_active() is false at any step, and returns the number of steps completed as WorkDone(steps_done). The server keeps a running total of the steps completed by all Work calls so far, and GetStats returns that total as Stats(work_steps). The grader calls a 20-step job of 100ms each with a 0.35-second deadline, then waits 2.5 seconds and checks the total.

Notes

Generate code from the contract

Generate Python code from /opt/fixtures/rt/grpc/meter.proto and place it in /root/rt/grpc. With the single line /opt/rt-lab/bin/python -m grpc_tools.protoc -I/opt/fixtures/rt/grpc --python_out=/root/rt/grpc --grpc_python_out=/root/rt/grpc /opt/fixtures/rt/grpc/meter.proto, meter_pb2.py (the messages) and meter_pb2_grpc.py (the service skeleton and stubs) are created. Open the two files and check what shape of call each of Countdown, Sum, and Chat was generated as.

Four shapes are separated by whether the stream in the rpc declaration is attached to the request side or the response side. unary_stream, stream_unary, and stream_stream in the generated code are those shapes. You do not edit the proto file.

Server streaming — send as things are produced

In /root/rt/grpc/server.py, create Meter, which inherits meter_pb2_grpc.MeterServicer, and serve(port). serve attaches Meter to a grpc.server with 16 threads, listens on 127.0.0.1:port, and waits until it ends. Countdown sends Ticks one at a time from request.n down to 1, and rests request.interval_ms milliseconds after each one. The grader measures when the first Tick arrives.

If you yield from a generator, gRPC sends them out one at a time. If you collect everything into a list and return it, the first message arrives at the same moment as the last message, and then there is no point in having made it streaming.

Client streaming — receive to the end and answer once

Add Sum to Meter. Read the incoming Num stream to the end and return the count and sum once as Total(count, sum). If the stream ends with nothing having arrived, it is Total(count=0, sum=0).

The request iterator ends only when the client closes the stream (half-close). Waiting for the end is the contract of this shape.

Bidirectional streaming — answer before receiving everything

Add Chat to Meter. For each incoming Line, immediately return a Line with its text converted to uppercase. The grader sends one line, and sends the next line only after receiving the answer. And it sets a 3-second deadline.

If you first collect all the requests, as in list(request_iterator), the client waits for an answer and the server waits for the requests to end, so they stop each other. Yield immediately inside the loop that runs over the iterator.

A mid-stream error does not erase what was already received

Fix Countdown so that if request.fail_at is not 0, it ends with context.abort(grpc.StatusCode.ABORTED, "fail_at") after sending fail_at Ticks. And in /root/rt/grpc/client.py, create read_countdown(target, n, interval_ms, fail_at, timeout=10). Open a channel to target, call Countdown, and return a tuple of the list of received values and the name of the status code it ended with ("OK" if normal). The grader calls it against the reference server.

The status code arrives after the messages, in the trailers. So the error pops out as a grpc.RpcError in the middle of iteration, and the messages received before it are already in your hands. If you receive with a single list(...), you lose them. e.code().name is the status code name.

Put the deadline on the call

In /root/rt/grpc/client.py, add countdown_deadline(target, n, interval_ms, timeout). Call Countdown with a deadline of timeout seconds, and return a tuple of the list of received values, the status code name, and the seconds the call took. The deadline must be set through the timeout argument of the stub call.

If you set the timeout on the call, it is conveyed to the server through the grpc-timeout header too. If you time separately on the client with a clock and stop, the server keeps working without knowing the deadline. When the deadline passes, the status is DEADLINE_EXCEEDED.

When the deadline passes, the server's work stops too

Add Work and GetStats to Meter. Work works request.steps times for request.step_ms milliseconds each, but stops if context.is_active() is false at any step, and returns the number of steps completed as WorkDone(steps_done). The server keeps a running total of the steps completed by all Work calls so far, and GetStats returns that total as Stats(work_steps). The grader calls a 20-step job of 100ms each with a 0.35-second deadline, then waits 2.5 seconds and checks the total.

The fact that the client received a deadline overrun does not stop the server's thread. The server code has to check for itself. If it does not check, it uses CPU and DB connections to the end for a result nobody will receive. Several threads raise the total, so use a lock.