From 810d918d364bd3d47dabde5899632bafa4fb9aed Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Mon, 20 Jul 2026 01:13:11 +0200 Subject: [PATCH] fix(core/quic): make the control-stream read cancel-safe MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `io::read_msg` frames a message with two `quinn::RecvStream::read_exact` calls, and quinn documents `read_exact` as explicitly NOT cancel-safe: the bytes it has already taken out of the stream live only in the future's own buffer and nothing puts them back on drop. Both long-lived control loops drive that read from a `tokio::select!` arm — the client pump alongside `ctrl_rx.recv()` and the resync tick, the host alongside probe/reconfig/clip-offer channels — and neither uses `biased;`, so any sibling that becomes ready ends the iteration and drops a partially-progressed read. `clock_sync` has the same shape via `tokio::time::timeout`, which can fire mid-frame before the session even starts. A control frame only has to straddle two wakeups for this to bite: a ClipOffer carries up to 16 kinds x 128 bytes of MIME, ~2 KB, which exceeds one QUIC packet and is subject to the pacer; any frame whose second half is lost or reordered does it too. Losing the consumed length prefix misaligns the stream permanently — the next read takes two payload bytes as a length, so Reconfigured, ProbeResult, BitrateChanged, ClockEcho and ClipState all decode as garbage and are silently dropped, and a bogus length up to 64 KiB parks the read forever. Mode switches, adaptive bitrate, mid-stream clock resync and clipboard are dead for the rest of the session; only a reconnect recovers, and the log shows at most one `warn!`. Add `io::MsgReader`, which keeps the frame in progress in the reader rather than the future and reads via quinn's cancel-safe `read`, and switch the three cancelling sites to it (client control loop, host control loop, clock_sync). The sequential handshake/pairing callers keep the plain `read_msg`, whose doc comment now states the constraint. No wire bytes and no ABI change — only how the same length-prefixed frames are assembled. Tests: a frame split across two wakeups with the read cancelled in between must resume and leave the following frame correctly framed (confirmed to fail — it hangs on the desynced stream — against the old behavior), plus a zero-length frame round-trip. Co-Authored-By: Claude Opus 4.8 --- crates/punktfunk-core/src/client/pump.rs | 10 ++- crates/punktfunk-core/src/quic/clock.rs | 4 +- crates/punktfunk-core/src/quic/io.rs | 76 +++++++++++++++++++ crates/punktfunk-core/src/quic/tests.rs | 83 ++++++++++++++++++++- crates/punktfunk-host/src/native/control.rs | 8 +- 5 files changed, 173 insertions(+), 8 deletions(-) diff --git a/crates/punktfunk-core/src/client/pump.rs b/crates/punktfunk-core/src/client/pump.rs index 56c4d23c..9c864eca 100644 --- a/crates/punktfunk-core/src/client/pump.rs +++ b/crates/punktfunk-core/src/client/pump.rs @@ -94,10 +94,14 @@ pub(super) async fn run_pump(args: WorkerArgs) { // as `PunktfunkError::Rejected` instead of the generic transport error the failed // read produces — the difference between "not accepted" and the actual cause. let handshake = async { - let (mut send, mut recv) = conn + let (mut send, recv) = conn .open_bi() .await .map_err(|e| PunktfunkError::Io(std::io::Error::other(e.to_string())))?; + // Frame every read on this stream through the resumable reader: the control loop + // below drives it from a `select!` arm and `clock_sync` wraps it in a timeout, and a + // partial frame lost to either would misalign the stream for the whole session. + let mut recv = io::MsgReader::new(recv); io::write_msg( &mut send, @@ -136,7 +140,7 @@ pub(super) async fn run_pump(args: WorkerArgs) { .encode(), ) .await?; - let welcome = Welcome::decode(&io::read_msg(&mut recv).await?)?; + let welcome = Welcome::decode(&recv.read_msg().await?)?; if welcome.compositor != CompositorPref::Auto { tracing::info!( compositor = welcome.compositor.as_str(), @@ -462,7 +466,7 @@ pub(super) async fn run_pump(args: WorkerArgs) { break; } } - msg = io::read_msg(&mut ctrl_recv) => { + msg = ctrl_recv.read_msg() => { let Ok(msg) = msg else { break }; // stream closed if let Ok(ack) = Reconfigured::decode(&msg) { if ack.accepted { diff --git a/crates/punktfunk-core/src/quic/clock.rs b/crates/punktfunk-core/src/quic/clock.rs index 1a22b7f1..a48ffd56 100644 --- a/crates/punktfunk-core/src/quic/clock.rs +++ b/crates/punktfunk-core/src/quic/clock.rs @@ -38,7 +38,7 @@ pub struct ClockSkew { /// with, so the offset aligns a client receive instant to the host's capture clock. pub async fn clock_sync( send: &mut quinn::SendStream, - recv: &mut quinn::RecvStream, + recv: &mut io::MsgReader, ) -> Option { use std::time::Duration; const ROUNDS: usize = 8; @@ -50,7 +50,7 @@ pub async fn clock_sync( if io::write_msg(send, &probe).await.is_err() { break; } - let read = tokio::time::timeout(read_timeout, io::read_msg(recv)).await; + let read = tokio::time::timeout(read_timeout, recv.read_msg()).await; let echo = match read { Ok(Ok(b)) => match ClockEcho::decode(&b) { Ok(e) => e, diff --git a/crates/punktfunk-core/src/quic/io.rs b/crates/punktfunk-core/src/quic/io.rs index be49aefe..1a6b8bee 100644 --- a/crates/punktfunk-core/src/quic/io.rs +++ b/crates/punktfunk-core/src/quic/io.rs @@ -1,6 +1,14 @@ //! Length-prefixed framing for QUIC control-stream messages: a `u16` length header followed by the //! payload, bounded at 64 KiB (control messages are tiny). /// Read one framed message (bounded at 64 KiB — control messages are tiny). +/// +/// **Not cancel-safe**: it frames with two `quinn::RecvStream::read_exact` calls, and quinn +/// documents `read_exact` as not cancel-safe (the bytes it has already taken out of the stream +/// live only in the future's own buffer, and nothing puts them back on drop). Dropping a +/// partially-progressed future therefore destroys the bytes it consumed and misaligns every +/// subsequent read on that stream. Use it only where the read runs to completion — the sequential +/// handshake/pairing exchanges. Anything driving a read from a `select!` arm or a +/// `tokio::time::timeout` must use [`MsgReader`] instead. pub async fn read_msg(recv: &mut quinn::RecvStream) -> std::io::Result> { let mut len = [0u8; 2]; recv.read_exact(&mut len) @@ -14,6 +22,74 @@ pub async fn read_msg(recv: &mut quinn::RecvStream) -> std::io::Result> Ok(buf) } +/// Cancel-safe framed reader for a long-lived control stream. +/// +/// Keeps the frame in progress in `buf` rather than inside the read future, so dropping the future +/// — which both control loops do on every iteration where a sibling `select!` arm wins, and which +/// [`clock_sync`](super::clock_sync) does on a read timeout — resumes instead of losing bytes. +/// With the plain [`read_msg`] a control frame that straddles two wakeups (a ~2 KB `ClipOffer` +/// exceeds one QUIC packet; so does any frame whose second half is lost or reordered) left the +/// stream permanently misaligned: the next read took two payload bytes as a length, every later +/// message decoded as garbage and was silently ignored, and a bogus 64 KiB length parked the read +/// forever — killing mode switches, adaptive bitrate, clock re-sync and clipboard for the rest of +/// the session with nothing but a `warn!` in the log. +pub struct MsgReader { + recv: quinn::RecvStream, + /// The frame in progress, length prefix included. + buf: Vec, + /// Bytes `buf` must reach: 2 while reading the prefix, then `2 + payload length`. + need: usize, +} + +impl MsgReader { + pub fn new(recv: quinn::RecvStream) -> Self { + MsgReader { + recv, + buf: Vec::new(), + need: 2, + } + } + + /// Read one framed message. Cancel-safe: dropping the future keeps the partial frame, so the + /// next call resumes where this one stopped. + pub async fn read_msg(&mut self) -> std::io::Result> { + loop { + while self.buf.len() < self.need { + let mut chunk = [0u8; 2048]; + let want = (self.need - self.buf.len()).min(chunk.len()); + // `read` IS cancel-safe: it only reports bytes it hands back, and they are + // committed to `self.buf` before the next await point. + match self + .recv + .read(&mut chunk[..want]) + .await + .map_err(std::io::Error::other)? + { + Some(n) => self.buf.extend_from_slice(&chunk[..n]), + None => { + return Err(std::io::Error::new( + std::io::ErrorKind::UnexpectedEof, + "control stream finished mid-frame", + )) + } + } + } + if self.need == 2 { + self.need = 2 + u16::from_le_bytes([self.buf[0], self.buf[1]]) as usize; + if self.need == 2 { + self.buf.clear(); + return Ok(Vec::new()); // zero-length frame + } + } else { + let msg = self.buf.split_off(2); + self.buf.clear(); + self.need = 2; + return Ok(msg); + } + } + } +} + /// Write one framed message. pub async fn write_msg(send: &mut quinn::SendStream, payload: &[u8]) -> std::io::Result<()> { send.write_all(&super::frame(payload)) diff --git a/crates/punktfunk-core/src/quic/tests.rs b/crates/punktfunk-core/src/quic/tests.rs index 502e0e0d..dc019d5e 100644 --- a/crates/punktfunk-core/src/quic/tests.rs +++ b/crates/punktfunk-core/src/quic/tests.rs @@ -1585,7 +1585,7 @@ mod clip_loopback { /// Stand up two loopback quinn endpoints, connect, and return /// `(server_ep, client_ep, host_conn, client_conn)`. Both endpoints are returned so the caller /// keeps them in scope — dropping a `quinn::Endpoint` tears down its connections. - async fn connect_pair() -> ( + pub(super) async fn connect_pair() -> ( quinn::Endpoint, quinn::Endpoint, quinn::Connection, @@ -1731,3 +1731,84 @@ mod clip_loopback { let _host_conn = holder.await.unwrap(); } } + +/// The control stream is read from a `select!` arm on both peers, so the read future is dropped +/// routinely — and quinn documents `read_exact` (what `io::read_msg` uses) as NOT cancel-safe. +/// [`io::MsgReader`] must survive that: the partial frame lives in the reader, not the future. +mod ctrl_framing { + use super::clip_loopback::connect_pair; + use super::*; + use crate::quic::io; + + /// A frame whose halves land in different wakeups, with the read cancelled in between, must + /// still be delivered whole — and the NEXT frame must decode correctly too. Without a + /// resumable reader the consumed length prefix is lost, the following read takes two payload + /// bytes as a length, and every later control message is garbage for the rest of the session. + #[tokio::test] + async fn cancelled_mid_frame_read_resumes_without_desync() { + let (_server_ep, _client_ep, host_conn, client_conn) = connect_pair().await; + + let first = b"the-frame-that-straddles-two-wakeups".to_vec(); + let second = b"the-frame-after-it".to_vec(); + let (f1, f2) = (first.clone(), second.clone()); + + let writer = tokio::spawn(async move { + let (mut send, _recv) = host_conn.open_bi().await.expect("open bi"); + let framed = crate::quic::frame(&f1); + // Length prefix + only part of the payload, then a real pause: this is the ClipOffer + // -sized frame split across two QUIC packets that made the bug reachable. + let split = 2 + f1.len() / 3; + send.write_all(&framed[..split]).await.expect("write head"); + tokio::time::sleep(std::time::Duration::from_millis(120)).await; + send.write_all(&framed[split..]).await.expect("write tail"); + send.write_all(&crate::quic::frame(&f2)) + .await + .expect("write second"); + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + host_conn + }); + + let (_send, recv) = client_conn.accept_bi().await.expect("accept bi"); + let mut reader = io::MsgReader::new(recv); + + // Cancel mid-frame — exactly what a sibling `select!` arm does. + let cancelled = + tokio::time::timeout(std::time::Duration::from_millis(30), reader.read_msg()).await; + assert!( + cancelled.is_err(), + "the head-only frame must not complete yet (test setup)" + ); + + let got = tokio::time::timeout(std::time::Duration::from_secs(5), reader.read_msg()) + .await + .expect("first frame must arrive after resuming") + .expect("first frame reads cleanly"); + assert_eq!(got, first, "the cancelled read must resume, not lose bytes"); + + let got2 = tokio::time::timeout(std::time::Duration::from_secs(5), reader.read_msg()) + .await + .expect("second frame must arrive") + .expect("second frame reads cleanly"); + assert_eq!(got2, second, "stream must still be framed correctly"); + + let _host_conn = writer.await.unwrap(); + } + + /// A zero-length frame is a legal encoding and must not stall the reader or eat the next one. + #[tokio::test] + async fn zero_length_frame_round_trips() { + let (_server_ep, _client_ep, host_conn, client_conn) = connect_pair().await; + let writer = tokio::spawn(async move { + let (mut send, _recv) = host_conn.open_bi().await.expect("open bi"); + send.write_all(&crate::quic::frame(&[])).await.unwrap(); + send.write_all(&crate::quic::frame(b"after")).await.unwrap(); + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + host_conn + }); + let (_send, recv) = client_conn.accept_bi().await.expect("accept bi"); + let mut reader = io::MsgReader::new(recv); + assert!(reader.read_msg().await.unwrap().is_empty()); + assert_eq!(reader.read_msg().await.unwrap(), b"after"); + let _host_conn = writer.await.unwrap(); + } +} diff --git a/crates/punktfunk-host/src/native/control.rs b/crates/punktfunk-host/src/native/control.rs index 8cd7fa4c..59aac14c 100644 --- a/crates/punktfunk-host/src/native/control.rs +++ b/crates/punktfunk-host/src/native/control.rs @@ -16,7 +16,7 @@ use punktfunk_core::quic::{ClipControl, ClipOffer, ClipState}; #[allow(clippy::too_many_arguments)] pub(super) async fn run( mut ctrl_send: quinn::SendStream, - mut ctrl_recv: quinn::RecvStream, + ctrl_recv: quinn::RecvStream, initial_mode: punktfunk_core::Mode, codec: crate::encode::Codec, live_reconfig_ok: bool, @@ -47,9 +47,13 @@ pub(super) async fn run( // coalesces a well-behaved resize drag; compliant clients self-limit to ≥ 1 s). const MIN_SWITCH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(500); let mut last_accepted_switch: Option = None; + // Resumable framing: this read is one arm of a `select!` whose siblings fire on every probe + // result / reconfigure / clip offer, so the read future is dropped routinely. `io::read_msg` + // would lose the partial frame and misalign the stream for the rest of the session. + let mut ctrl_reader = io::MsgReader::new(ctrl_recv); loop { tokio::select! { - msg = io::read_msg(&mut ctrl_recv) => { + msg = ctrl_reader.read_msg() => { let Ok(msg) = msg else { break }; // stream closed if let Ok(req) = Reconfigure::decode(&msg) { let now = std::time::Instant::now();