diff --git a/src/connections.rs b/src/connections.rs index e57227b..29bf43f 100644 --- a/src/connections.rs +++ b/src/connections.rs @@ -9,11 +9,12 @@ use std::{ use thiserror::Error; use crate::{ - PeerMessage, + PeerConfig, PeerMessage, listener::{DATAGRAM_SIZE, Datagram, ListenerError}, package::Package, }; +/// Represents an error during sending a package for any number of reason #[derive(Debug, Error)] pub enum SendError { #[error("Serialization Error: {0}")] @@ -26,6 +27,7 @@ pub enum SendError { SendError(std::io::Error), } +/// Represents a general error which can occur specifically to a connection #[derive(Debug, Error)] pub enum ConnectionError { #[error(transparent)] @@ -36,21 +38,27 @@ pub enum ConnectionError { TimedOut, } -pub struct ConnectionManager { +/// Manages connections for the Peer +pub(crate) struct ConnectionManager { + /// Wraps UdpSocket with some helper methods pub udp: UdpWrapper, connections: HashMap, closing_since: Option, + config: PeerConfig, } impl ConnectionManager { - pub fn new(socket: Arc) -> ConnectionManager { + /// Creates a new ConnectionManager with the given socket and config + pub fn new(socket: Arc, config: PeerConfig) -> ConnectionManager { ConnectionManager { udp: UdpWrapper::new(socket), connections: HashMap::new(), closing_since: None, + config, } } + /// Start closing all outbound connections. pub fn close(&mut self) { self.closing_since = Some(Instant::now()); for (addr, _) in self.connections.clone() { @@ -58,11 +66,15 @@ impl ConnectionManager { } } + /// Attempt to connect to a remote peer pub fn connect_to(&mut self, addr: SocketAddr) { self.connections .insert(addr, Connection::from(addr, ConnectionState::Connecting)); } + /// Send possible "pings". This also does some other busywork that is + /// necessary to be run as often as possible, such as closing the udp socket + /// itself, and "hello" or "close" messages. pub fn send_pings(&mut self) { let now = Instant::now(); @@ -71,7 +83,7 @@ impl ConnectionManager { { if self.connections.len() == 0 { self.udp.close(); - } else if now - closing_since > Duration::from_millis(1000) { + } else if now - closing_since > self.config.disconnect_timeout * 2 { self.udp.close(); } } @@ -79,7 +91,7 @@ impl ConnectionManager { if let Some(_) = &self.udp.socket { for (addr, conn) in &mut self.connections { let duration = now - conn.last_sent_ping; - if duration > Duration::from_millis(100) { + if duration > self.config.ping_interval { match conn.state { ConnectionState::ReceivingConnection | ConnectionState::Connecting => { match self.udp.send_to(*addr, Package::Hello) { @@ -124,6 +136,7 @@ impl ConnectionManager { } } + /// Close the connection to another peer at the given address pub fn close_connection(&mut self, addr: &SocketAddr) { if let Some(conn) = self.connections.get_mut(addr) { if conn.state != ConnectionState::Closing { @@ -133,6 +146,7 @@ impl ConnectionManager { } } + /// Close the connection to another peer, but with a specified error. pub fn error_connection(&mut self, addr: &SocketAddr, error: ConnectionError) { if let Some(conn) = self.connections.get_mut(addr) { if conn.state != ConnectionState::Error { @@ -143,6 +157,7 @@ impl ConnectionManager { } } + /// Handle Package from a remote peer #[must_use] pub fn handle_package(&mut self, package: Package, addr: &SocketAddr) -> Vec { let mut messages = Vec::new(); @@ -220,15 +235,18 @@ impl ConnectionManager { messages } + /// Close any connections that have not sent a ping for at least "timeout"s + /// duration. pub fn close_silent_connections(&mut self) { let now = Instant::now(); for (addr, conn) in self.connections.clone() { - if (now - conn.last_recv_ping) > Duration::from_millis(2000) { + if (now - conn.last_recv_ping) > self.config.timeout { self.error_connection(&addr, ConnectionError::TimedOut); } } } + /// Remove any connections that have closed or errored. pub fn purge_old_connections(&mut self) -> Vec { let mut messages = Vec::new(); @@ -237,7 +255,7 @@ impl ConnectionManager { self.connections.retain(|_, conn| { let retain = match conn.state { ConnectionState::ReceivingClosing => { - if (now - conn.last_recv_close) > Duration::from_millis(500) { + if (now - conn.last_recv_close) > self.config.disconnect_timeout { messages.push(PeerMessage::Disconnected(conn.clone(), None)); false } else { @@ -257,7 +275,8 @@ impl ConnectionManager { } } -pub struct UdpWrapper { +/// Wrap UdpSocket with some helper methods +pub(crate) struct UdpWrapper { socket: Option>, } @@ -268,14 +287,17 @@ impl UdpWrapper { } } + /// Close the method, that is to say, drop it. pub fn close(&mut self) { self.socket = None; } + /// Is the socket closed pub fn is_closed(&self) -> bool { self.socket.is_none() } + /// Send a Package to a remote peer at addr fn send_to(&self, addr: SocketAddr, package: Package) -> Result<(), SendError> { if let Some(socket) = &self.socket { let mut buf = Vec::new(); @@ -295,8 +317,12 @@ impl UdpWrapper { } } +/// Represents a single connection from the peer to another peer pub struct Connection { + /// Address of the connection pub address: SocketAddr, + /// "ping", as in the time between the last sent ping and the last received + /// pong after it. pub ping: Duration, last_sent_ping: Instant, @@ -304,6 +330,7 @@ pub struct Connection { last_recv_close: Instant, closing_since: Instant, + /// What state is the connection currently in? pub state: ConnectionState, error: Option, } @@ -338,15 +365,24 @@ impl Connection { } } +/// State of the connection #[derive(PartialEq, Eq, PartialOrd, Ord, Clone, Copy)] pub enum ConnectionState { + /// Connection is currently attempting to connect to another peer Connecting, + /// Connection is attepting to connect to another peer, and has received at + /// least one "hello" in response. ConnectingNearlyReady, + /// Connection is receiving a connection from another peer ReceivingConnection, + /// Connection has been properly established Connected, + /// Connection is being closed Closing, + /// Connection is being closed by the other peer ReceivingClosing, + /// Connection has ended up in an erronous state Error, } diff --git a/src/lib.rs b/src/lib.rs index 0a97c02..d378a66 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -7,6 +7,7 @@ use std::{ mpsc::{Receiver, channel}, }, thread::spawn, + time::Duration, }; use thiserror::*; @@ -20,6 +21,7 @@ pub mod connections; pub(crate) mod listener; pub(crate) mod package; +/// Error for the Peer #[derive(Error, Debug)] pub enum PeerError { #[error("Error binding to port: {0}")] @@ -28,12 +30,61 @@ pub enum PeerError { ListenerError(std::io::Error), } +/// Possible messages from occurring events from the peer pub enum PeerMessage { + /// A new connection has been connected NewConnection(Connection), + /// An existing connection has disconnected, with an optional error Disconnected(Connection, Option), + /// The Peer has been closed Closed, } +/// Optional configuration available for peers, mainly configuration of +/// different intervals +#[derive(Debug, Clone)] +pub struct PeerConfig { + ping_interval: Duration, + timeout: Duration, + disconnect_timeout: Duration, +} + +impl Default for PeerConfig { + fn default() -> Self { + Self { + ping_interval: Duration::from_millis(100), + timeout: Duration::from_millis(2000), + disconnect_timeout: Duration::from_millis(500), + } + } +} + +impl PeerConfig { + /// Set the desired interval between "ping"s + pub fn with_ping_interval(self, interval: Duration) -> PeerConfig { + PeerConfig { + ping_interval: interval, + ..self + } + } + + /// Set the length of timeout; that is to say the duration of time which is + /// acceptable to occur between pings before a Timeout-error occurs + pub fn with_timeout(self, timeout: Duration) -> PeerConfig { + PeerConfig { timeout, ..self } + } + + /// Set the length of time that is awaited since the last "Closed" message + /// during disconnection before disconnection actually occurs. Should be + /// significantly higher than ping interval. + pub fn with_disconnect(self, timeout: Duration) -> PeerConfig { + PeerConfig { + disconnect_timeout: timeout, + ..self + } + } +} + pub struct Peer { connection_mgr: ConnectionManager, closed: Arc, @@ -44,7 +95,7 @@ pub struct 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) -> Result { + pub fn listen(port: Option, config: PeerConfig) -> Result { let socket = Arc::new( UdpSocket::bind(SocketAddr::from(([0, 0, 0, 0], port.unwrap_or(0)))) .map_err(|e| PeerError::BindError(e))?, @@ -66,17 +117,21 @@ impl Peer { }); Ok(Peer { - connection_mgr: ConnectionManager::new(socket), + connection_mgr: ConnectionManager::new(socket, config.clone()), closed, receiver, messages: VecDeque::new(), }) } + /// Connect the peer to the given address. pub fn connect_to(&mut self, addr: SocketAddr) { self.connection_mgr.connect_to(addr); } + /// Send pings and poll for events from the listener thread. Returns + /// possible events, such as connections and messages. Should be called as + /// often as possible. #[must_use] pub fn poll(&mut self) -> Result, PeerError> { if self.connection_mgr.udp.is_closed() { @@ -125,6 +180,7 @@ impl Peer { Ok(self.messages.pop_front()) } + /// Close the Peer, disconnecting all connections. pub fn close(&mut self) { self.connection_mgr.close(); } diff --git a/src/listener.rs b/src/listener.rs index 8f73023..2343d71 100644 --- a/src/listener.rs +++ b/src/listener.rs @@ -11,6 +11,7 @@ use crate::package::Package; pub const DATAGRAM_SIZE: usize = 102400; pub type Datagram = [u8; DATAGRAM_SIZE]; +/// Listener thread for Peer's UdpSocket pub struct Listener { socket: Arc, sender: Sender, @@ -21,6 +22,7 @@ impl Listener { Listener { socket, sender } } + /// Read from the socket, parse packages and emit messages. pub fn poll(&mut self) { let mut buffer: Datagram = [0; _]; match self.socket.recv_from(&mut buffer) {