ibrahim-tcp 1.0.0

High-performance lightweight TCP protocol with reliability, compression, and fault tolerance
Documentation
use crate::{AckManager, Packet};
use std::{io, net::SocketAddr, sync::Arc};
use tokio::{
    io::{AsyncReadExt, AsyncWriteExt},
    net::{TcpListener, TcpStream},
    sync::Mutex,
    time::{sleep, Duration},
};

#[derive(Clone)]
pub struct Server {
    peers: Arc<Mutex<Vec<(SocketAddr, TcpStream)>>>,
}

impl Server {
    pub async fn new(addr: &str) -> io::Result<Self> {
        let listener = TcpListener::bind(addr).await?;
        let peers = Arc::new(Mutex::new(Vec::new()));
        let peers_accept = peers.clone();

        tokio::spawn(async move {
            loop {
                if let Ok((stream, addr)) = listener.accept().await {
                    println!("[+] {} connected", addr);
                    peers_accept.lock().await.push((addr, stream));
                }
            }
        });

        Ok(Self { peers })
    }

    pub async fn broadcast(&self, pkt: &Packet) {
        let data = pkt.serialize();
        let mut peers = self.peers.lock().await;
        let mut i = 0;
        while i < peers.len() {
            let (_, stream) = &mut peers[i];
            if stream.write_all(&data).await.is_err() {
                peers.remove(i);
            } else {
                i += 1;
            }
        }
    }
}

#[derive(Clone)]
pub struct Client {
    addr: String,
    ack: AckManager,
}

impl Client {
    pub fn new(addr: &str) -> Self {
        Self {
            addr: addr.into(),
            ack: AckManager::new(),
        }
    }

    pub async fn run(&self, name: &str) -> io::Result<()> {
        loop {
            match TcpStream::connect(&self.addr).await {
                Ok(mut stream) => {
                    println!("[{}] connected", name);
                    let mut id = 1u64;
                    let mut buf = vec![0u8; 8192];

                    loop {
                        let payload = format!(
                            "{}|health={} bullets={} pos=({}, {})",
                            name,
                            100 - (id % 10),
                            30 - (id % 5),
                            id % 100,
                            (id * 3) % 100
                        );

                        let pkt = Packet::new(id, payload.clone(), true);
                        if let Err(e) = stream.write_all(&pkt.serialize()).await {
                            eprintln!("[{}] send error: {}", name, e);
                            break;
                        }
                        self.ack.track(&pkt);

                        if let Ok(n) = stream.read(&mut buf).await {
                            if n > 0 {
                                if let Ok(Some(resp)) = Packet::deserialize(
                                    &mut std::io::Cursor::new(buf[..n].to_vec()),
                                    true,
                                ) {
                                    if resp.ack {
                                        self.ack.confirm(resp.id);
                                    } else {
                                        println!("[{}] recv: {}", name, resp.data);
                                    }
                                }
                            }
                        }

                        for r in self.ack.retries() {
                            let _ = stream.write_all(&r.serialize()).await;
                        }

                        id += 1;
                        sleep(Duration::from_secs(1)).await;
                    }
                }
                Err(e) => {
                    eprintln!("[{}] reconnecting: {}", name, e);
                    sleep(Duration::from_secs(2)).await;
                }
            }
        }
    }
}