diff --git a/README.md b/README.md index f6314c5..56900b6 100644 --- a/README.md +++ b/README.md @@ -303,6 +303,14 @@ Credentials and session/capability tokens are redacted from diagnostics, and offline fake login/simulator tests cover success, rejection, redirect, timeout, and cancellation cleanup. The login contract is documented with the network lifecycle in [`docs/network-manager.md`](docs/network-manager.md). +Region handoff and logout now follow the same packet-level lifecycle as the +golden implementation: a promoted simulator receives `UseCircuitCode` and +`CompleteAgentMovement`, while blocking, asynchronous, and nonblocking logout +send `LogoutRequest`, validate `LogoutReply`, preserve callback ordering, and +perform bounded idempotent teardown. Shutdown cancels login/logout work before +closing every simulator and worker, clears session and capability secrets, and +allows the manager to reconnect cleanly afterward. Loopback fake-server tests +cover handoff, reconnect, reply, timeout, cancellation, and repeated shutdown. Seed-cap discovery and `EventQueueGet` are native Rust too: the full reference capability list is posted as LLSD/XML, accepted URIs are rate-categorized, and bounded long polls preserve ack IDs, reconnect retries, shutdown `done`, and diff --git a/api/SHIM-COVERAGE.md b/api/SHIM-COVERAGE.md index e4858e3..7d1f6fa 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: 77 types / 14,065 members; remaining surface is callable failure-only shims | +| `LibreMetaverse` | 2,711 | 27,281 | native implementation: 78 types / 14,072 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/src/generated.rs b/crates/libremetaverse/src/generated.rs index b1ed24f..fb1acde 100644 --- a/crates/libremetaverse/src/generated.rs +++ b/crates/libremetaverse/src/generated.rs @@ -16206,18 +16206,9 @@ impl LoadUrlEventArgs { } /// C# type: `T:LibreMetaverse.LoggedOutEventArgs`. -pub struct LoggedOutEventArgs { - /// C# member: `F:LibreMetaverse.LoggedOutEventArgs.InventoryItems`. - pub inventory_items: Vec, -} -impl LoggedOutEventArgs { - /// C# member: `M:LibreMetaverse.LoggedOutEventArgs.#ctor(System.Collections.Generic.List{LibreMetaverse.UUID})`. - pub fn new(inventory_items: Vec) -> Result { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.LoggedOutEventArgs.#ctor(System.Collections.Generic.List{LibreMetaverse.UUID})", - ) - } -} +/// C# member: `F:LibreMetaverse.LoggedOutEventArgs.InventoryItems`. +/// C# member: `M:LibreMetaverse.LoggedOutEventArgs.#ctor(System.Collections.Generic.List{LibreMetaverse.UUID})`. +pub use crate::network_manager::LoggedOutEventArgs; /// C# type: `T:LibreMetaverse.Logger`. pub struct Logger; @@ -17640,7 +17631,8 @@ impl NetworkManager { &self, handler: libremetaverse_types::compat::EventHandler, ) -> libremetaverse_types::compat::Subscription { - libremetaverse_types::unimplemented_api!("E:LibreMetaverse.NetworkManager.LoggedOut") + /* native network implementation */ + self.native_subscribe_logged_out(handler) } /// C# member: `E:LibreMetaverse.NetworkManager.LoginProgress`. pub fn subscribe_login_progress( @@ -17712,7 +17704,8 @@ impl NetworkManager { } /// C# member: `M:LibreMetaverse.NetworkManager.BeginLogout`. pub fn begin_logout(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.NetworkManager.BeginLogout") + /* native network implementation */ + self.native_begin_logout() } /// C# member: `M:LibreMetaverse.NetworkManager.Connect(System.Net.IPAddress,System.UInt16,System.UInt64,System.Boolean,System.Uri)`. pub fn connect_with_ip_address_u_int16_u_int64_boolean_uri( @@ -18035,16 +18028,16 @@ impl NetworkManager { } /// C# member: `M:LibreMetaverse.NetworkManager.Logout`. pub fn logout_with_method(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.NetworkManager.Logout") + /* native network implementation */ + self.native_logout() } /// C# member: `M:LibreMetaverse.NetworkManager.LogoutAsync(System.Threading.CancellationToken)`. pub async fn logout_with_cancellation_token( &self, cancellation_token: Option, ) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented( - "M:LibreMetaverse.NetworkManager.LogoutAsync(System.Threading.CancellationToken)", - ) + /* native network implementation */ + self.native_logout_async(cancellation_token).await } /// C# member: `M:LibreMetaverse.NetworkManager.RefreshMac`. pub fn refresh_mac() -> Result<(), crate::Error> { @@ -18102,7 +18095,8 @@ impl NetworkManager { } /// C# member: `M:LibreMetaverse.NetworkManager.RequestLogout`. pub fn request_logout(&self) -> Result<(), crate::Error> { - libremetaverse_types::not_implemented("M:LibreMetaverse.NetworkManager.RequestLogout") + /* native network implementation */ + self.native_request_logout() } /// C# member: `M:LibreMetaverse.NetworkManager.SendPacket(LibreMetaverse.Packets.Packet)`. pub fn send_packet_with_packet( diff --git a/crates/libremetaverse/src/network_manager.rs b/crates/libremetaverse/src/network_manager.rs index 27aef40..d66df2a 100644 --- a/crates/libremetaverse/src/network_manager.rs +++ b/crates/libremetaverse/src/network_manager.rs @@ -13,9 +13,10 @@ use crate::interfaces::IMessage; use crate::packets::{ - AgentPausePacket, AgentResumePacket, CloseCircuitPacket, CompletePingCheckPacket, - GenericStreamingMessagePacket, KickUserPacket, Packet, PacketType, RegionHandshakePacket, - RegionHandshakeReplyPacket, StartPingCheckPacket, UseCircuitCodePacket, + AgentPausePacket, AgentResumePacket, CloseCircuitPacket, CompleteAgentMovementPacket, + CompletePingCheckPacket, GenericStreamingMessagePacket, KickUserPacket, LogoutReplyPacket, + LogoutRequestPacket, Packet, PacketType, RegionHandshakePacket, RegionHandshakeReplyPacket, + StartPingCheckPacket, UseCircuitCodePacket, }; use crate::udp_transport::{UDPBase, UDPPacketBuffer, UdpPacketHandler, UdpTransportConfig}; use crate::{ @@ -32,7 +33,7 @@ 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::atomic::{AtomicBool, AtomicI32, AtomicU8, 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}; @@ -95,6 +96,63 @@ fn positive_millisecond_duration(value: i32) -> Duration { Duration::from_millis(u64::try_from(value.max(1)).unwrap_or(1)) } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum LogoutWaitOutcome { + Reply, + Timeout, + Cancelled, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[repr(u8)] +enum NetworkLifecycleState { + Disconnected, + Connecting, + Connected, + LoggingOut, + ShuttingDown, +} + +impl NetworkLifecycleState { + fn from_u8(value: u8) -> Self { + match value { + 1 => Self::Connecting, + 2 => Self::Connected, + 3 => Self::LoggingOut, + 4 => Self::ShuttingDown, + _ => Self::Disconnected, + } + } +} + +fn wait_for_logout( + receiver: &Receiver<()>, + timeout: Duration, + cancellations: &[CancellationToken], +) -> LogoutWaitOutcome { + let started = std::time::Instant::now(); + loop { + if receiver.try_recv().is_ok() { + return LogoutWaitOutcome::Reply; + } + if cancellations + .iter() + .any(CancellationToken::is_cancellation_requested) + { + return LogoutWaitOutcome::Cancelled; + } + let remaining = timeout.saturating_sub(started.elapsed()); + if remaining.is_zero() { + return LogoutWaitOutcome::Timeout; + } + match receiver.recv_timeout(remaining.min(Duration::from_millis(20))) { + Ok(()) => return LogoutWaitOutcome::Reply, + Err(RecvTimeoutError::Timeout) => {} + Err(RecvTimeoutError::Disconnected) => return LogoutWaitOutcome::Cancelled, + } + } +} + 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() { @@ -287,6 +345,18 @@ pub struct SimChangedEventArgs { previous_simulator: Option, } +/// Inventory item identifiers returned by a successful logout handshake. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct LoggedOutEventArgs { + pub inventory_items: Vec, +} + +impl LoggedOutEventArgs { + pub const fn new(inventory_items: Vec) -> Result { + Ok(Self { inventory_items }) + } +} + impl SimChangedEventArgs { pub fn new(previous_simulator: Option) -> Result { Ok(Self { previous_simulator }) @@ -1141,9 +1211,12 @@ impl Simulator { *write(&self.session_id) = *read(&manager.session_id); } - pub fn native_connect(&self, _move_to_sim: bool) -> Result { + pub fn native_connect(&self, move_to_sim: bool) -> Result { if self.native_is_connected() { self.native_use_circuit_code(true)?; + if move_to_sim { + self.native_complete_agent_movement()?; + } return Ok(true); } let mut transport_thread = mutex(&self.transport_thread); @@ -1195,6 +1268,9 @@ impl Simulator { self.handshake_complete.store(false, Ordering::Release); drop(transport_thread); self.native_use_circuit_code(true)?; + if move_to_sim { + self.native_complete_agent_movement()?; + } let timeout = positive_millisecond_duration(self.client.settings_ref().timing.login_timeout); let guard = mutex(&self.handshake_wait_lock); @@ -1500,6 +1576,20 @@ impl Simulator { packet.to_bytes_with_method() } + fn native_complete_agent_movement(&self) -> Result<(), Error> { + let mut packet = CompleteAgentMovementPacket::new_with_constructor()?; + packet.agent_data.agent_id = *read(&self.agent_id); + packet.agent_data.session_id = *read(&self.session_id); + packet.agent_data.circuit_code = self.circuit_code.load(Ordering::Acquire); + let data = packet.to_bytes_with_method()?; + self.native_send_packet_data( + data.clone(), + i32::try_from(data.len()).map_err(|_| Error::Argument)?, + PacketType::CompleteAgentMovement, + false, + ) + } + #[must_use] pub fn native_is_connected(&self) -> bool { self.connected.load(Ordering::Acquire) @@ -1647,6 +1737,7 @@ struct NetworkEvents { disconnected: EventRegistry, event_queue_running: EventRegistry, generic_streaming_message: EventRegistry, + logged_out: EventRegistry, login_progress: EventRegistry, packet_sent: EventRegistry, sim_changed: EventRegistry, @@ -1688,14 +1779,25 @@ pub(crate) struct NetworkManagerInner { login_seed_capability: RwLock>, login_response: LoginResponseData, login_runtime: Mutex, + logout_cancellation: Mutex, inbox_count: AtomicI32, outbox_count: AtomicI32, workers: Mutex>, disconnect_lock: Mutex<()>, - shutting_down: AtomicBool, + lifecycle_state: AtomicU8, + logout_task_count: AtomicI32, + last_logout_error: Mutex>, } impl NetworkManagerInner { + fn lifecycle_state(&self) -> NetworkLifecycleState { + NetworkLifecycleState::from_u8(self.lifecycle_state.load(Ordering::Acquire)) + } + + fn set_lifecycle_state(&self, state: NetworkLifecycleState) { + self.lifecycle_state.store(state as u8, Ordering::Release); + } + fn ensure_workers(self: &Arc) -> Result<(), Error> { let mut workers = mutex(&self.workers); if workers.is_some() { @@ -1808,8 +1910,8 @@ impl NetworkManagerInner { 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.agent_data.agent_id = *read(&self.agent_id); + reply.agent_data.session_id = *read(&self.session_id); reply.region_info.flags = 0x1 | 0x2 | 0x4; let Ok(data) = reply.to_bytes_with_method() else { return; @@ -1875,6 +1977,51 @@ impl NetworkManagerInner { false, ); } + PacketType::LogoutReply => { + let Some(raw_data) = raw_data else { + return; + }; + let Ok(mut incoming) = LogoutReplyPacket::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.agent_data.agent_id != *read(&self.agent_id) + || incoming.agent_data.session_id != *read(&self.session_id) + { + return; + } + let Ok(args) = LoggedOutEventArgs::new( + incoming + .inventory_data + .into_iter() + .map(|item| item.item_id) + .collect(), + ) else { + return; + }; + self.events.logged_out.emit_with(|| args.clone()); + let _ = self.shutdown( + crate::NetworkManagerDisconnectType::ClientInitiated, + "ClientInitiated".to_owned(), + ); + } PacketType::DisableSimulator => { let _ = self.disconnect_simulator(simulator.clone(), false); } @@ -2035,7 +2182,9 @@ impl NetworkManagerInner { } fn keepalive_tick(&self) { - if !self.connected.load(Ordering::Acquire) { + if !self.connected.load(Ordering::Acquire) + || self.lifecycle_state() != NetworkLifecycleState::Connected + { return; } let current = read(&self.current_sim).clone(); @@ -2063,17 +2212,16 @@ impl NetworkManagerInner { } 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 - }; + let transition = mutex(&self.disconnect_lock); + let previous = read(&self.current_sim).clone(); + if previous.as_ref() == Some(&simulator) { + drop(transition); + simulator.native_set_seed_caps(seed_caps, Some(false))?; + return Ok(()); + } simulator.native_set_seed_caps(seed_caps, Some(previous.as_ref() != Some(&simulator)))?; + *write(&self.current_sim) = Some(simulator.clone()); + drop(transition); let args = SimChangedEventArgs::new(previous)?; self.events.sim_changed.emit_with(|| args.clone()); Ok(()) @@ -2084,9 +2232,31 @@ impl NetworkManagerInner { reason: crate::NetworkManagerDisconnectType, message: String, ) -> Result<(), Error> { - if self.shutting_down.swap(true, Ordering::AcqRel) { - return Ok(()); + loop { + let state = self.lifecycle_state(); + if matches!( + state, + NetworkLifecycleState::Disconnected | NetworkLifecycleState::ShuttingDown + ) { + return Ok(()); + } + if self + .lifecycle_state + .compare_exchange( + state as u8, + NetworkLifecycleState::ShuttingDown as u8, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_ok() + { + break; + } } + if let Some((_, cancellation)) = mutex(&self.login_runtime).active.take() { + cancellation.cancel(); + } + mutex(&self.logout_cancellation).cancel(); let current = write(&self.current_sim).take(); let mut simulators: Vec<_> = self .simulators @@ -2105,20 +2275,29 @@ impl NetworkManagerInner { 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()); + if let Ok(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(()) + self.circuit_code.store(0, Ordering::Release); + *write(&self.agent_id) = UUID::zero(); + *write(&self.session_id) = UUID::zero(); + *write(&self.login_seed_capability) = None; + self.login_response.clear_secrets(); + let event_result = DisconnectedEventArgs::new(reason, message); + if let Ok(args) = &event_result { + self.events.disconnected.emit_with(|| args.clone()); + } + self.set_lifecycle_state(NetworkLifecycleState::Disconnected); + event_result.map(|_| ()) } } impl Drop for NetworkManagerInner { fn drop(&mut self) { + mutex(&self.logout_cancellation).cancel(); self.stop_workers(); for simulator in self.simulators.drain() { let _ = simulator.native_disconnect(false); @@ -2179,11 +2358,14 @@ impl NetworkManager { login_seed_capability: RwLock::new(None), login_response: login_response.clone(), login_runtime: Mutex::new(LoginRuntimeState::default()), + logout_cancellation: Mutex::new(CancellationTokenSource::new()), inbox_count: AtomicI32::new(0), outbox_count: AtomicI32::new(0), workers: Mutex::new(None), disconnect_lock: Mutex::new(()), - shutting_down: AtomicBool::new(false), + lifecycle_state: AtomicU8::new(NetworkLifecycleState::Disconnected as u8), + logout_task_count: AtomicI32::new(0), + last_logout_error: Mutex::new(None), }); let weak_inner = Arc::downgrade(&inner); inner.caps_events.register_event( @@ -2239,6 +2421,13 @@ impl NetworkManager { self.inner.events.login_progress.subscribe(handler) } + pub fn native_subscribe_logged_out( + &self, + handler: EventHandler, + ) -> Subscription { + self.inner.events.logged_out.subscribe(handler) + } + fn update_login_status(&self, status: LoginStatus, message: String) { if status == LoginStatus::Failed { *write(&self.inner.agent_id) = UUID::zero(); @@ -2840,6 +3029,8 @@ impl NetworkManager { } if self.simulators.is_empty() { self.inner.connected.store(false, Ordering::Release); + self.inner + .set_lifecycle_state(NetworkLifecycleState::Disconnected); } } @@ -2982,21 +3173,31 @@ impl NetworkManager { size_x: u32, size_y: u32, ) -> Result, Error> { + if matches!( + self.inner.lifecycle_state(), + NetworkLifecycleState::LoggingOut | NetworkLifecycleState::ShuttingDown + ) { + return Err(Error::InvalidOperation); + } 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), - )?); + let existing = self.native_find_endpoint(endpoint); + let newly_created = existing.is_none(); + let simulator = existing.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 + .set_lifecycle_state(NetworkLifecycleState::Connecting); self.inner.connected.store(true, Ordering::Release); + self.inner + .set_lifecycle_state(NetworkLifecycleState::Connected); let args = SimConnectingEventArgs::new(simulator.clone())?; self.inner.events.sim_connecting.emit_with(|| args.clone()); if args.cancel() { @@ -3015,13 +3216,20 @@ impl NetworkManager { return Err(error); } } - if set_default { - self.inner.set_current_sim(simulator.clone(), seed_caps)?; + if set_default + && let Err(error) = self.inner.set_current_sim(simulator.clone(), seed_caps) + { + if newly_created { + let _ = simulator.native_disconnect(false); + self.simulators.remove(&simulator); + } + return Err(error); } 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)?; + simulator.native_complete_agent_movement()?; self.inner.set_current_sim(simulator.clone(), seed_caps)?; } Ok(Some(simulator)) @@ -3098,6 +3306,206 @@ impl NetworkManager { simulator.native_send_packet(packet) } + fn send_logout_request_if_needed(&self) -> Result<(), Error> { + let Some(simulator) = self.native_current_sim() else { + return Ok(()); + }; + if !self.native_connected() { + return Ok(()); + } + loop { + match self.inner.lifecycle_state() { + NetworkLifecycleState::Connected => { + if self + .inner + .lifecycle_state + .compare_exchange( + NetworkLifecycleState::Connected as u8, + NetworkLifecycleState::LoggingOut as u8, + Ordering::AcqRel, + Ordering::Acquire, + ) + .is_ok() + { + break; + } + } + NetworkLifecycleState::LoggingOut + | NetworkLifecycleState::Disconnected + | NetworkLifecycleState::Connecting + | NetworkLifecycleState::ShuttingDown => return Ok(()), + } + } + { + let mut cancellation = mutex(&self.inner.logout_cancellation); + if cancellation.token().is_cancellation_requested() { + *cancellation = CancellationTokenSource::new(); + } + } + let result = (|| { + let mut packet = LogoutRequestPacket::new_with_constructor()?; + packet.agent_data.agent_id = *read(&self.inner.agent_id); + packet.agent_data.session_id = *read(&self.inner.session_id); + let data = packet.to_bytes_with_method()?; + simulator.native_send_packet_data( + data.clone(), + i32::try_from(data.len()).map_err(|_| Error::Argument)?, + PacketType::LogoutRequest, + false, + ) + })(); + if result.is_err() { + let _ = self.inner.lifecycle_state.compare_exchange( + NetworkLifecycleState::LoggingOut as u8, + NetworkLifecycleState::Connected as u8, + Ordering::AcqRel, + Ordering::Acquire, + ); + } + result + } + + fn logout_signal(&self) -> (Receiver<()>, Subscription) { + let (sender, receiver) = sync_channel(1); + let subscription = self.native_subscribe_logged_out(Arc::new(move |_| { + let _ = sender.try_send(()); + })); + (receiver, subscription) + } + + fn finish_logout_wait( + &self, + outcome: LogoutWaitOutcome, + emit_timeout_logout: bool, + ) -> Result<(), Error> { + match outcome { + LogoutWaitOutcome::Reply => { + self.native_shutdown( + crate::NetworkManagerDisconnectType::ClientInitiated, + "ClientInitiated".to_owned(), + )?; + let deadline = std::time::Instant::now() + + positive_millisecond_duration( + self.inner.client.settings_ref().timing.logout_timeout, + ); + while self.inner.lifecycle_state() != NetworkLifecycleState::Disconnected + && std::time::Instant::now() < deadline + { + thread::sleep(Duration::from_millis(1)); + } + Ok(()) + } + LogoutWaitOutcome::Timeout => { + self.native_shutdown( + crate::NetworkManagerDisconnectType::NetworkTimeout, + "NetworkTimeout".to_owned(), + )?; + if emit_timeout_logout { + let args = LoggedOutEventArgs::new(Vec::new())?; + self.inner.events.logged_out.emit_with(|| args.clone()); + } + Ok(()) + } + LogoutWaitOutcome::Cancelled => { + let _ = self.inner.lifecycle_state.compare_exchange( + NetworkLifecycleState::LoggingOut as u8, + NetworkLifecycleState::Connected as u8, + Ordering::AcqRel, + Ordering::Acquire, + ); + Err(Error::Cancelled) + } + } + } + + fn logout_blocking(&self, emit_timeout_logout: bool) -> Result<(), Error> { + if !self.native_connected() || self.native_current_sim().is_none() { + return Ok(()); + } + let (receiver, _subscription) = self.logout_signal(); + self.send_logout_request_if_needed()?; + let timeout = + positive_millisecond_duration(self.inner.client.settings_ref().timing.logout_timeout); + let lifecycle_cancellation = mutex(&self.inner.logout_cancellation).token(); + self.finish_logout_wait( + wait_for_logout(&receiver, timeout, &[lifecycle_cancellation]), + emit_timeout_logout, + ) + } + + pub fn native_begin_logout(&self) -> Result<(), Error> { + if !self.native_connected() || self.native_current_sim().is_none() { + return Ok(()); + } + *mutex(&self.inner.last_logout_error) = None; + let manager = Self { + login_response_data: Some(self.inner.login_response.clone()), + simulators: self.simulators.clone(), + inner: Arc::clone(&self.inner), + }; + self.inner.logout_task_count.fetch_add(1, Ordering::AcqRel); + let result = thread::Builder::new() + .name("libremetaverse-logout".to_owned()) + .spawn(move || { + if let Err(error) = manager.logout_blocking(true) { + *mutex(&manager.inner.last_logout_error) = Some(error); + } + manager + .inner + .logout_task_count + .fetch_sub(1, Ordering::AcqRel); + }) + .map(|_| ()) + .map_err(|_| Error::InvalidOperation); + if result.is_err() { + self.inner.logout_task_count.fetch_sub(1, Ordering::AcqRel); + *mutex(&self.inner.last_logout_error) = Some(Error::InvalidOperation); + } + result + } + + pub fn native_logout(&self) -> Result<(), Error> { + self.logout_blocking(false) + } + + pub async fn native_logout_async( + &self, + cancellation_token: Option, + ) -> Result<(), Error> { + if !self.native_connected() || self.native_current_sim().is_none() { + return Ok(()); + } + let (receiver, _subscription) = self.logout_signal(); + self.send_logout_request_if_needed()?; + let timeout = + positive_millisecond_duration(self.inner.client.settings_ref().timing.logout_timeout); + let cancellations = vec![ + cancellation_token.unwrap_or_default(), + mutex(&self.inner.logout_cancellation).token(), + ]; + let outcome = run_blocking_compat("libremetaverse-logout-wait", move || { + wait_for_logout(&receiver, timeout, &cancellations) + }) + .await?; + self.finish_logout_wait(outcome, false) + } + + pub fn native_request_logout(&self) -> Result<(), Error> { + self.send_logout_request_if_needed() + } + + /// Number of nonblocking logout waiters still owned by this manager. + #[must_use] + pub fn pending_logout_tasks(&self) -> i32 { + self.inner.logout_task_count.load(Ordering::Acquire) + } + + /// Last error produced by a detached [`Self::native_begin_logout`] waiter. + #[must_use] + pub fn last_logout_error(&self) -> Option { + *mutex(&self.inner.last_logout_error) + } + pub fn native_shutdown( &self, reason: crate::NetworkManagerDisconnectType, @@ -3137,6 +3545,10 @@ impl NetworkManager { pub fn native_set_connected(&self, value: bool) { self.inner.connected.store(value, Ordering::Release); + if value && self.inner.lifecycle_state() == NetworkLifecycleState::Disconnected { + self.inner + .set_lifecycle_state(NetworkLifecycleState::Connected); + } } #[must_use] diff --git a/crates/libremetaverse/tests/network_manager.rs b/crates/libremetaverse/tests/network_manager.rs index aac9fc9..2122bd1 100644 --- a/crates/libremetaverse/tests/network_manager.rs +++ b/crates/libremetaverse/tests/network_manager.rs @@ -3,8 +3,10 @@ use libremetaverse::messages::linden::{ EnableSimulatorMessage, EnableSimulatorMessageSimulatorInfoBlock, }; use libremetaverse::packets::{ - DisableSimulatorPacket, Packet, PacketAckPacket, PacketAckPacketPacketsBlock, PacketType, - RegionHandshakePacket, RegionHandshakeReplyPacket, StartPingCheckPacket, UseCircuitCodePacket, + CompleteAgentMovementPacket, DisableSimulatorPacket, LogoutReplyPacket, + LogoutReplyPacketInventoryDataBlock, LogoutRequestPacket, Packet, PacketAckPacket, + PacketAckPacketPacketsBlock, PacketType, RegionHandshakePacket, RegionHandshakeReplyPacket, + StartPingCheckPacket, UseCircuitCodePacket, }; use libremetaverse::{ CapsEventDictionary, CapsEventQueueCallback, GridClient, Helpers, HttpCapsClient, LoginState, @@ -20,7 +22,7 @@ use std::collections::{BTreeMap, HashMap}; use std::net::{IpAddr, Ipv4Addr, SocketAddr, UdpSocket}; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::mpsc::{Receiver, Sender, channel}; -use std::sync::{Arc, Barrier, Mutex, Weak}; +use std::sync::{Arc, Barrier, Mutex, MutexGuard, OnceLock, Weak}; use std::thread::{self, JoinHandle}; use std::time::{Duration, Instant}; @@ -73,6 +75,14 @@ fn wait_until(mut condition: impl FnMut() -> bool) { } } +fn network_test_guard() -> MutexGuard<'static, ()> { + static NETWORK_TEST_LOCK: OnceLock> = OnceLock::new(); + NETWORK_TEST_LOCK + .get_or_init(|| Mutex::new(())) + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + #[test] fn packet_callbacks_preserve_filtering_order_async_policy_and_reentrancy() { let client = GridClient::new().expect("client"); @@ -355,6 +365,9 @@ enum ServerCommand { Ping, Disable, ExpectEconomy, + Handoff, + ExpectLogout, + LogoutReply(Vec), Stop, } @@ -364,6 +377,8 @@ enum ServerReport { Circuit(u32), HandshakeReply(u32), Economy, + Logout(UUID, UUID), + Movement(u32), Ping(u8), } @@ -516,6 +531,93 @@ fn spawn_fake_server() -> FakeServer { break; } }, + ServerCommand::Handoff => { + loop { + let (length, _) = socket.recv_from(&mut buffer).expect("UseCircuitCode"); + if packet_type(&buffer[..length]) != Some(PacketType::UseCircuitCode) { + continue; + } + let mut position = 0; + let use_circuit = UseCircuitCodePacket::new_with_bytes_int32( + buffer[..length].to_vec(), + &mut position, + ) + .unwrap(); + report_sender + .send(ServerReport::Circuit(use_circuit.circuit_code.code)) + .unwrap(); + let sequence = u32::from_be_bytes(buffer[1..5].try_into().unwrap()); + socket.send_to(&ack(sequence), client_endpoint).unwrap(); + break; + } + loop { + let (length, _) = socket + .recv_from(&mut buffer) + .expect("CompleteAgentMovement"); + if packet_type(&buffer[..length]) != Some(PacketType::CompleteAgentMovement) + { + continue; + } + let mut position = 0; + let movement = CompleteAgentMovementPacket::new_with_bytes_int32( + buffer[..length].to_vec(), + &mut position, + ) + .unwrap(); + report_sender + .send(ServerReport::Movement(movement.agent_data.circuit_code)) + .unwrap(); + break; + } + } + ServerCommand::ExpectLogout => loop { + let (length, _) = socket.recv_from(&mut buffer).expect("LogoutRequest"); + if packet_type(&buffer[..length]) != Some(PacketType::LogoutRequest) { + continue; + } + let mut position = 0; + let request = LogoutRequestPacket::new_with_bytes_int32( + buffer[..length].to_vec(), + &mut position, + ) + .unwrap(); + report_sender + .send(ServerReport::Logout( + request.agent_data.agent_id, + request.agent_data.session_id, + )) + .unwrap(); + break; + }, + ServerCommand::LogoutReply(inventory_items) => loop { + let (length, _) = socket.recv_from(&mut buffer).expect("LogoutRequest"); + if packet_type(&buffer[..length]) != Some(PacketType::LogoutRequest) { + continue; + } + let mut position = 0; + let request = LogoutRequestPacket::new_with_bytes_int32( + buffer[..length].to_vec(), + &mut position, + ) + .unwrap(); + report_sender + .send(ServerReport::Logout( + request.agent_data.agent_id, + request.agent_data.session_id, + )) + .unwrap(); + let mut reply = LogoutReplyPacket::new_with_constructor().unwrap(); + reply.agent_data.agent_id = request.agent_data.agent_id; + reply.agent_data.session_id = request.agent_data.session_id; + reply.inventory_data = inventory_items + .into_iter() + .map(|item_id| LogoutReplyPacketInventoryDataBlock { item_id }) + .collect(); + let mut bytes = reply.to_bytes_with_method().unwrap(); + bytes[0] &= !(Helpers::MSG_RELIABLE | Helpers::MSG_ZEROCODED); + socket.send_to(&bytes, client_endpoint).unwrap(); + break; + }, ServerCommand::Stop => break, } } @@ -592,6 +694,7 @@ fn assert_reuse_preserves_connected_override( #[test] fn fake_server_drives_circuit_handshake_ping_disable_and_disconnect_reasons() { + let _network_guard = network_test_guard(); let server = spawn_fake_server(); let client = GridClient::new().expect("client"); let mut manager = NetworkManager::new(client).expect("manager"); @@ -689,6 +792,7 @@ fn fake_server_drives_circuit_handshake_ping_disable_and_disconnect_reasons() { #[test] fn concurrent_disconnect_emits_each_transition_once() { + let _network_guard = network_test_guard(); const CALLERS: usize = 12; let server = spawn_fake_server(); let client = GridClient::new().expect("client"); @@ -741,6 +845,7 @@ fn concurrent_disconnect_emits_each_transition_once() { #[test] fn sim_disconnected_callback_can_reenter_disconnect_without_deadlock() { + let _network_guard = network_test_guard(); let server = spawn_fake_server(); let client = GridClient::new().expect("client"); let mut manager = NetworkManager::new(client).expect("manager"); @@ -778,6 +883,7 @@ fn sim_disconnected_callback_can_reenter_disconnect_without_deadlock() { #[test] fn keepalive_requires_two_silent_intervals_before_network_timeout() { + let _network_guard = network_test_guard(); let server = spawn_fake_server(); let client = GridClient::new().expect("client"); let mut manager = NetworkManager::new(client).expect("manager"); @@ -821,6 +927,7 @@ fn keepalive_requires_two_silent_intervals_before_network_timeout() { #[test] fn enable_simulator_caps_event_adds_each_new_region_once() { + let _network_guard = network_test_guard(); let first_server = spawn_fake_server(); let second_server = spawn_fake_server(); let mut client = GridClient::new().expect("client"); @@ -879,8 +986,446 @@ fn enable_simulator_caps_event_adds_each_new_region_once() { .unwrap(); } +#[test] +fn connected_region_handoff_preserves_old_sim_and_sends_complete_agent_movement() { + let _network_guard = network_test_guard(); + let first_server = spawn_fake_server(); + let second_server = spawn_fake_server(); + let client = GridClient::new().expect("client"); + let mut manager = NetworkManager::new(client).expect("manager"); + manager.set_circuit_code(0x7788); + let first = manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + first_server.endpoint, + 100, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + let second = manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + second_server.endpoint, + 200, + false, + None, + 512, + 256, + ) + .unwrap() + .unwrap(); + for server in [&first_server, &second_server] { + assert_eq!( + server.reports.recv_timeout(Duration::from_secs(1)).unwrap(), + ServerReport::Circuit(0x7788) + ); + assert_eq!( + server.reports.recv_timeout(Duration::from_secs(1)).unwrap(), + ServerReport::HandshakeReply(0x1 | 0x2 | 0x4) + ); + } + + let (changed_sender, changed_receiver) = channel(); + let _changed = manager.subscribe_sim_changed(Arc::new(move |args| { + changed_sender.send(args.previous_simulator()).unwrap(); + })); + second_server.commands.send(ServerCommand::Handoff).unwrap(); + let selected = manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + second_server.endpoint, + 200, + true, + None, + 512, + 256, + ) + .unwrap() + .unwrap(); + + assert_eq!(selected, second); + assert_eq!(manager.current_sim(), Some(second.clone())); + assert_eq!( + changed_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + Some(first.clone()) + ); + assert_eq!( + second_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + ServerReport::Circuit(0x7788) + ); + assert_eq!( + second_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + ServerReport::Movement(0x7788) + ); + assert!(first.connected()); + assert!(second.connected()); + assert_eq!(manager.simulators.len(), 2); + + manager + .shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated) + .unwrap(); +} + +#[test] +fn manager_reconnects_with_fresh_workers_after_complete_shutdown() { + let _network_guard = network_test_guard(); + let first_server = spawn_fake_server(); + let second_server = spawn_fake_server(); + let client = GridClient::new().expect("client"); + let mut manager = NetworkManager::new(client).expect("manager"); + manager.set_circuit_code(0x7878); + + manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + first_server.endpoint, + 250, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + assert_eq!( + first_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + ServerReport::Circuit(0x7878) + ); + assert_eq!( + first_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + ServerReport::HandshakeReply(0x1 | 0x2 | 0x4) + ); + manager + .shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated) + .unwrap(); + assert!(!manager.connected()); + assert!(manager.current_sim().is_none()); + assert!(manager.simulators.is_empty()); + + manager.set_circuit_code(0x7979); + let reconnected = manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + second_server.endpoint, + 251, + true, + None, + 512, + 256, + ) + .unwrap() + .unwrap(); + assert_eq!( + second_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + ServerReport::Circuit(0x7979) + ); + assert_eq!( + second_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + ServerReport::HandshakeReply(0x1 | 0x2 | 0x4) + ); + assert!(manager.connected()); + assert_eq!(manager.current_sim(), Some(reconnected)); + assert_eq!(manager.simulators.len(), 1); + + manager + .shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated) + .unwrap(); +} + +#[test] +fn logout_reply_orders_events_returns_inventory_and_tears_down_once() { + let _network_guard = network_test_guard(); + let server = spawn_fake_server(); + let client = GridClient::new().expect("client"); + let mut manager = NetworkManager::new(client).expect("manager"); + manager.set_circuit_code(0x8899); + manager.set_login_seed_capability(Some(Uri("https://caps.invalid/secret".into()))); + let login_response = manager + .login_response_data + .as_mut() + .expect("login response storage"); + login_response.set_session_id( + UUID::new_with_string("aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa".into()).unwrap(), + ); + login_response.set_secure_session_id( + UUID::new_with_string("bbbbbbbb-bbbb-bbbb-bbbb-bbbbbbbbbbbb".into()).unwrap(), + ); + login_response.set_mfa_hash("secret-mfa".into()); + login_response.set_seed_capability("https://caps.invalid/secret".into()); + login_response.set_next_url("https://login.invalid/secret-token".into()); + manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + server.endpoint, + 300, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + assert_eq!( + server.reports.recv_timeout(Duration::from_secs(1)).unwrap(), + ServerReport::Circuit(0x8899) + ); + assert_eq!( + server.reports.recv_timeout(Duration::from_secs(1)).unwrap(), + ServerReport::HandshakeReply(0x1 | 0x2 | 0x4) + ); + + let order = Arc::new(Mutex::new(Vec::new())); + let logged_order = Arc::clone(&order); + let (inventory_sender, inventory_receiver) = channel(); + let _logged_out = manager.subscribe_logged_out(Arc::new(move |args| { + logged_order.lock().unwrap().push("logged_out"); + inventory_sender.send(args.inventory_items).unwrap(); + })); + let sim_order = Arc::clone(&order); + let _sim_disconnected = manager.subscribe_sim_disconnected(Arc::new(move |_| { + sim_order.lock().unwrap().push("sim_disconnected"); + })); + let disconnected_order = Arc::clone(&order); + let disconnected_count = Arc::new(AtomicUsize::new(0)); + let count = Arc::clone(&disconnected_count); + let _disconnected = manager.subscribe_disconnected(Arc::new(move |_| { + count.fetch_add(1, Ordering::Relaxed); + disconnected_order.lock().unwrap().push("disconnected"); + })); + let inventory = vec![ + UUID::new_with_string("11111111-1111-1111-1111-111111111111".into()).unwrap(), + UUID::new_with_string("22222222-2222-2222-2222-222222222222".into()).unwrap(), + ]; + server + .commands + .send(ServerCommand::LogoutReply(inventory.clone())) + .unwrap(); + + manager.logout_with_method().unwrap(); + + assert_eq!( + server.reports.recv_timeout(Duration::from_secs(1)).unwrap(), + ServerReport::Logout(UUID::zero(), UUID::zero()) + ); + assert_eq!( + inventory_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + inventory + ); + assert_eq!( + *order.lock().unwrap(), + ["logged_out", "sim_disconnected", "disconnected"] + ); + assert!(!manager.connected()); + assert!(manager.current_sim().is_none()); + assert!(manager.simulators.is_empty()); + assert_eq!(manager.circuit_code(), 0); + assert!(manager.login_seed_capability().is_none()); + let login_response = manager + .login_response_data + .as_ref() + .expect("login response storage"); + assert_eq!(login_response.session_id(), UUID::zero()); + assert_eq!(login_response.secure_session_id(), UUID::zero()); + assert!(login_response.mfa_hash().is_empty()); + assert!(login_response.seed_capability().is_empty()); + assert!(login_response.next_url().is_empty()); + + manager + .shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated) + .unwrap(); + assert_eq!(disconnected_count.load(Ordering::Relaxed), 1); +} + +#[test] +fn async_logout_timeout_and_cancellation_have_distinct_reference_results() { + let _network_guard = network_test_guard(); + let timeout_server = spawn_fake_server(); + let mut timeout_client = GridClient::new().expect("client"); + timeout_client.settings().timing().logout_timeout = 25; + let mut timeout_manager = NetworkManager::new(timeout_client).expect("manager"); + timeout_manager.set_circuit_code(0x9911); + timeout_manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + timeout_server.endpoint, + 400, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + let _ = timeout_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + let _ = timeout_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + let (reason_sender, reason_receiver) = channel(); + let _disconnected = timeout_manager.subscribe_disconnected(Arc::new(move |args| { + reason_sender.send(args.reason()).unwrap(); + })); + + block_on(timeout_manager.logout_with_cancellation_token(None)).unwrap(); + + assert_eq!( + reason_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + NetworkManagerDisconnectType::NetworkTimeout + ); + assert!(!timeout_manager.connected()); + + let cancel_server = spawn_fake_server(); + let mut cancel_client = GridClient::new().expect("client"); + cancel_client.settings().timing().logout_timeout = 1_000; + let mut cancel_manager = NetworkManager::new(cancel_client).expect("manager"); + cancel_manager.set_circuit_code(0x9922); + cancel_manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + cancel_server.endpoint, + 500, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + let _ = cancel_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + let _ = cancel_server + .reports + .recv_timeout(Duration::from_secs(1)) + .unwrap(); + let cancellation = CancellationTokenSource::new(); + cancellation.cancel(); + + let result = + block_on(cancel_manager.logout_with_cancellation_token(Some(cancellation.token()))); + + assert_eq!(result, Err(libremetaverse::Error::Cancelled)); + assert!(cancel_manager.connected()); + cancel_manager + .shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated) + .unwrap(); +} + +#[test] +fn begin_logout_times_out_nonblocking_and_emits_logged_out_after_teardown() { + let _network_guard = network_test_guard(); + let server = spawn_fake_server(); + let mut client = GridClient::new().expect("client"); + client.settings().timing().logout_timeout = 250; + let mut manager = NetworkManager::new(client).expect("manager"); + manager.set_circuit_code(0x9933); + manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + server.endpoint, + 600, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + let _ = server.reports.recv_timeout(Duration::from_secs(1)).unwrap(); + let _ = server.reports.recv_timeout(Duration::from_secs(1)).unwrap(); + let order = Arc::new(Mutex::new(Vec::new())); + let disconnected_order = Arc::clone(&order); + let _disconnected = manager.subscribe_disconnected(Arc::new(move |_| { + disconnected_order.lock().unwrap().push("disconnected"); + })); + let logged_order = Arc::clone(&order); + let (logged_sender, logged_receiver) = channel(); + let _logged = manager.subscribe_logged_out(Arc::new(move |args| { + logged_order.lock().unwrap().push("logged_out"); + logged_sender.send(args.inventory_items).unwrap(); + })); + + let started = Instant::now(); + manager.begin_logout().unwrap(); + assert!(started.elapsed() < Duration::from_millis(100)); + assert_eq!( + logged_receiver + .recv_timeout(Duration::from_secs(1)) + .unwrap(), + Vec::::new() + ); + assert_eq!(*order.lock().unwrap(), ["disconnected", "logged_out"]); + assert!(!manager.connected()); + assert!(manager.simulators.is_empty()); +} + +#[test] +fn explicit_shutdown_cancels_a_pending_begin_logout_without_retaining_client() { + let _network_guard = network_test_guard(); + let server = spawn_fake_server(); + let mut client = GridClient::new().expect("client"); + client.settings().timing().logout_timeout = 5_000; + let mut manager = NetworkManager::new(client.clone()).expect("manager"); + manager.set_circuit_code(0x9944); + manager + .connect_with_ip_end_point_u_int64_boolean_uri_u_int32_u_int32( + server.endpoint, + 700, + true, + None, + 256, + 256, + ) + .unwrap() + .unwrap(); + let _ = server.reports.recv_timeout(Duration::from_secs(1)).unwrap(); + let _ = server.reports.recv_timeout(Duration::from_secs(1)).unwrap(); + + server.commands.send(ServerCommand::ExpectLogout).unwrap(); + manager.begin_logout().unwrap(); + assert_eq!( + server.reports.recv_timeout(Duration::from_secs(1)).unwrap(), + ServerReport::Logout(UUID::zero(), UUID::zero()) + ); + assert_eq!(manager.pending_logout_tasks(), 1); + manager + .shutdown_with_disconnect_type(NetworkManagerDisconnectType::ClientInitiated) + .unwrap(); + wait_until(|| manager.pending_logout_tasks() == 0); + assert_eq!( + manager.last_logout_error(), + Some(libremetaverse::Error::Cancelled) + ); + assert!(!manager.connected()); + assert!(manager.simulators.is_empty()); +} + #[test] fn async_connect_and_shutdown_do_not_require_an_ambient_tokio_runtime() { + let _network_guard = network_test_guard(); let server = spawn_fake_server(); let client = GridClient::new().expect("client"); let mut manager = NetworkManager::new(client).expect("manager"); @@ -911,6 +1456,7 @@ fn async_connect_and_shutdown_do_not_require_an_ambient_tokio_runtime() { #[test] fn handshake_timeout_preserves_reference_connect_result_and_shutdown_tears_down_current_sim() { + let _network_guard = network_test_guard(); let (endpoint, server) = spawn_ack_only_server(); let mut client = GridClient::new().expect("client"); client.settings().timing().login_timeout = 25; @@ -993,6 +1539,7 @@ fn successful_login_response(endpoint: SocketAddr, seed: &str) -> OSD { #[tokio::test(flavor = "current_thread")] async fn login_hashes_credentials_parses_response_and_hands_off_session_to_udp_and_caps() { + let _network_guard = network_test_guard(); let server = spawn_fake_server(); let requests = Arc::new(Mutex::new(Vec::::new())); let recorded = Arc::clone(&requests); diff --git a/docs/network-manager.md b/docs/network-manager.md index f904e8c..09d6c19 100644 --- a/docs/network-manager.md +++ b/docs/network-manager.md @@ -66,6 +66,35 @@ the last simulator reports `SimShutdown`. Explicit shutdown preserves the caller-provided reason/message and sends `CloseCircuit` only for client or network-timeout shutdowns. +Promoting an already-connected simulator is a controlled handoff. The manager +first sends a reliable `UseCircuitCode` followed by the agent/session/circuit +identified `CompleteAgentMovement`, installs the new seed capability, swaps +`CurrentSim`, and finally raises `SimChanged` with the previous simulator. The +old simulator remains tracked until an explicit disconnect so multi-simulator +traffic can continue. A completed shutdown leaves the manager reusable: a later +connection creates fresh workers and transports rather than retaining canceled +ones. Connection attempts made while logout or teardown is still active fail +with `InvalidOperation`, preventing a reentrant reconnect from racing resource +cleanup. + +`RequestLogout`, blocking `Logout`, runtime-neutral `LogoutAsync`, and +nonblocking `BeginLogout` all send the native `LogoutRequest` packet at most +once per active handshake. A matching `LogoutReply` is accepted only for the +current agent and session, raises `LoggedOut` with the returned inventory IDs, +and then performs a client-initiated shutdown. Blocking calls wait for that +ordered teardown; asynchronous calls additionally observe caller cancellation. +`BeginLogout` raises an empty `LoggedOut` only after a bounded network-timeout +teardown when no reply arrives, matching the reference ordering. + +Teardown is idempotent and ordered. It cancels active login and logout waits, +disconnects non-current simulators before the current simulator, stops manager +workers, clears connection/circuit/session/seed state and parsed login secrets, +then raises one `Disconnected` event and marks shutdown complete. Reentrant and +repeated shutdown calls do not repeat events. Nonblocking logout waiters observe +the lifecycle cancellation and terminate rather than retaining the manager or +client; their last detached error remains available through +`NetworkManager::last_logout_error` for diagnostics. + ## Login behavior The native login path constructs the reference LLSD request, including the C# @@ -92,6 +121,8 @@ cargo test -p libremetaverse --test network_manager They cover login hashing, redirects, cancellation and initial handoff alongside ACK-gated circuit setup, typed handshake reply/state, ping response, CAPS enable, -UDP disable, current/seed selection, callback ordering and filtering, RAII -unregistration, runtime-neutral async calls, concurrent registration/removal, -reentrant/concurrent disconnect, and keepalive timeout reasons. +UDP disable, current/seed selection, controlled region handoff, reconnect after +teardown, all three logout wait modes, reply validation and event order, callback +ordering and filtering, RAII unregistration, runtime-neutral async calls, +concurrent registration/removal, reentrant/concurrent disconnect, and keepalive +timeout reasons. diff --git a/tools/generate_api_shims.py b/tools/generate_api_shims.py index 9f2dbb7..b533e3e 100644 --- a/tools/generate_api_shims.py +++ b/tools/generate_api_shims.py @@ -58,6 +58,7 @@ NATIVE_TYPES = { "T:LibreMetaverse.LoginResponseData": "crate::login::LoginResponseData", "T:LibreMetaverse.LoginResponseData.ParseMessage": "crate::login::LoginResponseDataParseMessage", "T:LibreMetaverse.LoginResponseData.ParseResult": "crate::login::LoginResponseDataParseResult", + "T:LibreMetaverse.LoggedOutEventArgs": "crate::network_manager::LoggedOutEventArgs", "T:LibreMetaverse.MediaEntry": "crate::message_decoder::MediaEntry", "T:LibreMetaverse.Caps.EventQueueCallback": "crate::network_manager::CapsEventQueueCallback", "T:LibreMetaverse.CapsEventDictionary": "crate::network_manager::CapsEventDictionary", @@ -226,6 +227,8 @@ NATIVE_MEMBER_BODIES = { "self.native_subscribe_event_queue_running(handler)", "E:LibreMetaverse.NetworkManager.GenericStreamingMessage": "self.native_subscribe_generic_streaming_message(handler)", + "E:LibreMetaverse.NetworkManager.LoggedOut": + "self.native_subscribe_logged_out(handler)", "E:LibreMetaverse.NetworkManager.LoginProgress": "self.native_subscribe_login_progress(handler)", "E:LibreMetaverse.NetworkManager.PacketSent": @@ -244,6 +247,8 @@ NATIVE_MEMBER_BODIES = { "self.native_abort_login()", "M:LibreMetaverse.NetworkManager.BeginLogin(LibreMetaverse.LoginParams)": "self.native_begin_login(login_params)", + "M:LibreMetaverse.NetworkManager.BeginLogout": + "self.native_begin_logout()", "M:LibreMetaverse.NetworkManager.Connect(System.Net.IPAddress,System.UInt16,System.UInt64,System.Boolean,System.Uri)": "self.native_connect(std::net::SocketAddr::new(ip, port), handle, set_default, seedcaps, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_X, crate::network_manager::Simulator::DEFAULT_REGION_SIZE_Y)", "M:LibreMetaverse.NetworkManager.Connect(System.Net.IPAddress,System.UInt16,System.UInt64,System.Boolean,System.Uri,System.UInt32,System.UInt32)": @@ -310,10 +315,16 @@ NATIVE_MEMBER_BODIES = { "self.native_login_response_async(login_params, cancellation_token).await", "M:LibreMetaverse.NetworkManager.LoginWithResponseAsync(System.String,System.String,System.String,System.String,System.String,System.String,System.Threading.CancellationToken)": "self.native_login_strings_response(first_name, last_name, password, channel, start, version, cancellation_token).await", + "M:LibreMetaverse.NetworkManager.Logout": + "self.native_logout()", + "M:LibreMetaverse.NetworkManager.LogoutAsync(System.Threading.CancellationToken)": + "self.native_logout_async(cancellation_token).await", "M:LibreMetaverse.NetworkManager.RegisterLoginResponseCallback(LibreMetaverse.NetworkManager.LoginResponseCallback)": "self.native_register_login_response_callback(callback, None)", "M:LibreMetaverse.NetworkManager.RegisterLoginResponseCallback(LibreMetaverse.NetworkManager.LoginResponseCallback,System.String[])": "self.native_register_login_response_callback(callback, options)", + "M:LibreMetaverse.NetworkManager.RequestLogout": + "self.native_request_logout()", "M:LibreMetaverse.NetworkManager.StartLocation(System.String,System.Int32,System.Int32,System.Int32)": "Ok(crate::login::start_location(&sim, x, y, z))", "M:LibreMetaverse.NetworkManager.UnregisterLoginResponseCallback(LibreMetaverse.NetworkManager.LoginResponseCallback)":