From f3429ff4e86bb68b1fe42f908aaf561b3a9f3976 Mon Sep 17 00:00:00 2001 From: Sofia Date: Thu, 16 Jul 2026 01:59:19 +0300 Subject: [PATCH] Iron out some bugs --- src/connections.rs | 67 +++++++++++++++++++++++++++++++--------------- src/lib.rs | 17 +++++++----- src/listener.rs | 1 + 3 files changed, 57 insertions(+), 28 deletions(-) diff --git a/src/connections.rs b/src/connections.rs index 94c7748..18fdf93 100644 --- a/src/connections.rs +++ b/src/connections.rs @@ -37,7 +37,7 @@ pub enum ConnectionError { } pub struct ConnectionManager { - udp: UdpWrapper, + pub udp: UdpWrapper, connections: HashMap, closing_since: Option, } @@ -68,8 +68,10 @@ impl ConnectionManager { if let Some(closing_since) = self.closing_since { if self.connections.len() == 0 { + println!("UDP closed"); self.udp.close(); } else if now - closing_since > Duration::from_millis(1000) { + println!("UDP closed"); self.udp.close(); } } @@ -103,16 +105,18 @@ impl ConnectionManager { } } } - ConnectionState::Closing => match self.udp.send_to(*addr, Package::Close) { - Ok(_) => { - conn.last_sent_ping = now; + ConnectionState::Closing | ConnectionState::ReceivingClosing => { + match self.udp.send_to(*addr, Package::Close) { + Ok(_) => { + conn.last_sent_ping = now; + } + Err(err) => { + conn.closing_since = Instant::now(); + conn.state = ConnectionState::Error; + conn.error = Some(ConnectionError::SendError(err)); + } } - Err(err) => { - conn.closing_since = Instant::now(); - conn.state = ConnectionState::Error; - conn.error = Some(ConnectionError::SendError(err)); - } - }, + } _ => {} } } @@ -188,13 +192,18 @@ impl ConnectionManager { Package::Close => { let remove = if let Some(conn) = self.connections.get_mut(addr) { match conn.state { - ConnectionState::Closing => true, + ConnectionState::Closing => { + messages + .push(PeerMessage::Disconnected(conn.clone(), conn.error.take())); + true + } ConnectionState::Error => false, ConnectionState::ReceivingClosing => { conn.last_recv_close = Instant::now(); false } _ => { + conn.last_recv_close = Instant::now(); conn.state = ConnectionState::ReceivingClosing; false } @@ -224,21 +233,27 @@ impl ConnectionManager { let mut messages = Vec::new(); let now = Instant::now(); - for (addr, conn) in self.connections.clone() { - match conn.state { + + self.connections.retain(|addr, conn| { + let retain = match conn.state { ConnectionState::ReceivingClosing => { if (now - conn.last_recv_close) > Duration::from_millis(500) { - self.connections.remove(&addr); - messages.push(PeerMessage::Disconnected(conn, None)); + messages.push(PeerMessage::Disconnected(conn.clone(), None)); + false + } else { + true } } ConnectionState::Error => { - self.connections.remove(&addr); - messages.push(PeerMessage::Disconnected(conn.clone(), conn.error)); + messages.push(PeerMessage::Disconnected(conn.clone(), conn.error.take())); + false } - _ => {} - } - } + _ => true, + }; + retain + }); + + for (addr, conn) in &self.connections {} messages } @@ -259,6 +274,10 @@ impl UdpWrapper { self.socket = None; } + pub fn is_closed(&self) -> bool { + self.socket.is_none() + } + fn send_to(&self, addr: SocketAddr, package: Package) -> Result<(), SendError> { if let Some(socket) = &self.socket { let mut buf = Vec::new(); @@ -268,10 +287,14 @@ impl UdpWrapper { if buf.len() > DATAGRAM_SIZE { return Err(SendError::DatagramTooLarge(buf.len())); } - socket + let res = socket .send_to(&buf, addr) .map_err(|e| SendError::SendError(e)) - .map(|_| ()) + .map(|_| ()); + if let Ok(_) = res { + println!("Sent: {:?}", package); + } + res } else { Err(SendError::SocketDisconnected) } diff --git a/src/lib.rs b/src/lib.rs index 5264177..8b55807 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -35,7 +35,7 @@ pub enum PeerMessage { pub struct Peer { connection_mgr: ConnectionManager, - closing: Arc, + closed: Arc, receiver: Receiver, messages: VecDeque, } @@ -48,17 +48,17 @@ impl Peer { UdpSocket::bind(SocketAddr::from(([0, 0, 0, 0], port.unwrap_or(0)))) .map_err(|e| PeerError::BindError(e))?, ); - let closing = Arc::new(AtomicBool::new(false)); + let closed = Arc::new(AtomicBool::new(false)); let (sender, receiver) = channel(); spawn({ let socket = socket.clone(); - let closing = closing.clone(); + let closed = closed.clone(); move || { let mut listener = Listener::new(socket, sender); - while !closing.load(Ordering::Relaxed) { + while !closed.load(Ordering::Relaxed) { listener.poll(); } } @@ -66,7 +66,7 @@ impl Peer { Ok(Peer { connection_mgr: ConnectionManager::new(socket), - closing, + closed, receiver, messages: VecDeque::new(), }) @@ -78,6 +78,10 @@ impl Peer { #[must_use] pub fn poll(&mut self) -> Result, PeerError> { + if self.connection_mgr.udp.is_closed() { + self.closed.store(true, Ordering::Relaxed); + } + self.connection_mgr.send_pings(); let mut msg; @@ -95,12 +99,14 @@ impl Peer { .extend(self.connection_mgr.handle_package(package, &socket_addr)); } ListenerMessage::PackageError(listener_error, socket_addr) => { + println!("Error: {}", listener_error); self.connection_mgr.error_connection( &socket_addr, ConnectionError::ListenerError(listener_error), ); } ListenerMessage::Error(error) => { + println!("Error {}", error); self.close(); return Err(PeerError::ListenerError(error)); } @@ -116,7 +122,6 @@ impl Peer { } pub fn close(&mut self) { - self.closing.store(true, Ordering::Relaxed); self.connection_mgr.close(); } } diff --git a/src/listener.rs b/src/listener.rs index 8f73023..bea256c 100644 --- a/src/listener.rs +++ b/src/listener.rs @@ -28,6 +28,7 @@ impl Listener { 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); + println!("Received: {:?}", res); match res { Ok(pkg) => self.sender.send(ListenerMessage::Package(pkg, from_addr)), Err(_) => self.sender.send(ListenerMessage::PackageError(