Skip to main content

onlyne_wire/
frame.rs

1//! Length-prefixed JSON framing for Onlyne (decision D8).
2//!
3//! A frame is a `u32 big-endian byte length` followed by that many bytes of
4//! UTF-8 JSON. One connection carries request, response, and event frames side
5//! by side. This crate moves opaque serialisable values and holds no protocol
6//! semantics, so `onlyne-proto` stays free of tokio.
7//!
8//! Failure modes are attached to an [`io::Error`] kind so callers can map them
9//! onto a wire error code:
10//!
11//! - [`TooLarge`] -> `frame_too_large`
12//! - undecodable JSON ([`io::ErrorKind::InvalidData`]) -> `bad_frame`
13//! - [`UnexpectedEof`] -> the peer dropped the connection mid-frame
14//! - a clean end of stream at a frame boundary -> [`read_frame`] yields `Ok(None)`
15//!
16//! The decoder never buffers more than one frame body at a time, and the ceiling
17//! is applied before any body byte is read: the announced length is checked
18//! against [`MAX_FRAME_BYTES`], and [`fill`] refuses a destination larger than
19//! [`MAX_PARTIAL_FRAME_BYTES`]. A sender that dribbles bytes therefore holds the
20//! reader's memory to the announced length it was allowed to announce, so
21//! `MAX_FRAME_BYTES` is the binding limit and `MAX_PARTIAL_FRAME_BYTES` is the
22//! outer wall behind it.
23
24use serde::Serialize;
25use serde::de::DeserializeOwned;
26use std::fmt;
27use std::io::{Error, ErrorKind, Result};
28use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
29
30/// Hard ceiling for one frame body. Larger payloads are rejected before any byte
31/// reaches the stream, which keeps the connection framing consistent.
32pub const MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
33
34/// Outer ceiling for the bytes one decode holds: the body buffer, on top of the
35/// four header bytes that name it.
36///
37/// [`MAX_FRAME_BYTES`] is the limit a peer actually meets: an announced length
38/// over it fails before the body is read. This ceiling sits behind that check and
39/// is enforced in [`fill`], the one place stream bytes land in memory, so the
40/// accumulation bound stays true even if the per-frame ceiling is ever raised or a
41/// new caller hands `fill` a buffer sized from untrusted input. Reading a frame
42/// cannot buffer past this many bytes, however slowly the sender dribbles them.
43pub const MAX_PARTIAL_FRAME_BYTES: usize = 16 * 1024 * 1024;
44
45// The per-frame gate must stay inside the accumulation wall, or the tighter check
46// stops being the one a peer meets.
47const _: () = assert!(MAX_FRAME_BYTES + 4 <= MAX_PARTIAL_FRAME_BYTES);
48
49/// A frame ran past a ceiling: [`MAX_FRAME_BYTES`] for one body, announced or
50/// serialised, or [`MAX_PARTIAL_FRAME_BYTES`] when the decoder's own accumulation
51/// bound is what a read ran into.
52#[derive(Debug, Clone, Copy, PartialEq, Eq)]
53pub struct TooLarge {
54    /// Length the peer announced, or the length the local side serialised to.
55    pub len: u64,
56    /// The ceiling that was crossed: [`MAX_FRAME_BYTES`] for a frame body, or
57    /// [`MAX_PARTIAL_FRAME_BYTES`] when the decoder's own accumulation bound is
58    /// what a read ran into.
59    pub max: usize,
60}
61
62impl fmt::Display for TooLarge {
63    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
64        write!(
65            f,
66            "frame of {} bytes exceeds {max}",
67            self.len,
68            max = self.max
69        )
70    }
71}
72
73impl std::error::Error for TooLarge {}
74
75/// The stream ended in the middle of a frame: fewer bytes arrived than the
76/// length prefix announced.
77#[derive(Debug, Clone, Copy, PartialEq, Eq)]
78pub struct UnexpectedEof {
79    /// Bytes the length prefix promised.
80    pub expected: u64,
81    /// Bytes that actually arrived.
82    pub got: u64,
83}
84
85impl fmt::Display for UnexpectedEof {
86    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
87        write!(
88            f,
89            "stream ended mid-frame: {} of {} bytes arrived",
90            self.got, self.expected
91        )
92    }
93}
94
95impl std::error::Error for UnexpectedEof {}
96
97fn too_large(len: u64, max: usize) -> Error {
98    Error::new(ErrorKind::InvalidInput, TooLarge { len, max })
99}
100
101/// Serialise `value`, prepend its length, then flush the pair as one frame.
102///
103/// The body is buffered in full first, so an oversize payload fails without
104/// putting a partial frame on the wire.
105pub async fn write_frame<W, T>(w: &mut W, value: &T) -> Result<()>
106where
107    W: AsyncWrite + Unpin,
108    T: Serialize + ?Sized,
109{
110    let body = serde_json::to_vec(value).map_err(Error::other)?;
111    let len = body.len();
112    if len > MAX_FRAME_BYTES {
113        return Err(too_large(len as u64, MAX_FRAME_BYTES));
114    }
115    let prefix = u32::try_from(len).map_err(|_| too_large(len as u64, MAX_FRAME_BYTES))?;
116    let mut buf = Vec::with_capacity(4 + len);
117    buf.extend_from_slice(&prefix.to_be_bytes());
118    buf.extend_from_slice(&body);
119    w.write_all(&buf).await?;
120    w.flush().await
121}
122
123/// Read one frame.
124///
125/// `Ok(None)` means the peer closed the stream while the reader sat at a frame
126/// boundary.
127pub async fn read_frame<R, T>(r: &mut R) -> Result<Option<T>>
128where
129    R: AsyncRead + Unpin,
130    T: DeserializeOwned,
131{
132    let mut header = [0u8; 4];
133    if !fill(r, &mut header, 4).await? {
134        return Ok(None);
135    }
136    let len = u32::from_be_bytes(header) as u64;
137    if len > MAX_FRAME_BYTES as u64 {
138        return Err(too_large(len, MAX_FRAME_BYTES));
139    }
140    let mut body = vec![0u8; len as usize];
141    if !body.is_empty() && !fill(r, &mut body, len).await? {
142        return Err(Error::new(
143            ErrorKind::UnexpectedEof,
144            UnexpectedEof {
145                expected: len,
146                got: 0,
147            },
148        ));
149    }
150    serde_json::from_slice::<T>(&body)
151        .map(Some)
152        .map_err(|e| Error::new(ErrorKind::InvalidData, format!("bad frame json: {e}")))
153}
154
155/// Read until `dst` is full.
156///
157/// `Ok(true)` means the buffer filled. `Ok(false)` means the stream ended with
158/// nothing read at all, which is a clean close at a frame boundary. Ending
159/// partway through yields [`UnexpectedEof`] against `announced`, the total byte
160/// count the caller promised for this read.
161///
162/// `dst` must fit the accumulation ceiling: a destination larger than
163/// [`MAX_PARTIAL_FRAME_BYTES`] is refused with [`TooLarge`] before a byte is read,
164/// which is what caps the memory one decode holds, whatever the peer announces and
165/// however slowly it sends them.
166async fn fill<R>(r: &mut R, dst: &mut [u8], announced: u64) -> Result<bool>
167where
168    R: AsyncRead + Unpin,
169{
170    if dst.len() > MAX_PARTIAL_FRAME_BYTES {
171        return Err(too_large(dst.len() as u64, MAX_PARTIAL_FRAME_BYTES));
172    }
173    let mut got = 0usize;
174    while got < dst.len() {
175        match r.read(&mut dst[got..]).await {
176            Ok(0) => {
177                if got == 0 {
178                    return Ok(false);
179                }
180                return Err(Error::new(
181                    ErrorKind::UnexpectedEof,
182                    UnexpectedEof {
183                        expected: announced,
184                        got: got as u64,
185                    },
186                ));
187            }
188            Ok(n) => got += n,
189            Err(e) if e.kind() == ErrorKind::Interrupted => continue,
190            Err(e) => return Err(e),
191        }
192    }
193    Ok(true)
194}
195
196/// Bytes one [`FrameReader::next`] asks the stream for before it knows a length.
197const READ_CHUNK: usize = 64 * 1024;
198
199/// A frame decoder whose partial progress outlives the future that reads it.
200///
201/// [`read_frame`] keeps the bytes it has read inside its own future, so a
202/// caller that drops the future mid-frame — `tokio::select!` choosing another
203/// branch — loses them and every later frame boundary with them. This reader
204/// keeps them in `buf` instead: [`FrameReader::next`] awaits only
205/// [`AsyncReadExt::read_buf`], which is cancel-safe, so dropping a `next`
206/// future discards nothing, and the next call resumes where the stream is.
207#[derive(Debug, Default)]
208pub struct FrameReader {
209    buf: Vec<u8>,
210}
211
212impl FrameReader {
213    pub fn new() -> Self {
214        Self::default()
215    }
216
217    /// Read one frame. Cancel-safe: see the type's documentation.
218    ///
219    /// `Ok(None)` means the peer closed the stream at a frame boundary.
220    pub async fn next<R, T>(&mut self, r: &mut R) -> Result<Option<T>>
221    where
222        R: AsyncRead + Unpin,
223        T: DeserializeOwned,
224    {
225        loop {
226            if let Some(end) = self.frame_end()? {
227                let decoded = serde_json::from_slice::<T>(&self.buf[4..end]);
228                self.buf.drain(..end);
229                return decoded.map(Some).map_err(|e| {
230                    Error::new(ErrorKind::InvalidData, format!("bad frame json: {e}"))
231                });
232            }
233            self.buf.reserve(READ_CHUNK);
234            if r.read_buf(&mut self.buf).await? == 0 {
235                if self.buf.is_empty() {
236                    return Ok(None);
237                }
238                // Same accounting as `read_frame`: a short header counts against
239                // its 4 bytes, a short body against the announced length.
240                let got = self.buf.len() as u64;
241                let eof = match self.announced() {
242                    Some(len) => UnexpectedEof {
243                        expected: len,
244                        got: got - 4,
245                    },
246                    None => UnexpectedEof { expected: 4, got },
247                };
248                return Err(Error::new(ErrorKind::UnexpectedEof, eof));
249            }
250        }
251    }
252
253    fn announced(&self) -> Option<u64> {
254        let header: [u8; 4] = self.buf.get(..4)?.try_into().ok()?;
255        Some(u32::from_be_bytes(header) as u64)
256    }
257
258    /// The end offset of a complete frame at the front of the buffer, refusing
259    /// an oversize length before any of its body is read.
260    fn frame_end(&self) -> Result<Option<usize>> {
261        let Some(len) = self.announced() else {
262            return Ok(None);
263        };
264        if len > MAX_FRAME_BYTES as u64 {
265            return Err(too_large(len, MAX_FRAME_BYTES));
266        }
267        let end = 4 + len as usize;
268        Ok((self.buf.len() >= end).then_some(end))
269    }
270}
271
272/// `true` when `err` reports a frame over [`MAX_FRAME_BYTES`].
273pub fn is_too_large(err: &Error) -> bool {
274    err.get_ref()
275        .and_then(|r| r.downcast_ref::<TooLarge>())
276        .is_some()
277}
278
279/// `true` when `err` reports undecodable frame content.
280pub fn is_bad_frame(err: &Error) -> bool {
281    err.kind() == ErrorKind::InvalidData
282}