Properly dispose of listener
This commit is contained in:
parent
a947f3a8cb
commit
3d1c07c2a9
12
src/lib.rs
12
src/lib.rs
@ -108,6 +108,13 @@ pub struct Peer<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send
|
|||||||
receiver: Receiver<ListenerMessage<T>>,
|
receiver: Receiver<ListenerMessage<T>>,
|
||||||
messages: VecDeque<PeerMessage<T>>,
|
messages: VecDeque<PeerMessage<T>>,
|
||||||
}
|
}
|
||||||
|
impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> Drop
|
||||||
|
for Peer<T>
|
||||||
|
{
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.closed.store(true, Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> Peer<T> {
|
impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> Peer<T> {
|
||||||
/// Bind given port to listen to incoming messages. Creates a new Peer that
|
/// Bind given port to listen to incoming messages. Creates a new Peer that
|
||||||
@ -124,6 +131,11 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
spawn({
|
spawn({
|
||||||
let socket = socket.clone();
|
let socket = socket.clone();
|
||||||
let closed = closed.clone();
|
let closed = closed.clone();
|
||||||
|
|
||||||
|
if let Err(e) = socket.set_nonblocking(true) {
|
||||||
|
sender.send(ListenerMessage::Error(e)).ok();
|
||||||
|
}
|
||||||
|
|
||||||
move || {
|
move || {
|
||||||
let mut listener = Listener::new(socket, sender);
|
let mut listener = Listener::new(socket, sender);
|
||||||
|
|
||||||
|
|||||||
@ -1,5 +1,5 @@
|
|||||||
use std::{
|
use std::{
|
||||||
io::Cursor,
|
io::{self, Cursor},
|
||||||
marker::PhantomData,
|
marker::PhantomData,
|
||||||
net::{SocketAddr, UdpSocket},
|
net::{SocketAddr, UdpSocket},
|
||||||
sync::{Arc, mpsc::Sender},
|
sync::{Arc, mpsc::Sender},
|
||||||
@ -49,6 +49,7 @@ impl<'a, T: Clone + Serialize + DeserializeOwned> Listener<T> {
|
|||||||
}
|
}
|
||||||
.ok();
|
.ok();
|
||||||
}
|
}
|
||||||
|
Err(ref err) if err.kind() == io::ErrorKind::WouldBlock => {}
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
self.sender.send(ListenerMessage::Error(err)).ok();
|
self.sender.send(ListenerMessage::Error(err)).ok();
|
||||||
}
|
}
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user