Add simulated ping
This commit is contained in:
parent
a2e25d1c3f
commit
39fb40e0a0
@ -1,5 +1,5 @@
|
|||||||
use std::{
|
use std::{
|
||||||
collections::HashMap,
|
collections::{HashMap, VecDeque},
|
||||||
io::Cursor,
|
io::Cursor,
|
||||||
net::{SocketAddr, UdpSocket},
|
net::{SocketAddr, UdpSocket},
|
||||||
sync::Arc,
|
sync::Arc,
|
||||||
@ -13,7 +13,7 @@ use thiserror::Error;
|
|||||||
use crate::{
|
use crate::{
|
||||||
PeerConfig, PeerMessage,
|
PeerConfig, PeerMessage,
|
||||||
listener::{DATAGRAM_SIZE, ListenerError},
|
listener::{DATAGRAM_SIZE, ListenerError},
|
||||||
package::{CloseReason, Message, Messages, Package},
|
package::{CloseReason, DelayedPackage, Message, Messages, Package},
|
||||||
stats::NetStats,
|
stats::NetStats,
|
||||||
};
|
};
|
||||||
|
|
||||||
@ -46,7 +46,7 @@ pub(crate) struct ConnectionManager<
|
|||||||
T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static,
|
T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static,
|
||||||
> {
|
> {
|
||||||
/// Wraps UdpSocket with some helper methods
|
/// Wraps UdpSocket with some helper methods
|
||||||
pub udp: UdpWrapper,
|
pub udp: UdpWrapper<T>,
|
||||||
connections: HashMap<SocketAddr, Connection<T>>,
|
connections: HashMap<SocketAddr, Connection<T>>,
|
||||||
closing_since: Option<Instant>,
|
closing_since: Option<Instant>,
|
||||||
config: PeerConfig,
|
config: PeerConfig,
|
||||||
@ -60,7 +60,7 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
/// Creates a new ConnectionManager with the given socket and config
|
/// Creates a new ConnectionManager with the given socket and config
|
||||||
pub fn new(socket: Arc<UdpSocket>, config: PeerConfig) -> ConnectionManager<T> {
|
pub fn new(socket: Arc<UdpSocket>, config: PeerConfig) -> ConnectionManager<T> {
|
||||||
ConnectionManager {
|
ConnectionManager {
|
||||||
udp: UdpWrapper::new(socket),
|
udp: UdpWrapper::new(socket, config.additional_ping),
|
||||||
connections: HashMap::new(),
|
connections: HashMap::new(),
|
||||||
closing_since: None,
|
closing_since: None,
|
||||||
config,
|
config,
|
||||||
@ -177,6 +177,13 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
match self.udp.send_delayed_messages() {
|
||||||
|
Ok(_) => {}
|
||||||
|
Err(err) => {
|
||||||
|
println!("Error: {}", err);
|
||||||
|
self.close();
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn send(&mut self, addr: &SocketAddr, message: T, reliable: bool) {
|
pub fn send(&mut self, addr: &SocketAddr, message: T, reliable: bool) {
|
||||||
@ -453,14 +460,22 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Wrap UdpSocket with some helper methods
|
/// Wrap UdpSocket with some helper methods
|
||||||
pub(crate) struct UdpWrapper {
|
pub(crate) struct UdpWrapper<
|
||||||
|
T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static,
|
||||||
|
> {
|
||||||
socket: Option<Arc<UdpSocket>>,
|
socket: Option<Arc<UdpSocket>>,
|
||||||
|
delayed_queue: VecDeque<DelayedPackage<T>>,
|
||||||
|
additional_ping: Duration,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl UdpWrapper {
|
impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static>
|
||||||
pub fn new(socket: Arc<UdpSocket>) -> UdpWrapper {
|
UdpWrapper<T>
|
||||||
|
{
|
||||||
|
pub fn new(socket: Arc<UdpSocket>, additional_ping: Duration) -> UdpWrapper<T> {
|
||||||
UdpWrapper {
|
UdpWrapper {
|
||||||
socket: Some(socket),
|
socket: Some(socket),
|
||||||
|
delayed_queue: VecDeque::new(),
|
||||||
|
additional_ping,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -475,16 +490,32 @@ impl UdpWrapper {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Send a Package to a remote peer at addr
|
/// Send a Package to a remote peer at addr
|
||||||
fn send_to<T: Serialize>(
|
fn send_to(&mut self, addr: SocketAddr, package: Package<T>) -> Result<usize, SendError> {
|
||||||
&self,
|
let delayed = DelayedPackage {
|
||||||
addr: SocketAddr,
|
package,
|
||||||
package: Package<T>,
|
created: Instant::now(),
|
||||||
) -> Result<usize, SendError> {
|
addr,
|
||||||
|
size: 0,
|
||||||
|
};
|
||||||
|
if Instant::now() < (delayed.created + self.additional_ping) {
|
||||||
|
let mut buf = Vec::new();
|
||||||
|
let cursor = Cursor::new(&mut buf);
|
||||||
|
let encoder = GzEncoder::new(cursor, Compression::fast());
|
||||||
|
ciborium::into_writer(&delayed.package, encoder)
|
||||||
|
.map_err(|err| SendError::SerializationError(err))?;
|
||||||
|
if buf.len() > DATAGRAM_SIZE {
|
||||||
|
return Err(SendError::DatagramTooLarge(buf.len()));
|
||||||
|
}
|
||||||
|
|
||||||
|
self.delayed_queue.push_back(delayed);
|
||||||
|
return Ok(buf.len());
|
||||||
|
}
|
||||||
|
|
||||||
if let Some(socket) = &self.socket {
|
if let Some(socket) = &self.socket {
|
||||||
let mut buf = Vec::new();
|
let mut buf = Vec::new();
|
||||||
let cursor = Cursor::new(&mut buf);
|
let cursor = Cursor::new(&mut buf);
|
||||||
let encoder = GzEncoder::new(cursor, Compression::fast());
|
let encoder = GzEncoder::new(cursor, Compression::fast());
|
||||||
ciborium::into_writer(&package, encoder)
|
ciborium::into_writer(&delayed.package, encoder)
|
||||||
.map_err(|err| SendError::SerializationError(err))?;
|
.map_err(|err| SendError::SerializationError(err))?;
|
||||||
if buf.len() > DATAGRAM_SIZE {
|
if buf.len() > DATAGRAM_SIZE {
|
||||||
return Err(SendError::DatagramTooLarge(buf.len()));
|
return Err(SendError::DatagramTooLarge(buf.len()));
|
||||||
@ -497,6 +528,32 @@ impl UdpWrapper {
|
|||||||
Err(SendError::SocketDisconnected)
|
Err(SendError::SocketDisconnected)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn send_delayed_messages(&mut self) -> Result<(), SendError> {
|
||||||
|
let queue_iter = self.delayed_queue.clone().into_iter();
|
||||||
|
let treshold = Instant::now() - self.additional_ping;
|
||||||
|
let messages = queue_iter.take_while(|m| treshold >= m.created);
|
||||||
|
self.delayed_queue.retain(|m| treshold < m.created);
|
||||||
|
|
||||||
|
if let Some(socket) = &self.socket {
|
||||||
|
for message in messages {
|
||||||
|
let mut buf = Vec::new();
|
||||||
|
let cursor = Cursor::new(&mut buf);
|
||||||
|
let encoder = GzEncoder::new(cursor, Compression::fast());
|
||||||
|
ciborium::into_writer(&message.package, encoder)
|
||||||
|
.map_err(|err| SendError::SerializationError(err))?;
|
||||||
|
if buf.len() > DATAGRAM_SIZE {
|
||||||
|
return Err(SendError::DatagramTooLarge(buf.len()));
|
||||||
|
}
|
||||||
|
socket
|
||||||
|
.send_to(&buf, message.addr)
|
||||||
|
.map_err(|e| SendError::SendError(e))?;
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
} else {
|
||||||
|
Err(SendError::SocketDisconnected)
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Represents a single connection from the peer to another peer
|
/// Represents a single connection from the peer to another peer
|
||||||
|
|||||||
45
src/lib.rs
45
src/lib.rs
@ -7,7 +7,7 @@ use std::{
|
|||||||
mpsc::{Receiver, channel},
|
mpsc::{Receiver, channel},
|
||||||
},
|
},
|
||||||
thread::spawn,
|
thread::spawn,
|
||||||
time::Duration,
|
time::{Duration, Instant},
|
||||||
};
|
};
|
||||||
|
|
||||||
use serde::{Serialize, de::DeserializeOwned};
|
use serde::{Serialize, de::DeserializeOwned};
|
||||||
@ -16,7 +16,7 @@ use thiserror::*;
|
|||||||
use crate::{
|
use crate::{
|
||||||
connections::{Connection, ConnectionError, ConnectionManager},
|
connections::{Connection, ConnectionError, ConnectionManager},
|
||||||
listener::{Listener, ListenerMessage},
|
listener::{Listener, ListenerMessage},
|
||||||
package::CloseReason,
|
package::{CloseReason, DelayedPackage, Package},
|
||||||
stats::NetStats,
|
stats::NetStats,
|
||||||
};
|
};
|
||||||
|
|
||||||
@ -55,6 +55,7 @@ pub struct PeerConfig {
|
|||||||
disconnect_timeout: Duration,
|
disconnect_timeout: Duration,
|
||||||
message_retry: Duration,
|
message_retry: Duration,
|
||||||
identifier: String,
|
identifier: String,
|
||||||
|
additional_ping: Duration,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for PeerConfig {
|
impl Default for PeerConfig {
|
||||||
@ -65,6 +66,7 @@ impl Default for PeerConfig {
|
|||||||
disconnect_timeout: Duration::from_millis(500),
|
disconnect_timeout: Duration::from_millis(500),
|
||||||
message_retry: Duration::from_millis(100),
|
message_retry: Duration::from_millis(100),
|
||||||
identifier: "ExamplePeer".to_string(),
|
identifier: "ExamplePeer".to_string(),
|
||||||
|
additional_ping: Duration::from_nanos(0),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -112,12 +114,23 @@ impl PeerConfig {
|
|||||||
..self
|
..self
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Sets additional ping for this peer, useful for debugging high-latency
|
||||||
|
/// connections.
|
||||||
|
pub fn with_additional_ping(self, ping: Duration) -> PeerConfig {
|
||||||
|
PeerConfig {
|
||||||
|
additional_ping: ping,
|
||||||
|
..self
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub struct Peer<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> {
|
pub struct Peer<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> {
|
||||||
connection_mgr: ConnectionManager<T>,
|
connection_mgr: ConnectionManager<T>,
|
||||||
|
config: PeerConfig,
|
||||||
closed: Arc<AtomicBool>,
|
closed: Arc<AtomicBool>,
|
||||||
receiver: Receiver<ListenerMessage<T>>,
|
receiver: Receiver<ListenerMessage<T>>,
|
||||||
|
delayed_messages: VecDeque<DelayedPackage<T>>,
|
||||||
messages: VecDeque<PeerMessage<T>>,
|
messages: VecDeque<PeerMessage<T>>,
|
||||||
}
|
}
|
||||||
impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> Drop
|
impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static> Drop
|
||||||
@ -159,9 +172,11 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
|
|
||||||
Ok(Peer {
|
Ok(Peer {
|
||||||
connection_mgr: ConnectionManager::new(socket, config.clone()),
|
connection_mgr: ConnectionManager::new(socket, config.clone()),
|
||||||
|
config: config,
|
||||||
closed,
|
closed,
|
||||||
receiver,
|
receiver,
|
||||||
messages: VecDeque::new(),
|
messages: VecDeque::new(),
|
||||||
|
delayed_messages: VecDeque::new(),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -185,6 +200,22 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
self.connection_mgr.send_pings();
|
self.connection_mgr.send_pings();
|
||||||
self.connection_mgr.send_queued_messages();
|
self.connection_mgr.send_queued_messages();
|
||||||
|
|
||||||
|
// Process delayed messages
|
||||||
|
let now = Instant::now();
|
||||||
|
for i in (0..self.delayed_messages.len()).rev() {
|
||||||
|
if let Some(msg) = self.delayed_messages.get(i) {
|
||||||
|
if (msg.created + self.config.additional_ping) > now {
|
||||||
|
if let Some(msg) = self.delayed_messages.remove(i) {
|
||||||
|
self.messages.extend(self.connection_mgr.handle_package(
|
||||||
|
msg.package,
|
||||||
|
&msg.addr,
|
||||||
|
msg.size,
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
let mut msg;
|
let mut msg;
|
||||||
while {
|
while {
|
||||||
msg = match self.receiver.try_recv() {
|
msg = match self.receiver.try_recv() {
|
||||||
@ -196,11 +227,21 @@ impl<T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + '
|
|||||||
} {
|
} {
|
||||||
match msg.unwrap() {
|
match msg.unwrap() {
|
||||||
ListenerMessage::Package(package, socket_addr, bytes) => {
|
ListenerMessage::Package(package, socket_addr, bytes) => {
|
||||||
|
let received = Instant::now();
|
||||||
|
if Instant::now() > (received + self.config.additional_ping) {
|
||||||
self.messages.extend(self.connection_mgr.handle_package(
|
self.messages.extend(self.connection_mgr.handle_package(
|
||||||
package,
|
package,
|
||||||
&socket_addr,
|
&socket_addr,
|
||||||
bytes,
|
bytes,
|
||||||
));
|
));
|
||||||
|
} else {
|
||||||
|
self.delayed_messages.push_back(DelayedPackage {
|
||||||
|
package,
|
||||||
|
created: received,
|
||||||
|
addr: socket_addr,
|
||||||
|
size: bytes,
|
||||||
|
});
|
||||||
|
}
|
||||||
}
|
}
|
||||||
ListenerMessage::PackageError(listener_error, socket_addr) => {
|
ListenerMessage::PackageError(listener_error, socket_addr) => {
|
||||||
println!("Error: {}", listener_error);
|
println!("Error: {}", listener_error);
|
||||||
|
|||||||
@ -1,8 +1,10 @@
|
|||||||
use std::{
|
use std::{
|
||||||
|
collections::VecDeque,
|
||||||
io::{self, Cursor},
|
io::{self, Cursor},
|
||||||
marker::PhantomData,
|
marker::PhantomData,
|
||||||
net::{SocketAddr, UdpSocket},
|
net::{SocketAddr, UdpSocket},
|
||||||
sync::{Arc, mpsc::Sender},
|
sync::{Arc, mpsc::Sender},
|
||||||
|
time::{Duration, Instant},
|
||||||
};
|
};
|
||||||
|
|
||||||
use flate2::read::GzDecoder;
|
use flate2::read::GzDecoder;
|
||||||
|
|||||||
@ -1,7 +1,9 @@
|
|||||||
use serde::{Deserialize, Serialize};
|
use std::{net::SocketAddr, time::Instant};
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize)]
|
use serde::{Deserialize, Serialize, de::DeserializeOwned};
|
||||||
pub(crate) enum Package<T> {
|
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
pub(crate) enum Package<T: Clone> {
|
||||||
Hello(String),
|
Hello(String),
|
||||||
Ping,
|
Ping,
|
||||||
Pong,
|
Pong,
|
||||||
@ -33,7 +35,7 @@ impl std::fmt::Display for CloseReason {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize)]
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
pub(crate) struct Messages<T> {
|
pub(crate) struct Messages<T> {
|
||||||
pub ack: u64,
|
pub ack: u64,
|
||||||
pub messages_sent: usize,
|
pub messages_sent: usize,
|
||||||
@ -46,3 +48,13 @@ pub(crate) struct Message<T> {
|
|||||||
pub message_id: u64,
|
pub message_id: u64,
|
||||||
pub message: T,
|
pub message: T,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub(crate) struct DelayedPackage<
|
||||||
|
T: std::fmt::Debug + Clone + Serialize + DeserializeOwned + Send + Sync + 'static,
|
||||||
|
> {
|
||||||
|
pub package: Package<T>,
|
||||||
|
pub created: Instant,
|
||||||
|
pub addr: SocketAddr,
|
||||||
|
pub size: usize,
|
||||||
|
}
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user