From c84f761783db31118cdd151f893a0834163d8af2 Mon Sep 17 00:00:00 2001 From: l0ng-ai <24760907+l0ng-ai@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:17:40 +0800 Subject: [PATCH] fix(mobile): read frames cancel-safely so a slow link can't desync a stream (#1093) read_frame awaits the header and the payload separately and keeps nothing between them, so dropping it mid-frame discards the bytes it has read and leaves the stream inside a payload. The next read takes terminal output for a length: "Frame of 2086478377 bytes exceeds the limit" is `)"]|` off a pane running jq, seen on a phone over 4G. Three callers drop it routinely: - the phone app batches output with timeout(16ms, PaneReader::next), which fires whenever a large frame is still arriving on a slow link; - the gateway's control and pane loops race it in select! against outgoing events and output, so a paste while a pane prints can lose input bytes. Add FrameReader to tty7-mobile-proto: it keeps a partial frame in a Decoder and only awaits a single read(), which takes nothing when dropped. Use it in PaneReader, ControlReceiver and both gateway loops. read_frame stays for the one-shot reads whose stream is abandoned on a timeout, and says so. --- crates/tty7-gateway/src/serve.rs | 20 +++-- crates/tty7-mobile-client/src/lib.rs | 30 +++++--- crates/tty7-mobile-proto/src/lib.rs | 108 ++++++++++++++++++++++++++- 3 files changed, 140 insertions(+), 18 deletions(-) diff --git a/crates/tty7-gateway/src/serve.rs b/crates/tty7-gateway/src/serve.rs index d1704063..54b0a995 100644 --- a/crates/tty7-gateway/src/serve.rs +++ b/crates/tty7-gateway/src/serve.rs @@ -23,9 +23,9 @@ use tty7_core::core::machine::Machine; use tty7_core::daemon::control::PaneAgentState; use tty7_core::daemon::protocol::{DaemonMsg, LeaseRequest, WinSize}; use tty7_mobile_proto::{ - ControlEvent, ControlRequest, Diff, Frame, GridSize, MAX_DIFF, MAX_UPLOAD, Open, OpenReply, - PROTOCOL_VERSION, PaneEvent, PaneRequest, RemoteView, TabCreated, Tree, Uploaded, read_frame, - write_bytes, write_msg, + ControlEvent, ControlRequest, Diff, Frame, FrameReader, GridSize, MAX_DIFF, MAX_UPLOAD, Open, + OpenReply, PROTOCOL_VERSION, PaneEvent, PaneRequest, RemoteView, TabCreated, Tree, Uploaded, + read_frame, write_bytes, write_msg, }; use crate::state::State; @@ -466,7 +466,7 @@ async fn finish(mut send: SendStream) { async fn control_stream( mut send: SendStream, - mut recv: RecvStream, + recv: RecvStream, backend: Arc, ) -> io::Result<()> { let (events_tx, mut events) = mpsc::channel::(4); @@ -475,9 +475,12 @@ async fn control_stream( .name("gateway-tree".into()) .spawn(move || watch_tree(backend, events_tx, refresh_rx))?; + // A `select!` drops whichever arm loses, so the read has to survive + // being dropped mid-frame. + let mut frames = FrameReader::new(recv); loop { tokio::select! { - frame = read_frame(&mut recv) => match frame? { + frame = frames.next() => match frame? { Some(frame) => match frame.msg::()? { ControlRequest::Refresh => { let _ = refresh_tx.send(()); @@ -604,7 +607,7 @@ async fn pane_stream( by: String, mut feed: Box, mut send: SendStream, - mut recv: RecvStream, + recv: RecvStream, backend: Arc, ) -> io::Result<()> { let (down_tx, mut down) = mpsc::channel::(PANE_BACKLOG); @@ -674,9 +677,12 @@ async fn pane_stream( } })?; + // A `select!` drops whichever arm loses — every time output goes down — + // so the read has to survive being dropped mid-frame. + let mut frames = FrameReader::new(recv); loop { tokio::select! { - frame = read_frame(&mut recv) => match frame? { + frame = frames.next() => match frame? { Some(Frame::Bytes(bytes)) => { let _ = input_tx.send(Up::Keys(bytes)); } diff --git a/crates/tty7-mobile-client/src/lib.rs b/crates/tty7-mobile-client/src/lib.rs index 9d5cc0ed..fef776ed 100644 --- a/crates/tty7-mobile-client/src/lib.rs +++ b/crates/tty7-mobile-client/src/lib.rs @@ -14,9 +14,9 @@ use iroh::{Endpoint, EndpointAddr, EndpointId, RelayUrl, SecretKey}; use iroh_mdns_address_lookup::MdnsAddressLookup; use serde::{Deserialize, Serialize}; use tty7_mobile_proto::{ - ALPN, ControlEvent, ControlRequest, Diff, Frame, GridSize, MAX_UPLOAD, MDNS_SERVICE, Open, - OpenReply, PROTOCOL_VERSION, PairCode, PaneEvent, PaneRequest, TabCreated, Uploaded, - read_frame, write_bytes, write_msg, + ALPN, ControlEvent, ControlRequest, Diff, Frame, FrameReader, GridSize, MAX_UPLOAD, + MDNS_SERVICE, Open, OpenReply, PROTOCOL_VERSION, PairCode, PaneEvent, PaneRequest, TabCreated, + Uploaded, read_frame, write_bytes, write_msg, }; /// How long to wait for a gateway to answer an [`Open`]. @@ -185,7 +185,12 @@ impl Session { machine: machine.map(str::to_string), }; let (send, recv) = open(&self.conn, &ask).await?; - Ok((PaneWriter { send }, PaneReader { recv })) + Ok(( + PaneWriter { send }, + PaneReader { + frames: FrameReader::new(recv), + }, + )) } } @@ -364,7 +369,9 @@ impl ControlStream { pub fn split(self) -> (ControlSender, ControlReceiver) { ( ControlSender { send: self.send }, - ControlReceiver { recv: self.recv }, + ControlReceiver { + frames: FrameReader::new(self.recv), + }, ) } } @@ -380,13 +387,15 @@ impl ControlSender { } } +/// Reads with a [`FrameReader`], so a caller may race [`Self::next`] +/// against anything else without tearing a frame. pub struct ControlReceiver { - recv: RecvStream, + frames: FrameReader, } impl ControlReceiver { pub async fn next(&mut self) -> Result> { - match read_frame(&mut self.recv).await? { + match self.frames.next().await? { Some(frame) => Ok(Some(frame.msg()?)), None => Ok(None), } @@ -419,14 +428,17 @@ pub enum PaneItem { Event(PaneEvent), } +/// Reads with a [`FrameReader`], so the app may put a deadline on +/// [`Self::next`] (it batches output by one) without tearing a frame. pub struct PaneReader { - recv: RecvStream, + frames: FrameReader, } impl PaneReader { /// The next output chunk or event, `None` once the stream has ended. + /// Cancel-safe. pub async fn next(&mut self) -> Result> { - match read_frame(&mut self.recv).await? { + match self.frames.next().await? { Some(Frame::Bytes(bytes)) => Ok(Some(PaneItem::Output(bytes))), Some(frame) => Ok(Some(PaneItem::Event(frame.msg()?))), None => Ok(None), diff --git a/crates/tty7-mobile-proto/src/lib.rs b/crates/tty7-mobile-proto/src/lib.rs index 9a334fd7..b1de52b4 100644 --- a/crates/tty7-mobile-proto/src/lib.rs +++ b/crates/tty7-mobile-proto/src/lib.rs @@ -508,8 +508,53 @@ mod io_async { use super::*; use tokio::io::{AsyncRead, AsyncReadExt as _, AsyncWrite, AsyncWriteExt as _}; + /// Reads frames off a stream, keeping a frame it has only part of in its + /// own buffer. That makes [`FrameReader::next`] cancel-safe: dropping it in + /// a `select!` or at a `timeout` loses nothing, where dropping + /// [`read_frame`] mid-frame throws away the bytes it has read and leaves + /// the stream in the middle of a payload — whose bytes the next read then + /// takes for a length. + pub struct FrameReader { + inner: R, + decoder: Decoder, + chunk: Box<[u8]>, + } + + impl FrameReader { + pub fn new(inner: R) -> Self { + FrameReader { + inner, + decoder: Decoder::default(), + chunk: vec![0; 64 << 10].into_boxed_slice(), + } + } + + /// The next frame. `Ok(None)` is a clean end of stream between frames; + /// an end in the middle of one is an error. Cancel-safe. + pub async fn next(&mut self) -> io::Result> { + loop { + if let Some(frame) = self.decoder.next_frame()? { + return Ok(Some(frame)); + } + // The only await: a read that is dropped before it completes + // has taken nothing off the stream. + let n = self.inner.read(&mut self.chunk).await?; + if n == 0 { + return if self.decoder.buf.is_empty() { + Ok(None) + } else { + Err(io::ErrorKind::UnexpectedEof.into()) + }; + } + self.decoder.push(&self.chunk[..n]); + } + } + } + /// Reads one frame. `Ok(None)` is a clean end of stream between frames; an - /// end in the middle of one is an error. + /// end in the middle of one is an error. Not cancel-safe: anything that + /// may drop it before it finishes — a `select!` arm, a `timeout` the + /// stream outlives — wants a [`FrameReader`]. pub async fn read_frame(r: &mut R) -> io::Result> { let mut header = [0u8; 5]; let mut got = 0; @@ -543,7 +588,7 @@ mod io_async { } #[cfg(feature = "tokio")] -pub use io_async::{read_frame, write_bytes, write_msg}; +pub use io_async::{FrameReader, read_frame, write_bytes, write_msg}; fn invalid(e: impl Into>) -> io::Error { io::Error::new(io::ErrorKind::InvalidData, e) @@ -587,6 +632,65 @@ mod tests { } } + /// A read dropped with half a frame in hand — a phone's output batching + /// timing out on a slow link, a gateway's `select!` taking the other arm — + /// must leave the next read at the start of that frame, not in the middle + /// of its payload reading terminal output as a length. + #[tokio::test] + async fn a_frame_reader_dropped_mid_frame_loses_nothing() { + use tokio::io::AsyncWriteExt as _; + + let (mut far, near) = tokio::io::duplex(1 << 16); + let mut reader = FrameReader::new(near); + let output = encode_bytes(b"[]|\"\\(.name)=\\(.conclusion)\""); + let (head, tail) = output.split_at(9); + + far.write_all(head).await.unwrap(); + // Poll once, so the reader takes what has arrived, then drop it. + tokio::select! { + biased; + _ = reader.next() => panic!("half a frame read as a whole one"), + () = std::future::ready(()) => {} + } + far.write_all(tail).await.unwrap(); + far.write_all(&encode_msg(&PaneEvent::Size { cols: 80, rows: 24 })) + .await + .unwrap(); + drop(far); + + assert_eq!( + reader.next().await.unwrap(), + Some(Frame::Bytes(b"[]|\"\\(.name)=\\(.conclusion)\"".to_vec())) + ); + assert_eq!( + reader + .next() + .await + .unwrap() + .unwrap() + .msg::() + .unwrap(), + PaneEvent::Size { cols: 80, rows: 24 } + ); + assert_eq!(reader.next().await.unwrap(), None); + } + + #[tokio::test] + async fn a_frame_reader_refuses_a_stream_that_ends_mid_frame() { + use tokio::io::AsyncWriteExt as _; + + let (mut far, near) = tokio::io::duplex(1 << 16); + let mut reader = FrameReader::new(near); + far.write_all(&encode_bytes(b"cut short")[..7]) + .await + .unwrap(); + drop(far); + assert_eq!( + reader.next().await.unwrap_err().kind(), + io::ErrorKind::UnexpectedEof + ); + } + #[test] fn oversized_and_unknown_frames_are_refused() { let mut dec = Decoder::default();