fix(core/quic): make the control-stream read cancel-safe
`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 <noreply@anthropic.com>
This commit is contained in:
@@ -94,10 +94,14 @@ pub(super) async fn run_pump(args: WorkerArgs) {
|
|||||||
// as `PunktfunkError::Rejected` instead of the generic transport error the failed
|
// as `PunktfunkError::Rejected` instead of the generic transport error the failed
|
||||||
// read produces — the difference between "not accepted" and the actual cause.
|
// read produces — the difference between "not accepted" and the actual cause.
|
||||||
let handshake = async {
|
let handshake = async {
|
||||||
let (mut send, mut recv) = conn
|
let (mut send, recv) = conn
|
||||||
.open_bi()
|
.open_bi()
|
||||||
.await
|
.await
|
||||||
.map_err(|e| PunktfunkError::Io(std::io::Error::other(e.to_string())))?;
|
.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(
|
io::write_msg(
|
||||||
&mut send,
|
&mut send,
|
||||||
@@ -136,7 +140,7 @@ pub(super) async fn run_pump(args: WorkerArgs) {
|
|||||||
.encode(),
|
.encode(),
|
||||||
)
|
)
|
||||||
.await?;
|
.await?;
|
||||||
let welcome = Welcome::decode(&io::read_msg(&mut recv).await?)?;
|
let welcome = Welcome::decode(&recv.read_msg().await?)?;
|
||||||
if welcome.compositor != CompositorPref::Auto {
|
if welcome.compositor != CompositorPref::Auto {
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
compositor = welcome.compositor.as_str(),
|
compositor = welcome.compositor.as_str(),
|
||||||
@@ -462,7 +466,7 @@ pub(super) async fn run_pump(args: WorkerArgs) {
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
msg = io::read_msg(&mut ctrl_recv) => {
|
msg = ctrl_recv.read_msg() => {
|
||||||
let Ok(msg) = msg else { break }; // stream closed
|
let Ok(msg) = msg else { break }; // stream closed
|
||||||
if let Ok(ack) = Reconfigured::decode(&msg) {
|
if let Ok(ack) = Reconfigured::decode(&msg) {
|
||||||
if ack.accepted {
|
if ack.accepted {
|
||||||
|
|||||||
@@ -38,7 +38,7 @@ pub struct ClockSkew {
|
|||||||
/// with, so the offset aligns a client receive instant to the host's capture clock.
|
/// with, so the offset aligns a client receive instant to the host's capture clock.
|
||||||
pub async fn clock_sync(
|
pub async fn clock_sync(
|
||||||
send: &mut quinn::SendStream,
|
send: &mut quinn::SendStream,
|
||||||
recv: &mut quinn::RecvStream,
|
recv: &mut io::MsgReader,
|
||||||
) -> Option<ClockSkew> {
|
) -> Option<ClockSkew> {
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
const ROUNDS: usize = 8;
|
const ROUNDS: usize = 8;
|
||||||
@@ -50,7 +50,7 @@ pub async fn clock_sync(
|
|||||||
if io::write_msg(send, &probe).await.is_err() {
|
if io::write_msg(send, &probe).await.is_err() {
|
||||||
break;
|
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 {
|
let echo = match read {
|
||||||
Ok(Ok(b)) => match ClockEcho::decode(&b) {
|
Ok(Ok(b)) => match ClockEcho::decode(&b) {
|
||||||
Ok(e) => e,
|
Ok(e) => e,
|
||||||
|
|||||||
@@ -1,6 +1,14 @@
|
|||||||
//! Length-prefixed framing for QUIC control-stream messages: a `u16` length header followed by the
|
//! Length-prefixed framing for QUIC control-stream messages: a `u16` length header followed by the
|
||||||
//! payload, bounded at 64 KiB (control messages are tiny).
|
//! payload, bounded at 64 KiB (control messages are tiny).
|
||||||
/// Read one framed message (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<Vec<u8>> {
|
pub async fn read_msg(recv: &mut quinn::RecvStream) -> std::io::Result<Vec<u8>> {
|
||||||
let mut len = [0u8; 2];
|
let mut len = [0u8; 2];
|
||||||
recv.read_exact(&mut len)
|
recv.read_exact(&mut len)
|
||||||
@@ -14,6 +22,74 @@ pub async fn read_msg(recv: &mut quinn::RecvStream) -> std::io::Result<Vec<u8>>
|
|||||||
Ok(buf)
|
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<u8>,
|
||||||
|
/// 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<Vec<u8>> {
|
||||||
|
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.
|
/// Write one framed message.
|
||||||
pub async fn write_msg(send: &mut quinn::SendStream, payload: &[u8]) -> std::io::Result<()> {
|
pub async fn write_msg(send: &mut quinn::SendStream, payload: &[u8]) -> std::io::Result<()> {
|
||||||
send.write_all(&super::frame(payload))
|
send.write_all(&super::frame(payload))
|
||||||
|
|||||||
@@ -1585,7 +1585,7 @@ mod clip_loopback {
|
|||||||
/// Stand up two loopback quinn endpoints, connect, and return
|
/// Stand up two loopback quinn endpoints, connect, and return
|
||||||
/// `(server_ep, client_ep, host_conn, client_conn)`. Both endpoints are returned so the caller
|
/// `(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.
|
/// 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::Endpoint,
|
quinn::Endpoint,
|
||||||
quinn::Connection,
|
quinn::Connection,
|
||||||
@@ -1731,3 +1731,84 @@ mod clip_loopback {
|
|||||||
let _host_conn = holder.await.unwrap();
|
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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -16,7 +16,7 @@ use punktfunk_core::quic::{ClipControl, ClipOffer, ClipState};
|
|||||||
#[allow(clippy::too_many_arguments)]
|
#[allow(clippy::too_many_arguments)]
|
||||||
pub(super) async fn run(
|
pub(super) async fn run(
|
||||||
mut ctrl_send: quinn::SendStream,
|
mut ctrl_send: quinn::SendStream,
|
||||||
mut ctrl_recv: quinn::RecvStream,
|
ctrl_recv: quinn::RecvStream,
|
||||||
initial_mode: punktfunk_core::Mode,
|
initial_mode: punktfunk_core::Mode,
|
||||||
codec: crate::encode::Codec,
|
codec: crate::encode::Codec,
|
||||||
live_reconfig_ok: bool,
|
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).
|
// 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);
|
const MIN_SWITCH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(500);
|
||||||
let mut last_accepted_switch: Option<std::time::Instant> = None;
|
let mut last_accepted_switch: Option<std::time::Instant> = 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 {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
msg = io::read_msg(&mut ctrl_recv) => {
|
msg = ctrl_reader.read_msg() => {
|
||||||
let Ok(msg) = msg else { break }; // stream closed
|
let Ok(msg) = msg else { break }; // stream closed
|
||||||
if let Ok(req) = Reconfigure::decode(&msg) {
|
if let Ok(req) = Reconfigure::decode(&msg) {
|
||||||
let now = std::time::Instant::now();
|
let now = std::time::Instant::now();
|
||||||
|
|||||||
Reference in New Issue
Block a user