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:
10
README.md
10
README.md
@@ -272,3 +272,13 @@ zerocode expansion, ACK state, and reliable windows have explicit limits;
|
|||||||
linked cancellation and final drop release every socket task. The executor,
|
linked cancellation and final drop release every socket task. The executor,
|
||||||
wire, backpressure, retry, and security contracts are documented in
|
wire, backpressure, retry, and security contracts are documented in
|
||||||
[`docs/udp-transport.md`](docs/udp-transport.md).
|
[`docs/udp-transport.md`](docs/udp-transport.md).
|
||||||
|
The native network manager now layers ordered packet and CAPS callback
|
||||||
|
registries, typed RAII subscriptions, synchronized simulator collections,
|
||||||
|
cancellable connection events, ACK-gated circuit setup, current-simulator and
|
||||||
|
seed-capability selection, CAPS `EnableSimulator`, UDP `DisableSimulator`, ping
|
||||||
|
handling, typed region-handshake replies, two-interval keepalive detection, and
|
||||||
|
deterministic disconnect reasons on that transport. Callback lists are
|
||||||
|
snapshotted before invocation, and manager/transport workers retain no dropped
|
||||||
|
client owner. The ownership,
|
||||||
|
dispatch, lifecycle, and fake-server verification contracts are documented in
|
||||||
|
[`docs/network-manager.md`](docs/network-manager.md).
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ Generated by `python3 tools/generate_api_shims.py`; do not edit by hand.
|
|||||||
|
|
||||||
| Assembly | Types | Members | Status |
|
| Assembly | Types | Members | Status |
|
||||||
|---|---:|---:|---|
|
|---|---:|---:|---|
|
||||||
| `LibreMetaverse` | 2,711 | 27,281 | native implementation: 32 types / 13,459 members; remaining surface is callable failure-only shims |
|
| `LibreMetaverse` | 2,711 | 27,281 | native implementation: 47 types / 13,615 members; remaining surface is callable failure-only shims |
|
||||||
| `LibreMetaverse.Imaging.Abstractions` | 3 | 20 | native implementation: 3 types / 20 members; no generated shims remain |
|
| `LibreMetaverse.Imaging.Abstractions` | 3 | 20 | native implementation: 3 types / 20 members; no generated shims remain |
|
||||||
| `LibreMetaverse.Imaging.Skia` | 1 | 3 | native implementation: 1 type / 3 members; no generated shims remain |
|
| `LibreMetaverse.Imaging.Skia` | 1 | 3 | native implementation: 1 type / 3 members; no generated shims remain |
|
||||||
| `LibreMetaverse.LslTools` | 164 | 768 | callable failure-only shim |
|
| `LibreMetaverse.LslTools` | 164 | 768 | callable failure-only shim |
|
||||||
|
|||||||
@@ -466,8 +466,61 @@ impl Drop for CancellationRegistration {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Default, Eq, Hash, PartialEq)]
|
/// RAII guard for a mapped event subscription.
|
||||||
pub struct Subscription;
|
///
|
||||||
|
/// Dropping or explicitly closing the guard invokes the removal callback once.
|
||||||
|
/// The callback normally owns only a weak reference to the event registry, so
|
||||||
|
/// an abandoned subscription cannot keep its publisher alive.
|
||||||
|
#[must_use = "dropping the subscription immediately unregisters the handler"]
|
||||||
|
pub struct Subscription {
|
||||||
|
close: Option<Box<dyn FnOnce() + Send + Sync>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Subscription {
|
||||||
|
pub fn new(close: impl FnOnce() + Send + Sync + 'static) -> Self {
|
||||||
|
Self {
|
||||||
|
close: Some(Box::new(close)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn close(mut self) {
|
||||||
|
if let Some(close) = self.close.take() {
|
||||||
|
close();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[must_use]
|
||||||
|
pub const fn is_active(&self) -> bool {
|
||||||
|
self.close.is_some()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub const fn detached() -> Self {
|
||||||
|
Self { close: None }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for Subscription {
|
||||||
|
fn default() -> Self {
|
||||||
|
Self::detached()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl fmt::Debug for Subscription {
|
||||||
|
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||||
|
formatter
|
||||||
|
.debug_struct("Subscription")
|
||||||
|
.field("active", &self.is_active())
|
||||||
|
.finish()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for Subscription {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if let Some(close) = self.close.take() {
|
||||||
|
close();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub type EventHandler<T> = Arc<dyn Fn(T) + Send + Sync>;
|
pub type EventHandler<T> = Arc<dyn Fn(T) + Send + Sync>;
|
||||||
|
|
||||||
|
|||||||
@@ -223,12 +223,23 @@ impl ClientRuntime {
|
|||||||
/// Construction is side-effect free: it creates no runtime, task, socket, or
|
/// Construction is side-effect free: it creates no runtime, task, socket, or
|
||||||
/// HTTP request. Network-facing services are attached explicitly and are owned
|
/// HTTP request. Network-facing services are attached explicitly and are owned
|
||||||
/// until ordered shutdown or drop.
|
/// until ordered shutdown or drop.
|
||||||
|
#[derive(Clone)]
|
||||||
pub struct GridClient {
|
pub struct GridClient {
|
||||||
pub(crate) settings: Settings,
|
pub(crate) settings: Settings,
|
||||||
pub(crate) time_provider: TimeProvider,
|
pub(crate) time_provider: TimeProvider,
|
||||||
runtime: Arc<ClientRuntime>,
|
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 {
|
impl GridClient {
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn builder() -> GridClientBuilder {
|
pub fn builder() -> GridClientBuilder {
|
||||||
@@ -301,6 +312,11 @@ impl GridClient {
|
|||||||
pub(crate) fn shutdown(&self) -> Result<(), ClientCoreError> {
|
pub(crate) fn shutdown(&self) -> Result<(), ClientCoreError> {
|
||||||
self.runtime.shutdown()
|
self.runtime.shutdown()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) fn retention_probe(&self) -> ClientRetentionProbe {
|
||||||
|
ClientRetentionProbe(Arc::downgrade(&self.runtime))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl fmt::Debug for GridClient {
|
impl fmt::Debug for GridClient {
|
||||||
@@ -325,9 +341,11 @@ impl fmt::Debug for GridClient {
|
|||||||
|
|
||||||
impl Drop for GridClient {
|
impl Drop for GridClient {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
|
if Arc::strong_count(&self.runtime) == 1 {
|
||||||
let _ = self.runtime.shutdown();
|
let _ = self.runtime.shutdown();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl crate::IGridClient for GridClient {
|
impl crate::IGridClient for GridClient {
|
||||||
fn settings(&self) -> Settings {
|
fn settings(&self) -> Settings {
|
||||||
@@ -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]
|
#[must_use]
|
||||||
pub fn connection_mut(&mut self) -> &mut ConnectionSettings {
|
pub fn connection_mut(&mut self) -> &mut ConnectionSettings {
|
||||||
&mut self.connection
|
&mut self.connection
|
||||||
|
|||||||
File diff suppressed because it is too large
Load Diff
@@ -13,6 +13,7 @@ mod generated;
|
|||||||
mod genepool_catalog;
|
mod genepool_catalog;
|
||||||
mod j2k;
|
mod j2k;
|
||||||
mod message_codec;
|
mod message_codec;
|
||||||
|
mod network_manager;
|
||||||
#[rustfmt::skip]
|
#[rustfmt::skip]
|
||||||
pub mod packet_catalog;
|
pub mod packet_catalog;
|
||||||
mod packet_wire;
|
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_structured_data as structured_data;
|
||||||
pub use libremetaverse_types as types;
|
pub use libremetaverse_types as types;
|
||||||
pub use libremetaverse_types::Error;
|
pub use libremetaverse_types::Error;
|
||||||
|
pub use network_manager::SimulatorCollection;
|
||||||
pub use udp_transport::{
|
pub use udp_transport::{
|
||||||
AgentThrottleSender, UdpPacketHandler, UdpThrottleCategory, UdpTransportConfig,
|
AgentThrottleSender, UdpPacketHandler, UdpThrottleCategory, UdpTransportConfig,
|
||||||
UdpTransportError, UdpTransportStats,
|
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(
|
pub(crate) fn decode_header(
|
||||||
bytes: &[u8],
|
bytes: &[u8],
|
||||||
position: &mut i32,
|
position: &mut i32,
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ use std::fmt;
|
|||||||
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
|
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
|
||||||
use std::panic::{AssertUnwindSafe, catch_unwind};
|
use std::panic::{AssertUnwindSafe, catch_unwind};
|
||||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
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::sync::{Arc, Mutex};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
use tokio::net::UdpSocket;
|
use tokio::net::UdpSocket;
|
||||||
@@ -966,6 +967,7 @@ impl UDPBase {
|
|||||||
packet_type,
|
packet_type,
|
||||||
do_zerocode,
|
do_zerocode,
|
||||||
response: response_sender,
|
response: response_sender,
|
||||||
|
acknowledgement: None,
|
||||||
};
|
};
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
result = sender.send(command) => {
|
result = sender.send(command) => {
|
||||||
@@ -999,6 +1001,7 @@ impl UDPBase {
|
|||||||
packet_type,
|
packet_type,
|
||||||
do_zerocode,
|
do_zerocode,
|
||||||
response,
|
response,
|
||||||
|
acknowledgement: None,
|
||||||
})
|
})
|
||||||
.map_err(|error| {
|
.map_err(|error| {
|
||||||
if matches!(error, mpsc::error::TrySendError::Full(_)) {
|
if matches!(error, mpsc::error::TrySendError::Full(_)) {
|
||||||
@@ -1014,6 +1017,64 @@ impl UDPBase {
|
|||||||
Ok(receiver)
|
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> {
|
pub fn update_throttle(&self, throttle: AgentThrottle) -> Result<(), UdpTransportError> {
|
||||||
let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?;
|
let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?;
|
||||||
sender
|
sender
|
||||||
@@ -1059,6 +1120,10 @@ enum CoordinatorCommand {
|
|||||||
packet_type: PacketType,
|
packet_type: PacketType,
|
||||||
do_zerocode: bool,
|
do_zerocode: bool,
|
||||||
response: oneshot::Sender<Result<u32, UdpTransportError>>,
|
response: oneshot::Sender<Result<u32, UdpTransportError>>,
|
||||||
|
acknowledgement: Option<AckSender<()>>,
|
||||||
|
},
|
||||||
|
FlushAcks {
|
||||||
|
destination: SocketAddr,
|
||||||
},
|
},
|
||||||
UpdateThrottle(AgentThrottle),
|
UpdateThrottle(AgentThrottle),
|
||||||
}
|
}
|
||||||
@@ -1076,6 +1141,7 @@ struct ReliablePacket {
|
|||||||
category: UdpThrottleCategory,
|
category: UdpThrottleCategory,
|
||||||
last_sent: Instant,
|
last_sent: Instant,
|
||||||
resend_count: u32,
|
resend_count: u32,
|
||||||
|
acknowledgement: Option<AckSender<()>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
struct PeerState {
|
struct PeerState {
|
||||||
@@ -1167,6 +1233,7 @@ impl CoordinatorState {
|
|||||||
destination: SocketAddr,
|
destination: SocketAddr,
|
||||||
packet_type: PacketType,
|
packet_type: PacketType,
|
||||||
do_zerocode: bool,
|
do_zerocode: bool,
|
||||||
|
acknowledgement: Option<AckSender<()>>,
|
||||||
) -> Result<(WriteCommand, u32), UdpTransportError> {
|
) -> Result<(WriteCommand, u32), UdpTransportError> {
|
||||||
let mut data = encode_for_transport(data, do_zerocode, self.config.protocol_mtu)?;
|
let mut data = encode_for_transport(data, do_zerocode, self.config.protocol_mtu)?;
|
||||||
if data.len() < 6 {
|
if data.len() < 6 {
|
||||||
@@ -1198,8 +1265,11 @@ impl CoordinatorState {
|
|||||||
category,
|
category,
|
||||||
last_sent: Instant::now(),
|
last_sent: Instant::now(),
|
||||||
resend_count: 0,
|
resend_count: 0,
|
||||||
|
acknowledgement,
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
} else if let Some(acknowledgement) = acknowledgement {
|
||||||
|
let _ = acknowledgement.try_send(());
|
||||||
}
|
}
|
||||||
Ok((WriteCommand::Datagram { buffer, category }, sequence))
|
Ok((WriteCommand::Datagram { buffer, category }, sequence))
|
||||||
}
|
}
|
||||||
@@ -1410,8 +1480,10 @@ async fn process_command(
|
|||||||
packet_type,
|
packet_type,
|
||||||
do_zerocode,
|
do_zerocode,
|
||||||
response,
|
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 {
|
let result = match result {
|
||||||
Ok((write, sequence)) => writer
|
Ok((write, sequence)) => writer
|
||||||
.send(write)
|
.send(write)
|
||||||
@@ -1422,6 +1494,9 @@ async fn process_command(
|
|||||||
};
|
};
|
||||||
let _ = response.send(result);
|
let _ = response.send(result);
|
||||||
}
|
}
|
||||||
|
CoordinatorCommand::FlushAcks { destination } => {
|
||||||
|
send_peer_acks(state, destination, writer).await;
|
||||||
|
}
|
||||||
CoordinatorCommand::UpdateThrottle(throttle) => {
|
CoordinatorCommand::UpdateThrottle(throttle) => {
|
||||||
let _ = writer.send(WriteCommand::UpdateThrottle(throttle)).await;
|
let _ = writer.send(WriteCommand::UpdateThrottle(throttle)).await;
|
||||||
}
|
}
|
||||||
@@ -1489,7 +1564,11 @@ async fn process_incoming(
|
|||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
for ack in appended_acks.into_iter().chain(standalone_acks) {
|
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
|
counters
|
||||||
.acknowledgements_received
|
.acknowledgements_received
|
||||||
.fetch_add(1, Ordering::Relaxed);
|
.fetch_add(1, Ordering::Relaxed);
|
||||||
@@ -1561,7 +1640,8 @@ async fn send_peer_acks(
|
|||||||
}
|
}
|
||||||
Err(_) => return,
|
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;
|
return;
|
||||||
};
|
};
|
||||||
if writer.send(write).await.is_ok() {
|
if writer.send(write).await.is_ok() {
|
||||||
@@ -1759,7 +1839,7 @@ mod tests {
|
|||||||
let destination: SocketAddr = "127.0.0.1:13000".parse().unwrap();
|
let destination: SocketAddr = "127.0.0.1:13000".parse().unwrap();
|
||||||
let packet = vec![Helpers::MSG_RELIABLE, 0, 0, 0, 0, 0, 1];
|
let packet = vec![Helpers::MSG_RELIABLE, 0, 0, 0, 0, 0, 1];
|
||||||
coordinator
|
coordinator
|
||||||
.prepare_packet(packet, destination, PacketType::ObjectUpdate, false)
|
.prepare_packet(packet, destination, PacketType::ObjectUpdate, false, None)
|
||||||
.unwrap();
|
.unwrap();
|
||||||
assert!(coordinator.take_resends().is_empty());
|
assert!(coordinator.take_resends().is_empty());
|
||||||
|
|
||||||
|
|||||||
899
crates/libremetaverse/tests/network_manager.rs
Normal file
899
crates/libremetaverse/tests/network_manager.rs
Normal file
@@ -0,0 +1,899 @@
|
|||||||
|
use libremetaverse::interfaces::IMessage;
|
||||||
|
use libremetaverse::messages::linden::{
|
||||||
|
EnableSimulatorMessage, EnableSimulatorMessageSimulatorInfoBlock,
|
||||||
|
};
|
||||||
|
use libremetaverse::packets::{
|
||||||
|
DisableSimulatorPacket, Packet, PacketAckPacket, PacketAckPacketPacketsBlock, PacketType,
|
||||||
|
RegionHandshakePacket, RegionHandshakeReplyPacket, StartPingCheckPacket, UseCircuitCodePacket,
|
||||||
|
};
|
||||||
|
use libremetaverse::{
|
||||||
|
CapsEventDictionary, CapsEventQueueCallback, GridClient, Helpers, NetworkManager,
|
||||||
|
NetworkManagerDisconnectType, PacketEventDictionary, Simulator,
|
||||||
|
};
|
||||||
|
use libremetaverse_types::compat::{EventHandler, Uri};
|
||||||
|
use std::net::{IpAddr, Ipv4Addr, SocketAddr, UdpSocket};
|
||||||
|
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
|
||||||
|
use std::sync::mpsc::{Receiver, Sender, channel};
|
||||||
|
use std::sync::{Arc, Barrier, Mutex, Weak};
|
||||||
|
use std::thread::{self, JoinHandle};
|
||||||
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
|
fn block_on<F: std::future::Future>(future: F) -> F::Output {
|
||||||
|
use std::task::{Context, Poll, Wake, Waker};
|
||||||
|
|
||||||
|
struct ThreadWake(thread::Thread);
|
||||||
|
|
||||||
|
impl Wake for ThreadWake {
|
||||||
|
fn wake(self: Arc<Self>) {
|
||||||
|
self.0.unpark();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let waker = Waker::from(Arc::new(ThreadWake(thread::current())));
|
||||||
|
let mut context = Context::from_waker(&waker);
|
||||||
|
let mut future = std::pin::pin!(future);
|
||||||
|
loop {
|
||||||
|
match future.as_mut().poll(&mut context) {
|
||||||
|
Poll::Ready(output) => return output,
|
||||||
|
Poll::Pending => thread::park(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn loopback(port: u16) -> SocketAddr {
|
||||||
|
SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), port)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn simulator(client: &GridClient) -> Simulator {
|
||||||
|
Simulator::new(
|
||||||
|
client.clone(),
|
||||||
|
loopback(13000),
|
||||||
|
0x1000,
|
||||||
|
Some(256),
|
||||||
|
Some(256),
|
||||||
|
)
|
||||||
|
.expect("simulator")
|
||||||
|
}
|
||||||
|
|
||||||
|
fn base_packet(packet_type: PacketType) -> Packet {
|
||||||
|
Packet::build_packet_with_packet_type(packet_type).expect("base packet")
|
||||||
|
}
|
||||||
|
|
||||||
|
fn wait_until(mut condition: impl FnMut() -> bool) {
|
||||||
|
let deadline = Instant::now() + Duration::from_secs(2);
|
||||||
|
while !condition() {
|
||||||
|
assert!(Instant::now() < deadline, "condition timed out");
|
||||||
|
thread::sleep(Duration::from_millis(2));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn packet_callbacks_preserve_filtering_order_async_policy_and_reentrancy() {
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let sim = simulator(&client);
|
||||||
|
let events = Arc::new(PacketEventDictionary::new(client).expect("packet events"));
|
||||||
|
let order = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
|
||||||
|
let default_order = Arc::clone(&order);
|
||||||
|
events
|
||||||
|
.register_event(
|
||||||
|
PacketType::Default,
|
||||||
|
Arc::new(move |_| default_order.lock().unwrap().push("default")),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
let specific_order = Arc::clone(&order);
|
||||||
|
events
|
||||||
|
.register_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
Arc::new(move |_| specific_order.lock().unwrap().push("specific")),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
base_packet(PacketType::StartPingCheck),
|
||||||
|
sim.clone(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(*order.lock().unwrap(), ["default", "specific"]);
|
||||||
|
|
||||||
|
order.lock().unwrap().clear();
|
||||||
|
events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::CompletePingCheck,
|
||||||
|
base_packet(PacketType::CompletePingCheck),
|
||||||
|
sim.clone(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(*order.lock().unwrap(), ["default"]);
|
||||||
|
|
||||||
|
let reentrant_count = Arc::new(AtomicUsize::new(0));
|
||||||
|
let weak_events: Weak<PacketEventDictionary> = Arc::downgrade(&events);
|
||||||
|
let count = Arc::clone(&reentrant_count);
|
||||||
|
let reentrant: EventHandler<_> = Arc::new(move |_| {
|
||||||
|
count.fetch_add(1, Ordering::Relaxed);
|
||||||
|
if let Some(events) = weak_events.upgrade() {
|
||||||
|
events
|
||||||
|
.register_event(PacketType::UseCircuitCode, Arc::new(|_| {}), false)
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
events
|
||||||
|
.register_event(PacketType::UseCircuitCode, reentrant, false)
|
||||||
|
.unwrap();
|
||||||
|
events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::UseCircuitCode,
|
||||||
|
base_packet(PacketType::UseCircuitCode),
|
||||||
|
sim.clone(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(reentrant_count.load(Ordering::Relaxed), 1);
|
||||||
|
|
||||||
|
let async_events = PacketEventDictionary::new(sim.client.clone()).unwrap();
|
||||||
|
let removed_async: EventHandler<_> = Arc::new(|_| {});
|
||||||
|
async_events
|
||||||
|
.register_event(PacketType::StartPingCheck, Arc::clone(&removed_async), true)
|
||||||
|
.unwrap();
|
||||||
|
let caller_thread = thread::current().id();
|
||||||
|
let (thread_sender, thread_receiver) = channel();
|
||||||
|
async_events
|
||||||
|
.register_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
Arc::new(move |_| thread_sender.send(thread::current().id()).unwrap()),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
async_events
|
||||||
|
.unregister_event(PacketType::StartPingCheck, removed_async)
|
||||||
|
.unwrap();
|
||||||
|
async_events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
base_packet(PacketType::StartPingCheck),
|
||||||
|
sim,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_ne!(
|
||||||
|
thread_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
caller_thread,
|
||||||
|
"C# keeps a packet type asynchronous after its first async registration"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn synchronous_specific_dispatch_suppresses_default_async_chain_like_reference() {
|
||||||
|
let mixed_events = PacketEventDictionary::new(GridClient::new().unwrap()).unwrap();
|
||||||
|
let (default_sender, default_receiver) = channel();
|
||||||
|
mixed_events
|
||||||
|
.register_event(
|
||||||
|
PacketType::Default,
|
||||||
|
Arc::new(move |_| default_sender.send(()).unwrap()),
|
||||||
|
true,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
let specific_calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let calls = Arc::clone(&specific_calls);
|
||||||
|
mixed_events
|
||||||
|
.register_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
Arc::new(move |_| {
|
||||||
|
calls.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
mixed_events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
base_packet(PacketType::StartPingCheck),
|
||||||
|
simulator(&mixed_events.client),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(specific_calls.load(Ordering::Relaxed), 1);
|
||||||
|
assert!(
|
||||||
|
default_receiver
|
||||||
|
.recv_timeout(Duration::from_millis(50))
|
||||||
|
.is_err(),
|
||||||
|
"C# returns after a synchronous specific chain without queuing the default async chain"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
struct TestCapsMessage;
|
||||||
|
impl IMessage for TestCapsMessage {}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn caps_callbacks_are_ordered_and_subscription_guards_unregister() {
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let sim = simulator(&client);
|
||||||
|
let events = CapsEventDictionary::new(client).expect("CAPS events");
|
||||||
|
let order = Arc::new(Mutex::new(Vec::new()));
|
||||||
|
|
||||||
|
let default_order = Arc::clone(&order);
|
||||||
|
let default = CapsEventQueueCallback::from_handler(move |_, _, _| {
|
||||||
|
default_order.lock().unwrap().push("default");
|
||||||
|
});
|
||||||
|
let specific_order = Arc::clone(&order);
|
||||||
|
let specific = CapsEventQueueCallback::from_handler(move |_, _, _| {
|
||||||
|
specific_order.lock().unwrap().push("specific");
|
||||||
|
});
|
||||||
|
let default_guard = events.subscribe(String::new(), default);
|
||||||
|
let specific_guard = events.subscribe("EnableSimulator".into(), specific);
|
||||||
|
|
||||||
|
events.invoke_raise_event("EnableSimulator", &TestCapsMessage, sim.clone());
|
||||||
|
assert_eq!(*order.lock().unwrap(), ["default", "specific"]);
|
||||||
|
|
||||||
|
drop(specific_guard);
|
||||||
|
order.lock().unwrap().clear();
|
||||||
|
events.invoke_raise_event("EnableSimulator", &TestCapsMessage, sim);
|
||||||
|
assert_eq!(*order.lock().unwrap(), ["default"]);
|
||||||
|
drop(default_guard);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn concurrent_registration_and_removal_is_deterministic() {
|
||||||
|
const CALLBACKS: usize = 16;
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let sim = simulator(&client);
|
||||||
|
let events = Arc::new(PacketEventDictionary::new(client).unwrap());
|
||||||
|
let barrier = Arc::new(Barrier::new(CALLBACKS + 1));
|
||||||
|
let calls = Arc::new(AtomicUsize::new(0));
|
||||||
|
let (guard_sender, guard_receiver) = channel();
|
||||||
|
|
||||||
|
thread::scope(|scope| {
|
||||||
|
for _ in 0..CALLBACKS {
|
||||||
|
let events = Arc::clone(&events);
|
||||||
|
let barrier = Arc::clone(&barrier);
|
||||||
|
let calls = Arc::clone(&calls);
|
||||||
|
let guard_sender = guard_sender.clone();
|
||||||
|
scope.spawn(move || {
|
||||||
|
let guard = events.subscribe(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
Arc::new(move |_| {
|
||||||
|
calls.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}),
|
||||||
|
false,
|
||||||
|
);
|
||||||
|
guard_sender.send(guard).unwrap();
|
||||||
|
barrier.wait();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
barrier.wait();
|
||||||
|
});
|
||||||
|
drop(guard_sender);
|
||||||
|
let guards: Vec<_> = guard_receiver.into_iter().collect();
|
||||||
|
|
||||||
|
events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
base_packet(PacketType::StartPingCheck),
|
||||||
|
sim.clone(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(calls.load(Ordering::Relaxed), CALLBACKS);
|
||||||
|
|
||||||
|
drop(guards);
|
||||||
|
calls.store(0, Ordering::Relaxed);
|
||||||
|
events
|
||||||
|
.invoke_raise_event(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
base_packet(PacketType::StartPingCheck),
|
||||||
|
sim,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(calls.load(Ordering::Relaxed), 0);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn sim_connecting_can_cancel_without_leaking_collection_state() {
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let manager = NetworkManager::new(client).expect("manager");
|
||||||
|
let guard = manager.subscribe_sim_connecting(Arc::new(|mut args| args.set_cancel(true)));
|
||||||
|
|
||||||
|
let connected = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
loopback(9),
|
||||||
|
0x2000,
|
||||||
|
true,
|
||||||
|
Some(Uri("http://caps.invalid/seed".into())),
|
||||||
|
512,
|
||||||
|
256,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert!(connected.is_none());
|
||||||
|
assert!(manager.simulators.is_empty());
|
||||||
|
assert!(
|
||||||
|
manager.connected(),
|
||||||
|
"C# marks the manager connected before raising the cancellable event"
|
||||||
|
);
|
||||||
|
assert!(manager.current_sim().is_none());
|
||||||
|
drop(guard);
|
||||||
|
}
|
||||||
|
|
||||||
|
enum ServerCommand {
|
||||||
|
Ping,
|
||||||
|
Disable,
|
||||||
|
Stop,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Eq, PartialEq)]
|
||||||
|
enum ServerReport {
|
||||||
|
Ack(u32),
|
||||||
|
Circuit(u32),
|
||||||
|
HandshakeReply(u32),
|
||||||
|
Ping(u8),
|
||||||
|
}
|
||||||
|
|
||||||
|
struct FakeServer {
|
||||||
|
endpoint: SocketAddr,
|
||||||
|
commands: Sender<ServerCommand>,
|
||||||
|
reports: Receiver<ServerReport>,
|
||||||
|
handle: Option<JoinHandle<()>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for FakeServer {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
let _ = self.commands.send(ServerCommand::Stop);
|
||||||
|
if let Some(handle) = self.handle.take() {
|
||||||
|
let _ = handle.join();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn packet_type(bytes: &[u8]) -> Option<PacketType> {
|
||||||
|
let mut packet_end = i32::try_from(bytes.len()).ok()?.checked_sub(1)?;
|
||||||
|
Packet::build_packet_with_bytes_int32_bytes(bytes.to_vec(), &mut packet_end, vec![0; 8192])
|
||||||
|
.ok()
|
||||||
|
.map(|packet| packet.type_)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn ack(sequence: u32) -> Vec<u8> {
|
||||||
|
let mut packet = PacketAckPacket::new_with_constructor().unwrap();
|
||||||
|
packet.packets = vec![PacketAckPacketPacketsBlock { id: sequence }];
|
||||||
|
let mut bytes = packet.to_bytes_with_method().unwrap();
|
||||||
|
bytes[0] &= !(Helpers::MSG_RELIABLE | Helpers::MSG_ZEROCODED);
|
||||||
|
bytes
|
||||||
|
}
|
||||||
|
|
||||||
|
fn serve_ping_command(
|
||||||
|
socket: &UdpSocket,
|
||||||
|
client_endpoint: SocketAddr,
|
||||||
|
report_sender: &Sender<ServerReport>,
|
||||||
|
buffer: &mut [u8],
|
||||||
|
) {
|
||||||
|
let mut ping = StartPingCheckPacket::new_with_constructor().unwrap();
|
||||||
|
ping.ping_id.ping_id = 42;
|
||||||
|
ping.ping_id.oldest_unacked = 1;
|
||||||
|
let mut bytes = ping.to_bytes_with_method().unwrap();
|
||||||
|
bytes[0] &= !Helpers::MSG_ZEROCODED;
|
||||||
|
bytes[0] |= Helpers::MSG_RELIABLE;
|
||||||
|
bytes[1..5].copy_from_slice(&99_u32.to_be_bytes());
|
||||||
|
socket.send_to(&bytes, client_endpoint).unwrap();
|
||||||
|
|
||||||
|
loop {
|
||||||
|
let (length, _) = socket.recv_from(buffer).expect("ping response");
|
||||||
|
match packet_type(&buffer[..length]) {
|
||||||
|
Some(PacketType::PacketAck) => {
|
||||||
|
let mut position = 0;
|
||||||
|
let response =
|
||||||
|
PacketAckPacket::new_with_bytes_int32(buffer[..length].to_vec(), &mut position)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(response.packets.len(), 1);
|
||||||
|
report_sender
|
||||||
|
.send(ServerReport::Ack(response.packets[0].id))
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
Some(PacketType::CompletePingCheck) => {
|
||||||
|
let mut position = 0;
|
||||||
|
let response =
|
||||||
|
libremetaverse::packets::CompletePingCheckPacket::new_with_bytes_int32(
|
||||||
|
buffer[..length].to_vec(),
|
||||||
|
&mut position,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
report_sender
|
||||||
|
.send(ServerReport::Ping(response.ping_id.ping_id))
|
||||||
|
.unwrap();
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn spawn_fake_server() -> FakeServer {
|
||||||
|
let socket = UdpSocket::bind(loopback(0)).expect("fake UDP server");
|
||||||
|
socket
|
||||||
|
.set_read_timeout(Some(Duration::from_secs(2)))
|
||||||
|
.unwrap();
|
||||||
|
let endpoint = socket.local_addr().unwrap();
|
||||||
|
let (command_sender, command_receiver) = channel();
|
||||||
|
let (report_sender, report_receiver) = channel();
|
||||||
|
let handle = thread::spawn(move || {
|
||||||
|
let mut buffer = [0_u8; 4096];
|
||||||
|
let (length, client_endpoint) = socket.recv_from(&mut buffer).expect("UseCircuitCode");
|
||||||
|
assert_eq!(
|
||||||
|
packet_type(&buffer[..length]),
|
||||||
|
Some(PacketType::UseCircuitCode)
|
||||||
|
);
|
||||||
|
let mut position = 0;
|
||||||
|
let use_circuit =
|
||||||
|
UseCircuitCodePacket::new_with_bytes_int32(buffer[..length].to_vec(), &mut position)
|
||||||
|
.expect("decode UseCircuitCode");
|
||||||
|
report_sender
|
||||||
|
.send(ServerReport::Circuit(use_circuit.circuit_code.code))
|
||||||
|
.unwrap();
|
||||||
|
let sequence = u32::from_be_bytes(buffer[1..5].try_into().unwrap());
|
||||||
|
socket.send_to(&ack(sequence), client_endpoint).unwrap();
|
||||||
|
|
||||||
|
let mut handshake = RegionHandshakePacket::new_with_constructor().unwrap();
|
||||||
|
handshake.region_info.sim_name = b"Fake Region".to_vec();
|
||||||
|
let mut bytes = handshake.to_bytes_with_method().unwrap();
|
||||||
|
bytes[0] &= !(Helpers::MSG_RELIABLE | Helpers::MSG_ZEROCODED);
|
||||||
|
socket.send_to(&bytes, client_endpoint).unwrap();
|
||||||
|
|
||||||
|
loop {
|
||||||
|
let (length, _) = socket.recv_from(&mut buffer).expect("handshake reply");
|
||||||
|
let received_type = packet_type(&buffer[..length]);
|
||||||
|
if received_type == Some(PacketType::RegionHandshakeReply) {
|
||||||
|
let mut position = 0;
|
||||||
|
let mut packet_end = i32::try_from(length).unwrap() - 1;
|
||||||
|
let mut zero_buffer = vec![0_u8; 8192];
|
||||||
|
let mut response = RegionHandshakeReplyPacket::new_with_constructor().unwrap();
|
||||||
|
response
|
||||||
|
.from_bytes_with_bytes_int32_int32_bytes(
|
||||||
|
buffer[..length].to_vec(),
|
||||||
|
&mut position,
|
||||||
|
&mut packet_end,
|
||||||
|
Some(&mut zero_buffer),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
report_sender
|
||||||
|
.send(ServerReport::HandshakeReply(response.region_info.flags))
|
||||||
|
.unwrap();
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
while let Ok(command) = command_receiver.recv() {
|
||||||
|
match command {
|
||||||
|
ServerCommand::Ping => {
|
||||||
|
serve_ping_command(&socket, client_endpoint, &report_sender, &mut buffer);
|
||||||
|
}
|
||||||
|
ServerCommand::Disable => {
|
||||||
|
let disable = DisableSimulatorPacket::new_with_constructor().unwrap();
|
||||||
|
let mut bytes = disable.to_bytes_with_method().unwrap();
|
||||||
|
bytes[0] &= !(Helpers::MSG_RELIABLE | Helpers::MSG_ZEROCODED);
|
||||||
|
socket.send_to(&bytes, client_endpoint).unwrap();
|
||||||
|
}
|
||||||
|
ServerCommand::Stop => break,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
FakeServer {
|
||||||
|
endpoint,
|
||||||
|
commands: command_sender,
|
||||||
|
reports: report_receiver,
|
||||||
|
handle: Some(handle),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn spawn_ack_only_server() -> (SocketAddr, JoinHandle<()>) {
|
||||||
|
let socket = UdpSocket::bind(loopback(0)).expect("ACK-only UDP server");
|
||||||
|
socket
|
||||||
|
.set_read_timeout(Some(Duration::from_secs(2)))
|
||||||
|
.unwrap();
|
||||||
|
let endpoint = socket.local_addr().unwrap();
|
||||||
|
let handle = thread::spawn(move || {
|
||||||
|
let mut buffer = [0_u8; 4096];
|
||||||
|
let (length, client_endpoint) = socket.recv_from(&mut buffer).expect("UseCircuitCode");
|
||||||
|
assert_eq!(
|
||||||
|
packet_type(&buffer[..length]),
|
||||||
|
Some(PacketType::UseCircuitCode)
|
||||||
|
);
|
||||||
|
let sequence = u32::from_be_bytes(buffer[1..5].try_into().unwrap());
|
||||||
|
socket.send_to(&ack(sequence), client_endpoint).unwrap();
|
||||||
|
});
|
||||||
|
(endpoint, handle)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn assert_ping_exchange(server: &FakeServer, packet_receiver: &Receiver<PacketType>) {
|
||||||
|
server.commands.send(ServerCommand::Ping).unwrap();
|
||||||
|
assert_eq!(
|
||||||
|
packet_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
PacketType::StartPingCheck
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
server.reports.recv_timeout(Duration::from_secs(1)).unwrap(),
|
||||||
|
ServerReport::Ack(99)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
server.reports.recv_timeout(Duration::from_secs(1)).unwrap(),
|
||||||
|
ServerReport::Ping(42)
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn assert_reuse_preserves_connected_override(
|
||||||
|
manager: &mut NetworkManager,
|
||||||
|
server: &FakeServer,
|
||||||
|
simulator: &Simulator,
|
||||||
|
) {
|
||||||
|
manager.set_connected(false);
|
||||||
|
let reused = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
server.endpoint,
|
||||||
|
simulator.handle,
|
||||||
|
false,
|
||||||
|
None,
|
||||||
|
simulator.size_x,
|
||||||
|
simulator.size_y,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(reused, simulator.clone());
|
||||||
|
assert!(
|
||||||
|
!manager.connected(),
|
||||||
|
"C# does not alter manager Connected when reusing a connected simulator"
|
||||||
|
);
|
||||||
|
manager.set_connected(true);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn fake_server_drives_circuit_handshake_ping_disable_and_disconnect_reasons() {
|
||||||
|
let server = spawn_fake_server();
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(0x1122_3344);
|
||||||
|
let (changed_sender, changed_receiver) = channel();
|
||||||
|
let (connected_sender, connected_receiver) = channel();
|
||||||
|
let (packet_sender, packet_receiver) = channel();
|
||||||
|
let (sim_disconnected_sender, sim_disconnected_receiver) = channel();
|
||||||
|
let (disconnected_sender, disconnected_receiver) = channel();
|
||||||
|
let guards = [
|
||||||
|
manager.subscribe_sim_changed(Arc::new(move |args| {
|
||||||
|
changed_sender
|
||||||
|
.send(args.previous_simulator().is_none())
|
||||||
|
.unwrap();
|
||||||
|
})),
|
||||||
|
manager.subscribe_sim_connected(Arc::new(move |args| {
|
||||||
|
connected_sender.send(args.simulator().handle).unwrap();
|
||||||
|
})),
|
||||||
|
manager.subscribe_sim_disconnected(Arc::new(move |args| {
|
||||||
|
sim_disconnected_sender
|
||||||
|
.send((args.reason(), args.was_connected()))
|
||||||
|
.unwrap();
|
||||||
|
})),
|
||||||
|
manager.subscribe_disconnected(Arc::new(move |args| {
|
||||||
|
disconnected_sender.send(args.reason()).unwrap();
|
||||||
|
})),
|
||||||
|
];
|
||||||
|
manager
|
||||||
|
.register_callback_with_packet_type_event_handler_boolean(
|
||||||
|
PacketType::StartPingCheck,
|
||||||
|
Arc::new(move |args| packet_sender.send(args.packet().type_).unwrap()),
|
||||||
|
false,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let seed = Uri("http://caps.invalid/seed".into());
|
||||||
|
let sim = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
server.endpoint,
|
||||||
|
0x1234_0000,
|
||||||
|
true,
|
||||||
|
Some(seed.clone()),
|
||||||
|
512,
|
||||||
|
256,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.expect("connected simulator");
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
server.reports.recv_timeout(Duration::from_secs(1)).unwrap(),
|
||||||
|
ServerReport::Circuit(0x1122_3344)
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
changed_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap()
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
connected_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
sim.handle
|
||||||
|
);
|
||||||
|
assert_eq!(manager.current_sim(), Some(sim.clone()));
|
||||||
|
assert_eq!(sim.seed_capability(), Some(seed));
|
||||||
|
assert_eq!(sim.size_x, 512);
|
||||||
|
assert_eq!(sim.size_y, 256);
|
||||||
|
wait_until(|| sim.handshake_complete());
|
||||||
|
assert_eq!(
|
||||||
|
server.reports.recv_timeout(Duration::from_secs(1)).unwrap(),
|
||||||
|
ServerReport::HandshakeReply(0x1 | 0x2 | 0x4)
|
||||||
|
);
|
||||||
|
|
||||||
|
assert_reuse_preserves_connected_override(&mut manager, &server, &sim);
|
||||||
|
|
||||||
|
assert_ping_exchange(&server, &packet_receiver);
|
||||||
|
|
||||||
|
server.commands.send(ServerCommand::Disable).unwrap();
|
||||||
|
assert_eq!(
|
||||||
|
sim_disconnected_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
(NetworkManagerDisconnectType::NetworkTimeout, true)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
disconnected_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
NetworkManagerDisconnectType::SimShutdown
|
||||||
|
);
|
||||||
|
assert!(manager.simulators.is_empty());
|
||||||
|
assert!(manager.current_sim().is_none());
|
||||||
|
assert!(!manager.connected());
|
||||||
|
drop(guards);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn concurrent_disconnect_emits_each_transition_once() {
|
||||||
|
const CALLERS: usize = 12;
|
||||||
|
let server = spawn_fake_server();
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(7);
|
||||||
|
let sim_disconnects = Arc::new(AtomicUsize::new(0));
|
||||||
|
let global_disconnects = Arc::new(AtomicUsize::new(0));
|
||||||
|
let sim_count = Arc::clone(&sim_disconnects);
|
||||||
|
let sim_guard = manager.subscribe_sim_disconnected(Arc::new(move |args| {
|
||||||
|
assert!(args.was_connected());
|
||||||
|
sim_count.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}));
|
||||||
|
let global_count = Arc::clone(&global_disconnects);
|
||||||
|
let global_guard = manager.subscribe_disconnected(Arc::new(move |_| {
|
||||||
|
global_count.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}));
|
||||||
|
let sim = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
server.endpoint,
|
||||||
|
9,
|
||||||
|
true,
|
||||||
|
None,
|
||||||
|
256,
|
||||||
|
256,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
let manager = Arc::new(manager);
|
||||||
|
let barrier = Arc::new(Barrier::new(CALLERS));
|
||||||
|
|
||||||
|
thread::scope(|scope| {
|
||||||
|
for _ in 0..CALLERS {
|
||||||
|
let manager = Arc::clone(&manager);
|
||||||
|
let simulator = sim.clone();
|
||||||
|
let barrier = Arc::clone(&barrier);
|
||||||
|
scope.spawn(move || {
|
||||||
|
barrier.wait();
|
||||||
|
manager.disconnect_sim(simulator, false).unwrap();
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
assert_eq!(sim_disconnects.load(Ordering::Relaxed), 1);
|
||||||
|
assert_eq!(global_disconnects.load(Ordering::Relaxed), 1);
|
||||||
|
assert!(manager.simulators.is_empty());
|
||||||
|
assert!(!manager.connected());
|
||||||
|
drop(sim_guard);
|
||||||
|
drop(global_guard);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn sim_disconnected_callback_can_reenter_disconnect_without_deadlock() {
|
||||||
|
let server = spawn_fake_server();
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(8);
|
||||||
|
let simulator = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
server.endpoint,
|
||||||
|
10,
|
||||||
|
true,
|
||||||
|
None,
|
||||||
|
256,
|
||||||
|
256,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
let manager = Arc::new(manager);
|
||||||
|
let weak_manager = Arc::downgrade(&manager);
|
||||||
|
let callback_completed = Arc::new(AtomicBool::new(false));
|
||||||
|
let completed = Arc::clone(&callback_completed);
|
||||||
|
let guard = manager.subscribe_sim_disconnected(Arc::new(move |args| {
|
||||||
|
let manager = weak_manager.upgrade().expect("manager during callback");
|
||||||
|
manager
|
||||||
|
.disconnect_sim(args.simulator(), false)
|
||||||
|
.expect("reentrant disconnect");
|
||||||
|
completed.store(true, Ordering::Release);
|
||||||
|
}));
|
||||||
|
|
||||||
|
manager.disconnect_sim(simulator, false).unwrap();
|
||||||
|
|
||||||
|
assert!(callback_completed.load(Ordering::Acquire));
|
||||||
|
assert!(manager.simulators.is_empty());
|
||||||
|
assert!(!manager.connected());
|
||||||
|
drop(guard);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn keepalive_requires_two_silent_intervals_before_network_timeout() {
|
||||||
|
let server = spawn_fake_server();
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(11);
|
||||||
|
let (disconnected_sender, disconnected_receiver) = channel();
|
||||||
|
let guard = manager.subscribe_disconnected(Arc::new(move |args| {
|
||||||
|
disconnected_sender.send(args.reason()).unwrap();
|
||||||
|
}));
|
||||||
|
let sim = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
server.endpoint,
|
||||||
|
12,
|
||||||
|
true,
|
||||||
|
None,
|
||||||
|
256,
|
||||||
|
256,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
wait_until(|| sim.handshake_complete());
|
||||||
|
|
||||||
|
manager.poll_keepalive();
|
||||||
|
assert!(manager.connected());
|
||||||
|
assert!(
|
||||||
|
disconnected_receiver
|
||||||
|
.recv_timeout(Duration::from_millis(50))
|
||||||
|
.is_err()
|
||||||
|
);
|
||||||
|
|
||||||
|
manager.poll_keepalive();
|
||||||
|
assert_eq!(
|
||||||
|
disconnected_receiver
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
NetworkManagerDisconnectType::NetworkTimeout
|
||||||
|
);
|
||||||
|
assert!(!manager.connected());
|
||||||
|
assert!(manager.current_sim().is_none());
|
||||||
|
drop(guard);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn enable_simulator_caps_event_adds_each_new_region_once() {
|
||||||
|
let first_server = spawn_fake_server();
|
||||||
|
let second_server = spawn_fake_server();
|
||||||
|
let mut client = GridClient::new().expect("client");
|
||||||
|
client.settings().agent_settings_mut().multiple_sims = true;
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(17);
|
||||||
|
let current = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
first_server.endpoint,
|
||||||
|
18,
|
||||||
|
true,
|
||||||
|
None,
|
||||||
|
256,
|
||||||
|
256,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
let mut message = EnableSimulatorMessage::new().unwrap();
|
||||||
|
let mut info = EnableSimulatorMessageSimulatorInfoBlock::new().unwrap();
|
||||||
|
info.ip = second_server.endpoint.ip();
|
||||||
|
info.port = i32::from(second_server.endpoint.port());
|
||||||
|
info.region_handle = 19;
|
||||||
|
info.region_size_x = 1024;
|
||||||
|
info.region_size_y = 512;
|
||||||
|
message.simulators = vec![info];
|
||||||
|
|
||||||
|
manager.dispatch_caps_event("EnableSimulator", &message, current.clone());
|
||||||
|
manager.dispatch_caps_event("EnableSimulator", &message, current);
|
||||||
|
|
||||||
|
assert_eq!(manager.simulators.len(), 2);
|
||||||
|
let added = manager
|
||||||
|
.find_simulator_with_ip_end_point(second_server.endpoint)
|
||||||
|
.unwrap()
|
||||||
|
.expect("enabled simulator");
|
||||||
|
assert_eq!(added.handle, 19);
|
||||||
|
assert_eq!(added.size_x, 1024);
|
||||||
|
assert_eq!(added.size_y, 512);
|
||||||
|
assert_eq!(
|
||||||
|
second_server
|
||||||
|
.reports
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
ServerReport::Circuit(17)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
second_server
|
||||||
|
.reports
|
||||||
|
.recv_timeout(Duration::from_secs(1))
|
||||||
|
.unwrap(),
|
||||||
|
ServerReport::HandshakeReply(0x1 | 0x2 | 0x4)
|
||||||
|
);
|
||||||
|
|
||||||
|
manager
|
||||||
|
.shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated)
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn async_connect_and_shutdown_do_not_require_an_ambient_tokio_runtime() {
|
||||||
|
let server = spawn_fake_server();
|
||||||
|
let client = GridClient::new().expect("client");
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(23);
|
||||||
|
|
||||||
|
let simulator = block_on(
|
||||||
|
manager.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32_20e51cc0(
|
||||||
|
server.endpoint,
|
||||||
|
24,
|
||||||
|
true,
|
||||||
|
None,
|
||||||
|
256,
|
||||||
|
256,
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.expect("async connection");
|
||||||
|
|
||||||
|
assert!(simulator.connected());
|
||||||
|
assert_eq!(manager.current_sim(), Some(simulator));
|
||||||
|
block_on(manager.shutdown_with_disconnect_type_string_7641243c(
|
||||||
|
NetworkManagerDisconnectType::ClientInitiated,
|
||||||
|
"ClientInitiated".into(),
|
||||||
|
))
|
||||||
|
.unwrap();
|
||||||
|
assert!(!manager.connected());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn handshake_timeout_preserves_reference_connect_result_and_shutdown_tears_down_current_sim() {
|
||||||
|
let (endpoint, server) = spawn_ack_only_server();
|
||||||
|
let mut client = GridClient::new().expect("client");
|
||||||
|
client.settings().timing().login_timeout = 25;
|
||||||
|
let mut manager = NetworkManager::new(client).expect("manager");
|
||||||
|
manager.set_circuit_code(25);
|
||||||
|
let sim_disconnects = Arc::new(AtomicUsize::new(0));
|
||||||
|
let disconnect_count = Arc::clone(&sim_disconnects);
|
||||||
|
let guard = manager.subscribe_sim_disconnected(Arc::new(move |_| {
|
||||||
|
disconnect_count.fetch_add(1, Ordering::Relaxed);
|
||||||
|
}));
|
||||||
|
|
||||||
|
let simulator = manager
|
||||||
|
.connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32(
|
||||||
|
endpoint, 26, true, None, 256, 256,
|
||||||
|
)
|
||||||
|
.unwrap()
|
||||||
|
.expect("C# returns the UDP-connected simulator after handshake timeout");
|
||||||
|
server.join().unwrap();
|
||||||
|
|
||||||
|
assert!(!simulator.handshake_complete());
|
||||||
|
assert!(manager.simulators.is_empty());
|
||||||
|
assert_eq!(manager.current_sim(), Some(simulator.clone()));
|
||||||
|
manager
|
||||||
|
.shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated)
|
||||||
|
.unwrap();
|
||||||
|
assert!(!simulator.connected());
|
||||||
|
assert_eq!(sim_disconnects.load(Ordering::Relaxed), 1);
|
||||||
|
drop(guard);
|
||||||
|
}
|
||||||
79
docs/network-manager.md
Normal file
79
docs/network-manager.md
Normal file
@@ -0,0 +1,79 @@
|
|||||||
|
# Native network manager and simulator lifecycle
|
||||||
|
|
||||||
|
`NetworkManager`, `Simulator`, `PacketEventDictionary`, and
|
||||||
|
`CapsEventDictionary` implement the native Rust boundary corresponding to the
|
||||||
|
LibreMetaverse connection and dispatch layer. They use the generated packet
|
||||||
|
codecs and the bounded `UDPBase` transport; they do not invoke .NET code or
|
||||||
|
start a helper process.
|
||||||
|
|
||||||
|
## Ownership and execution
|
||||||
|
|
||||||
|
A manager owns bounded incoming and outgoing channels plus one keepalive
|
||||||
|
worker. Each simulator owns a `UDPBase` transport running on a dedicated
|
||||||
|
current-thread Tokio executor so the mapped synchronous connection API does not
|
||||||
|
depend on the caller already having entered a runtime. Async connection methods
|
||||||
|
move blocking compatibility work to short-lived named standard threads and
|
||||||
|
return their result through a runtime-neutral one-shot future. Every worker
|
||||||
|
holds only a weak manager reference between operations. Dropping the last
|
||||||
|
manager and subscription guard therefore releases all manager-owned
|
||||||
|
`GridClient` handles, cancels the workers, joins their threads, and closes
|
||||||
|
simulator transports.
|
||||||
|
|
||||||
|
The public `SimulatorCollection` is a shared snapshot collection. Lookup,
|
||||||
|
addition, removal, current-simulator selection, and disconnect transitions are
|
||||||
|
synchronized. User callbacks receive cloned handles only after internal locks
|
||||||
|
have been released. `Subscription` is an RAII guard: dropping or closing it
|
||||||
|
removes exactly its registration through a weak registry reference.
|
||||||
|
|
||||||
|
## Dispatch semantics
|
||||||
|
|
||||||
|
Packet callbacks preserve registration order. `PacketType::Default` handlers
|
||||||
|
are selected before handlers for the concrete packet type. Once any callback
|
||||||
|
registered for a packet type requests asynchronous dispatch, that packet type
|
||||||
|
remains asynchronous after later removals, matching the conservative C#
|
||||||
|
delegate state. A bounded worker executes each asynchronous default/specific
|
||||||
|
batch in order. Callback panics are isolated and cannot terminate a processor
|
||||||
|
or poison a callback registry.
|
||||||
|
|
||||||
|
CAPS dispatch snapshots the default (`""`) and named callback chains, then
|
||||||
|
invokes them in that order without holding the registry lock. The built-in
|
||||||
|
`EnableSimulator` handler honors `Agent.MultipleSims`, ignores duplicate
|
||||||
|
endpoints, validates ports, and connects each newly advertised region with its
|
||||||
|
reported handle and dimensions. Incoming `DisableSimulator`, `KickUser`,
|
||||||
|
`RegionHandshake`, `StartPingCheck`, and generic-streaming UDP packets perform
|
||||||
|
their reference lifecycle actions before user dispatch. Ping replies preserve
|
||||||
|
the incoming ping identifier and immediately flush pending acknowledgements
|
||||||
|
when the simulator reports a nonzero `OldestUnacked` sequence.
|
||||||
|
|
||||||
|
## Connection and shutdown behavior
|
||||||
|
|
||||||
|
Connecting reuses an existing endpoint or inserts one simulator atomically,
|
||||||
|
starts the bounded manager workers, raises cancellable `SimConnecting`, sends a
|
||||||
|
reliable `UseCircuitCode`, and waits up to `Timing.LoginTimeout` for its actual
|
||||||
|
protocol ACK. Async compatibility methods use a runtime-neutral blocking
|
||||||
|
bridge, so they can be polled without an ambient Tokio runtime. Circuit-code
|
||||||
|
changes propagate to every tracked simulator. A default connection stores its
|
||||||
|
seed capability, updates `CurrentSim`, raises `SimChanged`, and then raises
|
||||||
|
`SimConnected`. A received region handshake is decoded and answered with the
|
||||||
|
reference `RegionHandshakeReply` flags before the simulator is marked complete.
|
||||||
|
|
||||||
|
The keepalive uses the same two-interval disconnect-candidate transition as the
|
||||||
|
C# timer: traffic clears the candidate, the first silent interval marks it and
|
||||||
|
sends a ping, and the second shuts the manager down with `NetworkTimeout`.
|
||||||
|
Concurrent disconnect attempts are serialized so one simulator and one global
|
||||||
|
transition are emitted. Per-simulator removal reports `NetworkTimeout`; losing
|
||||||
|
the last simulator reports `SimShutdown`. Explicit shutdown preserves the
|
||||||
|
caller-provided reason/message and sends `CloseCircuit` only for client or
|
||||||
|
network-timeout shutdowns.
|
||||||
|
|
||||||
|
The deterministic tests use loopback fake simulators and require no live grid:
|
||||||
|
|
||||||
|
```sh
|
||||||
|
cargo test -p libremetaverse --test network_manager
|
||||||
|
```
|
||||||
|
|
||||||
|
They cover ACK-gated circuit setup, typed handshake reply/state, ping response,
|
||||||
|
CAPS enable, UDP disable, current/seed selection, callback ordering and
|
||||||
|
filtering, RAII unregistration, runtime-neutral async calls, concurrent
|
||||||
|
registration/removal, reentrant/concurrent disconnect, and keepalive timeout
|
||||||
|
reasons.
|
||||||
@@ -38,6 +38,19 @@ TARGETS = {
|
|||||||
# implementations. The generated module keeps catalog markers and re-exports
|
# implementations. The generated module keeps catalog markers and re-exports
|
||||||
# the hand-written type so coverage remains deterministic.
|
# the hand-written type so coverage remains deterministic.
|
||||||
NATIVE_TYPES = {
|
NATIVE_TYPES = {
|
||||||
|
"T:LibreMetaverse.Caps.EventQueueCallback": "crate::network_manager::CapsEventQueueCallback",
|
||||||
|
"T:LibreMetaverse.CapsEventDictionary": "crate::network_manager::CapsEventDictionary",
|
||||||
|
"T:LibreMetaverse.DisconnectedEventArgs": "crate::network_manager::DisconnectedEventArgs",
|
||||||
|
"T:LibreMetaverse.GenericStreamingMessageEventArgs": "crate::network_manager::GenericStreamingMessageEventArgs",
|
||||||
|
"T:LibreMetaverse.NetworkManager.IncomingPacket": "crate::network_manager::NetworkManagerIncomingPacket",
|
||||||
|
"T:LibreMetaverse.NetworkManager.OutgoingPacket": "crate::network_manager::NetworkManagerOutgoingPacket",
|
||||||
|
"T:LibreMetaverse.PacketEventDictionary": "crate::network_manager::PacketEventDictionary",
|
||||||
|
"T:LibreMetaverse.PacketReceivedEventArgs": "crate::network_manager::PacketReceivedEventArgs",
|
||||||
|
"T:LibreMetaverse.PacketSentEventArgs": "crate::network_manager::PacketSentEventArgs",
|
||||||
|
"T:LibreMetaverse.SimChangedEventArgs": "crate::network_manager::SimChangedEventArgs",
|
||||||
|
"T:LibreMetaverse.SimConnectedEventArgs": "crate::network_manager::SimConnectedEventArgs",
|
||||||
|
"T:LibreMetaverse.SimConnectingEventArgs": "crate::network_manager::SimConnectingEventArgs",
|
||||||
|
"T:LibreMetaverse.SimDisconnectedEventArgs": "crate::network_manager::SimDisconnectedEventArgs",
|
||||||
"T:LibreMetaverse.ArchetypeParam": "crate::genepool_catalog::ArchetypeParam",
|
"T:LibreMetaverse.ArchetypeParam": "crate::genepool_catalog::ArchetypeParam",
|
||||||
"T:LibreMetaverse.AttentionData": "crate::attention_catalog::AttentionData",
|
"T:LibreMetaverse.AttentionData": "crate::attention_catalog::AttentionData",
|
||||||
"T:LibreMetaverse.AttentionSet": "crate::attention_catalog::AttentionSet",
|
"T:LibreMetaverse.AttentionSet": "crate::attention_catalog::AttentionSet",
|
||||||
@@ -132,7 +145,9 @@ NATIVE_TYPES = {
|
|||||||
# hand-written, while its fixed public methods remain generator-audited.
|
# hand-written, while its fixed public methods remain generator-audited.
|
||||||
NATIVE_DECLARATIONS = {
|
NATIVE_DECLARATIONS = {
|
||||||
"T:LibreMetaverse.GridClient": "crate::client_core::GridClient",
|
"T:LibreMetaverse.GridClient": "crate::client_core::GridClient",
|
||||||
|
"T:LibreMetaverse.NetworkManager": "crate::network_manager::NetworkManager",
|
||||||
"T:LibreMetaverse.Settings": "crate::client_core::Settings",
|
"T:LibreMetaverse.Settings": "crate::client_core::Settings",
|
||||||
|
"T:LibreMetaverse.Simulator": "crate::network_manager::Simulator",
|
||||||
"T:LibreMetaverse.StructuredData.OSDParser": "crate::model::OSDParser",
|
"T:LibreMetaverse.StructuredData.OSDParser": "crate::model::OSDParser",
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -163,6 +178,122 @@ NATIVE_MEMBER_BODIES = {
|
|||||||
"self.shutdown().map_err(Into::into)",
|
"self.shutdown().map_err(Into::into)",
|
||||||
"M:LibreMetaverse.GridClient.ToString":
|
"M:LibreMetaverse.GridClient.ToString":
|
||||||
"String::new()",
|
"String::new()",
|
||||||
|
"E:LibreMetaverse.NetworkManager.Disconnected":
|
||||||
|
"self.native_subscribe_disconnected(handler)",
|
||||||
|
"E:LibreMetaverse.NetworkManager.GenericStreamingMessage":
|
||||||
|
"self.native_subscribe_generic_streaming_message(handler)",
|
||||||
|
"E:LibreMetaverse.NetworkManager.PacketSent":
|
||||||
|
"self.native_subscribe_packet_sent(handler)",
|
||||||
|
"E:LibreMetaverse.NetworkManager.SimChanged":
|
||||||
|
"self.native_subscribe_sim_changed(handler)",
|
||||||
|
"E:LibreMetaverse.NetworkManager.SimConnected":
|
||||||
|
"self.native_subscribe_sim_connected(handler)",
|
||||||
|
"E:LibreMetaverse.NetworkManager.SimConnecting":
|
||||||
|
"self.native_subscribe_sim_connecting(handler)",
|
||||||
|
"E:LibreMetaverse.NetworkManager.SimDisconnected":
|
||||||
|
"self.native_subscribe_sim_disconnected(handler)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.#ctor(LibreMetaverse.GridClient)":
|
||||||
|
"crate::network_manager::NetworkManager::native_new(client)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.Connect(System.Net.IPAddress,System.UInt16,System.UInt64,System.Boolean,System.Uri)":
|
||||||
|
"self.native_connect(std::net::SocketAddr::new(ip, port), handle, set_default, seedcaps, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_X, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_Y)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.Connect(System.Net.IPAddress,System.UInt16,System.UInt64,System.Boolean,System.Uri,System.UInt32,System.UInt32)":
|
||||||
|
"self.native_connect(std::net::SocketAddr::new(ip, port), handle, set_default, seedcaps, size_x, size_y)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.Connect(System.Net.IPEndPoint,System.UInt64,System.Boolean,System.Uri)":
|
||||||
|
"self.native_connect(end_point, handle, set_default, seedcaps, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_X, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_Y)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.Connect(System.Net.IPEndPoint,System.UInt64,System.Boolean,System.Uri,System.UInt32,System.UInt32)":
|
||||||
|
"self.native_connect(end_point, handle, set_default, seedcaps, size_x, size_y)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.ConnectAsync(System.Net.IPAddress,System.UInt16,System.UInt64,System.Boolean,System.Uri)":
|
||||||
|
"self.native_connect_async(std::net::SocketAddr::new(ip, port), handle, set_default, seedcaps, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_X, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_Y).await",
|
||||||
|
"M:LibreMetaverse.NetworkManager.ConnectAsync(System.Net.IPEndPoint,System.UInt64,System.Boolean,System.Uri,System.UInt32,System.UInt32)":
|
||||||
|
"self.native_connect_async(end_point, handle, set_default, seedcaps, size_x, size_y).await",
|
||||||
|
"M:LibreMetaverse.NetworkManager.DisconnectSim(LibreMetaverse.Simulator,System.Boolean)":
|
||||||
|
"self.native_disconnect_sim(simulator, send_close_circuit)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.EnqueueIncoming(LibreMetaverse.NetworkManager.IncomingPacket)":
|
||||||
|
"self.native_enqueue_incoming(packet)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.EnqueueOutgoing(LibreMetaverse.NetworkManager.OutgoingPacket)":
|
||||||
|
"self.native_enqueue_outgoing(packet)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.FindSimulator(System.Net.IPEndPoint)":
|
||||||
|
"Ok(self.native_find_endpoint(end_point))",
|
||||||
|
"M:LibreMetaverse.NetworkManager.FindSimulator(System.UInt64)":
|
||||||
|
"Ok(self.native_find_handle(handle))",
|
||||||
|
"M:LibreMetaverse.NetworkManager.RegisterCallback(LibreMetaverse.Packets.PacketType,System.EventHandler{LibreMetaverse.PacketReceivedEventArgs})":
|
||||||
|
"self.native_register_callback(type_, callback, true)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.RegisterCallback(LibreMetaverse.Packets.PacketType,System.EventHandler{LibreMetaverse.PacketReceivedEventArgs},System.Boolean)":
|
||||||
|
"self.native_register_callback(type_, callback, is_async)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.RegisterEventCallback(System.String,LibreMetaverse.Caps.EventQueueCallback)":
|
||||||
|
"self.native_register_caps(caps_event, callback)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.SendPacket(LibreMetaverse.Packets.Packet)":
|
||||||
|
"self.native_send_packet(packet, None)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.SendPacket(LibreMetaverse.Packets.Packet,LibreMetaverse.Simulator)":
|
||||||
|
"self.native_send_packet(packet, simulator)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.Shutdown(LibreMetaverse.NetworkManager.DisconnectType)":
|
||||||
|
"self.native_shutdown(type_, format!(\"{type_:?}\"))",
|
||||||
|
"M:LibreMetaverse.NetworkManager.Shutdown(LibreMetaverse.NetworkManager.DisconnectType,System.String)":
|
||||||
|
"self.native_shutdown(type_, message)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.ShutdownAsync(LibreMetaverse.NetworkManager.DisconnectType,System.String)":
|
||||||
|
"self.native_shutdown_async(type_, message).await",
|
||||||
|
"M:LibreMetaverse.NetworkManager.UnregisterCallback(LibreMetaverse.Packets.PacketType,System.EventHandler{LibreMetaverse.PacketReceivedEventArgs})":
|
||||||
|
"self.native_unregister_callback(type_, callback)",
|
||||||
|
"M:LibreMetaverse.NetworkManager.UnregisterEventCallback(System.String,LibreMetaverse.Caps.EventQueueCallback)":
|
||||||
|
"self.native_unregister_caps(caps_event, callback)",
|
||||||
|
"P:LibreMetaverse.NetworkManager.CircuitCode":
|
||||||
|
"self.native_circuit_code()",
|
||||||
|
"P:LibreMetaverse.NetworkManager.CircuitCode#set":
|
||||||
|
"self.native_set_circuit_code(value)",
|
||||||
|
"P:LibreMetaverse.NetworkManager.Connected":
|
||||||
|
"self.native_connected()",
|
||||||
|
"P:LibreMetaverse.NetworkManager.Connected#set":
|
||||||
|
"self.native_set_connected(value)",
|
||||||
|
"P:LibreMetaverse.NetworkManager.CurrentSim":
|
||||||
|
"self.native_current_sim()",
|
||||||
|
"P:LibreMetaverse.NetworkManager.CurrentSim#set":
|
||||||
|
"self.native_set_current_sim(value)",
|
||||||
|
"P:LibreMetaverse.NetworkManager.InboxCount":
|
||||||
|
"self.native_inbox_count()",
|
||||||
|
"P:LibreMetaverse.NetworkManager.OutboxCount":
|
||||||
|
"self.native_outbox_count()",
|
||||||
|
"M:LibreMetaverse.Simulator.#ctor(LibreMetaverse.GridClient,System.Net.IPEndPoint,System.UInt64,System.UInt32,System.UInt32)":
|
||||||
|
"crate::network_manager::Simulator::native_new(client, address, handle, size_x, size_y)",
|
||||||
|
"M:LibreMetaverse.Simulator.Connect(System.Boolean)":
|
||||||
|
"self.native_connect(move_to_sim)",
|
||||||
|
"M:LibreMetaverse.Simulator.ConnectAsync(System.Boolean)":
|
||||||
|
"self.native_connect_async(move_to_sim).await",
|
||||||
|
"M:LibreMetaverse.Simulator.Disconnect(System.Boolean)":
|
||||||
|
"self.native_disconnect(send_close_circuit)",
|
||||||
|
"M:LibreMetaverse.Simulator.Dispose":
|
||||||
|
"self.native_dispose()",
|
||||||
|
"M:LibreMetaverse.Simulator.Equals(System.Object)":
|
||||||
|
"self.native_equals(obj)",
|
||||||
|
"M:LibreMetaverse.Simulator.GetHashCode":
|
||||||
|
"self.native_hash_code()",
|
||||||
|
"M:LibreMetaverse.Simulator.Pause":
|
||||||
|
"self.native_pause()",
|
||||||
|
"M:LibreMetaverse.Simulator.Resume":
|
||||||
|
"self.native_resume()",
|
||||||
|
"M:LibreMetaverse.Simulator.SendPacket(LibreMetaverse.Packets.Packet)":
|
||||||
|
"self.native_send_packet(packet)",
|
||||||
|
"M:LibreMetaverse.Simulator.SendPacketData(System.Byte[],System.Int32,LibreMetaverse.Packets.PacketType,System.Boolean)":
|
||||||
|
"self.native_send_packet_data(data, data_length, type_, do_zerocode)",
|
||||||
|
"M:LibreMetaverse.Simulator.SendPing":
|
||||||
|
"self.native_send_ping()",
|
||||||
|
"M:LibreMetaverse.Simulator.SetSeedCaps(System.Uri,System.Boolean)":
|
||||||
|
"self.native_set_seed_caps(seedcaps, changed_sim)",
|
||||||
|
"M:LibreMetaverse.Simulator.ToString":
|
||||||
|
"self.native_to_string()",
|
||||||
|
"M:LibreMetaverse.Simulator.UseCircuitCode(System.Boolean)":
|
||||||
|
"self.native_use_circuit_code(wait_for_ack)",
|
||||||
|
"M:LibreMetaverse.Simulator.UseCircuitCodeAsync(System.Boolean)":
|
||||||
|
"self.native_use_circuit_code_async(wait_for_ack).await",
|
||||||
|
"M:LibreMetaverse.Simulator.op_Equality(LibreMetaverse.Simulator,LibreMetaverse.Simulator)":
|
||||||
|
"lhs == rhs",
|
||||||
|
"M:LibreMetaverse.Simulator.op_Inequality(LibreMetaverse.Simulator,LibreMetaverse.Simulator)":
|
||||||
|
"lhs != rhs",
|
||||||
|
"P:LibreMetaverse.Simulator.Connected":
|
||||||
|
"self.native_is_connected()",
|
||||||
|
"P:LibreMetaverse.Simulator.HandshakeComplete":
|
||||||
|
"self.native_handshake_complete()",
|
||||||
|
"P:LibreMetaverse.Simulator.IPEndPoint":
|
||||||
|
"self.native_ip_end_point()",
|
||||||
"P:LibreMetaverse.GridClient.Settings":
|
"P:LibreMetaverse.GridClient.Settings":
|
||||||
"&mut self.settings",
|
"&mut self.settings",
|
||||||
"P:LibreMetaverse.GridClient.TimeProvider":
|
"P:LibreMetaverse.GridClient.TimeProvider":
|
||||||
@@ -337,6 +468,11 @@ CLIENT_CORE_NATIVE_OWNERS = (
|
|||||||
"WorldSettings",
|
"WorldSettings",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
NETWORK_NATIVE_OWNERS = (
|
||||||
|
"NetworkManager",
|
||||||
|
"Simulator",
|
||||||
|
)
|
||||||
|
|
||||||
NATIVE_OWNER_BOUNDS = {
|
NATIVE_OWNER_BOUNDS = {
|
||||||
"T:LibreMetaverse.CacheDictionary`2": {
|
"T:LibreMetaverse.CacheDictionary`2": {
|
||||||
"TKey": "TKey: Clone + PartialEq",
|
"TKey": "TKey: Clone + PartialEq",
|
||||||
@@ -477,6 +613,7 @@ NATIVE_MESSAGE_CODEC_TYPES = {
|
|||||||
OPEN_ENUM_IDS = {"T:LibreMetaverse.BakeType", "T:LibreMetaverse.HoleType"}
|
OPEN_ENUM_IDS = {"T:LibreMetaverse.BakeType", "T:LibreMetaverse.HoleType"}
|
||||||
TRAIT_SUPERTRAITS = {
|
TRAIT_SUPERTRAITS = {
|
||||||
"T:LibreMetaverse.IBakingTextureProvider": "std::any::Any",
|
"T:LibreMetaverse.IBakingTextureProvider": "std::any::Any",
|
||||||
|
"T:LibreMetaverse.Interfaces.IMessage": "std::any::Any",
|
||||||
"T:LibreMetaverse.Appearance.ICurrentOutfitPolicy": "Send + Sync",
|
"T:LibreMetaverse.Appearance.ICurrentOutfitPolicy": "Send + Sync",
|
||||||
}
|
}
|
||||||
CORE_NAMESPACES = (
|
CORE_NAMESPACES = (
|
||||||
@@ -1053,11 +1190,12 @@ def render_body(signature: str, member_id: str, error_model: str, asyncness: str
|
|||||||
if trait:
|
if trait:
|
||||||
signature = signature.removeprefix("pub ")
|
signature = signature.removeprefix("pub ")
|
||||||
if body := native_member_body(member_id):
|
if body := native_member_body(member_id):
|
||||||
marker = (
|
if any(f"LibreMetaverse.{owner}." in member_id for owner in CLIENT_CORE_NATIVE_OWNERS):
|
||||||
" /* native client-core implementation */"
|
marker = " /* native client-core implementation */"
|
||||||
if any(f"LibreMetaverse.{owner}." in member_id for owner in CLIENT_CORE_NATIVE_OWNERS)
|
elif any(f"LibreMetaverse.{owner}." in member_id for owner in NETWORK_NATIVE_OWNERS):
|
||||||
else ""
|
marker = " /* native network implementation */"
|
||||||
)
|
else:
|
||||||
|
marker = ""
|
||||||
return f" {signature} {{{marker} {body} }}"
|
return f" {signature} {{{marker} {body} }}"
|
||||||
failure = (
|
failure = (
|
||||||
f"libremetaverse_types::not_implemented({json.dumps(member_id)})"
|
f"libremetaverse_types::not_implemented({json.dumps(member_id)})"
|
||||||
|
|||||||
@@ -971,7 +971,10 @@ def validate_generated_shims() -> None:
|
|||||||
if default_types - {"UUID"}:
|
if default_types - {"UUID"}:
|
||||||
raise ValueError(f"generated shim derives a plausible Default: {path.relative_to(ROOT)}")
|
raise ValueError(f"generated shim derives a plausible Default: {path.relative_to(ROOT)}")
|
||||||
for body in function_bodies(text):
|
for body in function_bodies(text):
|
||||||
if "native client-core implementation" in body:
|
if (
|
||||||
|
"native client-core implementation" in body
|
||||||
|
or "native network implementation" in body
|
||||||
|
):
|
||||||
continue
|
continue
|
||||||
if FAKE_BODY.search(body):
|
if FAKE_BODY.search(body):
|
||||||
raise ValueError(f"generated shim returns a plausible fallback: {path.relative_to(ROOT)}")
|
raise ValueError(f"generated shim returns a plausible fallback: {path.relative_to(ROOT)}")
|
||||||
|
|||||||
Reference in New Issue
Block a user