Implement network lifecycle teardown and logout (#57)
This commit is contained in:
@@ -16206,18 +16206,9 @@ impl LoadUrlEventArgs {
|
||||
}
|
||||
|
||||
/// C# type: `T:LibreMetaverse.LoggedOutEventArgs`.
|
||||
pub struct LoggedOutEventArgs {
|
||||
/// C# member: `F:LibreMetaverse.LoggedOutEventArgs.InventoryItems`.
|
||||
pub inventory_items: Vec<libremetaverse_types::UUID>,
|
||||
}
|
||||
impl LoggedOutEventArgs {
|
||||
/// C# member: `M:LibreMetaverse.LoggedOutEventArgs.#ctor(System.Collections.Generic.List{LibreMetaverse.UUID})`.
|
||||
pub fn new(inventory_items: Vec<libremetaverse_types::UUID>) -> Result<Self, crate::Error> {
|
||||
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::LoggedOutEventArgs>,
|
||||
) -> 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<libremetaverse_types::compat::CancellationToken>,
|
||||
) -> 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(
|
||||
|
||||
@@ -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<Uri, Error> {
|
||||
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<Simulator>,
|
||||
}
|
||||
|
||||
/// Inventory item identifiers returned by a successful logout handshake.
|
||||
#[derive(Clone, Debug, Eq, PartialEq)]
|
||||
pub struct LoggedOutEventArgs {
|
||||
pub inventory_items: Vec<UUID>,
|
||||
}
|
||||
|
||||
impl LoggedOutEventArgs {
|
||||
pub const fn new(inventory_items: Vec<UUID>) -> Result<Self, Error> {
|
||||
Ok(Self { inventory_items })
|
||||
}
|
||||
}
|
||||
|
||||
impl SimChangedEventArgs {
|
||||
pub fn new(previous_simulator: Option<Simulator>) -> Result<Self, Error> {
|
||||
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<bool, Error> {
|
||||
pub fn native_connect(&self, move_to_sim: bool) -> Result<bool, Error> {
|
||||
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<DisconnectedEventArgs>,
|
||||
event_queue_running: EventRegistry<EventQueueRunningEventArgs>,
|
||||
generic_streaming_message: EventRegistry<GenericStreamingMessageEventArgs>,
|
||||
logged_out: EventRegistry<LoggedOutEventArgs>,
|
||||
login_progress: EventRegistry<LoginProgressEventArgs>,
|
||||
packet_sent: EventRegistry<PacketSentEventArgs>,
|
||||
sim_changed: EventRegistry<SimChangedEventArgs>,
|
||||
@@ -1688,14 +1779,25 @@ pub(crate) struct NetworkManagerInner {
|
||||
login_seed_capability: RwLock<Option<Uri>>,
|
||||
login_response: LoginResponseData,
|
||||
login_runtime: Mutex<LoginRuntimeState>,
|
||||
logout_cancellation: Mutex<CancellationTokenSource>,
|
||||
inbox_count: AtomicI32,
|
||||
outbox_count: AtomicI32,
|
||||
workers: Mutex<Option<NetworkWorkers>>,
|
||||
disconnect_lock: Mutex<()>,
|
||||
shutting_down: AtomicBool,
|
||||
lifecycle_state: AtomicU8,
|
||||
logout_task_count: AtomicI32,
|
||||
last_logout_error: Mutex<Option<Error>>,
|
||||
}
|
||||
|
||||
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<Self>) -> 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<Uri>) -> 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<LoggedOutEventArgs>,
|
||||
) -> 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<Option<Simulator>, 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<CancellationToken>,
|
||||
) -> 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<Error> {
|
||||
*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]
|
||||
|
||||
Reference in New Issue
Block a user