The interactive version of this post, with animations, is at notes.shvbsle.in/ttrpc-data-loss.
Before the AI-accelerated programming age, it was expected from a software stack that all of its stable components are written in one uniform language. In the case of containers, the ecosystem (containerd, shim, runc) was written in Golang. But in this new age, armed with LLMs, any new language can make lateral entry in this ecosystem and quickly build stable runtimes. I've been trying to introduce more Rust to this ecosystem.
One day I was watching logs from a pod that was running on a container runtime written in Rust and noticed that occasionally the log stream just abruptly ended. I deployed another pod that prints numbers from 1 to 100, and the logs sometimes abruptly ended on 99, but never on 97 or 98. I kept thinking that it was a bug in the runtime, but it turned out that the bug was hidden much deeper in the stack. It was in the implementation of the protocol<sup>1</sup>. This page contains my notes on the internals of ttrpc and ways in which this bug manifests.
Interactive: step through real ttrpc frames byte by byte.
rt-multi-thread). 05) removes the stream ID from the client's streams map.
Interactive: watch the race drop frame 100.
In code, it was this:
async fn handle_msg(&self, msg: GenMessage) {
let req_map = self.streams.clone();
tokio::spawn(async move {
if let Some(resp_tx) = get_resp_tx(req_map, &msg.header).await {
resp_tx
.send(Ok(msg))
.await
.unwrap_or_else(|_e| error!("The request has returned"));
}
});
}
src/asynchronous/client.rs @ f31f592 · issue #311
Stop spawning: the reader task handles each frame itself, in the order it read them.
async fn handle_msg(&self, msg: GenMessage) {
+ // Do not `tokio::spawn` per frame: a `FLAG_REMOTE_CLOSED` frame could
+ // then `remove` a stream from `req_map` before the preceding DATA
+ // frame's task looked it up, silently dropping the final payload.
+ // The read loop already awaits this per frame, so inline is correct.
let req_map = self.streams.clone();
- tokio::spawn(async move {
- if let Some(resp_tx) = get_resp_tx(req_map, &msg.header).await {
- resp_tx
- .send(Ok(msg))
- .await
- .unwrap_or_else(|_e| error!("The request has returned"));
- }
- });
+ if let Some(resp_tx) = get_resp_tx(req_map, &msg.header).await {
+ resp_tx
+ .send(Ok(msg))
+ .await
+ .unwrap_or_else(|_e| error!("The request has returned"));
+ }
}
Interactive: the reader handling every frame in order.
But this naive solution was NOT perfect. My "reality has a surprising amount of detail"<sup>4</sup> moment happened when the reviewer<sup>5</sup> made me aware of a different limitation of this fix.
There is another construct that we must think of: the bounded mpsc channel. Every call's frames reach its caller through a tokio mpsc channel that holds 100 messages, and send().await waits while it is full.
let (tx, rx): (ResultSender, ResultReceiver) = mpsc::channel(100);
src/asynchronous/client.rs @ f31f592, Client::new_stream
With the reader doing that send itself, one full channel stops the reader, and every other call on the connection waits behind it. Here is the reviewer's reproduction: a 200-frame stream that nobody reads, plus an unrelated unary call.
Interactive: the 200-frame reproduction.
Each call gets its own mailbox: an unbounded queue that only the reader writes to.
Interactive: the same reproduction with mailboxes.
The server also hands frames to tokio tasks, so why did only the client need this fix?
This is my PR that this note is based on: containerd/ttrpc-rust#312↩
"GRPC for low-memory environments." The README explains that grpc-go's memory overhead is a problem "when running a large number of services on a single machine". From containerd/ttrpc.↩
"The protocol does not include features for handling unreliable connections such as handshakes, resets, pings, or flow control." From the ttrpc protocol specification.↩
John Salvatier, Reality has a surprising amount of detail, 2017.↩
"The ordering bug is real, but this implementation introduces connection-wide head-of-line blocking. [...] I reproduced this by leaving a 200-frame stream unread and issuing an unrelated unary RPC: it times out with this patch." Tim Zhang, review on #312.↩
A stress test on the merged code saw about 1 in 2000 client-streaming calls with two adjacent frames swapped. None lost a frame.↩