From 73e2224c70f074b44d238b18462fadb6a37a6382 Mon Sep 17 00:00:00 2001 From: Sofia Date: Thu, 16 Jul 2026 04:22:03 +0300 Subject: [PATCH] Add Messages and Message --- src/connections.rs | 49 ++++++++++++++++++++++++++++++---------------- src/lib.rs | 21 ++++++++++---------- src/listener.rs | 24 +++++++++++++++-------- src/package.rs | 16 ++++++++++++++- 4 files changed, 74 insertions(+), 36 deletions(-) diff --git a/src/connections.rs b/src/connections.rs index 29bf43f..d8ca1a5 100644 --- a/src/connections.rs +++ b/src/connections.rs @@ -6,12 +6,13 @@ use std::{ time::{Duration, Instant}, }; +use serde::{Serialize, de::DeserializeOwned}; use thiserror::Error; use crate::{ PeerConfig, PeerMessage, - listener::{DATAGRAM_SIZE, Datagram, ListenerError}, - package::Package, + listener::{DATAGRAM_SIZE, ListenerError}, + package::{Message, Package}, }; /// Represents an error during sending a package for any number of reason @@ -39,17 +40,18 @@ pub enum ConnectionError { } /// Manages connections for the Peer -pub(crate) struct ConnectionManager { +pub(crate) struct ConnectionManager +{ /// Wraps UdpSocket with some helper methods pub udp: UdpWrapper, - connections: HashMap, + connections: HashMap>, closing_since: Option, config: PeerConfig, } -impl ConnectionManager { +impl ConnectionManager { /// Creates a new ConnectionManager with the given socket and config - pub fn new(socket: Arc, config: PeerConfig) -> ConnectionManager { + pub fn new(socket: Arc, config: PeerConfig) -> ConnectionManager { ConnectionManager { udp: UdpWrapper::new(socket), connections: HashMap::new(), @@ -94,7 +96,7 @@ impl ConnectionManager { if duration > self.config.ping_interval { match conn.state { ConnectionState::ReceivingConnection | ConnectionState::Connecting => { - match self.udp.send_to(*addr, Package::Hello) { + match self.udp.send_to(*addr, Package::::Hello) { Ok(_) => { conn.last_sent_ping = now; } @@ -106,7 +108,7 @@ impl ConnectionManager { } } ConnectionState::Connected | ConnectionState::ConnectingNearlyReady => { - match self.udp.send_to(*addr, Package::Ping) { + match self.udp.send_to(*addr, Package::::Ping) { Ok(_) => { conn.last_sent_ping = now; } @@ -118,7 +120,7 @@ impl ConnectionManager { } } ConnectionState::Closing | ConnectionState::ReceivingClosing => { - match self.udp.send_to(*addr, Package::Close) { + match self.udp.send_to(*addr, Package::::Close) { Ok(_) => { conn.last_sent_ping = now; } @@ -159,7 +161,11 @@ impl ConnectionManager { /// Handle Package from a remote peer #[must_use] - pub fn handle_package(&mut self, package: Package, addr: &SocketAddr) -> Vec { + pub fn handle_package( + &mut self, + package: Package, + addr: &SocketAddr, + ) -> Vec> { let mut messages = Vec::new(); match package { @@ -186,7 +192,7 @@ impl ConnectionManager { if conn.state == ConnectionState::Connected { conn.last_recv_ping = Instant::now(); - match self.udp.send_to(*addr, Package::Pong) { + match self.udp.send_to(*addr, Package::::Pong) { Ok(_) => {} Err(err) => { conn.closing_since = Instant::now(); @@ -230,6 +236,7 @@ impl ConnectionManager { self.connections.remove(addr); } } + Package::Messages(messages) => todo!(), } messages @@ -247,7 +254,7 @@ impl ConnectionManager { } /// Remove any connections that have closed or errored. - pub fn purge_old_connections(&mut self) -> Vec { + pub fn purge_old_connections(&mut self) -> Vec> { let mut messages = Vec::new(); let now = Instant::now(); @@ -298,7 +305,11 @@ impl UdpWrapper { } /// Send a Package to a remote peer at addr - fn send_to(&self, addr: SocketAddr, package: Package) -> Result<(), SendError> { + fn send_to( + &self, + addr: SocketAddr, + package: Package, + ) -> Result<(), SendError> { if let Some(socket) = &self.socket { let mut buf = Vec::new(); let cursor = Cursor::new(&mut buf); @@ -318,7 +329,7 @@ impl UdpWrapper { } /// Represents a single connection from the peer to another peer -pub struct Connection { +pub struct Connection { /// Address of the connection pub address: SocketAddr, /// "ping", as in the time between the last sent ping and the last received @@ -330,12 +341,14 @@ pub struct Connection { last_recv_close: Instant, closing_since: Instant, + reliable_queue: Vec>, + /// What state is the connection currently in? pub state: ConnectionState, error: Option, } -impl Clone for Connection { +impl Clone for Connection { fn clone(&self) -> Self { Self { address: self.address.clone(), @@ -344,14 +357,15 @@ impl Clone for Connection { last_recv_ping: self.last_recv_ping.clone(), last_recv_close: self.last_recv_close.clone(), closing_since: self.closing_since.clone(), + reliable_queue: self.reliable_queue.clone(), state: self.state.clone(), error: None, } } } -impl Connection { - pub fn from(address: SocketAddr, state: ConnectionState) -> Connection { +impl Connection { + pub fn from(address: SocketAddr, state: ConnectionState) -> Connection { Connection { address, ping: Duration::default(), @@ -359,6 +373,7 @@ impl Connection { last_recv_ping: Instant::now(), last_recv_close: Instant::now(), closing_since: Instant::now(), + reliable_queue: Vec::new(), state, error: None, } diff --git a/src/lib.rs b/src/lib.rs index d378a66..9501c55 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -10,6 +10,7 @@ use std::{ time::Duration, }; +use serde::{Serialize, de::DeserializeOwned}; use thiserror::*; use crate::{ @@ -31,11 +32,11 @@ pub enum PeerError { } /// Possible messages from occurring events from the peer -pub enum PeerMessage { +pub enum PeerMessage { /// A new connection has been connected - NewConnection(Connection), + NewConnection(Connection), /// An existing connection has disconnected, with an optional error - Disconnected(Connection, Option), + Disconnected(Connection, Option), /// The Peer has been closed Closed, } @@ -85,17 +86,17 @@ impl PeerConfig { } } -pub struct Peer { - connection_mgr: ConnectionManager, +pub struct Peer { + connection_mgr: ConnectionManager, closed: Arc, - receiver: Receiver, - messages: VecDeque, + receiver: Receiver>, + messages: VecDeque>, } -impl Peer { +impl Peer { /// Bind given port to listen to incoming messages. Creates a new Peer that /// can establish new connections. - pub fn listen(port: Option, config: PeerConfig) -> Result { + pub fn listen(port: Option, config: PeerConfig) -> Result, PeerError> { let socket = Arc::new( UdpSocket::bind(SocketAddr::from(([0, 0, 0, 0], port.unwrap_or(0)))) .map_err(|e| PeerError::BindError(e))?, @@ -133,7 +134,7 @@ impl Peer { /// possible events, such as connections and messages. Should be called as /// often as possible. #[must_use] - pub fn poll(&mut self) -> Result, PeerError> { + pub fn poll(&mut self) -> Result>, PeerError> { if self.connection_mgr.udp.is_closed() { if self.closed.load(Ordering::Acquire) { self.messages.push_back(PeerMessage::Closed); diff --git a/src/listener.rs b/src/listener.rs index 2343d71..57b1d27 100644 --- a/src/listener.rs +++ b/src/listener.rs @@ -1,9 +1,11 @@ use std::{ io::Cursor, + marker::PhantomData, net::{SocketAddr, UdpSocket}, sync::{Arc, mpsc::Sender}, }; +use serde::{Deserialize, Serialize, de::DeserializeOwned}; use thiserror::Error; use crate::package::Package; @@ -12,14 +14,19 @@ pub const DATAGRAM_SIZE: usize = 102400; pub type Datagram = [u8; DATAGRAM_SIZE]; /// Listener thread for Peer's UdpSocket -pub struct Listener { +pub struct Listener { socket: Arc, - sender: Sender, + sender: Sender>, + pd: PhantomData, } -impl Listener { - pub fn new(socket: Arc, sender: Sender) -> Listener { - Listener { socket, sender } +impl<'a, T: Clone + Serialize + DeserializeOwned> Listener { + pub fn new(socket: Arc, sender: Sender>) -> Listener { + Listener { + socket, + sender, + pd: PhantomData::default(), + } } /// Read from the socket, parse packages and emit messages. @@ -29,7 +36,8 @@ impl Listener { Ok((num_bytes, from_addr)) => { let mut ciborium_buf: Datagram = [0; _]; let bytes = Cursor::new(&mut buffer[..num_bytes]); - let res = ciborium::from_reader_with_buffer::(bytes, &mut ciborium_buf); + let res = + ciborium::from_reader_with_buffer::, _>(bytes, &mut ciborium_buf); match res { Ok(pkg) => self.sender.send(ListenerMessage::Package(pkg, from_addr)), Err(_) => self.sender.send(ListenerMessage::PackageError( @@ -52,8 +60,8 @@ pub enum ListenerError { ParseError, } -pub enum ListenerMessage { - Package(Package, SocketAddr), +pub enum ListenerMessage { + Package(Package, SocketAddr), PackageError(ListenerError, SocketAddr), Error(std::io::Error), } diff --git a/src/package.rs b/src/package.rs index 258e82d..a908af9 100644 --- a/src/package.rs +++ b/src/package.rs @@ -1,9 +1,23 @@ use serde::{Deserialize, Serialize}; #[derive(Debug, Serialize, Deserialize)] -pub enum Package { +pub enum Package { Hello, Ping, Pong, Close, + Messages(Messages), +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct Messages { + pub ack: u64, + pub reliable: Vec>, + pub unreliable: Vec>, +} + +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct Message { + pub message_id: u64, + pub message: T, }