//! Native network-manager, packet-dispatch, and simulator lifecycle runtime. //! //! The C# implementation uses reference-identity objects protected by several //! locks and channel processors. Rust models those objects as cheap `Arc` //! handles. Callback lists are snapshotted while locked and always invoked //! after the lock is released. #![allow(clippy::missing_errors_doc)] // Result shapes are fixed by the compatibility map. #![allow(clippy::needless_pass_by_value)] // Owned mapped parameters preserve the public API. #![allow(clippy::too_many_arguments)] // Login overloads preserve the mapped public API. #![allow(clippy::too_many_lines)] // Packet dispatch keeps the complete protocol switch together. #![allow(clippy::type_complexity)] // Delegate signatures are fixed by the compatibility map. use crate::interfaces::IMessage; use crate::packets::{ AgentPausePacket, AgentResumePacket, CloseCircuitPacket, CompletePingCheckPacket, GenericStreamingMessagePacket, KickUserPacket, Packet, PacketType, RegionHandshakePacket, RegionHandshakeReplyPacket, StartPingCheckPacket, UseCircuitCodePacket, }; use crate::udp_transport::{UDPBase, UDPPacketBuffer, UdpPacketHandler, UdpTransportConfig}; use crate::{ AccountLevelBenefits, Avatar, Caps, Error, GenericStreamingMethod, GridClient, Helpers, LoginCredential, LoginParams, LoginProgressEventArgs, LoginResponseData, LoginStatus, NetworkManagerLoginResponseCallback, Primitive, RegionFlags, RegionProtocols, SimAccess, SimulatorDataPool, SimulatorFeatures, SimulatorSimStats, TerrainPatch, }; use libremetaverse_structured_data::{OSD, OSDFormat, OSDMap, OSDParser}; use libremetaverse_types::compat::{ CancellationToken, CancellationTokenSource, EventHandler, Object, Subscription, Uri, }; use libremetaverse_types::{UUID, Vector2}; use std::collections::HashMap; use std::fmt; use std::panic::{AssertUnwindSafe, catch_unwind}; use std::sync::atomic::{AtomicBool, AtomicI32, AtomicU32, AtomicU64, Ordering}; use std::sync::mpsc::{Receiver, RecvTimeoutError, SyncSender, sync_channel}; use std::sync::{Arc, Condvar, Mutex, RwLock, Weak}; use std::thread::{self, JoinHandle}; use std::time::Duration; impl Clone for crate::packets::Header { fn clone(&self) -> Self { crate::packet_wire::clone_header(self) } } impl Clone for Packet { fn clone(&self) -> Self { Self { has_variable_blocks: self.has_variable_blocks, header: self.header.clone(), type_: self.type_, } } } fn mutex(value: &Mutex) -> std::sync::MutexGuard<'_, T> { value .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) } fn read(value: &RwLock) -> std::sync::RwLockReadGuard<'_, T> { value .read() .unwrap_or_else(std::sync::PoisonError::into_inner) } fn write(value: &RwLock) -> std::sync::RwLockWriteGuard<'_, T> { value .write() .unwrap_or_else(std::sync::PoisonError::into_inner) } fn invoke_safely(handler: &EventHandler, value: T) { let _ = catch_unwind(AssertUnwindSafe(|| handler(value))); } fn wait_for_cancellation( token: &libremetaverse_types::compat::CancellationToken, duration: Duration, ) -> bool { let started = std::time::Instant::now(); while !token.is_cancellation_requested() { let remaining = duration.saturating_sub(started.elapsed()); if remaining.is_zero() { return false; } thread::sleep(remaining.min(Duration::from_millis(50))); } true } fn positive_millisecond_duration(value: i32) -> Duration { Duration::from_millis(u64::try_from(value.max(1)).unwrap_or(1)) } fn validate_login_uri(value: &str) -> Result { let parsed = reqwest::Url::parse(value).map_err(|_| Error::Argument)?; if !matches!(parsed.scheme(), "http" | "https") || parsed.host_str().is_none() { return Err(Error::Argument); } Ok(Uri(parsed.to_string())) } struct LoginWireSecrets { password: String, token: String, mfa_hash: String, } impl Drop for LoginWireSecrets { fn drop(&mut self) { // Overwrite the live buffers before releasing them. This cannot erase // allocator copies, but it prevents the owned request state surviving // the login/redirect sequence. self.password .replace_range(.., &"\0".repeat(self.password.len())); self.token.replace_range(.., &"\0".repeat(self.token.len())); self.mfa_hash .replace_range(.., &"\0".repeat(self.mfa_hash.len())); self.password.clear(); self.token.clear(); self.mfa_hash.clear(); } } fn login_request(params: &LoginParams, secrets: &LoginWireSecrets) -> OSD { let mut map = HashMap::from([ ("first".to_owned(), OSD::String(params.first_name.clone())), ("last".to_owned(), OSD::String(params.last_name.clone())), ("passwd".to_owned(), OSD::String(secrets.password.clone())), ("start".to_owned(), OSD::String(params.start.clone())), ("channel".to_owned(), OSD::String(params.channel.clone())), ("version".to_owned(), OSD::String(params.version.clone())), ("platform".to_owned(), OSD::String(params.platform.clone())), ( "platform_version".to_owned(), OSD::String(params.platform_version.clone()), ), ("mac".to_owned(), OSD::String(params.mac.clone())), ("agree_to_tos".to_owned(), OSD::Boolean(params.agree_to_tos)), ( "read_critical".to_owned(), OSD::Boolean(params.read_critical), ), ( "viewer_digest".to_owned(), OSD::String(params.viewer_digest.clone()), ), ("id0".to_owned(), OSD::String(params.id0.clone())), ( "last_exec_event".to_owned(), OSD::Integer(params.last_exec_event as i32), ), ( "options".to_owned(), OSD::Array(params.options.iter().cloned().map(OSD::String).collect()), ), ]); if params.mfa_enabled { map.insert("token".to_owned(), OSD::String(secrets.token.clone())); map.insert("mfa_hash".to_owned(), OSD::String(secrets.mfa_hash.clone())); } OSD::Map(map) } async fn run_blocking_compat( name: &str, operation: impl FnOnce() -> T + Send + 'static, ) -> Result { let (sender, receiver) = tokio::sync::oneshot::channel(); thread::Builder::new() .name(name.to_owned()) .spawn(move || { let _ = sender.send(operation()); }) .map_err(|_| Error::InvalidOperation)?; receiver.await.map_err(|_| Error::Cancelled) } struct EventRegistryState { next_id: AtomicU64, handlers: Mutex)>>, } struct EventRegistry { state: Arc>, } impl Default for EventRegistry { fn default() -> Self { Self { state: Arc::new(EventRegistryState { next_id: AtomicU64::new(1), handlers: Mutex::new(Vec::new()), }), } } } impl EventRegistry { fn subscribe(&self, handler: EventHandler) -> Subscription { let id = self.state.next_id.fetch_add(1, Ordering::Relaxed); mutex(&self.state.handlers).push((id, handler)); let state = Arc::downgrade(&self.state); Subscription::new(move || { let Some(state) = state.upgrade() else { return; }; mutex(&state.handlers).retain(|(candidate, _)| *candidate != id); }) } fn emit_with(&self, mut value: impl FnMut() -> T) { let handlers: Vec<_> = mutex(&self.state.handlers) .iter() .map(|(_, handler)| Arc::clone(handler)) .collect(); for handler in handlers { invoke_safely(&handler, value()); } } } /// C# `PacketReceivedEventArgs` with shared simulator identity. #[derive(Clone)] pub struct PacketReceivedEventArgs { packet: Packet, simulator: Simulator, } impl PacketReceivedEventArgs { pub fn new(packet: Packet, simulator: Simulator) -> Result { Ok(Self { packet, simulator }) } #[must_use] pub fn packet(&self) -> Packet { self.packet.clone() } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } #[derive(Clone)] pub struct PacketSentEventArgs { data: Vec, sent_bytes: i32, simulator: Simulator, } impl PacketSentEventArgs { pub fn new(data: Vec, sent_bytes: i32, simulator: Simulator) -> Result { let sent = usize::try_from(sent_bytes).map_err(|_| Error::Argument)?; if sent > data.len() { return Err(Error::Argument); } Ok(Self { data, sent_bytes, simulator, }) } #[must_use] pub fn data(&self) -> Vec { self.data.clone() } #[must_use] pub const fn sent_bytes(&self) -> i32 { self.sent_bytes } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } #[derive(Clone)] pub struct SimChangedEventArgs { previous_simulator: Option, } impl SimChangedEventArgs { pub fn new(previous_simulator: Option) -> Result { Ok(Self { previous_simulator }) } #[must_use] pub fn previous_simulator(&self) -> Option { self.previous_simulator.clone() } } #[derive(Clone)] pub struct SimConnectedEventArgs { simulator: Simulator, } /// Event arguments raised when a simulator's `EventQueueGet` poller is live. #[derive(Clone)] pub struct EventQueueRunningEventArgs { simulator: Simulator, } impl EventQueueRunningEventArgs { pub const fn new(simulator: Simulator) -> Result { Ok(Self { simulator }) } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } impl SimConnectedEventArgs { pub fn new(simulator: Simulator) -> Result { Ok(Self { simulator }) } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } #[derive(Clone)] pub struct SimConnectingEventArgs { simulator: Simulator, cancel: Arc, } impl SimConnectingEventArgs { pub fn new(simulator: Simulator) -> Result { Ok(Self { simulator, cancel: Arc::new(AtomicBool::new(false)), }) } #[must_use] pub fn cancel(&self) -> bool { self.cancel.load(Ordering::Acquire) } pub fn set_cancel(&mut self, value: bool) { self.cancel.store(value, Ordering::Release); } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } #[derive(Clone)] pub struct SimDisconnectedEventArgs { simulator: Simulator, reason: crate::NetworkManagerDisconnectType, was_connected: bool, } impl SimDisconnectedEventArgs { pub fn new( simulator: Simulator, reason: crate::NetworkManagerDisconnectType, was_connected: bool, ) -> Result { Ok(Self { simulator, reason, was_connected, }) } #[must_use] pub const fn reason(&self) -> crate::NetworkManagerDisconnectType { self.reason } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } #[must_use] pub const fn was_connected(&self) -> bool { self.was_connected } } #[derive(Clone)] pub struct DisconnectedEventArgs { reason: crate::NetworkManagerDisconnectType, message: String, } impl DisconnectedEventArgs { pub fn new( reason: crate::NetworkManagerDisconnectType, message: String, ) -> Result { Ok(Self { reason, message }) } #[must_use] pub fn message(&self) -> String { self.message.clone() } #[must_use] pub const fn reason(&self) -> crate::NetworkManagerDisconnectType { self.reason } } #[derive(Clone)] pub struct GenericStreamingMessageEventArgs { simulator: Simulator, method: GenericStreamingMethod, data: Vec, } impl GenericStreamingMessageEventArgs { pub fn new( simulator: Simulator, method: GenericStreamingMethod, data: Vec, ) -> Result { Ok(Self { simulator, method, data, }) } #[must_use] pub fn data(&self) -> Vec { self.data.clone() } #[must_use] pub const fn method(&self) -> GenericStreamingMethod { self.method } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } #[derive(Clone, Default)] struct RegisteredPacketCallbacks { handlers: Vec>, is_async: bool, } struct AsyncPacketBatch { default_handlers: Vec>, specific_handlers: Vec>, packet: Packet, simulator: Simulator, } enum PacketWorkerCommand { Dispatch(AsyncPacketBatch), Stop, } struct PacketEventState { callbacks: Mutex>, async_sender: SyncSender, worker: Mutex>>, } impl Drop for PacketEventState { fn drop(&mut self) { let _ = self.async_sender.try_send(PacketWorkerCommand::Stop); if let Some(worker) = self .worker .get_mut() .unwrap_or_else(std::sync::PoisonError::into_inner) .take() && worker.thread().id() != thread::current().id() { let _ = worker.join(); } } } /// Ordered packet callback registry matching `PacketEventDictionary`. pub struct PacketEventDictionary { pub client: GridClient, state: Arc, } impl PacketEventDictionary { pub fn new(client: GridClient) -> Result { let (sender, receiver) = sync_channel(256); let state = Arc::new(PacketEventState { callbacks: Mutex::new(HashMap::new()), async_sender: sender, worker: Mutex::new(None), }); let weak = Arc::downgrade(&state); let worker = thread::Builder::new() .name("libremetaverse-packet-events".to_owned()) .spawn(move || packet_event_worker(receiver, weak)) .map_err(|_| Error::InvalidOperation)?; *mutex(&state.worker) = Some(worker); Ok(Self { client, state }) } pub fn register_event( &self, packet_type: PacketType, event_handler: EventHandler, is_async: bool, ) -> Result<(), Error> { let mut table = mutex(&self.state.callbacks); let callbacks = table.entry(packet_type).or_default(); callbacks.is_async |= is_async; callbacks.handlers.push(event_handler); Ok(()) } pub fn subscribe( &self, packet_type: PacketType, event_handler: EventHandler, is_async: bool, ) -> Subscription { let identity = Arc::as_ptr(&event_handler).cast::<()>() as usize; let _ = self.register_event(packet_type, event_handler, is_async); let state = Arc::downgrade(&self.state); Subscription::new(move || { let Some(state) = state.upgrade() else { return; }; let mut table = mutex(&state.callbacks); let Some(callbacks) = table.get_mut(&packet_type) else { return; }; if let Some(index) = callbacks .handlers .iter() .rposition(|handler| Arc::as_ptr(handler).cast::<()>() as usize == identity) { callbacks.handlers.remove(index); } if callbacks.handlers.is_empty() { table.remove(&packet_type); } }) } pub fn unregister_event( &self, packet_type: PacketType, event_handler: EventHandler, ) -> Result<(), Error> { let mut table = mutex(&self.state.callbacks); let Some(callbacks) = table.get_mut(&packet_type) else { return Ok(()); }; if let Some(index) = callbacks .handlers .iter() .rposition(|handler| Arc::ptr_eq(handler, &event_handler)) { callbacks.handlers.remove(index); } if callbacks.handlers.is_empty() { table.remove(&packet_type); } Ok(()) } pub fn invoke_raise_event( &self, packet_type: PacketType, packet: Packet, simulator: Simulator, ) -> Result<(), Error> { self.raise_event(packet_type, packet, simulator) } fn raise_event( &self, packet_type: PacketType, packet: Packet, simulator: Simulator, ) -> Result<(), Error> { let (defaults, specific) = { let table = mutex(&self.state.callbacks); ( table.get(&PacketType::Default).cloned().unwrap_or_default(), table.get(&packet_type).cloned().unwrap_or_default(), ) }; let default_async = defaults.is_async; let specific_async = specific.is_async; let args = PacketReceivedEventArgs::new(packet.clone(), simulator.clone())?; if !default_async { for handler in &defaults.handlers { invoke_safely(handler, args.clone()); } } if !specific_async { for handler in &specific.handlers { invoke_safely(handler, args.clone()); } // The reference implementation returns as soon as a concrete // packet type has synchronous handlers. This also deliberately // suppresses a default asynchronous chain for that dispatch. if !specific.handlers.is_empty() { return Ok(()); } } let default_handlers: Vec> = if default_async { defaults.handlers } else { Vec::new() }; let specific_handlers: Vec> = if specific_async { specific.handlers } else { Vec::new() }; if default_handlers.is_empty() && specific_handlers.is_empty() { return Ok(()); } self.state .async_sender .try_send(PacketWorkerCommand::Dispatch(AsyncPacketBatch { default_handlers, specific_handlers, packet, simulator, })) .map_err(|_| Error::InvalidOperation) } } fn packet_event_worker(receiver: Receiver, state: Weak) { while state.strong_count() != 0 { match receiver.recv_timeout(Duration::from_millis(100)) { Ok(PacketWorkerCommand::Dispatch(batch)) => { let Ok(args) = PacketReceivedEventArgs::new(batch.packet, batch.simulator) else { continue; }; for handler in batch .default_handlers .into_iter() .chain(batch.specific_handlers) { invoke_safely(&handler, args.clone()); } } Ok(PacketWorkerCommand::Stop) | Err(RecvTimeoutError::Disconnected) => break, Err(RecvTimeoutError::Timeout) => {} } } } type CapsHandler = dyn Fn(&str, &dyn IMessage, Simulator) + Send + Sync; #[derive(Clone)] pub struct CapsEventQueueCallback { id: u64, handler: Arc, } impl CapsEventQueueCallback { pub fn from_handler( handler: impl Fn(&str, &dyn IMessage, Simulator) + Send + Sync + 'static, ) -> Self { static NEXT_ID: AtomicU64 = AtomicU64::new(1); Self { id: NEXT_ID.fetch_add(1, Ordering::Relaxed), handler: Arc::new(handler), } } pub fn new(_object: Object, _method: isize) -> Result { Err(Error::InvalidOperation) } pub fn begin_invoke( &self, caps_key: String, message: Box, simulator: Simulator, callback: Box, object: Object, ) -> Result, Error> { (self.handler)(&caps_key, message.as_ref(), simulator); callback(&object); Ok(Box::new(object)) } pub fn end_invoke(&self, _result: Box) -> Result<(), Error> { Ok(()) } pub fn invoke( &self, caps_key: String, message: Box, simulator: Simulator, ) -> Result<(), Error> { (self.handler)(&caps_key, message.as_ref(), simulator); Ok(()) } fn invoke_ref(&self, caps_key: &str, message: &dyn IMessage, simulator: Simulator) { let _ = catch_unwind(AssertUnwindSafe(|| { (self.handler)(caps_key, message, simulator); })); } } pub struct CapsEventDictionary { pub client: GridClient, callbacks: Arc>>>, } impl CapsEventDictionary { pub fn new(client: GridClient) -> Result { Ok(Self { client, callbacks: Arc::new(Mutex::new(HashMap::new())), }) } pub fn register_event( &self, caps_event: String, event_handler: CapsEventQueueCallback, ) -> Result<(), Error> { mutex(&self.callbacks) .entry(caps_event) .or_default() .push(event_handler); Ok(()) } pub fn subscribe( &self, caps_event: String, event_handler: CapsEventQueueCallback, ) -> Subscription { let id = event_handler.id; let _ = self.register_event(caps_event.clone(), event_handler); let callbacks = Arc::downgrade(&self.callbacks); Subscription::new(move || { let Some(callbacks) = callbacks.upgrade() else { return; }; let mut table = mutex(&callbacks); let Some(entries) = table.get_mut(&caps_event) else { return; }; entries.retain(|entry| entry.id != id); if entries.is_empty() { table.remove(&caps_event); } }) } pub fn unregister_event( &self, caps_event: String, event_handler: CapsEventQueueCallback, ) -> Result<(), Error> { let mut table = mutex(&self.callbacks); if let Some(entries) = table.get_mut(&caps_event) { if let Some(index) = entries .iter() .rposition(|entry| entry.id == event_handler.id) { entries.remove(index); } if entries.is_empty() { table.remove(&caps_event); } } Ok(()) } pub fn invoke_raise_event( &self, caps_event: &str, message: &dyn IMessage, simulator: Simulator, ) { let (defaults, specific) = { let table = mutex(&self.callbacks); ( table.get("").cloned().unwrap_or_default(), table.get(caps_event).cloned().unwrap_or_default(), ) }; for callback in defaults.into_iter().chain(specific) { callback.invoke_ref(caps_event, message, simulator.clone()); } } } /// Mutable, thread-safe equivalent of the C# public simulator list. #[derive(Clone, Default)] pub struct SimulatorCollection(Arc>>); impl SimulatorCollection { #[must_use] pub fn snapshot(&self) -> Vec { read(&self.0).clone() } #[must_use] pub fn len(&self) -> usize { read(&self.0).len() } #[must_use] pub fn is_empty(&self) -> bool { self.len() == 0 } fn push_if_missing(&self, simulator: Simulator) -> Simulator { let mut simulators = write(&self.0); if let Some(existing) = simulators .iter() .find(|existing| existing.ip_end_point() == simulator.ip_end_point()) { return existing.clone(); } simulators.push(simulator.clone()); simulator } fn remove(&self, simulator: &Simulator) -> bool { let mut simulators = write(&self.0); let old_len = simulators.len(); simulators.retain(|candidate| candidate != simulator); old_len != simulators.len() } fn drain(&self) -> Vec { std::mem::take(&mut *write(&self.0)) } } impl fmt::Debug for SimulatorCollection { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter .debug_list() .entries(self.snapshot().iter()) .finish() } } struct TransportThread { shutdown: libremetaverse_types::compat::CancellationTokenSource, handle: JoinHandle<()>, } /// Shared data exposed through the C# `Simulator` field surface. pub struct SimulatorData { pub access: SimAccess, pub agent_movement_complete: bool, pub billable_factor: f32, pub cpu_class: i32, pub cpu_ratio: i32, pub client: GridClient, pub colo_location: String, pub data_pool: Option, pub features: SimulatorFeatures, pub flags: RegionFlags, pub global_to_local_id: RwLock>, pub handle: u64, pub id: UUID, pub is_estate_manager: bool, pub name: String, pub objects_avatars: RwLock>, pub objects_primitives: RwLock>, pub parcel_overlay: Vec, pub parcel_overlays_received: i32, pub product_name: String, pub product_sku: String, pub protocols: RegionProtocols, pub region_id: UUID, pub sequence: i32, pub sim_owner: UUID, pub sim_version: String, pub size_x: u32, pub size_y: u32, pub stats: SimulatorSimStats, pub terrain: Vec, pub terrain_base0: UUID, pub terrain_base1: UUID, pub terrain_base2: UUID, pub terrain_base3: UUID, pub terrain_detail0: UUID, pub terrain_detail1: UUID, pub terrain_detail2: UUID, pub terrain_detail3: UUID, pub terrain_height_range00: f32, pub terrain_height_range01: f32, pub terrain_height_range10: f32, pub terrain_height_range11: f32, pub terrain_start_height00: f32, pub terrain_start_height01: f32, pub terrain_start_height10: f32, pub terrain_start_height11: f32, pub water_height: f32, pub wind_speeds: Option>, endpoint: std::net::SocketAddr, connected: AtomicBool, handshake_complete: AtomicBool, handshake_wait: Condvar, handshake_wait_lock: Mutex<()>, disconnect_candidate: AtomicBool, circuit_code: AtomicU32, agent_id: RwLock, session_id: RwLock, seed_caps: RwLock>, caps_state: RwLock>>, event_queue_generation: Mutex, event_queue_wait: Condvar, manager: Mutex>, transport: UDPBase, transport_thread: Mutex>, ping_id: AtomicU32, pause_serial: AtomicU32, } impl Drop for SimulatorData { fn drop(&mut self) { if let Some(caps) = self .caps_state .get_mut() .unwrap_or_else(std::sync::PoisonError::into_inner) .take() { caps.disconnect(true); } let runtime = self .transport_thread .get_mut() .unwrap_or_else(std::sync::PoisonError::into_inner) .take(); if let Some(runtime) = runtime { runtime.shutdown.cancel(); if runtime.handle.thread().id() != thread::current().id() { let _ = runtime.handle.join(); } } } } /// Reference-identity simulator handle. pub struct Simulator { /// Capability client associated with this simulator. /// /// This remains on the value surface, rather than behind `Deref`, because /// existing translated code consumes the public C# field by value. pub caps: Option>, data: Arc, } impl Clone for Simulator { fn clone(&self) -> Self { let mut clone = self.native_clone_without_caps(); if let Some(inner) = read(&self.data.caps_state).clone() { clone.caps = Some(Box::new(Caps::native_from_inner( inner, clone.native_clone_without_caps(), ))); } clone } } impl std::ops::Deref for Simulator { type Target = SimulatorData; fn deref(&self) -> &Self::Target { &self.data } } impl PartialEq for Simulator { fn eq(&self, other: &Self) -> bool { self.ip_end_point() == other.ip_end_point() } } impl Eq for Simulator {} impl fmt::Debug for Simulator { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter .debug_struct("Simulator") .field("endpoint", &self.native_ip_end_point()) .field("handle", &self.handle) .field("connected", &self.native_is_connected()) .field("handshake_complete", &self.native_handshake_complete()) .finish_non_exhaustive() } } impl Simulator { pub fn native_new( client: GridClient, address: std::net::SocketAddr, handle: u64, size_x: Option, size_y: Option, ) -> Result { let cancellation = client.cancellation_token(); let simulator_reference = Arc::new(Mutex::new(Weak::new())); let handler: Arc = Arc::new(SimulatorUdpHandler { simulator: Arc::clone(&simulator_reference), }); let transport = UDPBase::client( address, UdpTransportConfig::default(), handler, cancellation, )?; let data = Arc::new(SimulatorData { access: SimAccess::UNKNOWN, agent_movement_complete: false, billable_factor: 0.0, cpu_class: 0, cpu_ratio: 0, client, colo_location: String::new(), data_pool: None, features: SimulatorFeatures, flags: RegionFlags(0), global_to_local_id: RwLock::new(HashMap::new()), handle, id: UUID::zero(), is_estate_manager: false, name: String::new(), objects_avatars: RwLock::new(HashMap::new()), objects_primitives: RwLock::new(HashMap::new()), parcel_overlay: vec![0; 4096], parcel_overlays_received: 0, product_name: String::new(), product_sku: String::new(), protocols: RegionProtocols(0), region_id: UUID::zero(), sequence: 0, sim_owner: UUID::zero(), sim_version: String::new(), size_x: size_x.unwrap_or(Self::DEFAULT_REGION_SIZE_X), size_y: size_y.unwrap_or(Self::DEFAULT_REGION_SIZE_Y), stats: SimulatorSimStats, terrain: Vec::new(), terrain_base0: UUID::zero(), terrain_base1: UUID::zero(), terrain_base2: UUID::zero(), terrain_base3: UUID::zero(), terrain_detail0: UUID::zero(), terrain_detail1: UUID::zero(), terrain_detail2: UUID::zero(), terrain_detail3: UUID::zero(), terrain_height_range00: 0.0, terrain_height_range01: 0.0, terrain_height_range10: 0.0, terrain_height_range11: 0.0, terrain_start_height00: 0.0, terrain_start_height01: 0.0, terrain_start_height10: 0.0, terrain_start_height11: 0.0, water_height: 0.0, wind_speeds: None, endpoint: address, connected: AtomicBool::new(false), handshake_complete: AtomicBool::new(false), handshake_wait: Condvar::new(), handshake_wait_lock: Mutex::new(()), disconnect_candidate: AtomicBool::new(false), circuit_code: AtomicU32::new(0), agent_id: RwLock::new(UUID::zero()), session_id: RwLock::new(UUID::zero()), seed_caps: RwLock::new(None), caps_state: RwLock::new(None), event_queue_generation: Mutex::new(0), event_queue_wait: Condvar::new(), manager: Mutex::new(Weak::new()), transport, transport_thread: Mutex::new(None), ping_id: AtomicU32::new(0), pause_serial: AtomicU32::new(0), }); *mutex(&simulator_reference) = Arc::downgrade(&data); Ok(Self { caps: None, data }) } pub(crate) fn native_clone_without_caps(&self) -> Self { Self { caps: None, data: Arc::clone(&self.data), } } pub(crate) fn native_data_weak(&self) -> Weak { Arc::downgrade(&self.data) } pub(crate) fn native_data_arc(&self) -> Arc { Arc::clone(&self.data) } pub(crate) fn native_from_weak(data: &Weak) -> Option { data.upgrade().map(|data| Self { caps: None, data }) } pub(crate) fn native_from_data(data: Arc) -> Self { Self { caps: None, data } } fn attach_manager(&self, manager: &Arc, circuit_code: u32) { *mutex(&self.manager) = Arc::downgrade(manager); self.circuit_code.store(circuit_code, Ordering::Release); *write(&self.agent_id) = *read(&manager.agent_id); *write(&self.session_id) = *read(&manager.session_id); } pub fn native_connect(&self, _move_to_sim: bool) -> Result { if self.native_is_connected() { self.native_use_circuit_code(true)?; return Ok(true); } let mut transport_thread = mutex(&self.transport_thread); if transport_thread.is_none() { let shutdown = libremetaverse_types::compat::CancellationTokenSource::default(); let shutdown_token = shutdown.token(); let transport = self.transport.clone(); let (ready_sender, ready_receiver) = sync_channel(1); let handle = thread::Builder::new() .name(format!("libremetaverse-udp-{}", self.endpoint)) .spawn(move || { let runtime = tokio::runtime::Builder::new_current_thread() .enable_all() .build(); let Ok(runtime) = runtime else { let _ = ready_sender.send(Err(crate::UdpTransportError::RuntimeUnavailable)); return; }; runtime.block_on(async move { let result = transport.start_transport(); let started = result.is_ok(); let _ = ready_sender.send(result); if started { shutdown_token.cancelled().await; transport.stop_async().await; } }); }) .map_err(|_| Error::InvalidOperation)?; let timeout = positive_millisecond_duration(self.client.settings_ref().timing.login_timeout); match ready_receiver.recv_timeout(timeout) { Ok(Ok(())) => { *transport_thread = Some(TransportThread { shutdown, handle }); } Ok(Err(error)) => { let _ = handle.join(); return Err(error.into()); } Err(_) => { shutdown.cancel(); let _ = handle.join(); return Ok(false); } } } self.connected.store(true, Ordering::Release); self.handshake_complete.store(false, Ordering::Release); drop(transport_thread); self.native_use_circuit_code(true)?; let timeout = positive_millisecond_duration(self.client.settings_ref().timing.login_timeout); let guard = mutex(&self.handshake_wait_lock); let (guard, _) = self .handshake_wait .wait_timeout_while(guard, timeout, |()| { !self.handshake_complete.load(Ordering::Acquire) && !self.client.cancellation_token().is_cancellation_requested() }) .unwrap_or_else(std::sync::PoisonError::into_inner); drop(guard); if !self.handshake_complete.load(Ordering::Acquire) && let Some(manager) = mutex(&self.manager).upgrade() { manager.simulators.remove(self); } Ok(true) } pub async fn native_connect_async(&self, move_to_sim: bool) -> Result { let simulator = self.clone(); run_blocking_compat("libremetaverse-simulator-connect", move || { simulator.native_connect(move_to_sim) }) .await? } pub fn native_disconnect(&self, send_close_circuit: bool) -> Result<(), Error> { if let Some(caps) = write(&self.caps_state).take() { caps.disconnect(true); } *write(&self.seed_caps) = None; self.disconnect_candidate.store(false, Ordering::Release); let was_connected = self.connected.swap(false, Ordering::AcqRel); self.handshake_complete.store(false, Ordering::Release); self.handshake_wait.notify_all(); if !was_connected { return Ok(()); } if send_close_circuit { let close = CloseCircuitPacket::new_with_constructor()?; let data = close.to_bytes_with_method()?; let _ = self.native_send_packet_data( data.clone(), i32::try_from(data.len()).map_err(|_| Error::Argument)?, PacketType::CloseCircuit, false, ); } if let Some(runtime) = mutex(&self.transport_thread).take() { runtime.shutdown.cancel(); if runtime.handle.thread().id() != thread::current().id() { let _ = runtime.handle.join(); } } Ok(()) } pub fn native_dispose(&self) -> Result<(), Error> { self.native_disconnect(false) } #[must_use] pub fn native_equals(&self, obj: Option) -> bool { obj.and_then(|value| value.downcast_arc::()) .is_some_and(|other| *self == *other) } #[must_use] pub fn native_hash_code(&self) -> i32 { let folded = self.handle ^ (self.handle >> 32); let low = u32::try_from(folded & u64::from(u32::MAX)).unwrap_or_default(); low.cast_signed() } pub fn native_pause(&self) -> Result<(), Error> { let mut packet = AgentPausePacket::new_with_constructor()?; packet.agent_data.agent_id = UUID::zero(); packet.agent_data.session_id = UUID::zero(); packet.agent_data.serial_num = self.pause_serial.fetch_add(1, Ordering::AcqRel); let data = packet.to_bytes_with_method()?; self.native_send_packet_data( data.clone(), i32::try_from(data.len()).map_err(|_| Error::Argument)?, PacketType::AgentPause, false, ) } pub fn native_resume(&self) -> Result<(), Error> { let mut packet = AgentResumePacket::new_with_constructor()?; packet.agent_data.agent_id = UUID::zero(); packet.agent_data.session_id = UUID::zero(); packet.agent_data.serial_num = self.pause_serial.fetch_add(1, Ordering::AcqRel); let data = packet.to_bytes_with_method()?; self.native_send_packet_data( data.clone(), i32::try_from(data.len()).map_err(|_| Error::Argument)?, PacketType::AgentResume, false, ) } pub fn native_send_packet(&self, packet: Packet) -> Result<(), Error> { let data = crate::packet_wire::encode_base_packet(&packet)?; self.native_send_packet_data( data.clone(), i32::try_from(data.len()).map_err(|_| Error::Argument)?, packet.type_, packet.header.zerocoded, ) } pub fn native_send_packet_data( &self, data: Vec, data_length: i32, packet_type: PacketType, do_zerocode: bool, ) -> Result<(), Error> { let length = usize::try_from(data_length).map_err(|_| Error::Argument)?; let payload = data.get(..length).ok_or(Error::Argument)?.to_vec(); self.transport .try_send_packet(payload, self.endpoint, packet_type, do_zerocode) .map(|_| ()) .map_err(Into::into) } pub fn native_send_ping(&self) -> Result<(), Error> { let mut ping = StartPingCheckPacket::new_with_constructor()?; ping.ping_id.ping_id = self.ping_id.fetch_add(1, Ordering::Relaxed).to_le_bytes()[0]; ping.ping_id.oldest_unacked = 0; let data = ping.to_bytes_with_method()?; self.native_send_packet_data( data.clone(), i32::try_from(data.len()).map_err(|_| Error::Argument)?, PacketType::StartPingCheck, false, ) } pub fn native_set_seed_caps( &self, seed_caps: Option, changed_sim: Option, ) -> Result<(), Error> { if !changed_sim.unwrap_or(false) && *read(&self.seed_caps) == seed_caps { return Ok(()); } if let Some(previous) = write(&self.caps_state).take() { previous.disconnect(true); } write(&self.seed_caps).clone_from(&seed_caps); if let Some(seed_caps) = seed_caps { let caps = Caps::native_create(self, seed_caps)?; *write(&self.caps_state) = Some(caps); } Ok(()) } pub fn native_is_event_queue_running( &self, block_until_running: Option, ) -> Result { if self.caps_running() { return Ok(true); } if !block_until_running.unwrap_or(false) { return Ok(false); } let mut generation = mutex(&self.event_queue_generation); for delay in [1_u64, 2, 4, 8] { if self.caps_running() || self.current_sim_caps_running() { return Ok(true); } self.start_current_event_queue(); let observed = *generation; let (next, _) = self .event_queue_wait .wait_timeout_while(generation, Duration::from_secs(delay), |value| { *value == observed && !self.client.cancellation_token().is_cancellation_requested() }) .unwrap_or_else(std::sync::PoisonError::into_inner); generation = next; if *generation != observed { break; } } Ok(self.caps_running() || self.current_sim_caps_running()) } fn caps_running(&self) -> bool { read(&self.caps_state).as_ref().is_some_and(|inner| { Caps::native_from_inner(Arc::clone(inner), self.native_clone_without_caps()) .is_event_queue_running() }) } fn current_sim_caps_running(&self) -> bool { mutex(&self.manager) .upgrade() .and_then(|manager| read(&manager.current_sim).clone()) .is_some_and(|simulator| simulator.caps_running()) } fn start_current_event_queue(&self) { let simulator = mutex(&self.manager) .upgrade() .and_then(|manager| read(&manager.current_sim).clone()) .unwrap_or_else(|| self.native_clone_without_caps()); if let Some(inner) = read(&simulator.caps_state).clone() { let caps = Caps::native_from_inner(inner, simulator.native_clone_without_caps()); if let Some(queue) = caps.event_queue() { let _ = queue.start(); } } } pub(crate) fn native_raise_event_queue_running(&self) { let mut generation = mutex(&self.event_queue_generation); *generation = generation.wrapping_add(1); drop(generation); self.event_queue_wait.notify_all(); if let Some(manager) = mutex(&self.manager).upgrade() { let args = EventQueueRunningEventArgs { simulator: self.clone(), }; manager .events .event_queue_running .emit_with(|| args.clone()); } } pub(crate) fn native_dispatch_caps_event(&self, event_name: &str, message: &dyn IMessage) { if let Some(manager) = mutex(&self.manager).upgrade() { manager .caps_events .invoke_raise_event(event_name, message, self.clone()); } } pub(crate) fn native_enqueue_caps_packet(&self, packet: Packet) -> Result<(), Error> { let manager = mutex(&self.manager) .upgrade() .ok_or(Error::InvalidOperation)?; manager.enqueue_incoming(NetworkManagerIncomingPacket { packet: Some(packet), simulator: Some(self.clone()), raw_data: None, }) } #[must_use] pub fn seed_capability(&self) -> Option { read(&self.seed_caps).clone() } pub fn native_use_circuit_code(&self, wait_for_ack: bool) -> Result<(), Error> { let data = self.use_circuit_code_bytes()?; if wait_for_ack { let receiver = self .transport .try_send_packet_wait_ack(data, self.endpoint, PacketType::UseCircuitCode, false) .map_err(Error::from)?; let timeout = positive_millisecond_duration(self.client.settings_ref().timing.login_timeout); let _ = receiver.recv_timeout(timeout); } else { self.transport .try_send_packet(data, self.endpoint, PacketType::UseCircuitCode, false) .map_err(Error::from)?; } Ok(()) } pub async fn native_use_circuit_code_async(&self, wait_for_ack: bool) -> Result { let data = self.use_circuit_code_bytes()?; if !wait_for_ack { self.transport .try_send_packet(data, self.endpoint, PacketType::UseCircuitCode, false) .map_err(Error::from)?; return Ok(true); } let receiver = self .transport .try_send_packet_wait_ack(data, self.endpoint, PacketType::UseCircuitCode, false) .map_err(Error::from)?; let timeout = positive_millisecond_duration(self.client.settings_ref().timing.login_timeout); run_blocking_compat("libremetaverse-circuit-ack", move || { receiver.recv_timeout(timeout).is_ok() }) .await } fn use_circuit_code_bytes(&self) -> Result, Error> { let mut packet = UseCircuitCodePacket::new_with_constructor()?; packet.circuit_code.code = self.circuit_code.load(Ordering::Acquire); packet.circuit_code.id = *read(&self.agent_id); packet.circuit_code.session_id = *read(&self.session_id); packet.to_bytes_with_method() } #[must_use] pub fn native_is_connected(&self) -> bool { self.connected.load(Ordering::Acquire) } #[cfg(test)] pub(crate) fn native_set_connected_for_tests(&self, connected: bool) { self.connected.store(connected, Ordering::Release); } #[must_use] pub fn native_handshake_complete(&self) -> bool { self.handshake_complete.load(Ordering::Acquire) } #[must_use] pub fn native_ip_end_point(&self) -> std::net::SocketAddr { self.data.endpoint } #[must_use] pub fn native_to_string(&self) -> String { if self.name.is_empty() { format!("({})", self.endpoint) } else { format!("{} ({})", self.name, self.endpoint) } } fn note_packet_received(&self, _packet_type: PacketType) { self.disconnect_candidate.store(false, Ordering::Release); } fn mark_disconnect_candidate(&self) -> bool { self.disconnect_candidate.swap(true, Ordering::AcqRel) } fn send_buffer(&self, buffer: UDPPacketBuffer) -> Result<(), Error> { self.transport.async_begin_send(buffer) } } struct SimulatorUdpHandler { simulator: Arc>>, } impl UdpPacketHandler for SimulatorUdpHandler { fn packet_received(&self, buffer: UDPPacketBuffer) { let Some(data) = mutex(&self.simulator).upgrade() else { return; }; let simulator = Simulator { caps: None, data }; let Ok(length) = usize::try_from(buffer.data_length) else { return; }; let Some(payload) = buffer.data.get(..length) else { return; }; let Ok(payload_length) = i32::try_from(payload.len()) else { return; }; let mut packet_end = payload_length - 1; let mut zero_buffer = vec![0_u8; UdpTransportConfig::default().max_decoded_packet_size]; let Ok(packet) = crate::packet_wire::build_packet_from_bytes(payload, &mut packet_end, &mut zero_buffer) else { return; }; simulator.note_packet_received(packet.type_); let manager = mutex(&simulator.manager).upgrade(); if let Some(manager) = manager { let _ = manager.enqueue_incoming(NetworkManagerIncomingPacket { simulator: Some(simulator), packet: Some(packet), raw_data: Some(payload.to_vec()), }); } } fn packet_sent(&self, buffer: UDPPacketBuffer, bytes_sent: usize) { let Some(data) = mutex(&self.simulator).upgrade() else { return; }; let simulator = Simulator { caps: None, data }; let manager = mutex(&simulator.manager).upgrade(); if let Some(manager) = manager { manager.raise_packet_sent(buffer.data, bytes_sent, simulator); } } } #[derive(Clone)] pub struct NetworkManagerIncomingPacket { pub packet: Option, pub simulator: Option, raw_data: Option>, } impl NetworkManagerIncomingPacket { pub fn new() -> Result { Ok(Self { packet: None, simulator: None, raw_data: None, }) } } #[derive(Clone)] pub struct NetworkManagerOutgoingPacket { pub buffer: UDPPacketBuffer, pub resend_count: i32, pub sequence_number: u32, pub simulator: Simulator, pub tick_count: i32, pub type_: PacketType, } impl NetworkManagerOutgoingPacket { pub fn new( simulator: Simulator, buffer: UDPPacketBuffer, type_: PacketType, ) -> Result { Ok(Self { simulator, buffer, sequence_number: 0, resend_count: 0, tick_count: 0, type_, }) } } struct NetworkWorkers { incoming: SyncSender, outgoing: SyncSender, cancellation: libremetaverse_types::compat::CancellationTokenSource, handles: Vec>, } #[derive(Default)] struct NetworkEvents { disconnected: EventRegistry, event_queue_running: EventRegistry, generic_streaming_message: EventRegistry, login_progress: EventRegistry, packet_sent: EventRegistry, sim_changed: EventRegistry, sim_connected: EventRegistry, sim_connecting: EventRegistry, sim_disconnected: EventRegistry, } struct LoginRuntimeState { next_generation: u64, active: Option<(u64, CancellationTokenSource)>, response_callbacks: Vec<(NetworkManagerLoginResponseCallback, Option>)>, } impl Default for LoginRuntimeState { fn default() -> Self { Self { next_generation: 1, active: None, response_callbacks: Vec::new(), } } } pub(crate) struct NetworkManagerInner { client: GridClient, simulators: SimulatorCollection, packet_events: PacketEventDictionary, caps_events: CapsEventDictionary, events: NetworkEvents, current_sim: RwLock>, connected: AtomicBool, circuit_code: AtomicU32, agent_id: RwLock, session_id: RwLock, login_status: RwLock, login_error_key: RwLock, login_message: RwLock, login_seed_capability: RwLock>, login_response: LoginResponseData, login_runtime: Mutex, inbox_count: AtomicI32, outbox_count: AtomicI32, workers: Mutex>, disconnect_lock: Mutex<()>, shutting_down: AtomicBool, } impl NetworkManagerInner { fn ensure_workers(self: &Arc) -> Result<(), Error> { let mut workers = mutex(&self.workers); if workers.is_some() { return Ok(()); } let capacity = usize::try_from(crate::Settings::udp_receive_queue_capacity()) .map_err(|_| Error::InvalidOperation)? .max(1); let (incoming_sender, incoming_receiver) = sync_channel::(capacity); let (outgoing_sender, outgoing_receiver) = sync_channel::(capacity); let cancellation = libremetaverse_types::compat::CancellationTokenSource::default(); let incoming_token = cancellation.token(); let outgoing_token = cancellation.token(); let keepalive_token = cancellation.token(); let incoming_weak = Arc::downgrade(self); let outgoing_weak = Arc::downgrade(self); let keepalive_weak = Arc::downgrade(self); let incoming_handle = thread::Builder::new() .name("libremetaverse-network-incoming".to_owned()) .spawn(move || { while !incoming_token.is_cancellation_requested() { match incoming_receiver.recv_timeout(Duration::from_millis(50)) { Ok(incoming) => { let Some(inner) = incoming_weak.upgrade() else { break; }; inner.inbox_count.fetch_sub(1, Ordering::AcqRel); let raw_data = incoming.raw_data; let (Some(packet), Some(simulator)) = (incoming.packet, incoming.simulator) else { continue; }; inner.process_internal_packet(&packet, &simulator, raw_data.as_deref()); let _ = inner .packet_events .raise_event(packet.type_, packet, simulator); } Err(RecvTimeoutError::Timeout) => {} Err(RecvTimeoutError::Disconnected) => break, } } }) .map_err(|_| Error::InvalidOperation)?; let outgoing_handle = thread::Builder::new() .name("libremetaverse-network-outgoing".to_owned()) .spawn(move || { while !outgoing_token.is_cancellation_requested() { match outgoing_receiver.recv_timeout(Duration::from_millis(50)) { Ok(outgoing) => { let Some(inner) = outgoing_weak.upgrade() else { break; }; inner.outbox_count.fetch_sub(1, Ordering::AcqRel); let _ = outgoing.simulator.send_buffer(outgoing.buffer); } Err(RecvTimeoutError::Timeout) => {} Err(RecvTimeoutError::Disconnected) => break, } } }) .map_err(|_| Error::InvalidOperation)?; let timeout = positive_millisecond_duration(self.client.settings_ref().timing.simulator_timeout); let keepalive_handle = thread::Builder::new() .name("libremetaverse-network-keepalive".to_owned()) .spawn(move || { while !wait_for_cancellation(&keepalive_token, timeout) { let Some(inner) = keepalive_weak.upgrade() else { break; }; inner.keepalive_tick(); } }) .map_err(|_| Error::InvalidOperation)?; *workers = Some(NetworkWorkers { incoming: incoming_sender, outgoing: outgoing_sender, cancellation, handles: vec![incoming_handle, outgoing_handle, keepalive_handle], }); Ok(()) } fn process_internal_packet( &self, packet: &Packet, simulator: &Simulator, raw_data: Option<&[u8]>, ) { match packet.type_ { PacketType::RegionHandshake => { let Some(raw_data) = raw_data else { return; }; let mut position = 0; let Ok(_incoming) = RegionHandshakePacket::new_with_bytes_int32(raw_data.to_vec(), &mut position) else { return; }; let Ok(mut reply) = RegionHandshakeReplyPacket::new_with_constructor() else { return; }; reply.agent_data.agent_id = UUID::zero(); reply.agent_data.session_id = UUID::zero(); reply.region_info.flags = 0x1 | 0x2 | 0x4; let Ok(data) = reply.to_bytes_with_method() else { return; }; let Ok(length) = i32::try_from(data.len()) else { return; }; if simulator .native_send_packet_data(data, length, PacketType::RegionHandshakeReply, true) .is_err() { return; } simulator.connected.store(true, Ordering::Release); simulator.handshake_complete.store(true, Ordering::Release); simulator.handshake_wait.notify_all(); } PacketType::StartPingCheck => { let Some(raw_data) = raw_data else { return; }; let Ok(mut incoming) = StartPingCheckPacket::new_with_constructor() else { return; }; let Ok(mut packet_end) = i32::try_from(raw_data.len()) else { return; }; packet_end -= 1; let mut position = 0; let mut zero_buffer = vec![0_u8; UdpTransportConfig::default().max_decoded_packet_size]; if incoming .from_bytes_with_bytes_int32_int32_bytes( raw_data.to_vec(), &mut position, &mut packet_end, Some(&mut zero_buffer), ) .is_err() { return; } if incoming.ping_id.oldest_unacked > 0 { let _ = simulator.transport.try_flush_acks(simulator.endpoint); } let Ok(mut response) = CompletePingCheckPacket::new_with_constructor() else { return; }; response.ping_id.ping_id = incoming.ping_id.ping_id; let Ok(mut data) = response.to_bytes_with_method() else { return; }; if let Some(flags) = data.first_mut() { *flags &= !Helpers::MSG_RELIABLE; } let Ok(length) = i32::try_from(data.len()) else { return; }; let _ = simulator.native_send_packet_data( data, length, PacketType::CompletePingCheck, false, ); } PacketType::DisableSimulator => { let _ = self.disconnect_simulator(simulator.clone(), false); } PacketType::KickUser => { let Some(raw_data) = raw_data else { return; }; let Ok(mut incoming) = KickUserPacket::new_with_constructor() else { return; }; let Ok(mut packet_end) = i32::try_from(raw_data.len()) else { return; }; packet_end -= 1; let mut position = 0; let mut zero_buffer = vec![0_u8; UdpTransportConfig::default().max_decoded_packet_size]; if incoming .from_bytes_with_bytes_int32_int32_bytes( raw_data.to_vec(), &mut position, &mut packet_end, Some(&mut zero_buffer), ) .is_err() { return; } let message = String::from_utf8_lossy(&incoming.user_info.reason) .trim_end_matches('\0') .to_owned(); let _ = self.shutdown( crate::NetworkManagerDisconnectType::ServerInitiated, message, ); } PacketType::GenericStreamingMessage => { let Some(raw_data) = raw_data else { return; }; let Ok(mut incoming) = GenericStreamingMessagePacket::new_with_constructor() else { return; }; let Ok(mut packet_end) = i32::try_from(raw_data.len()) else { return; }; packet_end -= 1; let mut position = 0; let mut zero_buffer = vec![0_u8; UdpTransportConfig::default().max_decoded_packet_size]; if incoming .from_bytes_with_bytes_int32_int32_bytes( raw_data.to_vec(), &mut position, &mut packet_end, Some(&mut zero_buffer), ) .is_err() { return; } let method = if incoming.method_data.method == GenericStreamingMethod::GltfMaterialOverride as u16 { GenericStreamingMethod::GltfMaterialOverride } else { GenericStreamingMethod::Unknown }; let Ok(args) = GenericStreamingMessageEventArgs::new( simulator.clone(), method, incoming.data_block.data, ) else { return; }; self.events .generic_streaming_message .emit_with(|| args.clone()); } _ => {} } } fn disconnect_simulator( &self, simulator: Simulator, send_close_circuit: bool, ) -> Result<(), Error> { let disconnect_guard = mutex(&self.disconnect_lock); let was_connected = simulator.connected(); simulator.native_disconnect(send_close_circuit)?; if !self.simulators.remove(&simulator) { return Ok(()); } if read(&self.current_sim).as_ref() == Some(&simulator) { *write(&self.current_sim) = None; } let args = SimDisconnectedEventArgs::new( simulator.clone(), crate::NetworkManagerDisconnectType::NetworkTimeout, was_connected, )?; drop(disconnect_guard); self.events.sim_disconnected.emit_with(|| args.clone()); if self.simulators.is_empty() { self.shutdown( crate::NetworkManagerDisconnectType::SimShutdown, format!( "You have disconnected from {}.", simulator.native_to_string() ), )?; } Ok(()) } fn enqueue_incoming(&self, incoming: NetworkManagerIncomingPacket) -> Result<(), Error> { let workers = mutex(&self.workers); let Some(workers) = workers.as_ref() else { return Err(Error::InvalidOperation); }; workers .incoming .try_send(incoming) .map_err(|_| Error::InvalidOperation)?; self.inbox_count.fetch_add(1, Ordering::AcqRel); Ok(()) } fn enqueue_outgoing(&self, outgoing: NetworkManagerOutgoingPacket) -> Result<(), Error> { let workers = mutex(&self.workers); let Some(workers) = workers.as_ref() else { return Err(Error::InvalidOperation); }; workers .outgoing .try_send(outgoing) .map_err(|_| Error::InvalidOperation)?; self.outbox_count.fetch_add(1, Ordering::AcqRel); Ok(()) } fn stop_workers(&self) { let workers = mutex(&self.workers).take(); let Some(workers) = workers else { return; }; workers.cancellation.cancel(); drop(workers.incoming); drop(workers.outgoing); for handle in workers.handles { if handle.thread().id() != thread::current().id() { let _ = handle.join(); } } self.inbox_count.store(0, Ordering::Release); self.outbox_count.store(0, Ordering::Release); } fn keepalive_tick(&self) { if !self.connected.load(Ordering::Acquire) { return; } let current = read(&self.current_sim).clone(); let Some(current) = current else { return; }; if current.mark_disconnect_candidate() { let _ = self.shutdown( crate::NetworkManagerDisconnectType::NetworkTimeout, "NetworkTimeout".to_owned(), ); } else { let _ = current.native_send_ping(); } } fn raise_packet_sent(&self, data: Vec, bytes_sent: usize, simulator: Simulator) { let Ok(bytes_sent) = i32::try_from(bytes_sent) else { return; }; let Ok(args) = PacketSentEventArgs::new(data, bytes_sent, simulator) else { return; }; self.events.packet_sent.emit_with(|| args.clone()); } fn set_current_sim(&self, simulator: Simulator, seed_caps: Option) -> Result<(), Error> { let previous = { let mut current = write(&self.current_sim); if current.as_ref() == Some(&simulator) { simulator.native_set_seed_caps(seed_caps, Some(false))?; return Ok(()); } let previous = current.clone(); *current = Some(simulator.clone()); previous }; simulator.native_set_seed_caps(seed_caps, Some(previous.as_ref() != Some(&simulator)))?; let args = SimChangedEventArgs::new(previous)?; self.events.sim_changed.emit_with(|| args.clone()); Ok(()) } fn shutdown( &self, reason: crate::NetworkManagerDisconnectType, message: String, ) -> Result<(), Error> { if self.shutting_down.swap(true, Ordering::AcqRel) { return Ok(()); } let current = write(&self.current_sim).take(); let mut simulators: Vec<_> = self .simulators .drain() .into_iter() .filter(|simulator| current.as_ref() != Some(simulator)) .collect(); if let Some(current) = current { simulators.push(current); } let send_close = matches!( reason, crate::NetworkManagerDisconnectType::ClientInitiated | crate::NetworkManagerDisconnectType::NetworkTimeout ); for simulator in simulators { let was_connected = simulator.connected(); let _ = simulator.native_disconnect(send_close); let args = SimDisconnectedEventArgs::new(simulator, reason, was_connected)?; self.events.sim_disconnected.emit_with(|| args.clone()); } self.stop_workers(); self.connected.store(false, Ordering::Release); let args = DisconnectedEventArgs::new(reason, message)?; self.events.disconnected.emit_with(|| args.clone()); self.shutting_down.store(false, Ordering::Release); Ok(()) } } impl Drop for NetworkManagerInner { fn drop(&mut self) { self.stop_workers(); for simulator in self.simulators.drain() { let _ = simulator.native_disconnect(false); } } } /// Native implementation of C# `NetworkManager`. pub struct NetworkManager { pub login_response_data: Option, pub simulators: SimulatorCollection, inner: Arc, } impl fmt::Debug for NetworkManager { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter .debug_struct("NetworkManager") .field("connected", &self.native_connected()) .field("circuit_code", &self.native_circuit_code()) .field("simulators", &self.simulators) .finish_non_exhaustive() } } impl NetworkManager { pub(crate) fn native_from_inner(inner: Arc) -> Self { Self { login_response_data: Some(inner.login_response.clone()), simulators: inner.simulators.clone(), inner, } } pub(crate) fn native_inner_weak(&self) -> Weak { Arc::downgrade(&self.inner) } pub fn native_new(client: GridClient) -> Result { let simulators = SimulatorCollection::default(); let packet_events = PacketEventDictionary::new(client.clone())?; let caps_events = CapsEventDictionary::new(client.clone())?; let login_response = LoginResponseData::new()?; let inner = Arc::new(NetworkManagerInner { client, simulators: simulators.clone(), packet_events, caps_events, events: NetworkEvents::default(), current_sim: RwLock::new(None), connected: AtomicBool::new(false), circuit_code: AtomicU32::new(0), agent_id: RwLock::new(UUID::zero()), session_id: RwLock::new(UUID::zero()), login_status: RwLock::new(LoginStatus::None), login_error_key: RwLock::new(String::new()), login_message: RwLock::new(String::new()), login_seed_capability: RwLock::new(None), login_response: login_response.clone(), login_runtime: Mutex::new(LoginRuntimeState::default()), inbox_count: AtomicI32::new(0), outbox_count: AtomicI32::new(0), workers: Mutex::new(None), disconnect_lock: Mutex::new(()), shutting_down: AtomicBool::new(false), }); let weak_inner = Arc::downgrade(&inner); inner.caps_events.register_event( "EnableSimulator".to_owned(), CapsEventQueueCallback::from_handler(move |_, message, _| { let Some(inner) = weak_inner.upgrade() else { return; }; if !inner.client.settings_ref().agent.multiple_sims { return; } let any_message = message as &dyn std::any::Any; let Some(message) = any_message.downcast_ref::() else { return; }; let manager = NetworkManager { login_response_data: Some(inner.login_response.clone()), simulators: inner.simulators.clone(), inner: Arc::clone(&inner), }; for simulator in &message.simulators { let Ok(port) = u16::try_from(simulator.port) else { continue; }; let endpoint = std::net::SocketAddr::new(simulator.ip, port); if manager.native_find_endpoint(endpoint).is_some() { continue; } let _ = manager.native_connect( endpoint, simulator.region_handle, false, None, simulator.region_size_x, simulator.region_size_y, ); } }), )?; Ok(Self { login_response_data: Some(login_response), simulators, inner, }) } pub fn native_subscribe_login_progress( &self, handler: EventHandler, ) -> Subscription { self.inner.events.login_progress.subscribe(handler) } fn update_login_status(&self, status: LoginStatus, message: String) { if status == LoginStatus::Failed { *write(&self.inner.agent_id) = UUID::zero(); *write(&self.inner.session_id) = UUID::zero(); *write(&self.inner.login_seed_capability) = None; self.inner.login_response.clear_secrets(); } *write(&self.inner.login_status) = status; write(&self.inner.login_message).clone_from(&message); let fail_reason = read(&self.inner.login_error_key).clone(); if let Ok(args) = LoginProgressEventArgs::new(status, message, fail_reason) { self.inner.events.login_progress.emit_with(|| args.clone()); } } fn set_login_error(&self, value: &str) { let mut error = write(&self.inner.login_error_key); error.clear(); error.push_str(value); } #[must_use] pub fn native_login_status_code(&self) -> LoginStatus { *read(&self.inner.login_status) } pub fn native_set_login_status_code(&self, value: LoginStatus) { *write(&self.inner.login_status) = value; } #[must_use] pub fn native_login_error_key(&self) -> String { read(&self.inner.login_error_key).clone() } pub fn native_set_login_error_key(&self, value: String) { *write(&self.inner.login_error_key) = value; } #[must_use] pub fn native_login_message(&self) -> String { read(&self.inner.login_message).clone() } pub fn native_set_login_message(&self, value: String) { *write(&self.inner.login_message) = value; } #[must_use] pub fn native_login_seed_capability(&self) -> Option { read(&self.inner.login_seed_capability).clone() } pub fn native_set_login_seed_capability(&self, value: Option) { *write(&self.inner.login_seed_capability) = value; } #[must_use] pub fn native_account_level_benefits(&self) -> Option { self.inner.login_response.account_level_benefits() } #[must_use] pub fn native_agent_appearance_service_url(&self) -> Option { let value = self.inner.login_response.agent_appearance_service_url(); (!value.is_empty()).then_some(value) } #[must_use] pub fn native_max_agent_groups(&self) -> i32 { let fallback = self.inner.login_response.max_agent_groups(); self.inner .login_response .account_level_benefits() .map_or(fallback, |benefits| { let limit = benefits.group_membership_limit(); if limit > 0 { limit } else { fallback } }) } pub fn native_default_login_params( &self, first_name: String, last_name: String, password: String, channel: String, version: String, ) -> Result { LoginParams::new_with_grid_client_string_string_string_string_string( self.inner.client.clone(), first_name, last_name, password, channel, version, ) } pub fn native_register_login_response_callback( &self, callback: NetworkManagerLoginResponseCallback, options: Option>, ) -> Result<(), Error> { let mut state = mutex(&self.inner.login_runtime); if state .response_callbacks .iter() .any(|(candidate, _)| candidate.ptr_eq(&callback)) { return Err(Error::Argument); } state.response_callbacks.push((callback, options)); Ok(()) } pub fn native_unregister_login_response_callback( &self, callback: &NetworkManagerLoginResponseCallback, ) -> Result<(), Error> { let mut state = mutex(&self.inner.login_runtime); state .response_callbacks .retain(|(candidate, _)| !candidate.ptr_eq(callback)); Ok(()) } fn reserve_login( &self, cancellation_token: Option, ) -> Result<(u64, CancellationTokenSource), Error> { let mut runtime = mutex(&self.inner.login_runtime); if runtime.active.is_some() { return Err(Error::InvalidOperation); } let caller = cancellation_token.unwrap_or_default(); let source = CancellationTokenSource::new_linked(&[caller, self.inner.client.cancellation_token()]); let generation = runtime.next_generation; runtime.next_generation = runtime.next_generation.wrapping_add(1).max(1); runtime.active = Some((generation, source.clone())); Ok((generation, source)) } fn release_login(&self, generation: u64) { let mut runtime = mutex(&self.inner.login_runtime); if runtime .active .as_ref() .is_some_and(|(active, _)| *active == generation) { runtime.active = None; } } pub fn native_abort_login(&self) -> Result<(), Error> { let active = mutex(&self.inner.login_runtime).active.take(); if let Some((_, source)) = active { source.cancel(); } self.update_login_status(LoginStatus::Failed, "Abort Requested".to_owned()); Ok(()) } pub fn native_begin_login(&self, login_params: LoginParams) -> Result<(), Error> { let (generation, source) = self.reserve_login(None)?; let manager = Self { login_response_data: Some(self.inner.login_response.clone()), simulators: self.simulators.clone(), inner: Arc::clone(&self.inner), }; thread::Builder::new() .name("libremetaverse-login".to_owned()) .spawn(move || { let runtime = tokio::runtime::Builder::new_current_thread() .enable_all() .build(); let result = match runtime { Ok(runtime) => { runtime.block_on(manager.perform_login(login_params, source.token())) } Err(_) => Err(Error::InvalidOperation), }; if let Err(error) = result { manager.login_transport_failure(&error); } manager.release_login(generation); }) .map_err(|_| { self.release_login(generation); Error::InvalidOperation })?; Ok(()) } pub async fn native_login( &self, login_params: LoginParams, cancellation_token: Option, ) -> Result { Ok(self .native_login_response_async(login_params, cancellation_token) .await? .is_some_and(|response| response.success())) } pub async fn native_login_response_async( &self, login_params: LoginParams, cancellation_token: Option, ) -> Result, Error> { let timeout = positive_millisecond_duration(login_params.timeout); let (generation, source) = self.reserve_login(cancellation_token)?; let token = source.token(); let result = tokio::select! { result = tokio::time::timeout(timeout, self.perform_login(login_params, token.clone())) => { if let Ok(result) = result { result } else { source.cancel(); self.set_login_error("timeout"); self.update_login_status(LoginStatus::Failed, "Login timed out".to_owned()); Ok(None) } } () = token.cancelled() => { self.set_login_error("canceled"); self.update_login_status(LoginStatus::Failed, "Canceled".to_owned()); Ok(None) } }; self.release_login(generation); match result { Ok(value) => Ok(value), Err(error) => { self.login_transport_failure(&error); Ok(None) } } } pub async fn native_login_credential( &self, credential: LoginCredential, channel: String, start: Option, version: String, cancellation_token: Option, ) -> Result { let mut login_params = LoginParams::new_with_grid_client_login_credential_string_string( self.inner.client.clone(), credential, channel, version, )?; if let Some(start) = start { login_params.start = start; } self.native_login(login_params, cancellation_token).await } pub async fn native_login_credential_response( &self, credential: LoginCredential, channel: String, start: Option, version: String, cancellation_token: Option, ) -> Result, Error> { let mut login_params = LoginParams::new_with_grid_client_login_credential_string_string( self.inner.client.clone(), credential, channel, version, )?; if let Some(start) = start { login_params.start = start; } self.native_login_response_async(login_params, cancellation_token) .await } pub async fn native_login_strings( &self, first_name: String, last_name: String, password: String, channel: String, start: String, version: String, cancellation_token: Option, ) -> Result { let mut params = self.native_default_login_params(first_name, last_name, password, channel, version)?; params.start = start; self.native_login(params, cancellation_token).await } pub async fn native_login_strings_response( &self, first_name: String, last_name: String, password: String, channel: String, start: String, version: String, cancellation_token: Option, ) -> Result, Error> { let mut params = self.native_default_login_params(first_name, last_name, password, channel, version)?; params.start = start; self.native_login_response_async(params, cancellation_token) .await } async fn perform_login( &self, mut params: LoginParams, cancellation_token: CancellationToken, ) -> Result, Error> { if params.password.len() != 35 && !params.password.starts_with("$1$") { params.password = libremetaverse_types::Utils::md5_with_string(params.password)?; } if params.channel.trim().is_empty() { params.channel = crate::Settings::user_agent(); } if params.version.trim().is_empty() { params.version = String::from("?.?.?"); } let wire_secrets = LoginWireSecrets { password: std::mem::take(&mut params.password), token: std::mem::take(&mut params.token), mfa_hash: std::mem::take(&mut params.mfa_hash), }; params.clear_secrets(); let requested_start = if params.login_location.trim().is_empty() { params.start.clone() } else { params.login_location.clone() }; params.start = crate::login::normalize_start(&requested_start)?; let callback_options: Vec = mutex(&self.inner.login_runtime) .response_callbacks .iter() .filter_map(|(_, options)| options.as_ref()) .flatten() .cloned() .collect(); for option in callback_options { if !params.options.contains(&option) { params.options.push(option); } } let mut redirects = 0_usize; loop { cancellation_token.throw_if_cancellation_requested()?; let login_uri = validate_login_uri(¶ms.uri)?; self.update_login_status( LoginStatus::ConnectingToLogin, format!( "Logging in as {} {}...", params.first_name, params.last_name ), ); let request = login_request(¶ms, &wire_secrets); let response = self .inner .client .native_http_caps_client() .post_with_uri_osd_format_osd_cancellation_token_i_progress( login_uri, OSDFormat::Xml, request, cancellation_token.clone(), None, ) .await; let (http_response, response_data) = match response { Ok(response) => response, Err(Error::Cancelled) => { self.set_login_error("canceled"); self.update_login_status(LoginStatus::Failed, "Login canceled".to_owned()); return Ok(None); } Err(error) => { self.login_transport_failure(&error); return Ok(None); } }; if !http_response.is_success_status_code() { self.set_login_error("bad response"); self.update_login_status( LoginStatus::Failed, format!("Login service returned HTTP {}", http_response.status_code), ); return Ok(None); } self.update_login_status( LoginStatus::ReadingResponse, "Reading login response...".to_owned(), ); let parsed = match OSDParser::deserialize_with_bytes(response_data) { Ok(OSD::Map(map)) => OSDMap::new_with_dictionary(map), _ => Err(Error::Argument), }; let Ok(parsed) = parsed else { self.set_login_error("bad response"); self.update_login_status( LoginStatus::Failed, "Empty or corrupt login response".to_owned(), ); return Ok(None); }; let has_login = parsed.get("login").is_some(); self.inner.login_response.parse_with_osd_map(parsed)?; let response = self.inner.login_response.clone(); if !has_login { self.update_login_status( LoginStatus::Failed, "Login parameter missing in the response".to_owned(), ); return Ok(None); } match response.login() { crate::LoginState::Indeterminate => { redirects += 1; if redirects > 10 { self.set_login_error("redirect loop"); self.update_login_status( LoginStatus::Failed, "Too many login redirects".to_owned(), ); return Ok(None); } self.update_login_status(LoginStatus::Redirecting, response.message()); params.uri = response.next_url(); validate_login_uri(¶ms.uri)?; let delay = Duration::from_secs( u64::try_from(response.next_duration().clamp(0, 300)).unwrap_or(0), ); tokio::select! { () = tokio::time::sleep(delay) => {} () = cancellation_token.cancelled() => { self.set_login_error("canceled"); self.update_login_status(LoginStatus::Failed, "Login canceled".to_owned()); return Ok(None); } } } crate::LoginState::True => { self.invoke_login_response(true, false, &response); if response.seed_capability().is_empty() { self.set_login_error("bad response"); self.update_login_status( LoginStatus::Failed, "Login response did not include a seed capability".to_owned(), ); return Ok(None); } if self.initialize_login_response( response.clone(), response.message(), "Unable to establish a UDP connection to the simulator", "Login server did not return a simulator address", )? { return Ok(Some(response)); } return Ok(None); } crate::LoginState::False => { let reason = response.reason(); self.set_login_error(if reason.is_empty() { "unknown" } else { &reason }); self.update_login_status(LoginStatus::Failed, response.message()); return Ok(None); } } } } fn invoke_login_response(&self, success: bool, redirect: bool, response: &LoginResponseData) { let callbacks: Vec<_> = mutex(&self.inner.login_runtime) .response_callbacks .iter() .map(|(callback, _)| callback.clone()) .collect(); for callback in callbacks { let _ = callback.invoke( success, redirect, response.message(), response.reason(), Some(response.clone()), ); } } fn login_transport_failure(&self, error: &Error) { let (key, message) = match error { Error::Cancelled => ("canceled", "Login canceled"), Error::HttpRequest => ("no connection", "Unable to contact login service"), Error::Argument => ("bad response", "Invalid login request or response"), _ => ("no connection", "Login failed"), }; self.set_login_error(key); self.update_login_status(LoginStatus::Failed, message.to_owned()); } pub fn native_login_response(&self, response: LoginResponseData) -> Result { self.initialize_login_response( response, "Login success".to_owned(), "Unable to connect to simulator", "Unable to connect to simulator", ) } fn initialize_login_response( &self, response: LoginResponseData, success_message: String, connect_failure_message: &str, address_failure_message: &str, ) -> Result { self.inner.login_response.copy_from(&response); self.inner .circuit_code .store(response.circuit_code().cast_unsigned(), Ordering::Release); *write(&self.inner.agent_id) = response.agent_id(); *write(&self.inner.session_id) = response.session_id(); let seed = if response.seed_capability().is_empty() { None } else { Some(validate_login_uri(&response.seed_capability())?) }; write(&self.inner.login_seed_capability).clone_from(&seed); self.update_login_status( LoginStatus::ConnectingToSim, "Connecting to simulator...".to_owned(), ); let Some(ip) = response.sim_ip() else { self.set_login_error("bad response"); self.update_login_status(LoginStatus::Failed, address_failure_message.to_owned()); return Ok(false); }; if response.sim_port() == 0 { self.set_login_error("bad response"); self.update_login_status(LoginStatus::Failed, address_failure_message.to_owned()); return Ok(false); } let handle = (u64::from(response.region_x()) << 32) | u64::from(response.region_y()); let endpoint = std::net::SocketAddr::new(ip, response.sim_port()); let Ok(connected) = self.native_connect( endpoint, handle, true, seed, response.region_size_x(), response.region_size_y(), ) else { self.cleanup_failed_login_connection(endpoint); self.set_login_error("no connection"); self.update_login_status(LoginStatus::Failed, connect_failure_message.to_owned()); return Ok(false); }; if let Some(simulator) = connected { let economy_request = Packet::build_packet_with_packet_type(PacketType::EconomyDataRequest); if economy_request .and_then(|packet| simulator.native_send_packet(packet)) .is_err() { self.cleanup_failed_login_connection(endpoint); self.set_login_error("no connection"); self.update_login_status(LoginStatus::Failed, connect_failure_message.to_owned()); return Ok(false); } self.update_login_status(LoginStatus::Success, success_message); Ok(true) } else { self.cleanup_failed_login_connection(endpoint); self.update_login_status(LoginStatus::Failed, connect_failure_message.to_owned()); Ok(false) } } fn cleanup_failed_login_connection(&self, endpoint: std::net::SocketAddr) { if let Some(simulator) = self.native_find_endpoint(endpoint) { let _ = simulator.native_disconnect(false); self.simulators.remove(&simulator); let mut current = write(&self.inner.current_sim); if current.as_ref() == Some(&simulator) { *current = None; } } if self.simulators.is_empty() { self.inner.connected.store(false, Ordering::Release); } } pub fn native_subscribe_disconnected( &self, handler: EventHandler, ) -> Subscription { self.inner.events.disconnected.subscribe(handler) } pub fn native_subscribe_event_queue_running( &self, handler: EventHandler, ) -> Subscription { self.inner.events.event_queue_running.subscribe(handler) } pub fn native_subscribe_generic_streaming_message( &self, handler: EventHandler, ) -> Subscription { self.inner .events .generic_streaming_message .subscribe(handler) } pub fn native_subscribe_packet_sent( &self, handler: EventHandler, ) -> Subscription { self.inner.events.packet_sent.subscribe(handler) } pub fn native_subscribe_sim_changed( &self, handler: EventHandler, ) -> Subscription { self.inner.events.sim_changed.subscribe(handler) } pub fn native_subscribe_sim_connected( &self, handler: EventHandler, ) -> Subscription { self.inner.events.sim_connected.subscribe(handler) } pub fn native_subscribe_sim_connecting( &self, handler: EventHandler, ) -> Subscription { self.inner.events.sim_connecting.subscribe(handler) } pub fn native_subscribe_sim_disconnected( &self, handler: EventHandler, ) -> Subscription { self.inner.events.sim_disconnected.subscribe(handler) } pub fn subscribe_packet( &self, packet_type: PacketType, handler: EventHandler, is_async: bool, ) -> Subscription { self.inner .packet_events .subscribe(packet_type, handler, is_async) } pub fn subscribe_caps( &self, caps_event: String, handler: CapsEventQueueCallback, ) -> Subscription { self.inner.caps_events.subscribe(caps_event, handler) } /// Dispatches a decoded CAPS event through the default and named handler /// chains. This is the native boundary used by a CAPS event-queue client. pub fn dispatch_caps_event( &self, caps_event: &str, message: &dyn IMessage, simulator: Simulator, ) { self.inner .caps_events .invoke_raise_event(caps_event, message, simulator); } pub fn native_register_callback( &self, packet_type: PacketType, callback: EventHandler, is_async: bool, ) -> Result<(), Error> { self.inner .packet_events .register_event(packet_type, callback, is_async) } pub fn native_unregister_callback( &self, packet_type: PacketType, callback: EventHandler, ) -> Result<(), Error> { self.inner .packet_events .unregister_event(packet_type, callback) } pub fn native_register_caps( &self, caps_event: String, callback: CapsEventQueueCallback, ) -> Result<(), Error> { self.inner.caps_events.register_event(caps_event, callback) } pub fn native_unregister_caps( &self, caps_event: String, callback: CapsEventQueueCallback, ) -> Result<(), Error> { self.inner .caps_events .unregister_event(caps_event, callback) } pub fn native_connect( &self, endpoint: std::net::SocketAddr, handle: u64, set_default: bool, seed_caps: Option, size_x: u32, size_y: u32, ) -> Result, Error> { self.inner.ensure_workers()?; let simulator = self .native_find_endpoint(endpoint) .unwrap_or(Simulator::native_new( self.inner.client.clone(), endpoint, handle, Some(size_x), Some(size_y), )?); simulator.attach_manager(&self.inner, self.inner.circuit_code.load(Ordering::Acquire)); let simulator = self.simulators.push_if_missing(simulator); if !simulator.connected() { self.inner.connected.store(true, Ordering::Release); let args = SimConnectingEventArgs::new(simulator.clone())?; self.inner.events.sim_connecting.emit_with(|| args.clone()); if args.cancel() { self.simulators.remove(&simulator); return Ok(None); } match simulator.native_connect(set_default) { Ok(true) => {} Ok(false) => { self.simulators.remove(&simulator); return Ok(None); } Err(error) => { let _ = simulator.native_disconnect(false); self.simulators.remove(&simulator); return Err(error); } } if set_default { self.inner.set_current_sim(simulator.clone(), seed_caps)?; } let args = SimConnectedEventArgs::new(simulator.clone())?; self.inner.events.sim_connected.emit_with(|| args.clone()); } else if set_default { simulator.native_use_circuit_code(true)?; self.inner.set_current_sim(simulator.clone(), seed_caps)?; } Ok(Some(simulator)) } pub async fn native_connect_async( &self, endpoint: std::net::SocketAddr, handle: u64, set_default: bool, seed_caps: Option, size_x: u32, size_y: u32, ) -> Result, Error> { let manager = Self { login_response_data: Some(self.inner.login_response.clone()), simulators: self.simulators.clone(), inner: Arc::clone(&self.inner), }; run_blocking_compat("libremetaverse-manager-connect", move || { manager.native_connect(endpoint, handle, set_default, seed_caps, size_x, size_y) }) .await? } pub fn native_disconnect_sim( &self, simulator: Simulator, send_close_circuit: bool, ) -> Result<(), Error> { self.inner .disconnect_simulator(simulator, send_close_circuit) } pub fn native_enqueue_incoming( &self, incoming: NetworkManagerIncomingPacket, ) -> Result<(), Error> { self.inner.enqueue_incoming(incoming) } pub fn native_enqueue_outgoing( &self, outgoing: NetworkManagerOutgoingPacket, ) -> Result<(), Error> { self.inner.enqueue_outgoing(outgoing) } #[must_use] pub fn native_find_endpoint(&self, endpoint: std::net::SocketAddr) -> Option { self.simulators .snapshot() .into_iter() .find(|simulator| simulator.ip_end_point() == endpoint) } #[must_use] pub fn native_find_handle(&self, handle: u64) -> Option { self.simulators .snapshot() .into_iter() .find(|simulator| simulator.handle == handle) } pub fn native_send_packet( &self, packet: Packet, simulator: Option, ) -> Result<(), Error> { let simulator = simulator .or_else(|| self.native_current_sim()) .or_else(|| self.simulators.snapshot().into_iter().next()) .ok_or(Error::InvalidOperation)?; simulator.native_send_packet(packet) } pub fn native_shutdown( &self, reason: crate::NetworkManagerDisconnectType, message: String, ) -> Result<(), Error> { self.inner.shutdown(reason, message) } pub async fn native_shutdown_async( &self, reason: crate::NetworkManagerDisconnectType, message: String, ) -> Result<(), Error> { let inner = Arc::clone(&self.inner); run_blocking_compat("libremetaverse-manager-shutdown", move || { inner.shutdown(reason, message) }) .await? } #[must_use] pub fn native_circuit_code(&self) -> u32 { self.inner.circuit_code.load(Ordering::Acquire) } pub fn native_set_circuit_code(&self, value: u32) { self.inner.circuit_code.store(value, Ordering::Release); for simulator in self.simulators.snapshot() { simulator.circuit_code.store(value, Ordering::Release); } } #[must_use] pub fn native_connected(&self) -> bool { self.inner.connected.load(Ordering::Acquire) } pub fn native_set_connected(&self, value: bool) { self.inner.connected.store(value, Ordering::Release); } #[must_use] pub fn native_current_sim(&self) -> Option { read(&self.inner.current_sim).clone() } pub fn native_set_current_sim(&self, value: Option) { *write(&self.inner.current_sim) = value; } #[must_use] pub fn native_inbox_count(&self) -> i32 { self.inner.inbox_count.load(Ordering::Acquire) } #[must_use] pub fn native_outbox_count(&self) -> i32 { self.inner.outbox_count.load(Ordering::Acquire) } /// Executes one deterministic keepalive transition without sleeping. pub fn poll_keepalive(&self) { self.inner.keepalive_tick(); } } #[cfg(test)] mod tests { use super::*; use libremetaverse_structured_data::{OSD, OSDMap}; use std::sync::mpsc; #[test] fn manager_workers_and_callback_registries_do_not_retain_the_client() { let client = GridClient::new().expect("client"); let probe = client.retention_probe(); let manager = NetworkManager::new(client.clone()).expect("network manager"); manager.inner.ensure_workers().expect("manager workers"); let subscription = manager.subscribe_packet(PacketType::Default, Arc::new(|_| {}), true); drop(subscription); drop(manager); drop(client); assert!(probe.is_released()); } #[test] fn caps_packet_fallback_uses_the_bounded_udp_dispatch_pipeline() { let client = GridClient::new().expect("client"); let manager = NetworkManager::new(client.clone()).expect("network manager"); manager.inner.ensure_workers().expect("workers"); let simulator = Simulator::new(client, "127.0.0.1:14000".parse().unwrap(), 10, None, None).unwrap(); simulator.attach_manager(&manager.inner, 0); manager.simulators.push_if_missing(simulator.clone()); let (sender, receiver) = mpsc::channel(); let _subscription = manager.subscribe_packet( PacketType::PacketAck, Arc::new(move |args| { sender .send((args.packet().type_, args.simulator().handle)) .unwrap(); }), false, ); let body = OSDMap::new_with_constructor().unwrap(); body.add_with_string_osd("Packets".to_owned(), OSD::Array(Vec::new())) .unwrap(); let packet = Packet::build_packet_with_string_osd_map("PacketAck".to_owned(), body) .unwrap() .expect("packet fallback"); simulator.native_enqueue_caps_packet(packet).unwrap(); assert_eq!( receiver.recv_timeout(Duration::from_secs(2)).unwrap(), (PacketType::PacketAck, 10) ); } #[test] fn event_queue_running_and_typed_caps_events_use_manager_registries() { let client = GridClient::new().expect("client"); let manager = NetworkManager::new(client.clone()).expect("network manager"); let simulator = Simulator::new(client, "127.0.0.1:14001".parse().unwrap(), 11, None, None).unwrap(); simulator.attach_manager(&manager.inner, 0); let (running_sender, running_receiver) = mpsc::channel(); let _running = manager.native_subscribe_event_queue_running(Arc::new(move |args| { running_sender.send(args.simulator().handle).unwrap(); })); simulator.native_raise_event_queue_running(); assert_eq!( running_receiver .recv_timeout(Duration::from_secs(1)) .unwrap(), 11 ); let (caps_sender, caps_receiver) = mpsc::channel(); let _caps = manager.subscribe_caps( "UpdateAgentLanguage".to_owned(), CapsEventQueueCallback::from_handler(move |name, message, simulator| { let typed = (message as &dyn std::any::Any) .downcast_ref::(); caps_sender .send((name.to_owned(), typed.is_some(), simulator.handle)) .unwrap(); }), ); let body = OSDMap::new_with_dictionary(HashMap::from([ ("language".to_owned(), OSD::String("en".to_owned())), ("language_is_public".to_owned(), OSD::Boolean(true)), ])) .unwrap(); let message = crate::message_decoder::decode_event("UpdateAgentLanguage", &body) .unwrap() .unwrap(); simulator.native_dispatch_caps_event("UpdateAgentLanguage", message.as_ref()); assert_eq!( caps_receiver.recv_timeout(Duration::from_secs(1)).unwrap(), ("UpdateAgentLanguage".to_owned(), true, 11) ); } }