From c7777d42f5eff02eab5a93df9b649af3d4d2c1b8 Mon Sep 17 00:00:00 2001 From: Chili Palmer Date: Sun, 9 Aug 2026 11:39:59 +0000 Subject: [PATCH] Implement native UDP transport reliability (#52) --- Cargo.lock | 53 + README.md | 10 + api/SHIM-COVERAGE.md | 2 +- crates/libremetaverse/Cargo.toml | 4 + crates/libremetaverse/src/client_core.rs | 41 + crates/libremetaverse/src/generated.rs | 244 +-- crates/libremetaverse/src/lib.rs | 5 + crates/libremetaverse/src/udp_transport.rs | 1780 ++++++++++++++++++ crates/libremetaverse/tests/udp_transport.rs | 511 +++++ docs/udp-transport.md | 77 + tests/api-compile/Cargo.lock | 68 + tests/compat/tests/core_runtime_shims.rs | 8 +- tools/generate_api_shims.py | 4 + 13 files changed, 2593 insertions(+), 214 deletions(-) create mode 100644 crates/libremetaverse/src/udp_transport.rs create mode 100644 crates/libremetaverse/tests/udp_transport.rs create mode 100644 docs/udp-transport.md diff --git a/Cargo.lock b/Cargo.lock index 0575e5d..a710707 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -369,6 +369,7 @@ dependencies = [ "libremetaverse-structured-data", "libremetaverse-types", "roxmltree", + "tokio", ] [[package]] @@ -575,6 +576,17 @@ dependencies = [ "simd-adler32", ] +[[package]] +name = "mio" +version = "1.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "30d65c71f1ce40ab09135ce117d742b9f8a19ff91a41a8b57ed50bc2de59c427" +dependencies = [ + "libc", + "wasi", + "windows-sys", +] + [[package]] name = "nom" version = "7.1.3" @@ -845,6 +857,16 @@ version = "0.4.12" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5" +[[package]] +name = "socket2" +version = "0.6.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3d1e2c7f27f8d4cb10542a02c49005dbd6e93095799d6f3be745fae9f8fedd4" +dependencies = [ + "libc", + "windows-sys", +] + [[package]] name = "stats_alloc" version = "0.1.10" @@ -884,6 +906,31 @@ dependencies = [ "xattr", ] +[[package]] +name = "tokio" +version = "1.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "202caea871b69668250d242070849eb495be178ed697a3e98aebce5bc81a0bed" +dependencies = [ + "libc", + "mio", + "pin-project-lite", + "socket2", + "tokio-macros", + "windows-sys", +] + +[[package]] +name = "tokio-macros" +version = "2.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + [[package]] name = "toml" version = "1.1.4+spec-1.1.0" @@ -958,6 +1005,12 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "wasi" +version = "0.11.1+wasi-snapshot-preview1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" + [[package]] name = "wasm-bindgen" version = "0.2.127" diff --git a/README.md b/README.md index b155ee2..6c0c11b 100644 --- a/README.md +++ b/README.md @@ -262,3 +262,13 @@ token and are shut down idempotently in network, manager, HTTP, then rate-limite order. Injected clocks support deterministic tests, while `Debug` output omits endpoint values and service internals. The ownership and executor requirements are documented in [`docs/client-core.md`](docs/client-core.md). +The native Tokio UDP layer now implements the C# packet buffers and throttle +encoding plus bounded socket receive, coordination, and single-writer tasks. +It assigns wrapping protocol sequences, aggregates and consumes ACKs, retries +reliable packets with duplicate suppression, preserves zerocoding and the +1,200-byte MTU contract, applies independent task/texture/asset token buckets, +and exposes payload-free transport statistics. All queues, peer state, +zerocode expansion, ACK state, and reliable windows have explicit limits; +linked cancellation and final drop release every socket task. The executor, +wire, backpressure, retry, and security contracts are documented in +[`docs/udp-transport.md`](docs/udp-transport.md). diff --git a/api/SHIM-COVERAGE.md b/api/SHIM-COVERAGE.md index e365e9c..25bf381 100644 --- a/api/SHIM-COVERAGE.md +++ b/api/SHIM-COVERAGE.md @@ -4,7 +4,7 @@ Generated by `python3 tools/generate_api_shims.py`; do not edit by hand. | Assembly | Types | Members | Status | |---|---:|---:|---| -| `LibreMetaverse` | 2,711 | 27,281 | native implementation: 28 types / 13,428 members; remaining surface is callable failure-only shims | +| `LibreMetaverse` | 2,711 | 27,281 | native implementation: 32 types / 13,459 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.Skia` | 1 | 3 | native implementation: 1 type / 3 members; no generated shims remain | | `LibreMetaverse.LslTools` | 164 | 768 | callable failure-only shim | diff --git a/crates/libremetaverse/Cargo.toml b/crates/libremetaverse/Cargo.toml index d1309e3..8b58a7e 100644 --- a/crates/libremetaverse/Cargo.toml +++ b/crates/libremetaverse/Cargo.toml @@ -18,6 +18,10 @@ libremetaverse-imaging = { path = "../libremetaverse-imaging" } libremetaverse-structured-data = { path = "../libremetaverse-structured-data" } libremetaverse-types = { path = "../libremetaverse-types" } roxmltree = "0.21.1" +tokio = { version = "1.47.1", features = ["macros", "net", "rt", "sync", "time"] } + +[dev-dependencies] +tokio = { version = "1.47.1", features = ["macros", "net", "rt-multi-thread", "sync", "test-util", "time"] } [lints] workspace = true diff --git a/crates/libremetaverse/src/client_core.rs b/crates/libremetaverse/src/client_core.rs index fe7faad..5a7ff25 100644 --- a/crates/libremetaverse/src/client_core.rs +++ b/crates/libremetaverse/src/client_core.rs @@ -113,6 +113,7 @@ struct ClientRuntime { state: AtomicU8, cancellation: CancellationTokenSource, services: Mutex>, + agent_throttle_sender: Mutex>>, shutdown_complete: Condvar, shutdown_wait: Mutex<()>, } @@ -132,6 +133,7 @@ impl ClientRuntime { }) .collect(), ), + agent_throttle_sender: Mutex::new(None), shutdown_complete: Condvar::new(), shutdown_wait: Mutex::new(()), } @@ -200,6 +202,10 @@ impl ClientRuntime { } } services.clear(); + self.agent_throttle_sender + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take(); let shutdown_wait = self .shutdown_wait .lock() @@ -257,6 +263,41 @@ impl GridClient { self.runtime.register(service) } + /// Installs the network-owned callback used by mapped `AgentThrottle.Set`. + /// + /// # Errors + /// + /// Registration fails after client shutdown begins. + pub fn set_agent_throttle_sender( + &mut self, + sender: Arc, + ) -> Result<(), ClientCoreError> { + let mut registered = self + .runtime + .agent_throttle_sender + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let state = self.lifecycle_state(); + if state != ClientLifecycleState::Active { + return Err(ClientCoreError::InvalidLifecycle { + operation: "register the agent throttle sender", + state, + }); + } + *registered = Some(sender); + Ok(()) + } + + pub(crate) fn agent_throttle_sender( + &self, + ) -> Option> { + self.runtime + .agent_throttle_sender + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + } + pub(crate) fn shutdown(&self) -> Result<(), ClientCoreError> { self.runtime.shutdown() } diff --git a/crates/libremetaverse/src/generated.rs b/crates/libremetaverse/src/generated.rs index 4f7a6b4..1e04b44 100644 --- a/crates/libremetaverse/src/generated.rs +++ b/crates/libremetaverse/src/generated.rs @@ -3333,102 +3333,20 @@ impl AgentState { } /// C# type: `T:LibreMetaverse.AgentThrottle`. -pub struct AgentThrottle; -impl AgentThrottle { - /// C# member: `M:LibreMetaverse.AgentThrottle.#ctor(LibreMetaverse.GridClient)`. - pub fn new_with_grid_client(client: libremetaverse::GridClient) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.AgentThrottle.#ctor(LibreMetaverse.GridClient)", - ) - } - /// C# member: `M:LibreMetaverse.AgentThrottle.#ctor(System.Byte[],System.Int32)`. - pub fn new_with_bytes_int32(data: Vec, pos: i32) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.AgentThrottle.#ctor(System.Byte[],System.Int32)", - ) - } - /// C# member: `M:LibreMetaverse.AgentThrottle.Set`. - pub fn set_with_method(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.AgentThrottle.Set") - } - /// C# member: `M:LibreMetaverse.AgentThrottle.Set(LibreMetaverse.Simulator)`. - pub fn set_with_simulator( - &self, - simulator: Option, - ) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.AgentThrottle.Set(LibreMetaverse.Simulator)", - ) - } - /// C# member: `M:LibreMetaverse.AgentThrottle.ToBytes`. - pub fn to_bytes(&self) -> Result, crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.AgentThrottle.ToBytes") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Asset`. - pub fn asset(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Asset") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Asset`. - pub fn set_asset(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Asset") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Cloud`. - pub fn cloud(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Cloud") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Cloud`. - pub fn set_cloud(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Cloud") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Land`. - pub fn land(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Land") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Land`. - pub fn set_land(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Land") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Resend`. - pub fn resend(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Resend") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Resend`. - pub fn set_resend(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Resend") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Task`. - pub fn task(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Task") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Task`. - pub fn set_task(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Task") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Texture`. - pub fn texture(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Texture") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Texture`. - pub fn set_texture(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Texture") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Total`. - pub fn total(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Total") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Total`. - pub fn set_total(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Total") - } - /// C# member: `P:LibreMetaverse.AgentThrottle.Wind`. - pub fn wind(&self) -> f32 { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Wind") - } - /// Setter for C# member: `P:LibreMetaverse.AgentThrottle.Wind`. - pub fn set_wind(&mut self, value: f32) { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.AgentThrottle.Wind") - } -} +/// C# member: `M:LibreMetaverse.AgentThrottle.#ctor(LibreMetaverse.GridClient)`. +/// C# member: `M:LibreMetaverse.AgentThrottle.#ctor(System.Byte[],System.Int32)`. +/// C# member: `M:LibreMetaverse.AgentThrottle.Set`. +/// C# member: `M:LibreMetaverse.AgentThrottle.Set(LibreMetaverse.Simulator)`. +/// C# member: `M:LibreMetaverse.AgentThrottle.ToBytes`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Asset`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Cloud`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Land`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Resend`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Task`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Texture`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Total`. +/// C# member: `P:LibreMetaverse.AgentThrottle.Wind`. +pub use crate::udp_transport::AgentThrottle; /// C# type: `T:LibreMetaverse.AgentWearablesReplyEventArgs`. pub struct AgentWearablesReplyEventArgs; @@ -13666,21 +13584,9 @@ pub enum ImageType { } /// C# type: `T:LibreMetaverse.IncomingPacketIDCollection`. -pub struct IncomingPacketIDCollection; -impl IncomingPacketIDCollection { - /// C# member: `M:LibreMetaverse.IncomingPacketIDCollection.#ctor(System.Int32)`. - pub fn new(capacity: i32) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.IncomingPacketIDCollection.#ctor(System.Int32)", - ) - } - /// C# member: `M:LibreMetaverse.IncomingPacketIDCollection.TryEnqueue(System.UInt32)`. - pub fn try_enqueue(&self, ack: u32) -> bool { - libremetaverse_types::unimplemented_api!( - "M:LibreMetaverse.IncomingPacketIDCollection.TryEnqueue(System.UInt32)" - ) - } -} +/// C# member: `M:LibreMetaverse.IncomingPacketIDCollection.#ctor(System.Int32)`. +/// C# member: `M:LibreMetaverse.IncomingPacketIDCollection.TryEnqueue(System.UInt32)`. +pub use crate::udp_transport::IncomingPacketIDCollection; /// C# type: `T:LibreMetaverse.InitiateDownloadEventArgs`. pub struct InitiateDownloadEventArgs; @@ -27995,106 +27901,26 @@ pub use crate::foliage_catalog::TreeDefinition; pub use crate::foliage_catalog::TreeDefinitions; /// C# type: `T:LibreMetaverse.UDPBase`. -pub struct UDPBase; -impl UDPBase { - /// C# member: `M:LibreMetaverse.UDPBase.AsyncBeginSend(LibreMetaverse.UDPPacketBuffer)`. - pub fn async_begin_send( - &self, - buf: libremetaverse::UDPPacketBuffer, - ) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPBase.AsyncBeginSend(LibreMetaverse.UDPPacketBuffer)", - ) - } - /// C# member: `M:LibreMetaverse.UDPBase.Start`. - pub fn start(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.UDPBase.Start") - } - /// C# member: `M:LibreMetaverse.UDPBase.Stop`. - pub fn stop(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.UDPBase.Stop") - } - /// C# member: `P:LibreMetaverse.UDPBase.IsRunning`. - pub fn is_running(&self) -> bool { - libremetaverse_types::unimplemented_api!("P:LibreMetaverse.UDPBase.IsRunning") - } -} +/// C# member: `M:LibreMetaverse.UDPBase.AsyncBeginSend(LibreMetaverse.UDPPacketBuffer)`. +/// C# member: `M:LibreMetaverse.UDPBase.Start`. +/// C# member: `M:LibreMetaverse.UDPBase.Stop`. +/// C# member: `P:LibreMetaverse.UDPBase.IsRunning`. +pub use crate::udp_transport::UDPBase; /// C# type: `T:LibreMetaverse.UDPPacketBuffer`. -pub struct UDPPacketBuffer { - /// C# member: `F:LibreMetaverse.UDPPacketBuffer.Data`. - pub data: Vec, - /// C# member: `F:LibreMetaverse.UDPPacketBuffer.DataLength`. - pub data_length: i32, - /// C# member: `F:LibreMetaverse.UDPPacketBuffer.RemoteEndPoint`. - pub remote_end_point: std::net::SocketAddr, -} -impl UDPPacketBuffer { - /// C# member: `F:LibreMetaverse.UDPPacketBuffer.DEFAULT_BUFFER_SIZE`. - pub const DEFAULT_BUFFER_SIZE: i32 = 4096; - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor`. - pub fn new_with_constructor() -> Result { - libremetaverse_types::not_implemented("M:LibreMetaverse.UDPPacketBuffer.#ctor") - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Byte[],System.Int32,System.Net.IPEndPoint,System.Int32)`. - pub fn new_with_bytes_int32_ip_end_point_int32( - buffer: Vec, - buffer_size: i32, - destination: std::net::SocketAddr, - category: i32, - ) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Byte[],System.Int32,System.Net.IPEndPoint,System.Int32)", - ) - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint)`. - pub fn new_with_ip_end_point(end_point: std::net::SocketAddr) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint)", - ) - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint,System.Byte[])`. - pub fn new_with_ip_end_point_bytes( - end_point: std::net::SocketAddr, - data: Vec, - ) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint,System.Byte[])", - ) - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint,System.Int32)`. - pub fn new_with_ip_end_point_int32( - end_point: std::net::SocketAddr, - buffer_size: i32, - ) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint,System.Int32)", - ) - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.CopyFrom(System.Array)`. - pub fn copy_from_with_array( - &self, - src: libremetaverse_types::compat::Array, - ) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPPacketBuffer.CopyFrom(System.Array)", - ) - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.CopyFrom(System.Array,System.Int32)`. - pub fn copy_from_with_array_int32( - &self, - src: libremetaverse_types::compat::Array, - length: i32, - ) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.UDPPacketBuffer.CopyFrom(System.Array,System.Int32)", - ) - } - /// C# member: `M:LibreMetaverse.UDPPacketBuffer.ResetEndpoint`. - pub fn reset_endpoint(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.UDPPacketBuffer.ResetEndpoint") - } -} +/// C# member: `F:LibreMetaverse.UDPPacketBuffer.DEFAULT_BUFFER_SIZE`. +/// C# member: `F:LibreMetaverse.UDPPacketBuffer.Data`. +/// C# member: `F:LibreMetaverse.UDPPacketBuffer.DataLength`. +/// C# member: `F:LibreMetaverse.UDPPacketBuffer.RemoteEndPoint`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Byte[],System.Int32,System.Net.IPEndPoint,System.Int32)`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint)`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint,System.Byte[])`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.#ctor(System.Net.IPEndPoint,System.Int32)`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.CopyFrom(System.Array)`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.CopyFrom(System.Array,System.Int32)`. +/// C# member: `M:LibreMetaverse.UDPPacketBuffer.ResetEndpoint`. +pub use crate::udp_transport::UDPPacketBuffer; /// C# type: `T:LibreMetaverse.UUIDNameReplyEventArgs`. pub struct UUIDNameReplyEventArgs; diff --git a/crates/libremetaverse/src/lib.rs b/crates/libremetaverse/src/lib.rs index e867b9b..3e04749 100644 --- a/crates/libremetaverse/src/lib.rs +++ b/crates/libremetaverse/src/lib.rs @@ -20,6 +20,7 @@ mod skeleton; #[rustfmt::skip] // Deterministic machine output is formatted by the pinned generator. mod skeleton_catalog; mod targa; +mod udp_transport; #[rustfmt::skip] // Deterministic machine output is formatted by the pinned generator. mod visual_catalog; @@ -204,6 +205,10 @@ 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 udp_transport::{ + AgentThrottleSender, UdpPacketHandler, UdpThrottleCategory, UdpTransportConfig, + UdpTransportError, UdpTransportStats, +}; /// Rust object-safe view of the C# `InventoryBase` inheritance hierarchy. pub trait InventoryObjectClass: std::any::Any { diff --git a/crates/libremetaverse/src/udp_transport.rs b/crates/libremetaverse/src/udp_transport.rs new file mode 100644 index 0000000..cabf28d --- /dev/null +++ b/crates/libremetaverse/src/udp_transport.rs @@ -0,0 +1,1780 @@ +//! Bounded Tokio UDP transport and the C#-compatible transport value types. + +#![allow(clippy::missing_errors_doc)] // Result shapes are fixed by the compatibility map. +#![allow(clippy::needless_pass_by_value)] // Owned arrays and objects preserve mapped signatures. + +use crate::packets::{Packet, PacketAckPacket, PacketAckPacketPacketsBlock, PacketType}; +use crate::{Error, GridClient, Helpers, Simulator}; +use libremetaverse_types::compat::{Array, CancellationToken, CancellationTokenSource, Object}; +use std::collections::{BTreeMap, HashMap, HashSet, VecDeque}; +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::{Arc, Mutex}; +use std::time::Duration; +use tokio::net::UdpSocket; +use tokio::sync::{mpsc, oneshot}; +use tokio::task::JoinHandle; +use tokio::time::{Instant, MissedTickBehavior}; + +const DEFAULT_DECODE_BUFFER_SIZE: usize = 8 * 1024; +const THROTTLE_PERIOD: Duration = Duration::from_millis(100); +const THROTTLE_MIN_BYTES_PER_PERIOD: usize = 200; +const THROTTLE_BURST_PERIODS: usize = 4; + +/// A failure at the native UDP transport boundary. +/// +/// The variants deliberately carry no datagram bytes, credentials, endpoint +/// query data, or operating-system error strings. +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub enum UdpTransportError { + InvalidConfiguration(&'static str), + InvalidBuffer, + MtuExceeded, + NotRunning, + AlreadyRunning, + Backpressure, + ReliableWindowFull, + TooManyPeers, + RuntimeUnavailable, + Cancelled, + Socket, +} + +impl fmt::Display for UdpTransportError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::InvalidConfiguration(field) => { + write!(formatter, "invalid UDP transport setting {field}") + } + Self::InvalidBuffer => formatter.write_str("invalid UDP packet buffer"), + Self::MtuExceeded => formatter.write_str("UDP payload exceeds the protocol MTU"), + Self::NotRunning => formatter.write_str("UDP transport is not running"), + Self::AlreadyRunning => formatter.write_str("UDP transport is already running"), + Self::Backpressure => formatter.write_str("UDP transport queue is full"), + Self::ReliableWindowFull => formatter.write_str("UDP reliable-send window is full"), + Self::TooManyPeers => formatter.write_str("UDP peer limit is reached"), + Self::RuntimeUnavailable => { + formatter.write_str("UDP transport requires a current Tokio runtime") + } + Self::Cancelled => formatter.write_str("UDP transport was cancelled"), + Self::Socket => formatter.write_str("UDP socket operation failed"), + } + } +} + +impl std::error::Error for UdpTransportError {} + +impl From for Error { + fn from(error: UdpTransportError) -> Self { + match error { + UdpTransportError::InvalidConfiguration(_) + | UdpTransportError::InvalidBuffer + | UdpTransportError::MtuExceeded => Self::Argument, + UdpTransportError::Cancelled => Self::Cancelled, + UdpTransportError::Socket => Self::Socket, + UdpTransportError::NotRunning + | UdpTransportError::AlreadyRunning + | UdpTransportError::Backpressure + | UdpTransportError::ReliableWindowFull + | UdpTransportError::TooManyPeers + | UdpTransportError::RuntimeUnavailable => Self::InvalidOperation, + } + } +} + +/// C# `UDPPacketBuffer`, with one owned allocation per queued datagram. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct UDPPacketBuffer { + pub data: Vec, + pub data_length: i32, + pub remote_end_point: SocketAddr, +} + +impl UDPPacketBuffer { + pub const DEFAULT_BUFFER_SIZE: i32 = 4096; + + pub fn new_with_constructor() -> Result { + Self::new_with_ip_end_point_int32( + SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0), + Self::DEFAULT_BUFFER_SIZE, + ) + } + + pub fn new_with_bytes_int32_ip_end_point_int32( + buffer: Vec, + buffer_size: i32, + destination: SocketAddr, + _category: i32, + ) -> Result { + let mut packet = Self::new_with_ip_end_point_int32(destination, buffer_size)?; + packet.copy_from_slice_with_length(&buffer, buffer_size)?; + packet.data_length = buffer_size; + Ok(packet) + } + + pub fn new_with_ip_end_point(end_point: SocketAddr) -> Result { + Self::new_with_ip_end_point_int32(end_point, Self::DEFAULT_BUFFER_SIZE) + } + + pub fn new_with_ip_end_point_bytes( + end_point: SocketAddr, + data: Vec, + ) -> Result { + Ok(Self { + data, + // The C# constructor adopts the supplied array but deliberately + // leaves the public DataLength field at its zero default. + data_length: 0, + remote_end_point: end_point, + }) + } + + pub fn new_with_ip_end_point_int32( + end_point: SocketAddr, + buffer_size: i32, + ) -> Result { + let buffer_size = usize::try_from(buffer_size).map_err(|_| Error::Argument)?; + let mut data = Vec::new(); + data.try_reserve_exact(buffer_size) + .map_err(|_| Error::InvalidOperation)?; + data.resize(buffer_size, 0); + Ok(Self { + data, + data_length: 0, + remote_end_point: end_point, + }) + } + + /// Copies a mapped CLR array into this buffer. + /// + /// `Object::Bytes` represents a boxed `byte[]`; a flat array of integer + /// objects is also accepted so the mapped `System.Array` remains useful. + pub fn copy_from_with_array(&mut self, src: Array) -> Result<(), Error> { + let bytes = mapped_array_bytes(src)?; + let length = i32::try_from(bytes.len()).map_err(|_| Error::Argument)?; + self.copy_from_slice_with_length(&bytes, length) + } + + pub fn copy_from_with_array_int32(&mut self, src: Array, length: i32) -> Result<(), Error> { + let bytes = mapped_array_bytes(src)?; + self.copy_from_slice_with_length(&bytes, length) + } + + pub fn copy_from_slice(&mut self, src: &[u8]) -> Result<(), Error> { + let length = i32::try_from(src.len()).map_err(|_| Error::Argument)?; + self.copy_from_slice_with_length(src, length) + } + + pub fn copy_from_slice_with_length(&mut self, src: &[u8], length: i32) -> Result<(), Error> { + let length = usize::try_from(length).map_err(|_| Error::Argument)?; + if length > src.len() || length > self.data.len() { + return Err(Error::IndexOutOfRange); + } + self.data[..length].copy_from_slice(&src[..length]); + Ok(()) + } + + pub fn reset_endpoint(&mut self) -> Result<(), Error> { + self.remote_end_point = match self.remote_end_point.ip() { + IpAddr::V4(_) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0), + IpAddr::V6(_) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0), + }; + Ok(()) + } + + fn payload(&self) -> Result<&[u8], UdpTransportError> { + let length = + usize::try_from(self.data_length).map_err(|_| UdpTransportError::InvalidBuffer)?; + self.data + .get(..length) + .ok_or(UdpTransportError::InvalidBuffer) + } +} + +fn mapped_array_bytes(src: Array) -> Result, Error> { + if let [Object::Bytes(bytes)] = src.0.as_slice() { + return Ok(bytes.clone()); + } + let mut bytes = Vec::new(); + bytes + .try_reserve_exact(src.0.len()) + .map_err(|_| Error::InvalidOperation)?; + for value in src.0 { + let byte = match value { + Object::Integer(value) => u8::try_from(value).map_err(|_| Error::Argument)?, + Object::UInteger(value) => u8::try_from(value).map_err(|_| Error::Argument)?, + _ => return Err(Error::Argument), + }; + bytes.push(byte); + } + Ok(bytes) +} + +struct PacketArchiveState { + items: Vec, + members: HashSet, + first: usize, + next: usize, +} + +/// Fixed-size duplicate archive matching `IncomingPacketIDCollection`. +pub struct IncomingPacketIDCollection { + capacity: usize, + state: Mutex, +} + +impl IncomingPacketIDCollection { + pub fn new(capacity: i32) -> Result { + let capacity = usize::try_from(capacity).map_err(|_| Error::Argument)?; + if capacity == 0 { + return Err(Error::Argument); + } + let mut items = Vec::new(); + items + .try_reserve_exact(capacity) + .map_err(|_| Error::InvalidOperation)?; + items.resize(capacity, 0); + Ok(Self { + capacity, + state: Mutex::new(PacketArchiveState { + items, + members: HashSet::with_capacity(capacity), + first: 0, + next: 0, + }), + }) + } + + pub fn try_enqueue(&self, ack: u32) -> bool { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if !state.members.insert(ack) { + return false; + } + let next = state.next; + state.items[next] = ack; + state.next = (next + 1) % self.capacity; + if state.next == state.first { + let first = state.first; + let removed = state.items[first]; + state.members.remove(&removed); + state.first = (first + 1) % self.capacity; + } + true + } +} + +impl fmt::Debug for IncomingPacketIDCollection { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + let state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + formatter + .debug_struct("IncomingPacketIDCollection") + .field("capacity", &self.capacity) + .field("len", &state.members.len()) + .finish() + } +} + +/// Injection boundary used by `AgentThrottle::Set` and the later network +/// manager composition layer. +pub trait AgentThrottleSender: Send + Sync { + fn send_throttle( + &self, + throttle_bytes: &[u8], + simulator: Option<&Simulator>, + ) -> Result<(), Error>; +} + +/// Exact seven-stream C# throttle values and little-endian wire encoding. +#[derive(Clone)] +pub struct AgentThrottle { + resend: f32, + land: f32, + wind: f32, + cloud: f32, + task: f32, + texture: f32, + asset: f32, + sender: Option>, +} + +impl AgentThrottle { + pub fn new_with_grid_client(client: GridClient) -> Result { + Ok(Self { + sender: client.agent_throttle_sender(), + ..Self::default() + }) + } + + pub fn new_with_bytes_int32(data: Vec, pos: i32) -> Result { + let pos = usize::try_from(pos).map_err(|_| Error::Argument)?; + let end = pos.checked_add(28).ok_or(Error::Argument)?; + let bytes = data.get(pos..end).ok_or(Error::IndexOutOfRange)?; + let mut value = Self::default(); + value.set_resend(read_f32(bytes, 0)?); + value.set_land(read_f32(bytes, 4)?); + value.set_wind(read_f32(bytes, 8)?); + value.set_cloud(read_f32(bytes, 12)?); + value.set_task(read_f32(bytes, 16)?); + value.set_texture(read_f32(bytes, 20)?); + value.set_asset(read_f32(bytes, 24)?); + Ok(value) + } + + #[must_use] + pub fn with_sender(mut self, sender: Arc) -> Self { + self.sender = Some(sender); + self + } + + pub fn set_with_method(&self) -> Result<(), Error> { + if let Some(sender) = &self.sender { + sender.send_throttle(&self.to_bytes()?, None)?; + } + Ok(()) + } + + pub fn set_with_simulator(&self, simulator: Option) -> Result<(), Error> { + if let (Some(sender), Some(simulator)) = (&self.sender, simulator.as_ref()) { + sender.send_throttle(&self.to_bytes()?, Some(simulator))?; + } + Ok(()) + } + + pub fn to_bytes(&self) -> Result, Error> { + let mut output = Vec::with_capacity(28); + for value in [ + self.resend, + self.land, + self.wind, + self.cloud, + self.task, + self.texture, + self.asset, + ] { + output.extend_from_slice(&value.to_le_bytes()); + } + Ok(output) + } + + #[must_use] + pub const fn asset(&self) -> f32 { + self.asset + } + + pub fn set_asset(&mut self, value: f32) { + self.asset = value.clamp(10_000.0, 220_000.0); + } + + #[must_use] + pub const fn cloud(&self) -> f32 { + self.cloud + } + + pub fn set_cloud(&mut self, value: f32) { + self.cloud = value.clamp(0.0, 34_000.0); + } + + #[must_use] + pub const fn land(&self) -> f32 { + self.land + } + + pub fn set_land(&mut self, value: f32) { + self.land = value.clamp(0.0, 170_000.0); + } + + #[must_use] + pub const fn resend(&self) -> f32 { + self.resend + } + + pub fn set_resend(&mut self, value: f32) { + self.resend = value.clamp(10_000.0, 150_000.0); + } + + #[must_use] + pub const fn task(&self) -> f32 { + self.task + } + + pub fn set_task(&mut self, value: f32) { + self.task = value.clamp(4_000.0, 1_338_000.0); + } + + #[must_use] + pub const fn texture(&self) -> f32 { + self.texture + } + + pub fn set_texture(&mut self, value: f32) { + self.texture = value.clamp(4_000.0, 446_000.0); + } + + #[must_use] + pub fn total(&self) -> f32 { + self.resend + self.land + self.wind + self.cloud + self.task + self.texture + self.asset + } + + pub fn set_total(&mut self, value: f32) { + self.set_resend(value * 0.1); + self.set_land(value * 0.52 / 3.0); + self.set_wind(value * 0.05); + self.set_cloud(value * 0.05); + self.set_task(value * 0.704 / 3.0); + self.set_texture(value * 0.704 / 3.0); + self.set_asset(value * 0.484 / 3.0); + } + + #[must_use] + pub const fn wind(&self) -> f32 { + self.wind + } + + pub fn set_wind(&mut self, value: f32) { + self.wind = value.clamp(0.0, 34_000.0); + } +} + +impl Default for AgentThrottle { + fn default() -> Self { + let mut value = Self { + resend: 0.0, + land: 0.0, + wind: 0.0, + cloud: 0.0, + task: 0.0, + texture: 0.0, + asset: 0.0, + sender: None, + }; + value.set_total(1_536_000.0); + value + } +} + +impl fmt::Debug for AgentThrottle { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("AgentThrottle") + .field("resend", &self.resend) + .field("land", &self.land) + .field("wind", &self.wind) + .field("cloud", &self.cloud) + .field("task", &self.task) + .field("texture", &self.texture) + .field("asset", &self.asset) + .finish_non_exhaustive() + } +} + +fn read_f32(bytes: &[u8], offset: usize) -> Result { + Ok(f32::from_le_bytes( + bytes + .get(offset..offset + 4) + .ok_or(Error::IndexOutOfRange)? + .try_into() + .map_err(|_| Error::IndexOutOfRange)?, + )) +} + +/// Outgoing categories used by the native token buckets. +#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] +pub enum UdpThrottleCategory { + Unthrottled, + Task, + Texture, + Asset, +} + +impl UdpThrottleCategory { + #[must_use] + pub const fn classify(packet_type: PacketType) -> Self { + match packet_type { + PacketType::UseCircuitCode + | PacketType::CompleteAgentMovement + | PacketType::AgentThrottle + | PacketType::LogoutRequest + | PacketType::PacketAck + | PacketType::StartPingCheck + | PacketType::CompletePingCheck + | PacketType::CloseCircuit => Self::Unthrottled, + PacketType::RequestImage => Self::Texture, + PacketType::TransferRequest | PacketType::AbortXfer => Self::Asset, + _ => Self::Task, + } + } +} + +/// Bounded transport policy. Defaults match the C# settings used by a client +/// simulator, with explicit limits for queues and per-peer reliable state. +#[derive(Clone, Debug)] +pub struct UdpTransportConfig { + pub receive_queue_capacity: usize, + pub command_queue_capacity: usize, + pub write_queue_capacity: usize, + pub packet_archive_size: usize, + pub pending_ack_capacity: usize, + pub max_pending_acks: usize, + pub reliable_window_capacity: usize, + pub max_resend_count: u32, + pub resend_timeout: Duration, + pub network_tick_interval: Duration, + pub max_datagram_size: usize, + pub max_decoded_packet_size: usize, + pub protocol_mtu: usize, + pub max_peers: usize, + pub throttle: AgentThrottle, +} + +impl Default for UdpTransportConfig { + fn default() -> Self { + Self { + receive_queue_capacity: 512, + command_queue_capacity: 512, + write_queue_capacity: 512, + packet_archive_size: 1000, + pending_ack_capacity: 255, + max_pending_acks: 10, + reliable_window_capacity: 1024, + max_resend_count: 3, + resend_timeout: Duration::from_secs(4), + network_tick_interval: Duration::from_millis(500), + max_datagram_size: usize::try_from(UDPPacketBuffer::DEFAULT_BUFFER_SIZE) + .unwrap_or(4096), + max_decoded_packet_size: DEFAULT_DECODE_BUFFER_SIZE, + protocol_mtu: usize::try_from(Packet::MTU).unwrap_or(1200), + max_peers: 32, + throttle: AgentThrottle::default(), + } + } +} + +impl UdpTransportConfig { + fn validate(&self) -> Result<(), UdpTransportError> { + for (name, value) in [ + ("receive_queue_capacity", self.receive_queue_capacity), + ("command_queue_capacity", self.command_queue_capacity), + ("write_queue_capacity", self.write_queue_capacity), + ("packet_archive_size", self.packet_archive_size), + ("pending_ack_capacity", self.pending_ack_capacity), + ("max_pending_acks", self.max_pending_acks), + ("reliable_window_capacity", self.reliable_window_capacity), + ("max_datagram_size", self.max_datagram_size), + ("max_decoded_packet_size", self.max_decoded_packet_size), + ("protocol_mtu", self.protocol_mtu), + ("max_peers", self.max_peers), + ] { + if value == 0 { + return Err(UdpTransportError::InvalidConfiguration(name)); + } + } + if self.pending_ack_capacity > usize::from(u8::MAX) { + return Err(UdpTransportError::InvalidConfiguration( + "pending_ack_capacity", + )); + } + if self.max_pending_acks > self.pending_ack_capacity { + return Err(UdpTransportError::InvalidConfiguration("max_pending_acks")); + } + if self.protocol_mtu < 10 || self.protocol_mtu > self.max_datagram_size { + return Err(UdpTransportError::InvalidConfiguration("protocol_mtu")); + } + if self.max_decoded_packet_size < self.max_datagram_size { + return Err(UdpTransportError::InvalidConfiguration( + "max_decoded_packet_size", + )); + } + if self.resend_timeout.is_zero() { + return Err(UdpTransportError::InvalidConfiguration("resend_timeout")); + } + if self.network_tick_interval.is_zero() { + return Err(UdpTransportError::InvalidConfiguration( + "network_tick_interval", + )); + } + Ok(()) + } +} + +/// Immutable, redacted snapshot of transport counters. +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +pub struct UdpTransportStats { + pub received_datagrams: u64, + pub received_bytes: u64, + pub sent_datagrams: u64, + pub sent_bytes: u64, + pub dropped_receive_queue: u64, + pub dropped_send_queue: u64, + pub malformed_datagrams: u64, + pub rejected_sources: u64, + pub duplicate_datagrams: u64, + pub out_of_order_datagrams: u64, + pub sequence_gaps: u64, + pub acknowledgements_received: u64, + pub acknowledgements_sent: u64, + pub resent_datagrams: u64, + pub failed_resends: u64, + pub socket_errors: u64, + pub handler_panics: u64, +} + +#[derive(Default)] +struct StatsCounters { + received_datagrams: AtomicU64, + received_bytes: AtomicU64, + sent_datagrams: AtomicU64, + sent_bytes: AtomicU64, + dropped_receive_queue: AtomicU64, + dropped_send_queue: AtomicU64, + malformed_datagrams: AtomicU64, + rejected_sources: AtomicU64, + duplicate_datagrams: AtomicU64, + out_of_order_datagrams: AtomicU64, + sequence_gaps: AtomicU64, + acknowledgements_received: AtomicU64, + acknowledgements_sent: AtomicU64, + resent_datagrams: AtomicU64, + failed_resends: AtomicU64, + socket_errors: AtomicU64, + handler_panics: AtomicU64, +} + +impl StatsCounters { + fn snapshot(&self) -> UdpTransportStats { + let load = |counter: &AtomicU64| counter.load(Ordering::Relaxed); + UdpTransportStats { + received_datagrams: load(&self.received_datagrams), + received_bytes: load(&self.received_bytes), + sent_datagrams: load(&self.sent_datagrams), + sent_bytes: load(&self.sent_bytes), + dropped_receive_queue: load(&self.dropped_receive_queue), + dropped_send_queue: load(&self.dropped_send_queue), + malformed_datagrams: load(&self.malformed_datagrams), + rejected_sources: load(&self.rejected_sources), + duplicate_datagrams: load(&self.duplicate_datagrams), + out_of_order_datagrams: load(&self.out_of_order_datagrams), + sequence_gaps: load(&self.sequence_gaps), + acknowledgements_received: load(&self.acknowledgements_received), + acknowledgements_sent: load(&self.acknowledgements_sent), + resent_datagrams: load(&self.resent_datagrams), + failed_resends: load(&self.failed_resends), + socket_errors: load(&self.socket_errors), + handler_panics: load(&self.handler_panics), + } + } +} + +/// Callback boundary corresponding to the protected methods on C# `UDPBase`. +pub trait UdpPacketHandler: Send + Sync + 'static { + fn packet_received(&self, _buffer: UDPPacketBuffer) {} + fn packet_sent(&self, _buffer: UDPPacketBuffer, _bytes_sent: usize) {} + fn packet_dropped(&self) {} +} + +struct NoopPacketHandler; +impl UdpPacketHandler for NoopPacketHandler {} + +struct RuntimeTasks { + cancellation: CancellationTokenSource, + commands: mpsc::Sender, + handles: Vec>, + local_address: SocketAddr, +} + +struct UdpInner { + bind_address: SocketAddr, + remote_end_point: Option, + config: UdpTransportConfig, + handler: Arc, + parent_cancellation: CancellationToken, + running: AtomicBool, + runtime: Mutex>, + stats: Arc, +} + +impl Drop for UdpInner { + fn drop(&mut self) { + if let Some(runtime) = self + .runtime + .get_mut() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take() + { + runtime.cancellation.cancel(); + for handle in runtime.handles { + handle.abort(); + } + } + } +} + +/// Runtime-owning native implementation of the C# UDP transport slice. +/// +/// Construction creates no socket and no runtime. `start` binds a nonblocking +/// socket and spawns tasks on the caller's current Tokio runtime. +#[derive(Clone)] +pub struct UDPBase { + inner: Arc, +} + +impl UDPBase { + pub fn client( + remote_end_point: SocketAddr, + config: UdpTransportConfig, + handler: Arc, + parent_cancellation: CancellationToken, + ) -> Result { + let bind_address = match remote_end_point.ip() { + IpAddr::V4(_) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0), + IpAddr::V6(_) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0), + }; + Self::new( + bind_address, + Some(remote_end_point), + config, + handler, + parent_cancellation, + ) + } + + pub fn server( + bind_address: SocketAddr, + config: UdpTransportConfig, + handler: Arc, + parent_cancellation: CancellationToken, + ) -> Result { + Self::new(bind_address, None, config, handler, parent_cancellation) + } + + pub fn client_with_defaults(remote_end_point: SocketAddr) -> Result { + Self::client( + remote_end_point, + UdpTransportConfig::default(), + Arc::new(NoopPacketHandler), + CancellationToken::default(), + ) + } + + fn new( + bind_address: SocketAddr, + remote_end_point: Option, + config: UdpTransportConfig, + handler: Arc, + parent_cancellation: CancellationToken, + ) -> Result { + config.validate()?; + Ok(Self { + inner: Arc::new(UdpInner { + bind_address, + remote_end_point, + config, + handler, + parent_cancellation, + running: AtomicBool::new(false), + runtime: Mutex::new(None), + stats: Arc::new(StatsCounters::default()), + }), + }) + } + + pub fn start(&self) -> Result<(), Error> { + self.start_transport().map_err(Into::into) + } + + pub fn start_transport(&self) -> Result<(), UdpTransportError> { + let mut runtime = self + .inner + .runtime + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if runtime.is_some() { + return Ok(()); + } + if self.inner.parent_cancellation.is_cancellation_requested() { + return Err(UdpTransportError::Cancelled); + } + let handle = tokio::runtime::Handle::try_current() + .map_err(|_| UdpTransportError::RuntimeUnavailable)?; + let std_socket = std::net::UdpSocket::bind(self.inner.bind_address) + .map_err(|_| UdpTransportError::Socket)?; + std_socket + .set_nonblocking(true) + .map_err(|_| UdpTransportError::Socket)?; + let local_address = std_socket + .local_addr() + .map_err(|_| UdpTransportError::Socket)?; + let socket = + Arc::new(UdpSocket::from_std(std_socket).map_err(|_| UdpTransportError::Socket)?); + let cancellation = CancellationTokenSource::new_linked(std::slice::from_ref( + &self.inner.parent_cancellation, + )); + let token = cancellation.token(); + let (raw_sender, raw_receiver) = mpsc::channel(self.inner.config.receive_queue_capacity); + let (command_sender, command_receiver) = + mpsc::channel(self.inner.config.command_queue_capacity); + let (write_sender, write_receiver) = mpsc::channel(self.inner.config.write_queue_capacity); + + let receiver_handle = handle.spawn(receive_loop( + Arc::clone(&socket), + raw_sender, + token.clone(), + self.inner.config.max_datagram_size, + Arc::clone(&self.inner.stats), + Arc::clone(&self.inner.handler), + )); + let coordinator_handle = handle.spawn(coordinator_loop( + raw_receiver, + command_receiver, + write_sender, + token.clone(), + self.inner.remote_end_point, + self.inner.config.clone(), + Arc::clone(&self.inner.stats), + Arc::clone(&self.inner.handler), + )); + let writer_handle = handle.spawn(writer_loop( + socket, + write_receiver, + token, + self.inner.config.throttle.clone(), + Arc::clone(&self.inner.stats), + Arc::clone(&self.inner.handler), + )); + *runtime = Some(RuntimeTasks { + cancellation, + commands: command_sender, + handles: vec![receiver_handle, coordinator_handle, writer_handle], + local_address, + }); + self.inner.running.store(true, Ordering::Release); + Ok(()) + } + + pub fn stop(&self) -> Result<(), Error> { + self.stop_transport(); + Ok(()) + } + + pub fn stop_transport(&self) { + let runtime = self + .inner + .runtime + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take(); + self.inner.running.store(false, Ordering::Release); + if let Some(runtime) = runtime { + runtime.cancellation.cancel(); + for handle in runtime.handles { + handle.abort(); + } + } + } + + pub async fn stop_async(&self) { + let runtime = self + .inner + .runtime + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take(); + self.inner.running.store(false, Ordering::Release); + if let Some(runtime) = runtime { + runtime.cancellation.cancel(); + drop(runtime.commands); + for handle in runtime.handles { + let _ = handle.await; + } + } + } + + #[must_use] + pub fn is_running(&self) -> bool { + self.inner.running.load(Ordering::Acquire) + } + + #[must_use] + pub fn local_address(&self) -> Option { + self.inner + .runtime + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .as_ref() + .map(|runtime| runtime.local_address) + } + + #[must_use] + pub fn stats(&self) -> UdpTransportStats { + self.inner.stats.snapshot() + } + + pub fn async_begin_send(&self, buf: UDPPacketBuffer) -> Result<(), Error> { + let Some(sender) = self.command_sender() else { + return Ok(()); + }; + sender + .try_send(CoordinatorCommand::Raw(buf)) + .map_err(|error| { + if matches!(error, mpsc::error::TrySendError::Full(_)) { + self.inner + .stats + .dropped_send_queue + .fetch_add(1, Ordering::Relaxed); + Error::InvalidOperation + } else { + Error::Cancelled + } + }) + } + + pub async fn send_packet( + &self, + data: Vec, + packet_type: PacketType, + do_zerocode: bool, + cancellation: CancellationToken, + ) -> Result { + let destination = self + .inner + .remote_end_point + .ok_or(UdpTransportError::InvalidBuffer)?; + self.send_packet_to(data, destination, packet_type, do_zerocode, cancellation) + .await + } + + pub async fn send_packet_to( + &self, + data: Vec, + destination: SocketAddr, + packet_type: PacketType, + do_zerocode: bool, + cancellation: CancellationToken, + ) -> Result { + let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?; + let (response_sender, response_receiver) = oneshot::channel(); + let command = CoordinatorCommand::Packet { + data, + destination, + packet_type, + do_zerocode, + response: response_sender, + }; + tokio::select! { + result = sender.send(command) => { + result.map_err(|_| UdpTransportError::Cancelled)?; + } + () = cancellation.cancelled() => return Err(UdpTransportError::Cancelled), + () = self.inner.parent_cancellation.cancelled() => { + return Err(UdpTransportError::Cancelled); + } + } + tokio::select! { + result = response_receiver => result.map_err(|_| UdpTransportError::Cancelled)?, + () = cancellation.cancelled() => Err(UdpTransportError::Cancelled), + () = self.inner.parent_cancellation.cancelled() => Err(UdpTransportError::Cancelled), + } + } + + pub fn try_send_packet( + &self, + data: Vec, + destination: SocketAddr, + packet_type: PacketType, + do_zerocode: bool, + ) -> Result>, UdpTransportError> { + let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?; + let (response, receiver) = oneshot::channel(); + sender + .try_send(CoordinatorCommand::Packet { + data, + destination, + packet_type, + do_zerocode, + response, + }) + .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) + } + + pub fn update_throttle(&self, throttle: AgentThrottle) -> Result<(), UdpTransportError> { + let sender = self.command_sender().ok_or(UdpTransportError::NotRunning)?; + sender + .try_send(CoordinatorCommand::UpdateThrottle(throttle)) + .map_err(|error| { + if matches!(error, mpsc::error::TrySendError::Full(_)) { + UdpTransportError::Backpressure + } else { + UdpTransportError::Cancelled + } + }) + } + + fn command_sender(&self) -> Option> { + self.inner + .runtime + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .as_ref() + .map(|runtime| runtime.commands.clone()) + } +} + +impl fmt::Debug for UDPBase { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("UDPBase") + .field( + "mode", + &self.inner.remote_end_point.map_or("server", |_| "client"), + ) + .field("is_running", &self.is_running()) + .field("stats", &self.stats()) + .finish_non_exhaustive() + } +} + +enum CoordinatorCommand { + Raw(UDPPacketBuffer), + Packet { + data: Vec, + destination: SocketAddr, + packet_type: PacketType, + do_zerocode: bool, + response: oneshot::Sender>, + }, + UpdateThrottle(AgentThrottle), +} + +enum WriteCommand { + Datagram { + buffer: UDPPacketBuffer, + category: UdpThrottleCategory, + }, + UpdateThrottle(AgentThrottle), +} + +struct ReliablePacket { + buffer: UDPPacketBuffer, + category: UdpThrottleCategory, + last_sent: Instant, + resend_count: u32, +} + +struct PeerState { + sequence: u32, + archive: IncomingPacketIDCollection, + pending_acks: VecDeque, + need_ack: BTreeMap, + latest_received: Option, +} + +impl PeerState { + fn new(config: &UdpTransportConfig) -> Result { + Ok(Self { + sequence: 0, + archive: IncomingPacketIDCollection::new( + i32::try_from(config.packet_archive_size) + .map_err(|_| UdpTransportError::InvalidConfiguration("packet_archive_size"))?, + ) + .map_err(|_| UdpTransportError::InvalidConfiguration("packet_archive_size"))?, + pending_acks: VecDeque::with_capacity(config.pending_ack_capacity), + need_ack: BTreeMap::new(), + latest_received: None, + }) + } + + fn next_sequence(&mut self) -> u32 { + // C# increments an Int32 and casts it to UInt32 for the four-byte + // header. The equivalent wire sequence spans all u32 values and wraps + // through zero. + self.sequence = self.sequence.wrapping_add(1); + self.sequence + } + + fn observe_sequence(&mut self, sequence: u32, stats: &StatsCounters) { + let Some(latest) = self.latest_received else { + self.latest_received = Some(sequence); + return; + }; + let expected = latest.wrapping_add(1); + if sequence == expected { + self.latest_received = Some(sequence); + return; + } + let forward = sequence.wrapping_sub(latest); + if forward != 0 && forward <= (u32::MAX / 2) + 1 { + stats.sequence_gaps.fetch_add(1, Ordering::Relaxed); + self.latest_received = Some(sequence); + } else { + stats.out_of_order_datagrams.fetch_add(1, Ordering::Relaxed); + } + } +} + +struct CoordinatorState { + peers: HashMap, + remote_end_point: Option, + config: UdpTransportConfig, + stats: Arc, +} + +impl CoordinatorState { + fn canonical_peer(&self, source: SocketAddr) -> Result { + if let Some(remote) = self.remote_end_point { + if remote.ip() != source.ip() { + return Err(UdpTransportError::InvalidBuffer); + } + Ok(remote) + } else { + Ok(source) + } + } + + fn peer_mut(&mut self, endpoint: SocketAddr) -> Result<&mut PeerState, UdpTransportError> { + if !self.peers.contains_key(&endpoint) { + if self.peers.len() >= self.config.max_peers { + return Err(UdpTransportError::TooManyPeers); + } + let peer = PeerState::new(&self.config)?; + self.peers.insert(endpoint, peer); + } + self.peers + .get_mut(&endpoint) + .ok_or(UdpTransportError::TooManyPeers) + } + + fn prepare_packet( + &mut self, + data: Vec, + destination: SocketAddr, + packet_type: PacketType, + do_zerocode: bool, + ) -> Result<(WriteCommand, u32), UdpTransportError> { + let mut data = encode_for_transport(data, do_zerocode, self.config.protocol_mtu)?; + if data.len() < 6 { + return Err(UdpTransportError::InvalidBuffer); + } + let reliable = data[0] & Helpers::MSG_RELIABLE != 0; + let category = UdpThrottleCategory::classify(packet_type); + let reliable_capacity = self.config.reliable_window_capacity; + let protocol_mtu = self.config.protocol_mtu; + let peer = self.peer_mut(destination)?; + if reliable && peer.need_ack.len() >= reliable_capacity { + return Err(UdpTransportError::ReliableWindowFull); + } + append_pending_acks(&mut data, &mut peer.pending_acks, protocol_mtu)?; + let sequence = peer.next_sequence(); + data[1..5].copy_from_slice(&sequence.to_be_bytes()); + let data_length = + i32::try_from(data.len()).map_err(|_| UdpTransportError::InvalidBuffer)?; + let buffer = UDPPacketBuffer { + data, + data_length, + remote_end_point: destination, + }; + if reliable { + peer.need_ack.insert( + sequence, + ReliablePacket { + buffer: buffer.clone(), + category, + last_sent: Instant::now(), + resend_count: 0, + }, + ); + } + Ok((WriteCommand::Datagram { buffer, category }, sequence)) + } + + fn take_resends(&mut self) -> Vec { + let now = Instant::now(); + let mut writes = Vec::new(); + for peer in self.peers.values_mut() { + let mut failed = Vec::new(); + for (&sequence, packet) in &mut peer.need_ack { + if now.duration_since(packet.last_sent) <= self.config.resend_timeout { + continue; + } + if packet.resend_count < self.config.max_resend_count { + packet.resend_count += 1; + packet.last_sent = now; + if let Some(flags) = packet.buffer.data.first_mut() { + *flags |= Helpers::MSG_RESENT; + } + self.stats.resent_datagrams.fetch_add(1, Ordering::Relaxed); + writes.push(WriteCommand::Datagram { + buffer: packet.buffer.clone(), + category: packet.category, + }); + } else { + failed.push(sequence); + } + } + for sequence in failed { + peer.need_ack.remove(&sequence); + self.stats.failed_resends.fetch_add(1, Ordering::Relaxed); + } + } + writes + } +} + +fn encode_for_transport( + mut data: Vec, + do_zerocode: bool, + mtu: usize, +) -> Result, UdpTransportError> { + if data.len() < 6 { + return Err(UdpTransportError::InvalidBuffer); + } + if data.len() > mtu { + return Err(UdpTransportError::MtuExceeded); + } + if !do_zerocode { + return Ok(data); + } + data[0] |= Helpers::MSG_ZEROCODED; + let mut encoded = vec![0_u8; mtu]; + let source_length = i32::try_from(data.len()).map_err(|_| UdpTransportError::InvalidBuffer)?; + match crate::packet_wire::zero_encode(Some(&data), source_length, Some(&mut encoded)) { + Ok(length) => { + let length = usize::try_from(length).map_err(|_| UdpTransportError::InvalidBuffer)?; + encoded.truncate(length); + Ok(encoded) + } + Err(Error::IndexOutOfRange) => { + data[0] &= !Helpers::MSG_ZEROCODED; + Ok(data) + } + Err(_) => Err(UdpTransportError::InvalidBuffer), + } +} + +fn append_pending_acks( + data: &mut Vec, + pending: &mut VecDeque, + mtu: usize, +) -> Result<(), UdpTransportError> { + if data + .first() + .is_some_and(|flags| flags & Helpers::MSG_APPENDED_ACKS != 0) + { + return Err(UdpTransportError::InvalidBuffer); + } + let mut count = 0_u8; + while data.len().checked_add(5).is_some_and(|length| length < mtu) { + let Some(ack) = pending.pop_front() else { + break; + }; + data.extend_from_slice(&ack.to_be_bytes()); + count = count.saturating_add(1); + if count == u8::MAX { + break; + } + } + if count != 0 { + data.push(count); + data[0] |= Helpers::MSG_APPENDED_ACKS; + } + Ok(()) +} + +async fn receive_loop( + socket: Arc, + sender: mpsc::Sender, + cancellation: CancellationToken, + max_datagram_size: usize, + stats: Arc, + handler: Arc, +) { + let mut storage = vec![0_u8; max_datagram_size.saturating_add(1)]; + loop { + let received = tokio::select! { + () = cancellation.cancelled() => break, + result = socket.recv_from(&mut storage) => result, + }; + let Ok((length, source)) = received else { + if !cancellation.is_cancellation_requested() { + stats.socket_errors.fetch_add(1, Ordering::Relaxed); + } + break; + }; + stats.received_datagrams.fetch_add(1, Ordering::Relaxed); + stats + .received_bytes + .fetch_add(u64::try_from(length).unwrap_or(u64::MAX), Ordering::Relaxed); + if length > max_datagram_size { + stats.malformed_datagrams.fetch_add(1, Ordering::Relaxed); + continue; + } + let packet = UDPPacketBuffer { + data: storage[..length].to_vec(), + data_length: i32::try_from(length).unwrap_or(i32::MAX), + remote_end_point: source, + }; + if sender.try_send(packet).is_err() { + stats.dropped_receive_queue.fetch_add(1, Ordering::Relaxed); + invoke_handler(&stats, || handler.packet_dropped()); + } + } +} + +#[allow(clippy::too_many_arguments)] // Each argument is one explicit task ownership boundary. +async fn coordinator_loop( + mut incoming: mpsc::Receiver, + mut commands: mpsc::Receiver, + writer: mpsc::Sender, + cancellation: CancellationToken, + remote_end_point: Option, + config: UdpTransportConfig, + stats: Arc, + handler: Arc, +) { + let mut coordinator = CoordinatorState { + peers: HashMap::new(), + remote_end_point, + config: config.clone(), + stats: Arc::clone(&stats), + }; + let mut tick = tokio::time::interval(config.network_tick_interval); + tick.set_missed_tick_behavior(MissedTickBehavior::Skip); + tick.tick().await; + loop { + tokio::select! { + biased; + () = cancellation.cancelled() => break, + Some(buffer) = incoming.recv() => { + process_incoming(&mut coordinator, buffer, &writer, &handler).await; + } + Some(command) = commands.recv() => { + process_command(&mut coordinator, command, &writer).await; + } + _ = tick.tick() => { + flush_pending_acks(&mut coordinator, &writer).await; + for resend in coordinator.take_resends() { + if writer.send(resend).await.is_err() { + return; + } + } + } + else => break, + } + } +} + +async fn process_command( + state: &mut CoordinatorState, + command: CoordinatorCommand, + writer: &mpsc::Sender, +) { + match command { + CoordinatorCommand::Raw(buffer) => { + let valid = buffer.payload().is_ok() + && usize::try_from(buffer.data_length) + .is_ok_and(|length| length <= state.config.max_datagram_size); + if !valid { + state + .stats + .dropped_send_queue + .fetch_add(1, Ordering::Relaxed); + return; + } + let _ = writer + .send(WriteCommand::Datagram { + buffer, + category: UdpThrottleCategory::Unthrottled, + }) + .await; + } + CoordinatorCommand::Packet { + data, + destination, + packet_type, + do_zerocode, + response, + } => { + let result = state.prepare_packet(data, destination, packet_type, do_zerocode); + let result = match result { + Ok((write, sequence)) => writer + .send(write) + .await + .map(|()| sequence) + .map_err(|_| UdpTransportError::Cancelled), + Err(error) => Err(error), + }; + let _ = response.send(result); + } + CoordinatorCommand::UpdateThrottle(throttle) => { + let _ = writer.send(WriteCommand::UpdateThrottle(throttle)).await; + } + } +} + +async fn process_incoming( + state: &mut CoordinatorState, + buffer: UDPPacketBuffer, + writer: &mpsc::Sender, + handler: &Arc, +) { + let Ok(endpoint) = state.canonical_peer(buffer.remote_end_point) else { + state.stats.rejected_sources.fetch_add(1, Ordering::Relaxed); + return; + }; + let Ok(payload) = buffer.payload() else { + state + .stats + .malformed_datagrams + .fetch_add(1, Ordering::Relaxed); + return; + }; + let Ok(payload_length) = i32::try_from(payload.len()) else { + state + .stats + .malformed_datagrams + .fetch_add(1, Ordering::Relaxed); + return; + }; + let mut packet_end = payload_length - 1; + let mut zero_buffer = vec![0_u8; state.config.max_decoded_packet_size]; + let Ok(packet) = + crate::packet_wire::build_packet_from_bytes(payload, &mut packet_end, &mut zero_buffer) + else { + state + .stats + .malformed_datagrams + .fetch_add(1, Ordering::Relaxed); + return; + }; + + let appended_acks = packet.header.ack_list.clone().unwrap_or_default(); + let mut standalone_acks = Vec::new(); + if packet.type_ == PacketType::PacketAck { + let mut position = 0; + if let Ok(ack_packet) = + PacketAckPacket::new_with_bytes_int32(payload.to_vec(), &mut position) + { + standalone_acks.extend(ack_packet.packets.into_iter().map(|block| block.id)); + } else { + state + .stats + .malformed_datagrams + .fetch_add(1, Ordering::Relaxed); + return; + } + } + + let pending_threshold = state.config.max_pending_acks; + let pending_capacity = state.config.pending_ack_capacity; + let counters = Arc::clone(&state.stats); + let Ok(peer) = state.peer_mut(endpoint) else { + counters.rejected_sources.fetch_add(1, Ordering::Relaxed); + return; + }; + for ack in appended_acks.into_iter().chain(standalone_acks) { + peer.need_ack.remove(&ack); + counters + .acknowledgements_received + .fetch_add(1, Ordering::Relaxed); + } + + if packet.header.reliable { + if peer.pending_acks.len() < pending_capacity { + peer.pending_acks.push_back(packet.header.sequence); + } + if !peer.archive.try_enqueue(packet.header.sequence) { + counters.duplicate_datagrams.fetch_add(1, Ordering::Relaxed); + if peer.pending_acks.len() >= pending_threshold { + send_peer_acks(state, endpoint, writer).await; + } + return; + } + peer.observe_sequence(packet.header.sequence, &counters); + } + + if state + .peers + .get(&endpoint) + .is_some_and(|peer| peer.pending_acks.len() >= pending_threshold) + { + send_peer_acks(state, endpoint, writer).await; + } + invoke_handler(&state.stats, || handler.packet_received(buffer)); +} + +async fn flush_pending_acks(state: &mut CoordinatorState, writer: &mpsc::Sender) { + let endpoints: Vec<_> = state + .peers + .iter() + .filter_map(|(endpoint, peer)| (!peer.pending_acks.is_empty()).then_some(*endpoint)) + .collect(); + for endpoint in endpoints { + send_peer_acks(state, endpoint, writer).await; + } +} + +async fn send_peer_acks( + state: &mut CoordinatorState, + endpoint: SocketAddr, + writer: &mpsc::Sender, +) { + let Some(peer) = state.peers.get_mut(&endpoint) else { + return; + }; + if peer.pending_acks.is_empty() { + return; + } + let mut blocks = Vec::with_capacity(peer.pending_acks.len()); + while let Some(id) = peer.pending_acks.pop_front() { + blocks.push(PacketAckPacketPacketsBlock { id }); + } + let ack_count = blocks.len(); + let Ok(mut packet) = PacketAckPacket::new_with_constructor() else { + return; + }; + packet.packets = blocks; + let bytes = match packet.to_bytes_with_method() { + Ok(mut bytes) => { + // C# `SendAcks` explicitly marks standalone PacketAck packets as + // unreliable before serialization. + if let Some(flags) = bytes.first_mut() { + *flags &= !Helpers::MSG_RELIABLE; + } + bytes + } + Err(_) => return, + }; + let Ok((write, _)) = state.prepare_packet(bytes, endpoint, PacketType::PacketAck, false) else { + return; + }; + if writer.send(write).await.is_ok() { + state.stats.acknowledgements_sent.fetch_add( + u64::try_from(ack_count).unwrap_or(u64::MAX), + Ordering::Relaxed, + ); + } +} + +fn invoke_handler(stats: &StatsCounters, callback: impl FnOnce()) { + if catch_unwind(AssertUnwindSafe(callback)).is_err() { + stats.handler_panics.fetch_add(1, Ordering::Relaxed); + } +} + +struct TokenBucket { + tokens_per_period: usize, + token_limit: usize, + available: usize, + last_replenishment: Instant, +} + +impl TokenBucket { + #[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)] + fn new(bits_per_second: f32) -> Self { + let calculated = + (f64::from(bits_per_second) / 8.0 * THROTTLE_PERIOD.as_secs_f64()).max(0.0) as usize; + let tokens_per_period = calculated.max(THROTTLE_MIN_BYTES_PER_PERIOD); + let token_limit = tokens_per_period.saturating_mul(THROTTLE_BURST_PERIODS); + Self { + tokens_per_period, + token_limit, + available: token_limit, + last_replenishment: Instant::now(), + } + } + + fn replenish(&mut self, now: Instant) { + let elapsed = now.duration_since(self.last_replenishment); + let periods = elapsed.as_nanos() / THROTTLE_PERIOD.as_nanos(); + if periods == 0 { + return; + } + let periods = usize::try_from(periods).unwrap_or(usize::MAX); + self.available = self + .available + .saturating_add(self.tokens_per_period.saturating_mul(periods)) + .min(self.token_limit); + let periods_u32 = u32::try_from(periods).unwrap_or(u32::MAX); + self.last_replenishment += THROTTLE_PERIOD.saturating_mul(periods_u32); + } + + async fn acquire( + &mut self, + amount: usize, + cancellation: &CancellationToken, + ) -> Result<(), UdpTransportError> { + let amount = amount.clamp(1, self.token_limit); + loop { + self.replenish(Instant::now()); + if self.available >= amount { + self.available -= amount; + return Ok(()); + } + let deficit = amount - self.available; + let periods = deficit.div_ceil(self.tokens_per_period); + let periods = u32::try_from(periods).unwrap_or(u32::MAX); + let deadline = self.last_replenishment + THROTTLE_PERIOD.saturating_mul(periods.max(1)); + tokio::select! { + () = cancellation.cancelled() => return Err(UdpTransportError::Cancelled), + () = tokio::time::sleep_until(deadline) => {} + } + } + } +} + +struct UdpThrottle { + task: TokenBucket, + texture: TokenBucket, + asset: TokenBucket, +} + +impl UdpThrottle { + fn new(throttle: &AgentThrottle) -> Self { + Self { + task: TokenBucket::new(throttle.task()), + texture: TokenBucket::new(throttle.texture()), + asset: TokenBucket::new(throttle.asset()), + } + } + + async fn acquire( + &mut self, + category: UdpThrottleCategory, + amount: usize, + cancellation: &CancellationToken, + ) -> Result<(), UdpTransportError> { + match category { + UdpThrottleCategory::Unthrottled => Ok(()), + UdpThrottleCategory::Task => self.task.acquire(amount, cancellation).await, + UdpThrottleCategory::Texture => self.texture.acquire(amount, cancellation).await, + UdpThrottleCategory::Asset => self.asset.acquire(amount, cancellation).await, + } + } +} + +async fn writer_loop( + socket: Arc, + mut receiver: mpsc::Receiver, + cancellation: CancellationToken, + initial_throttle: AgentThrottle, + stats: Arc, + handler: Arc, +) { + let mut throttle = UdpThrottle::new(&initial_throttle); + loop { + let command = tokio::select! { + biased; + () = cancellation.cancelled() => break, + command = receiver.recv() => command, + }; + match command { + Some(WriteCommand::UpdateThrottle(values)) => { + throttle = UdpThrottle::new(&values); + } + Some(WriteCommand::Datagram { buffer, category }) => { + let Ok(payload) = buffer.payload() else { + stats.dropped_send_queue.fetch_add(1, Ordering::Relaxed); + continue; + }; + if throttle + .acquire(category, payload.len(), &cancellation) + .await + .is_err() + { + break; + } + match socket.send_to(payload, buffer.remote_end_point).await { + Ok(bytes_sent) => { + stats.sent_datagrams.fetch_add(1, Ordering::Relaxed); + stats.sent_bytes.fetch_add( + u64::try_from(bytes_sent).unwrap_or(u64::MAX), + Ordering::Relaxed, + ); + invoke_handler(&stats, || handler.packet_sent(buffer, bytes_sent)); + } + Err(_) => { + stats.socket_errors.fetch_add(1, Ordering::Relaxed); + } + } + } + None => break, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn sequence_rolls_over_across_the_full_reference_header_width() { + let config = UdpTransportConfig::default(); + let mut peer = PeerState::new(&config).unwrap(); + peer.sequence = u32::MAX - 1; + assert_eq!(peer.next_sequence(), u32::MAX); + assert_eq!(peer.next_sequence(), 0); + assert_eq!(peer.next_sequence(), 1); + } + + #[tokio::test(start_paused = true)] + async fn fake_time_drives_token_replenishment_and_reliable_expiry() { + let cancellation = CancellationToken::default(); + let mut bucket = TokenBucket::new(4_000.0); + bucket.acquire(800, &cancellation).await.unwrap(); + let started = Instant::now(); + bucket.acquire(800, &cancellation).await.unwrap(); + assert_eq!( + Instant::now().duration_since(started), + Duration::from_millis(400) + ); + + let config = UdpTransportConfig { + resend_timeout: Duration::from_secs(4), + max_resend_count: 1, + ..UdpTransportConfig::default() + }; + let mut coordinator = CoordinatorState { + peers: HashMap::new(), + remote_end_point: None, + config, + stats: Arc::new(StatsCounters::default()), + }; + 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) + .unwrap(); + assert!(coordinator.take_resends().is_empty()); + + tokio::time::advance(Duration::from_millis(4_001)).await; + let resend = coordinator.take_resends(); + assert_eq!(resend.len(), 1); + let WriteCommand::Datagram { buffer, .. } = &resend[0] else { + panic!("expected resend datagram"); + }; + assert_ne!(buffer.data[0] & Helpers::MSG_RESENT, 0); + assert_eq!(u32::from_be_bytes(buffer.data[1..5].try_into().unwrap()), 1); + + tokio::time::advance(Duration::from_millis(4_001)).await; + assert!(coordinator.take_resends().is_empty()); + assert_eq!(coordinator.stats.failed_resends.load(Ordering::Relaxed), 1); + assert!(coordinator.peers[&destination].need_ack.is_empty()); + } +} diff --git a/crates/libremetaverse/tests/udp_transport.rs b/crates/libremetaverse/tests/udp_transport.rs new file mode 100644 index 0000000..3eb6eec --- /dev/null +++ b/crates/libremetaverse/tests/udp_transport.rs @@ -0,0 +1,511 @@ +use libremetaverse::packets::{ + AgentThrottlePacket, CompletePingCheckPacket, Packet, PacketAckPacket, + PacketAckPacketPacketsBlock, PacketType, UseCircuitCodePacket, +}; +use libremetaverse::{ + AgentThrottle, AgentThrottleSender, GridClient, Helpers, IncomingPacketIDCollection, Simulator, + UDPBase, UDPPacketBuffer, UdpPacketHandler, UdpTransportConfig, UdpTransportError, +}; +use libremetaverse_types::Error; +use libremetaverse_types::compat::{Array, CancellationToken, Object}; +use std::net::{IpAddr, Ipv4Addr, SocketAddr}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use tokio::net::UdpSocket; + +#[derive(Default)] +struct RecordingHandler { + received: Mutex>>, + sent: Mutex>>, + drops: Mutex, +} + +impl UdpPacketHandler for RecordingHandler { + fn packet_received(&self, buffer: UDPPacketBuffer) { + self.received + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .push(buffer.data[..usize::try_from(buffer.data_length).unwrap()].to_vec()); + } + + fn packet_sent(&self, buffer: UDPPacketBuffer, bytes_sent: usize) { + self.sent + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .push(buffer.data[..bytes_sent].to_vec()); + } + + fn packet_dropped(&self) { + *self + .drops + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) += 1; + } +} + +fn loopback() -> SocketAddr { + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0) +} + +fn reliable_packet() -> Vec { + let mut packet = UseCircuitCodePacket::new_with_constructor().expect("packet constructor"); + packet.circuit_code.code = 0x1122_3344; + packet.to_bytes_with_method().expect("packet bytes") +} + +fn packet_ack(sequence: u32) -> Vec { + let mut packet = PacketAckPacket::new_with_constructor().expect("ACK constructor"); + packet.packets = vec![PacketAckPacketPacketsBlock { id: sequence }]; + let mut bytes = packet.to_bytes_with_method().expect("ACK bytes"); + bytes[0] &= !Helpers::MSG_RELIABLE; + bytes +} + +fn complete_ping_packet() -> Vec { + let mut packet = CompletePingCheckPacket::new_with_constructor().expect("ping constructor"); + packet.ping_id.ping_id = 7; + packet.to_bytes_with_method().expect("ping bytes") +} + +async fn wait_until(mut condition: impl FnMut() -> bool) { + tokio::time::timeout(Duration::from_secs(2), async { + while !condition() { + tokio::task::yield_now().await; + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .expect("condition timed out"); +} + +#[test] +fn packet_buffer_constructors_and_copy_match_the_reference() { + let endpoint = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 13000); + let mut packet = UDPPacketBuffer::new_with_ip_end_point_int32(endpoint, 4).unwrap(); + assert_eq!(packet.data, vec![0; 4]); + assert_eq!(packet.data_length, 0); + + packet + .copy_from_with_array(Array(vec![Object::Bytes(vec![1, 2, 3, 4])])) + .unwrap(); + assert_eq!(packet.data, vec![1, 2, 3, 4]); + assert_eq!(packet.data_length, 0); + assert_eq!( + packet.copy_from_slice(&[1, 2, 3, 4, 5]), + Err(Error::IndexOutOfRange) + ); + packet.reset_endpoint().unwrap(); + assert_eq!(packet.remote_end_point, "0.0.0.0:0".parse().unwrap()); + + let borrowed = vec![9, 8, 7]; + let packet = UDPPacketBuffer::new_with_ip_end_point_bytes(endpoint, borrowed).unwrap(); + assert_eq!(packet.data, vec![9, 8, 7]); + assert_eq!(packet.data_length, 0); +} + +#[test] +fn incoming_packet_archive_is_bounded_and_rejects_duplicates() { + let archive = IncomingPacketIDCollection::new(3).unwrap(); + assert!(archive.try_enqueue(1)); + assert!(!archive.try_enqueue(1)); + assert!(archive.try_enqueue(2)); + // The C# ring keeps one slot empty, so adding the third item evicts 1. + assert!(archive.try_enqueue(3)); + assert!(archive.try_enqueue(1)); + assert_eq!( + IncomingPacketIDCollection::new(0).unwrap_err(), + Error::Argument + ); +} + +#[test] +fn agent_throttle_clamps_and_round_trips_golden_little_endian_values() { + let throttle = AgentThrottle::default(); + assert_eq!(throttle.resend(), 150_000.0); + assert_eq!(throttle.land(), 170_000.0); + assert_eq!(throttle.wind(), 34_000.0); + assert_eq!(throttle.cloud(), 34_000.0); + assert_eq!(throttle.task(), 360_448.0); + assert_eq!(throttle.texture(), 360_448.0); + assert_eq!(throttle.asset(), 220_000.0); + assert_eq!(throttle.total(), 1_328_896.0); + + let bytes = throttle.to_bytes().unwrap(); + assert_eq!(bytes.len(), 28); + let decoded = AgentThrottle::new_with_bytes_int32(bytes, 0).unwrap(); + assert_eq!(decoded.to_bytes().unwrap(), throttle.to_bytes().unwrap()); + + let mut limits = AgentThrottle::default(); + limits.set_total(-1.0); + assert_eq!(limits.resend(), 10_000.0); + assert_eq!(limits.task(), 4_000.0); + assert_eq!(limits.land(), 0.0); +} + +#[derive(Default)] +struct RecordingThrottleSender(Mutex>>); + +impl AgentThrottleSender for RecordingThrottleSender { + fn send_throttle( + &self, + throttle_bytes: &[u8], + _simulator: Option<&Simulator>, + ) -> Result<(), Error> { + self.0.lock().unwrap().push(throttle_bytes.to_vec()); + Ok(()) + } +} + +#[test] +fn mapped_agent_throttle_set_uses_the_client_network_binding() { + let sender = Arc::new(RecordingThrottleSender::default()); + let mut client = GridClient::new().unwrap(); + client.set_agent_throttle_sender(sender.clone()).unwrap(); + let throttle = AgentThrottle::new_with_grid_client(client).unwrap(); + throttle.set_with_method().unwrap(); + assert_eq!( + sender.0.lock().unwrap().as_slice(), + &[throttle.to_bytes().unwrap()] + ); +} + +#[test] +fn construction_is_runtime_neutral_and_start_requires_an_injected_runtime() { + let transport = UDPBase::client_with_defaults("127.0.0.1:13000".parse().unwrap()).unwrap(); + assert!(!transport.is_running()); + assert_eq!(transport.start(), Err(Error::InvalidOperation)); + assert!(!transport.is_running()); +} + +#[tokio::test(flavor = "current_thread")] +async fn bounded_command_queue_reports_backpressure_without_unbounded_buffering() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let config = UdpTransportConfig { + command_queue_capacity: 1, + ..UdpTransportConfig::default() + }; + let transport = UDPBase::client( + server.local_addr().unwrap(), + config, + Arc::new(RecordingHandler::default()), + CancellationToken::default(), + ) + .unwrap(); + transport.start().unwrap(); + + // A current-thread runtime does not poll spawned tasks until this test + // yields, making the capacity-one boundary deterministic. + transport + .try_send_packet( + complete_ping_packet(), + server.local_addr().unwrap(), + PacketType::CompletePingCheck, + false, + ) + .unwrap(); + assert!(matches!( + transport.try_send_packet( + complete_ping_packet(), + server.local_addr().unwrap(), + PacketType::CompletePingCheck, + false, + ), + Err(UdpTransportError::Backpressure) + )); + assert_eq!(transport.stats().dropped_send_queue, 1); + transport.stop_async().await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn local_socket_send_assigns_sequence_and_ack_stops_resends() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let server_address = server.local_addr().unwrap(); + let handler = Arc::new(RecordingHandler::default()); + let config = UdpTransportConfig { + resend_timeout: Duration::from_millis(60), + network_tick_interval: Duration::from_millis(10), + ..UdpTransportConfig::default() + }; + let transport = UDPBase::client( + server_address, + config, + handler.clone(), + CancellationToken::default(), + ) + .unwrap(); + transport.start().unwrap(); + + let sequence = transport + .send_packet( + reliable_packet(), + PacketType::UseCircuitCode, + false, + CancellationToken::default(), + ) + .await + .unwrap(); + assert_eq!(sequence, 1); + let mut receive = [0_u8; 4096]; + let (length, client_address) = + tokio::time::timeout(Duration::from_secs(2), server.recv_from(&mut receive)) + .await + .unwrap() + .unwrap(); + assert_eq!(u32::from_be_bytes(receive[1..5].try_into().unwrap()), 1); + assert_ne!(receive[0] & Helpers::MSG_RELIABLE, 0); + server + .send_to(&packet_ack(sequence), client_address) + .await + .unwrap(); + wait_until(|| transport.stats().acknowledgements_received == 1).await; + tokio::time::sleep(Duration::from_millis(100)).await; + assert_eq!(transport.stats().resent_datagrams, 0); + assert_eq!(handler.sent.lock().unwrap().len(), 1); + assert_eq!(length, handler.sent.lock().unwrap()[0].len()); + + transport.stop_async().await; + assert!(!transport.is_running()); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn reliable_packets_resend_with_same_sequence_then_expire() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let config = UdpTransportConfig { + resend_timeout: Duration::from_millis(20), + network_tick_interval: Duration::from_millis(5), + max_resend_count: 2, + ..UdpTransportConfig::default() + }; + let transport = UDPBase::client( + server.local_addr().unwrap(), + config, + Arc::new(RecordingHandler::default()), + CancellationToken::default(), + ) + .unwrap(); + transport.start().unwrap(); + transport + .send_packet( + reliable_packet(), + PacketType::UseCircuitCode, + false, + CancellationToken::default(), + ) + .await + .unwrap(); + + let mut datagrams = Vec::new(); + let mut receive = [0_u8; 4096]; + for _ in 0..3 { + let (length, _) = + tokio::time::timeout(Duration::from_secs(1), server.recv_from(&mut receive)) + .await + .unwrap() + .unwrap(); + datagrams.push(receive[..length].to_vec()); + } + assert_eq!( + datagrams + .iter() + .map(|bytes| u32::from_be_bytes(bytes[1..5].try_into().unwrap())) + .collect::>(), + vec![1, 1, 1] + ); + assert_eq!(datagrams[0][0] & Helpers::MSG_RESENT, 0); + assert_ne!(datagrams[1][0] & Helpers::MSG_RESENT, 0); + wait_until(|| transport.stats().failed_resends == 1).await; + assert_eq!(transport.stats().resent_datagrams, 2); + transport.stop_async().await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn incoming_reliable_packets_are_acked_and_duplicates_are_not_dispatched() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let handler = Arc::new(RecordingHandler::default()); + let config = UdpTransportConfig { + max_pending_acks: 2, + network_tick_interval: Duration::from_secs(1), + ..UdpTransportConfig::default() + }; + let transport = UDPBase::client( + server.local_addr().unwrap(), + config, + handler.clone(), + CancellationToken::default(), + ) + .unwrap(); + transport.start().unwrap(); + let client_address = transport.local_address().unwrap(); + + let mut first = reliable_packet(); + first[1..5].copy_from_slice(&10_u32.to_be_bytes()); + let mut second = reliable_packet(); + second[1..5].copy_from_slice(&12_u32.to_be_bytes()); + let mut late = reliable_packet(); + late[1..5].copy_from_slice(&11_u32.to_be_bytes()); + server.send_to(&first, client_address).await.unwrap(); + server.send_to(&first, client_address).await.unwrap(); + server.send_to(&second, client_address).await.unwrap(); + server.send_to(&late, client_address).await.unwrap(); + + let mut receive = [0_u8; 4096]; + let mut ack_groups = Vec::new(); + for _ in 0..2 { + let (length, _) = + tokio::time::timeout(Duration::from_secs(2), server.recv_from(&mut receive)) + .await + .unwrap() + .unwrap(); + let mut position = 0; + let ack = PacketAckPacket::new_with_bytes_int32(receive[..length].to_vec(), &mut position) + .unwrap(); + ack_groups.push(ack.packets.iter().map(|block| block.id).collect::>()); + } + assert_eq!(ack_groups, vec![vec![10, 10], vec![12, 11]]); + wait_until(|| handler.received.lock().unwrap().len() == 3).await; + assert_eq!(transport.stats().duplicate_datagrams, 1); + assert_eq!(transport.stats().sequence_gaps, 1); + assert_eq!(transport.stats().out_of_order_datagrams, 1); + assert_eq!(transport.stats().acknowledgements_sent, 4); + transport.stop_async().await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn pending_ack_is_piggybacked_before_the_mtu_boundary() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let handler = Arc::new(RecordingHandler::default()); + let config = UdpTransportConfig { + network_tick_interval: Duration::from_secs(10), + ..UdpTransportConfig::default() + }; + let transport = UDPBase::client( + server.local_addr().unwrap(), + config, + handler.clone(), + CancellationToken::default(), + ) + .unwrap(); + transport.start().unwrap(); + + let mut reliable = reliable_packet(); + reliable[1..5].copy_from_slice(&20_u32.to_be_bytes()); + server + .send_to(&reliable, transport.local_address().unwrap()) + .await + .unwrap(); + wait_until(|| handler.received.lock().unwrap().len() == 1).await; + + transport + .send_packet( + complete_ping_packet(), + PacketType::CompletePingCheck, + false, + CancellationToken::default(), + ) + .await + .unwrap(); + let mut receive = [0_u8; 4096]; + let (length, _) = server.recv_from(&mut receive).await.unwrap(); + assert_ne!(receive[0] & Helpers::MSG_APPENDED_ACKS, 0); + let mut end = i32::try_from(length).unwrap() - 1; + let decoded = Packet::build_packet_with_bytes_int32_bytes( + receive[..length].to_vec(), + &mut end, + vec![0; 8192], + ) + .unwrap(); + assert_eq!(decoded.type_, PacketType::CompletePingCheck); + assert_eq!(decoded.header.ack_list, Some(vec![20])); + transport.stop_async().await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn zerocoding_mtu_malformed_input_and_cancellation_are_bounded() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let handler = Arc::new(RecordingHandler::default()); + let transport = UDPBase::client( + server.local_addr().unwrap(), + UdpTransportConfig::default(), + handler.clone(), + CancellationToken::default(), + ) + .unwrap(); + transport.start().unwrap(); + let client_address = transport.local_address().unwrap(); + + let mut packet = AgentThrottlePacket::new_with_constructor().unwrap(); + packet.throttle.throttles = vec![0; 28]; + let raw = packet.to_bytes_with_method().unwrap(); + transport + .send_packet( + raw.clone(), + PacketType::AgentThrottle, + true, + CancellationToken::default(), + ) + .await + .unwrap(); + let mut receive = [0_u8; 4096]; + let (length, _) = server.recv_from(&mut receive).await.unwrap(); + assert_ne!(receive[0] & Helpers::MSG_ZEROCODED, 0); + assert!(length < raw.len()); + let mut end = i32::try_from(length).unwrap() - 1; + Packet::build_packet_with_bytes_int32_bytes( + receive[..length].to_vec(), + &mut end, + vec![0; 8192], + ) + .unwrap(); + + assert_eq!( + transport + .send_packet( + vec![0; 1201], + PacketType::AgentThrottle, + false, + CancellationToken::default(), + ) + .await, + Err(UdpTransportError::MtuExceeded) + ); + + for malformed in [vec![], vec![0x40], vec![0x40, 0, 0, 0, 1, 0, 0xfe]] { + server.send_to(&malformed, client_address).await.unwrap(); + } + wait_until(|| transport.stats().malformed_datagrams >= 3).await; + assert!(handler.received.lock().unwrap().is_empty()); + + transport.stop_async().await; + assert!(!transport.is_running()); + assert_eq!( + transport + .send_packet( + reliable_packet(), + PacketType::UseCircuitCode, + false, + CancellationToken::default(), + ) + .await, + Err(UdpTransportError::NotRunning) + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn dropping_the_last_transport_handle_releases_its_socket_tasks() { + let server = UdpSocket::bind(loopback()).await.unwrap(); + let transport = UDPBase::client_with_defaults(server.local_addr().unwrap()).unwrap(); + transport.start().unwrap(); + let local_address = transport.local_address().unwrap(); + drop(transport); + + tokio::time::timeout(Duration::from_secs(2), async { + loop { + match UdpSocket::bind(local_address).await { + Ok(rebound) => break rebound, + Err(_) => tokio::time::sleep(Duration::from_millis(1)).await, + } + } + }) + .await + .expect("transport tasks retained the UDP socket after drop"); +} diff --git a/docs/udp-transport.md b/docs/udp-transport.md new file mode 100644 index 0000000..edb03f6 --- /dev/null +++ b/docs/udp-transport.md @@ -0,0 +1,77 @@ +# Native UDP transport + +`libremetaverse::UDPBase` is the native Tokio implementation of the +LibreMetaverse UDP transport slice. Construction is runtime-neutral: it does +not bind a socket, start a thread, or create an executor. Call `start` from an +entered Tokio runtime after constructing either a client transport with +`UDPBase::client` or a server transport with `UDPBase::server`. + +The implementation uses three bounded queues: + +- one socket receive task writes complete datagrams to the receive queue with + drop-on-full backpressure; +- one coordinator validates packet framing, owns per-peer sequence, duplicate, + pending-ACK, and reliable resend state, and writes prepared datagrams to a + bounded writer queue; +- one writer task is the only socket send owner and applies the per-category + token buckets before calling `send_to`. + +The default limits mirror the golden C# settings: 512 receive and command +entries, a 1,000-entry packet archive, ten ACKs before an immediate standalone +ACK, a 500 ms network tick, a 4 second resend timeout, and three retries. +Additional explicit limits bound the send queue, reliable window, ACK queue, +peer table, decoded zerocode buffer, and all allocations derived from incoming +datagrams. Client-mode source validation matches the reference by accepting +only the configured simulator IP address; replies continue to use the +configured endpoint. + +## Wire behavior + +Every newly prepared packet receives the next four-byte big-endian sequence, +using the full `u32` header width and wrapping through zero like the C# +`Interlocked.Increment(Int32)` plus `uint` cast. Reliable packets remain in the +bounded per-peer window until either an appended ACK or `PacketAck` block +removes them. Timed-out entries are retransmitted with `MSG_RESENT` and the +same sequence, then removed after the configured retry count. Reliable incoming +packets are ACKed even when duplicate; duplicate callbacks are suppressed. +Out-of-order packets are delivered immediately, as in the C# client, while +gaps and late arrivals are counted separately. + +Pending ACKs are appended in big-endian form while the strict reference +`data_length + 5 < MTU` condition holds. Remaining ACKs are emitted as an +unreliable `PacketAck`. Zerocoding occurs before ACK appending. If zerocoding +would expand a packet beyond the 1,200-byte protocol MTU, the zerocode flag is +cleared and the original packet is sent. The transport does not invent an +application fragmentation format: packet codecs must use their existing +`ToBytesMultiple` behavior, and an oversized unsplit outbound payload returns +`UdpTransportError::MtuExceeded`. + +Incoming storage is 4 KiB, matching `UDPPacketBuffer.DEFAULT_BUFFER_SIZE`, and +the receiver allocates one extra detection byte so an oversized datagram is +rejected rather than silently treated as complete. Packet parsing and +zerodecoding use fixed configured ceilings. Malformed or unknown packets are +counted and discarded without invoking application callbacks. + +## Throttling, cancellation, and diagnostics + +`AgentThrottle` reproduces all seven C# clamps and its 28-byte little-endian +wire layout. Outgoing control/handshake packets bypass throttling; task, +texture, and asset packets use independent 100 ms token buckets with the same +four-period burst and 200-byte minimum replenishment policy as the reference. +`update_throttle` swaps all three buckets in writer order. + +All tasks share a linked cancellation source. `stop_async` cancels and joins +the receive, coordinator, and writer tasks; `stop` and final drop cancel and +abort them for synchronous C# compatibility. A stopped transport can be +started again with fresh task state unless its parent token is cancelled. + +`UdpTransportStats` contains counts and byte totals only. The transport does +not log datagram bodies, credentials, capability URLs, or socket error text; +its `Debug` implementation also omits endpoint values. The implementation uses +only portable Rust and Tokio socket APIs on Linux, Windows, and macOS. + +The isolated loopback suite is: + +```sh +cargo test -p libremetaverse --test udp_transport --no-default-features +``` diff --git a/tests/api-compile/Cargo.lock b/tests/api-compile/Cargo.lock index b47076e..b5c8f2d 100644 --- a/tests/api-compile/Cargo.lock +++ b/tests/api-compile/Cargo.lock @@ -194,6 +194,7 @@ dependencies = [ "libremetaverse-structured-data", "libremetaverse-types", "roxmltree", + "tokio", ] [[package]] @@ -337,6 +338,17 @@ dependencies = [ "libremetaverse-voice-webrtc", ] +[[package]] +name = "mio" +version = "1.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "30d65c71f1ce40ab09135ce117d742b9f8a19ff91a41a8b57ed50bc2de59c427" +dependencies = [ + "libc", + "wasi", + "windows-sys", +] + [[package]] name = "once_cell" version = "1.21.4" @@ -458,6 +470,16 @@ version = "0.4.12" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5" +[[package]] +name = "socket2" +version = "0.6.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3d1e2c7f27f8d4cb10542a02c49005dbd6e93095799d6f3be745fae9f8fedd4" +dependencies = [ + "libc", + "windows-sys", +] + [[package]] name = "syn" version = "2.0.119" @@ -480,6 +502,31 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "tokio" +version = "1.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "202caea871b69668250d242070849eb495be178ed697a3e98aebce5bc81a0bed" +dependencies = [ + "libc", + "mio", + "pin-project-lite", + "socket2", + "tokio-macros", + "windows-sys", +] + +[[package]] +name = "tokio-macros" +version = "2.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + [[package]] name = "typenum" version = "1.20.1" @@ -509,6 +556,12 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "wasi" +version = "0.11.1+wasi-snapshot-preview1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" + [[package]] name = "wasm-bindgen" version = "0.2.127" @@ -554,6 +607,21 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "windows-link" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" + +[[package]] +name = "windows-sys" +version = "0.61.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc" +dependencies = [ + "windows-link", +] + [[package]] name = "zmij" version = "1.0.23" diff --git a/tests/compat/tests/core_runtime_shims.rs b/tests/compat/tests/core_runtime_shims.rs index dd9e8a0..bdd1df6 100644 --- a/tests/compat/tests/core_runtime_shims.rs +++ b/tests/compat/tests/core_runtime_shims.rs @@ -106,8 +106,8 @@ fn client_and_settings_construct_without_starting_runtime_work() { member_id(CapsRateLimiter::new_with_constructor()), "M:LibreMetaverse.CapsRateLimiter.#ctor" ); - assert_eq!( - member_id(UDPPacketBuffer::new_with_constructor()), - "M:LibreMetaverse.UDPPacketBuffer.#ctor" - ); + let udp_buffer = UDPPacketBuffer::new_with_constructor().expect("UDP packet buffer"); + assert_eq!(udp_buffer.data.len(), 4096); + assert_eq!(udp_buffer.data_length, 0); + assert_eq!(udp_buffer.remote_end_point, "0.0.0.0:0".parse().unwrap()); } diff --git a/tools/generate_api_shims.py b/tools/generate_api_shims.py index 3dcd5f1..0f3e26c 100644 --- a/tools/generate_api_shims.py +++ b/tools/generate_api_shims.py @@ -112,6 +112,10 @@ NATIVE_TYPES = { "T:LibreMetaverse.SkeletalBoneInfo": "crate::visual_catalog::SkeletalBoneInfo", "T:LibreMetaverse.TreeDefinition": "crate::foliage_catalog::TreeDefinition", "T:LibreMetaverse.TreeDefinitions": "crate::foliage_catalog::TreeDefinitions", + "T:LibreMetaverse.AgentThrottle": "crate::udp_transport::AgentThrottle", + "T:LibreMetaverse.IncomingPacketIDCollection": "crate::udp_transport::IncomingPacketIDCollection", + "T:LibreMetaverse.UDPBase": "crate::udp_transport::UDPBase", + "T:LibreMetaverse.UDPPacketBuffer": "crate::udp_transport::UDPPacketBuffer", "T:LibreMetaverse.VisualAlphaParam": "crate::visual_catalog::VisualAlphaParam", "T:LibreMetaverse.VisualColorParam": "crate::visual_catalog::VisualColorParam", "T:LibreMetaverse.VisualParam": "crate::visual_catalog::VisualParam",