bitbroker 0.1.0

A language agnostic message broker designed for real-time communication.
Documentation
use base64::prelude::*;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::TcpStream;

use crate::Result;

pub struct Connection {
    buffer: BufReader<TcpStream>,
}

impl Connection {
    pub async fn new(host: &str, port: usize) -> Result<Self> {
        let stream = TcpStream::connect(format!("{host}:{port}")).await?;
        Ok(Self::from(stream))
    }

    pub async fn read_message<T>(&mut self) -> Result<T>
    where
        T: From<Vec<u8>>,
    {
        let mut value = Vec::new();
        self.buffer.read_until(b'\n', &mut value).await?;
        value.pop();
        Ok(T::from(BASE64_STANDARD.decode(value)?))
    }

    pub async fn send_message<T>(&mut self, message: T) -> Result<()>
    where
        T: AsRef<[u8]>,
    {
        let mut encoded = BASE64_STANDARD.encode(message);
        encoded.push('\n');
        self.buffer.write_all(encoded.as_bytes()).await?;
        self.buffer.flush().await?;
        Ok(())
    }
}

impl From<TcpStream> for Connection {
    fn from(stream: TcpStream) -> Self {
        let buffer = BufReader::new(stream);
        Self { buffer }
    }
}