use scv_protocol::{ClientMessage, Frame, FrameDecoder, encode_frame};
use std::io;
use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite, AsyncWriteExt};
pub async fn read_frame<R>(reader: &mut R, decoder: &mut FrameDecoder) -> io::Result<Frame>
where
R: AsyncBufRead + Unpin,
{
loop {
let available = reader.fill_buf().await?;
if available.is_empty() {
return Ok(decoder.finish());
}
let step = decoder.feed(available);
reader.consume(step.consumed);
if let Some(frame) = step.frame {
return Ok(frame);
}
}
}
pub async fn write_message<W>(writer: &mut W, message: &ClientMessage) -> io::Result<()>
where
W: AsyncWrite + Unpin,
{
let bytes = encode_frame(message).map_err(io::Error::other)?;
writer.write_all(&bytes).await?;
writer.flush().await
}
#[derive(Debug)]
pub struct Connection<R, W> {
reader: R,
writer: W,
decoder: FrameDecoder,
}
impl<R, W> Connection<R, W>
where
R: AsyncBufRead + Unpin,
W: AsyncWrite + Unpin,
{
pub fn new(reader: R, writer: W, decoder: FrameDecoder) -> Self {
Self {
reader,
writer,
decoder,
}
}
pub async fn send(&mut self, message: &ClientMessage) -> io::Result<()> {
write_message(&mut self.writer, message).await
}
pub async fn read(&mut self) -> io::Result<Frame> {
read_frame(&mut self.reader, &mut self.decoder).await
}
#[cfg(test)]
pub(crate) fn decoder_mut(&mut self) -> &mut FrameDecoder {
&mut self.decoder
}
}
#[cfg(test)]
mod tests;