Iron out some bugs
This commit is contained in:
parent
d8781b55d0
commit
f3429ff4e8
@ -37,7 +37,7 @@ pub enum ConnectionError {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub struct ConnectionManager {
|
pub struct ConnectionManager {
|
||||||
udp: UdpWrapper,
|
pub udp: UdpWrapper,
|
||||||
connections: HashMap<SocketAddr, Connection>,
|
connections: HashMap<SocketAddr, Connection>,
|
||||||
closing_since: Option<Instant>,
|
closing_since: Option<Instant>,
|
||||||
}
|
}
|
||||||
@ -68,8 +68,10 @@ impl ConnectionManager {
|
|||||||
|
|
||||||
if let Some(closing_since) = self.closing_since {
|
if let Some(closing_since) = self.closing_since {
|
||||||
if self.connections.len() == 0 {
|
if self.connections.len() == 0 {
|
||||||
|
println!("UDP closed");
|
||||||
self.udp.close();
|
self.udp.close();
|
||||||
} else if now - closing_since > Duration::from_millis(1000) {
|
} else if now - closing_since > Duration::from_millis(1000) {
|
||||||
|
println!("UDP closed");
|
||||||
self.udp.close();
|
self.udp.close();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -103,7 +105,8 @@ impl ConnectionManager {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
ConnectionState::Closing => match self.udp.send_to(*addr, Package::Close) {
|
ConnectionState::Closing | ConnectionState::ReceivingClosing => {
|
||||||
|
match self.udp.send_to(*addr, Package::Close) {
|
||||||
Ok(_) => {
|
Ok(_) => {
|
||||||
conn.last_sent_ping = now;
|
conn.last_sent_ping = now;
|
||||||
}
|
}
|
||||||
@ -112,7 +115,8 @@ impl ConnectionManager {
|
|||||||
conn.state = ConnectionState::Error;
|
conn.state = ConnectionState::Error;
|
||||||
conn.error = Some(ConnectionError::SendError(err));
|
conn.error = Some(ConnectionError::SendError(err));
|
||||||
}
|
}
|
||||||
},
|
}
|
||||||
|
}
|
||||||
_ => {}
|
_ => {}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -188,13 +192,18 @@ impl ConnectionManager {
|
|||||||
Package::Close => {
|
Package::Close => {
|
||||||
let remove = if let Some(conn) = self.connections.get_mut(addr) {
|
let remove = if let Some(conn) = self.connections.get_mut(addr) {
|
||||||
match conn.state {
|
match conn.state {
|
||||||
ConnectionState::Closing => true,
|
ConnectionState::Closing => {
|
||||||
|
messages
|
||||||
|
.push(PeerMessage::Disconnected(conn.clone(), conn.error.take()));
|
||||||
|
true
|
||||||
|
}
|
||||||
ConnectionState::Error => false,
|
ConnectionState::Error => false,
|
||||||
ConnectionState::ReceivingClosing => {
|
ConnectionState::ReceivingClosing => {
|
||||||
conn.last_recv_close = Instant::now();
|
conn.last_recv_close = Instant::now();
|
||||||
false
|
false
|
||||||
}
|
}
|
||||||
_ => {
|
_ => {
|
||||||
|
conn.last_recv_close = Instant::now();
|
||||||
conn.state = ConnectionState::ReceivingClosing;
|
conn.state = ConnectionState::ReceivingClosing;
|
||||||
false
|
false
|
||||||
}
|
}
|
||||||
@ -224,21 +233,27 @@ impl ConnectionManager {
|
|||||||
let mut messages = Vec::new();
|
let mut messages = Vec::new();
|
||||||
|
|
||||||
let now = Instant::now();
|
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 => {
|
ConnectionState::ReceivingClosing => {
|
||||||
if (now - conn.last_recv_close) > Duration::from_millis(500) {
|
if (now - conn.last_recv_close) > Duration::from_millis(500) {
|
||||||
self.connections.remove(&addr);
|
messages.push(PeerMessage::Disconnected(conn.clone(), None));
|
||||||
messages.push(PeerMessage::Disconnected(conn, None));
|
false
|
||||||
|
} else {
|
||||||
|
true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
ConnectionState::Error => {
|
ConnectionState::Error => {
|
||||||
self.connections.remove(&addr);
|
messages.push(PeerMessage::Disconnected(conn.clone(), conn.error.take()));
|
||||||
messages.push(PeerMessage::Disconnected(conn.clone(), conn.error));
|
false
|
||||||
}
|
|
||||||
_ => {}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
_ => true,
|
||||||
|
};
|
||||||
|
retain
|
||||||
|
});
|
||||||
|
|
||||||
|
for (addr, conn) in &self.connections {}
|
||||||
|
|
||||||
messages
|
messages
|
||||||
}
|
}
|
||||||
@ -259,6 +274,10 @@ impl UdpWrapper {
|
|||||||
self.socket = None;
|
self.socket = None;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn is_closed(&self) -> bool {
|
||||||
|
self.socket.is_none()
|
||||||
|
}
|
||||||
|
|
||||||
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 {
|
if let Some(socket) = &self.socket {
|
||||||
let mut buf = Vec::new();
|
let mut buf = Vec::new();
|
||||||
@ -268,10 +287,14 @@ impl UdpWrapper {
|
|||||||
if buf.len() > DATAGRAM_SIZE {
|
if buf.len() > DATAGRAM_SIZE {
|
||||||
return Err(SendError::DatagramTooLarge(buf.len()));
|
return Err(SendError::DatagramTooLarge(buf.len()));
|
||||||
}
|
}
|
||||||
socket
|
let res = socket
|
||||||
.send_to(&buf, addr)
|
.send_to(&buf, addr)
|
||||||
.map_err(|e| SendError::SendError(e))
|
.map_err(|e| SendError::SendError(e))
|
||||||
.map(|_| ())
|
.map(|_| ());
|
||||||
|
if let Ok(_) = res {
|
||||||
|
println!("Sent: {:?}", package);
|
||||||
|
}
|
||||||
|
res
|
||||||
} else {
|
} else {
|
||||||
Err(SendError::SocketDisconnected)
|
Err(SendError::SocketDisconnected)
|
||||||
}
|
}
|
||||||
|
|||||||
17
src/lib.rs
17
src/lib.rs
@ -35,7 +35,7 @@ pub enum PeerMessage {
|
|||||||
|
|
||||||
pub struct Peer {
|
pub struct Peer {
|
||||||
connection_mgr: ConnectionManager,
|
connection_mgr: ConnectionManager,
|
||||||
closing: Arc<AtomicBool>,
|
closed: Arc<AtomicBool>,
|
||||||
receiver: Receiver<ListenerMessage>,
|
receiver: Receiver<ListenerMessage>,
|
||||||
messages: VecDeque<PeerMessage>,
|
messages: VecDeque<PeerMessage>,
|
||||||
}
|
}
|
||||||
@ -48,17 +48,17 @@ impl Peer {
|
|||||||
UdpSocket::bind(SocketAddr::from(([0, 0, 0, 0], port.unwrap_or(0))))
|
UdpSocket::bind(SocketAddr::from(([0, 0, 0, 0], port.unwrap_or(0))))
|
||||||
.map_err(|e| PeerError::BindError(e))?,
|
.map_err(|e| PeerError::BindError(e))?,
|
||||||
);
|
);
|
||||||
let closing = Arc::new(AtomicBool::new(false));
|
let closed = Arc::new(AtomicBool::new(false));
|
||||||
|
|
||||||
let (sender, receiver) = channel();
|
let (sender, receiver) = channel();
|
||||||
|
|
||||||
spawn({
|
spawn({
|
||||||
let socket = socket.clone();
|
let socket = socket.clone();
|
||||||
let closing = closing.clone();
|
let closed = closed.clone();
|
||||||
move || {
|
move || {
|
||||||
let mut listener = Listener::new(socket, sender);
|
let mut listener = Listener::new(socket, sender);
|
||||||
|
|
||||||
while !closing.load(Ordering::Relaxed) {
|
while !closed.load(Ordering::Relaxed) {
|
||||||
listener.poll();
|
listener.poll();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -66,7 +66,7 @@ impl Peer {
|
|||||||
|
|
||||||
Ok(Peer {
|
Ok(Peer {
|
||||||
connection_mgr: ConnectionManager::new(socket),
|
connection_mgr: ConnectionManager::new(socket),
|
||||||
closing,
|
closed,
|
||||||
receiver,
|
receiver,
|
||||||
messages: VecDeque::new(),
|
messages: VecDeque::new(),
|
||||||
})
|
})
|
||||||
@ -78,6 +78,10 @@ impl Peer {
|
|||||||
|
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn poll(&mut self) -> Result<Option<PeerMessage>, PeerError> {
|
pub fn poll(&mut self) -> Result<Option<PeerMessage>, PeerError> {
|
||||||
|
if self.connection_mgr.udp.is_closed() {
|
||||||
|
self.closed.store(true, Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
|
||||||
self.connection_mgr.send_pings();
|
self.connection_mgr.send_pings();
|
||||||
|
|
||||||
let mut msg;
|
let mut msg;
|
||||||
@ -95,12 +99,14 @@ impl Peer {
|
|||||||
.extend(self.connection_mgr.handle_package(package, &socket_addr));
|
.extend(self.connection_mgr.handle_package(package, &socket_addr));
|
||||||
}
|
}
|
||||||
ListenerMessage::PackageError(listener_error, socket_addr) => {
|
ListenerMessage::PackageError(listener_error, socket_addr) => {
|
||||||
|
println!("Error: {}", listener_error);
|
||||||
self.connection_mgr.error_connection(
|
self.connection_mgr.error_connection(
|
||||||
&socket_addr,
|
&socket_addr,
|
||||||
ConnectionError::ListenerError(listener_error),
|
ConnectionError::ListenerError(listener_error),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
ListenerMessage::Error(error) => {
|
ListenerMessage::Error(error) => {
|
||||||
|
println!("Error {}", error);
|
||||||
self.close();
|
self.close();
|
||||||
return Err(PeerError::ListenerError(error));
|
return Err(PeerError::ListenerError(error));
|
||||||
}
|
}
|
||||||
@ -116,7 +122,6 @@ impl Peer {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn close(&mut self) {
|
pub fn close(&mut self) {
|
||||||
self.closing.store(true, Ordering::Relaxed);
|
|
||||||
self.connection_mgr.close();
|
self.connection_mgr.close();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@ -28,6 +28,7 @@ impl Listener {
|
|||||||
let mut ciborium_buf: Datagram = [0; _];
|
let mut ciborium_buf: Datagram = [0; _];
|
||||||
let bytes = Cursor::new(&mut buffer[..num_bytes]);
|
let bytes = Cursor::new(&mut buffer[..num_bytes]);
|
||||||
let res = ciborium::from_reader_with_buffer::<Package, _>(bytes, &mut ciborium_buf);
|
let res = ciborium::from_reader_with_buffer::<Package, _>(bytes, &mut ciborium_buf);
|
||||||
|
println!("Received: {:?}", res);
|
||||||
match res {
|
match res {
|
||||||
Ok(pkg) => self.sender.send(ListenerMessage::Package(pkg, from_addr)),
|
Ok(pkg) => self.sender.send(ListenerMessage::Package(pkg, from_addr)),
|
||||||
Err(_) => self.sender.send(ListenerMessage::PackageError(
|
Err(_) => self.sender.send(ListenerMessage::PackageError(
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user