Implement network manager simulator lifecycle (#53)
All checks were successful
Native code generation / deterministic (push) Successful in 8m55s
Imaging and meshing gate / native (push) Successful in 2m45s
JPEG 2000 feature / linux (push) Successful in 1m41s
Native Rust workspace compile / compile (push) Successful in 1m57s
Skia feature / linux (push) Successful in 31m12s
All checks were successful
Native code generation / deterministic (push) Successful in 8m55s
Imaging and meshing gate / native (push) Successful in 2m45s
JPEG 2000 feature / linux (push) Successful in 1m41s
Native Rust workspace compile / compile (push) Successful in 1m57s
Skia feature / linux (push) Successful in 31m12s
This commit is contained in:
@@ -223,12 +223,23 @@ impl ClientRuntime {
|
||||
/// Construction is side-effect free: it creates no runtime, task, socket, or
|
||||
/// HTTP request. Network-facing services are attached explicitly and are owned
|
||||
/// until ordered shutdown or drop.
|
||||
#[derive(Clone)]
|
||||
pub struct GridClient {
|
||||
pub(crate) settings: Settings,
|
||||
pub(crate) time_provider: TimeProvider,
|
||||
runtime: Arc<ClientRuntime>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) struct ClientRetentionProbe(std::sync::Weak<ClientRuntime>);
|
||||
|
||||
#[cfg(test)]
|
||||
impl ClientRetentionProbe {
|
||||
pub(crate) fn is_released(&self) -> bool {
|
||||
self.0.upgrade().is_none()
|
||||
}
|
||||
}
|
||||
|
||||
impl GridClient {
|
||||
#[must_use]
|
||||
pub fn builder() -> GridClientBuilder {
|
||||
@@ -301,6 +312,11 @@ impl GridClient {
|
||||
pub(crate) fn shutdown(&self) -> Result<(), ClientCoreError> {
|
||||
self.runtime.shutdown()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn retention_probe(&self) -> ClientRetentionProbe {
|
||||
ClientRetentionProbe(Arc::downgrade(&self.runtime))
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Debug for GridClient {
|
||||
@@ -325,7 +341,9 @@ impl fmt::Debug for GridClient {
|
||||
|
||||
impl Drop for GridClient {
|
||||
fn drop(&mut self) {
|
||||
let _ = self.runtime.shutdown();
|
||||
if Arc::strong_count(&self.runtime) == 1 {
|
||||
let _ = self.runtime.shutdown();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -432,6 +450,11 @@ impl Settings {
|
||||
}
|
||||
}
|
||||
|
||||
/// Mutable access to the mapped agent settings for native composition.
|
||||
pub fn agent_settings_mut(&mut self) -> &mut AgentSettings {
|
||||
&mut self.agent
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn connection_mut(&mut self) -> &mut ConnectionSettings {
|
||||
&mut self.connection
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -13,6 +13,7 @@ mod generated;
|
||||
mod genepool_catalog;
|
||||
mod j2k;
|
||||
mod message_codec;
|
||||
mod network_manager;
|
||||
#[rustfmt::skip]
|
||||
pub mod packet_catalog;
|
||||
mod packet_wire;
|
||||
@@ -205,6 +206,7 @@ pub use libremetaverse_imaging as imaging_abstractions;
|
||||
pub use libremetaverse_structured_data as structured_data;
|
||||
pub use libremetaverse_types as types;
|
||||
pub use libremetaverse_types::Error;
|
||||
pub use network_manager::SimulatorCollection;
|
||||
pub use udp_transport::{
|
||||
AgentThrottleSender, UdpPacketHandler, UdpThrottleCategory, UdpTransportConfig,
|
||||
UdpTransportError, UdpTransportStats,
|
||||
|
||||
2272
crates/libremetaverse/src/network_manager.rs
Normal file
2272
crates/libremetaverse/src/network_manager.rs
Normal file
File diff suppressed because it is too large
Load Diff
@@ -367,6 +367,16 @@ pub(crate) fn ack_length(header: &Header) -> Result<usize, Error> {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn encode_base_packet(packet: &Packet) -> Result<Vec<u8>, Error> {
|
||||
let capacity = header_length(packet.header.frequency)
|
||||
.checked_add(ack_length(&packet.header)?)
|
||||
.ok_or(Error::Argument)?;
|
||||
let mut writer = WireWriter::with_capacity(capacity)?;
|
||||
encode_header(&packet.header, &mut writer)?;
|
||||
encode_acks(&packet.header, &mut writer)?;
|
||||
Ok(writer.into_inner())
|
||||
}
|
||||
|
||||
pub(crate) fn decode_header(
|
||||
bytes: &[u8],
|
||||
position: &mut i32,
|
||||
|
||||
@@ -11,6 +11,7 @@ use std::fmt;
|
||||
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
|
||||
use std::panic::{AssertUnwindSafe, catch_unwind};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||
use std::sync::mpsc::{Receiver as AckReceiver, SyncSender as AckSender, sync_channel};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Duration;
|
||||
use tokio::net::UdpSocket;
|
||||
@@ -966,6 +967,7 @@ impl UDPBase {
|
||||
packet_type,
|
||||
do_zerocode,
|
||||
response: response_sender,
|
||||
acknowledgement: None,
|
||||
};
|
||||
tokio::select! {
|
||||
result = sender.send(command) => {
|
||||
@@ -999,6 +1001,7 @@ impl UDPBase {
|
||||
packet_type,
|
||||
do_zerocode,
|
||||
response,
|
||||
acknowledgement: None,
|
||||
})
|
||||
.map_err(|error| {
|
||||
if matches!(error, mpsc::error::TrySendError::Full(_)) {
|
||||
@@ -1014,6 +1017,64 @@ impl UDPBase {
|
||||
Ok(receiver)
|
||||
}
|
||||
|
||||
/// Queues a packet and resolves only after its reliable sequence is
|
||||
/// acknowledged by the peer. The receiver disconnects if the reliable
|
||||
/// packet is dropped or the transport shuts down.
|
||||
pub(crate) fn try_send_packet_wait_ack(
|
||||
&self,
|
||||
data: Vec<u8>,
|
||||
destination: SocketAddr,
|
||||
packet_type: PacketType,
|
||||
do_zerocode: bool,
|
||||
) -> Result<AckReceiver<()>, UdpTransportError> {
|
||||
let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?;
|
||||
let (response, _queued) = oneshot::channel();
|
||||
let (acknowledgement, receiver) = sync_channel(1);
|
||||
sender
|
||||
.try_send(CoordinatorCommand::Packet {
|
||||
data,
|
||||
destination,
|
||||
packet_type,
|
||||
do_zerocode,
|
||||
response,
|
||||
acknowledgement: Some(acknowledgement),
|
||||
})
|
||||
.map_err(|error| {
|
||||
if matches!(error, mpsc::error::TrySendError::Full(_)) {
|
||||
self.inner
|
||||
.stats
|
||||
.dropped_send_queue
|
||||
.fetch_add(1, Ordering::Relaxed);
|
||||
UdpTransportError::Backpressure
|
||||
} else {
|
||||
UdpTransportError::Cancelled
|
||||
}
|
||||
})?;
|
||||
Ok(receiver)
|
||||
}
|
||||
|
||||
/// Queues an immediate flush of acknowledgements pending for `destination`.
|
||||
///
|
||||
/// The network manager uses this for the protocol's `OldestUnacked`
|
||||
/// request in `StartPingCheck`, matching `Simulator.SendAcks()` without
|
||||
/// exposing coordinator state or blocking the packet callback thread.
|
||||
pub(crate) fn try_flush_acks(&self, destination: SocketAddr) -> Result<(), UdpTransportError> {
|
||||
let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?;
|
||||
sender
|
||||
.try_send(CoordinatorCommand::FlushAcks { destination })
|
||||
.map_err(|error| {
|
||||
if matches!(error, mpsc::error::TrySendError::Full(_)) {
|
||||
self.inner
|
||||
.stats
|
||||
.dropped_send_queue
|
||||
.fetch_add(1, Ordering::Relaxed);
|
||||
UdpTransportError::Backpressure
|
||||
} else {
|
||||
UdpTransportError::Cancelled
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
pub fn update_throttle(&self, throttle: AgentThrottle) -> Result<(), UdpTransportError> {
|
||||
let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?;
|
||||
sender
|
||||
@@ -1059,6 +1120,10 @@ enum CoordinatorCommand {
|
||||
packet_type: PacketType,
|
||||
do_zerocode: bool,
|
||||
response: oneshot::Sender<Result<u32, UdpTransportError>>,
|
||||
acknowledgement: Option<AckSender<()>>,
|
||||
},
|
||||
FlushAcks {
|
||||
destination: SocketAddr,
|
||||
},
|
||||
UpdateThrottle(AgentThrottle),
|
||||
}
|
||||
@@ -1076,6 +1141,7 @@ struct ReliablePacket {
|
||||
category: UdpThrottleCategory,
|
||||
last_sent: Instant,
|
||||
resend_count: u32,
|
||||
acknowledgement: Option<AckSender<()>>,
|
||||
}
|
||||
|
||||
struct PeerState {
|
||||
@@ -1167,6 +1233,7 @@ impl CoordinatorState {
|
||||
destination: SocketAddr,
|
||||
packet_type: PacketType,
|
||||
do_zerocode: bool,
|
||||
acknowledgement: Option<AckSender<()>>,
|
||||
) -> Result<(WriteCommand, u32), UdpTransportError> {
|
||||
let mut data = encode_for_transport(data, do_zerocode, self.config.protocol_mtu)?;
|
||||
if data.len() < 6 {
|
||||
@@ -1198,8 +1265,11 @@ impl CoordinatorState {
|
||||
category,
|
||||
last_sent: Instant::now(),
|
||||
resend_count: 0,
|
||||
acknowledgement,
|
||||
},
|
||||
);
|
||||
} else if let Some(acknowledgement) = acknowledgement {
|
||||
let _ = acknowledgement.try_send(());
|
||||
}
|
||||
Ok((WriteCommand::Datagram { buffer, category }, sequence))
|
||||
}
|
||||
@@ -1410,8 +1480,10 @@ async fn process_command(
|
||||
packet_type,
|
||||
do_zerocode,
|
||||
response,
|
||||
acknowledgement,
|
||||
} => {
|
||||
let result = state.prepare_packet(data, destination, packet_type, do_zerocode);
|
||||
let result =
|
||||
state.prepare_packet(data, destination, packet_type, do_zerocode, acknowledgement);
|
||||
let result = match result {
|
||||
Ok((write, sequence)) => writer
|
||||
.send(write)
|
||||
@@ -1422,6 +1494,9 @@ async fn process_command(
|
||||
};
|
||||
let _ = response.send(result);
|
||||
}
|
||||
CoordinatorCommand::FlushAcks { destination } => {
|
||||
send_peer_acks(state, destination, writer).await;
|
||||
}
|
||||
CoordinatorCommand::UpdateThrottle(throttle) => {
|
||||
let _ = writer.send(WriteCommand::UpdateThrottle(throttle)).await;
|
||||
}
|
||||
@@ -1489,7 +1564,11 @@ async fn process_incoming(
|
||||
return;
|
||||
};
|
||||
for ack in appended_acks.into_iter().chain(standalone_acks) {
|
||||
peer.need_ack.remove(&ack);
|
||||
if let Some(mut reliable) = peer.need_ack.remove(&ack)
|
||||
&& let Some(acknowledgement) = reliable.acknowledgement.take()
|
||||
{
|
||||
let _ = acknowledgement.try_send(());
|
||||
}
|
||||
counters
|
||||
.acknowledgements_received
|
||||
.fetch_add(1, Ordering::Relaxed);
|
||||
@@ -1561,7 +1640,8 @@ async fn send_peer_acks(
|
||||
}
|
||||
Err(_) => return,
|
||||
};
|
||||
let Ok((write, _)) = state.prepare_packet(bytes, endpoint, PacketType::PacketAck, false) else {
|
||||
let Ok((write, _)) = state.prepare_packet(bytes, endpoint, PacketType::PacketAck, false, None)
|
||||
else {
|
||||
return;
|
||||
};
|
||||
if writer.send(write).await.is_ok() {
|
||||
@@ -1759,7 +1839,7 @@ mod tests {
|
||||
let destination: SocketAddr = "127.0.0.1:13000".parse().unwrap();
|
||||
let packet = vec![Helpers::MSG_RELIABLE, 0, 0, 0, 0, 0, 1];
|
||||
coordinator
|
||||
.prepare_packet(packet, destination, PacketType::ObjectUpdate, false)
|
||||
.prepare_packet(packet, destination, PacketType::ObjectUpdate, false, None)
|
||||
.unwrap();
|
||||
assert!(coordinator.take_resends().is_empty());
|
||||
|
||||
|
||||
Reference in New Issue
Block a user