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.
This commit is contained in:
l0ng-ai
2026-10-04 16:17:40 +08:00
committed by GitHub
parent 6b73242b1c
commit c84f761783
3 changed files with 140 additions and 18 deletions
+13 -7
View File
@@ -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<dyn Backend>,
) -> io::Result<()> {
let (events_tx, mut events) = mpsc::channel::<ControlEvent>(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>()? {
ControlRequest::Refresh => {
let _ = refresh_tx.send(());
@@ -604,7 +607,7 @@ async fn pane_stream(
by: String,
mut feed: Box<dyn PaneFeed>,
mut send: SendStream,
mut recv: RecvStream,
recv: RecvStream,
backend: Arc<dyn Backend>,
) -> io::Result<()> {
let (down_tx, mut down) = mpsc::channel::<Down>(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));
}
+21 -9
View File
@@ -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<RecvStream>,
}
impl ControlReceiver {
pub async fn next(&mut self) -> Result<Option<ControlEvent>> {
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<RecvStream>,
}
impl PaneReader {
/// The next output chunk or event, `None` once the stream has ended.
/// Cancel-safe.
pub async fn next(&mut self) -> Result<Option<PaneItem>> {
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),
+106 -2
View File
@@ -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<R> {
inner: R,
decoder: Decoder,
chunk: Box<[u8]>,
}
impl<R: AsyncRead + Unpin> FrameReader<R> {
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<Option<Frame>> {
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: AsyncRead + Unpin>(r: &mut R) -> io::Result<Option<Frame>> {
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<Box<dyn std::error::Error + Send + Sync>>) -> 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::<PaneEvent>()
.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();