Skip to main content

scv_client/
connection.rs

1//! Framed JSON over a byte stream: the reading loop around
2//! [`FrameDecoder`], and a [`Connection`] that pairs it with a writer.
3
4use scv_protocol::{ClientMessage, Frame, FrameDecoder, encode_frame};
5use std::io;
6use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite, AsyncWriteExt};
7
8/// Read the next frame from `reader`.
9///
10/// Cancel-safe: bytes are consumed from `reader` only as `decoder` takes
11/// them, and a partial line stays in `decoder`, so a read abandoned in a
12/// `select!` resumes where it stopped.
13pub async fn read_frame<R>(reader: &mut R, decoder: &mut FrameDecoder) -> io::Result<Frame>
14where
15    R: AsyncBufRead + Unpin,
16{
17    loop {
18        let available = reader.fill_buf().await?;
19        if available.is_empty() {
20            return Ok(decoder.finish());
21        }
22        let step = decoder.feed(available);
23        reader.consume(step.consumed);
24        if let Some(frame) = step.frame {
25            return Ok(frame);
26        }
27    }
28}
29
30/// Write `message` as one frame and flush it.
31pub async fn write_message<W>(writer: &mut W, message: &ClientMessage) -> io::Result<()>
32where
33    W: AsyncWrite + Unpin,
34{
35    let bytes = encode_frame(message).map_err(io::Error::other)?;
36    writer.write_all(&bytes).await?;
37    writer.flush().await
38}
39
40/// A client's side of one protocol connection: a buffered reader, a writer,
41/// and the decoder that bounds what the server may send.
42#[derive(Debug)]
43pub struct Connection<R, W> {
44    reader: R,
45    writer: W,
46    decoder: FrameDecoder,
47}
48
49impl<R, W> Connection<R, W>
50where
51    R: AsyncBufRead + Unpin,
52    W: AsyncWrite + Unpin,
53{
54    /// Wrap a stream's halves; `decoder` sets the frame limit.
55    pub fn new(reader: R, writer: W, decoder: FrameDecoder) -> Self {
56        Self {
57            reader,
58            writer,
59            decoder,
60        }
61    }
62
63    /// Send one message.
64    pub async fn send(&mut self, message: &ClientMessage) -> io::Result<()> {
65        write_message(&mut self.writer, message).await
66    }
67
68    /// The next frame. Cancel-safe, like [`read_frame`].
69    pub async fn read(&mut self) -> io::Result<Frame> {
70        read_frame(&mut self.reader, &mut self.decoder).await
71    }
72
73    /// The decoder, to change its limit once the peer declares one.
74    #[cfg(test)]
75    pub(crate) fn decoder_mut(&mut self) -> &mut FrameDecoder {
76        &mut self.decoder
77    }
78}
79
80#[cfg(test)]
81mod tests;