bitbroker 0.1.0

A language agnostic message broker designed for real-time communication.
Documentation
mod listener;
mod setup;

use std::collections::HashMap;

use tracing::info;

use setup::BrokerSetup;
use listener::ConnectionListener;

use crate::{proto::*, BitBrokerMessage, Connection, Error, Result};

pub async fn run_broker(num_players: usize, host: &str, port: usize) -> Result<()> {
    let mut broker = Broker::setup(num_players, host, port).await?;
    broker.initialise().await?;
    broker.run().await
}

struct Broker {
    players: HashMap<u8, Connection>,
    game: Connection,
}

impl Broker {
    async fn setup(num_players: usize, host: &str, port: usize) -> Result<Self> {
        Self::try_from(BrokerSetup::new(num_players).run(host, port).await?)
    }

    async fn initialise(&mut self) -> Result<()> {
        let mut batched_info = InitialMessageBatchProto::default();

        for (player, connection) in self.players.iter_mut() {
            let message: BitBrokerMessage = connection.read_message().await?;
            let mut info: InitialMessageProto = message.body_proto()?;
            info.bot_index = *player as i32;
            batched_info.initial_messages.push(info)
        }

        let message = BitBrokerMessage::new_proto(0x00, 0x11, batched_info);
        self.game.send_message(message).await
    }

    async fn run(&mut self) -> Result<()> {
        info!("Running the game");
        loop {
            self.game
                .send_message(BitBrokerMessage::trigger_message())
                .await?;
            let response: BitBrokerMessage = self.game.read_message().await?;
            match response.codes() {
                (0x00, 0xFE) => {
                    self.handle_game_end().await?;
                    return Ok(());
                }
                (0x01, player) => self.handler_player_action(player, response.body()).await?,
                _ => return Err(Error::from(response)),
            }
        }
    }

    async fn handle_game_end(&mut self) -> Result<()> {
        info!("Handling end of game");
        self.game
            .send_message(BitBrokerMessage::termination_message())
            .await?;

        for (_, connection) in self.players.iter_mut() {
            connection
                .send_message(BitBrokerMessage::termination_message())
                .await?
        }
        Ok(())
    }

    async fn handler_player_action(&mut self, player: &u8, body: impl Into<Vec<u8>>) -> Result<()> {
        let connection = self
            .players
            .get_mut(player)
            .ok_or(Error::BrokerError("Missing player"))?;
        connection
            .send_message(BitBrokerMessage::new(0x02, 0x05, body.into()))
            .await?;
        let response: BitBrokerMessage = connection.read_message().await?;
        self.game.send_message(response).await
    }
}

impl TryFrom<BrokerSetup> for Broker {
    type Error = Error;

    fn try_from(value: BrokerSetup) -> std::result::Result<Self, Self::Error> {
        Ok(Self {
            players: value.players,
            game: value.game.ok_or(Error::BrokerError("Game not connected"))?,
        })
    }
}