//! Seed-capability discovery and capability event routing. #![allow(clippy::missing_errors_doc)] // Public Result shapes are fixed by the compatibility map. #![allow(clippy::needless_pass_by_value)] // Mapped APIs preserve owned CLR argument shapes. use crate::event_queue::{ EventQueueClient, EventQueueClientConnectedCallback, EventQueueClientEventCallback, }; use crate::packets::Packet; use crate::{Error, Simulator}; use libremetaverse_structured_data::{OSD, OSDArray, OSDFormat, OSDMap, OSDParser}; use libremetaverse_types::compat::{ CancellationToken, CancellationTokenSource, EventHandler, Subscription, Uri, }; use std::collections::HashMap; use std::fmt; use std::panic::{AssertUnwindSafe, catch_unwind}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex, RwLock, Weak}; use std::thread::{self, JoinHandle}; use std::time::Duration; const INITIAL_SEED_RETRY: Duration = Duration::from_secs(1); const MAXIMUM_SEED_RETRY: Duration = Duration::from_secs(30); 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) } struct CapabilitiesEventState { next_id: AtomicU64, handlers: Mutex)>>, } impl Default for CapabilitiesEventState { fn default() -> Self { Self { next_id: AtomicU64::new(1), handlers: Mutex::new(Vec::new()), } } } #[derive(Clone)] pub struct CapabilitiesReceivedEventArgs { simulator: Simulator, } impl CapabilitiesReceivedEventArgs { pub const fn new(simulator: Simulator) -> Result { Ok(Self { simulator }) } #[must_use] pub fn simulator(&self) -> Simulator { self.simulator.clone() } } struct SeedTask { cancellation: CancellationTokenSource, handle: JoinHandle<()>, } pub(crate) struct CapsInner { simulator: Weak, seed_uri: Uri, caps: RwLock>, event_queue: Mutex>, seed_task: Mutex>, events: Arc, seed_request_completed: AtomicBool, disconnected: AtomicBool, } impl CapsInner { fn simulator(&self) -> Option { Simulator::native_from_weak(&self.simulator) } pub(crate) fn disconnect(&self, immediate: bool) { if self.disconnected.swap(true, Ordering::AcqRel) { return; } if let Some(task) = mutex(&self.seed_task).take() { task.cancellation.cancel(); if task.handle.thread().id() != thread::current().id() { let _ = task.handle.join(); } } if let Some(queue) = mutex(&self.event_queue).take() { let _ = queue.stop(immediate); let _ = queue.dispose(); } } fn emit_capabilities_received(&self, simulator: Simulator) { let handlers: Vec<_> = mutex(&self.events.handlers) .iter() .map(|(_, handler)| Arc::clone(handler)) .collect(); for handler in handlers { let args = CapabilitiesReceivedEventArgs { simulator: simulator.clone(), }; let _ = catch_unwind(AssertUnwindSafe(|| handler(args))); } } fn route_event(&self, event_name: String, body: OSDMap) { let Some(simulator) = self.simulator() else { return; }; match crate::message_decoder::decode_event(&event_name, &body) { Ok(Some(message)) => { simulator.native_dispatch_caps_event(&event_name, message.as_ref()); } Ok(None) | Err(_) => { if let Ok(Some(packet)) = Packet::build_packet_with_string_osd_map(event_name, body) { let _ = simulator.native_enqueue_caps_packet(packet); } } } } } impl Drop for CapsInner { fn drop(&mut self) { self.disconnected.store(true, Ordering::Release); if let Some(task) = self .seed_task .get_mut() .unwrap_or_else(std::sync::PoisonError::into_inner) .take() { task.cancellation.cancel(); if task.handle.thread().id() != thread::current().id() { let _ = task.handle.join(); } } if let Some(queue) = self .event_queue .get_mut() .unwrap_or_else(std::sync::PoisonError::into_inner) .take() { let _ = queue.dispose(); } } } /// Capability collection associated with a simulator. pub struct Caps { pub simulator: Box, inner: Arc, } impl Clone for Caps { fn clone(&self) -> Self { Self { simulator: Box::new(self.simulator.native_clone_without_caps()), inner: Arc::clone(&self.inner), } } } impl Caps { pub(crate) fn native_create( simulator: &Simulator, seed_uri: Uri, ) -> Result, Error> { if !is_http_uri(&seed_uri) { return Err(Error::Argument); } let inner = Arc::new(CapsInner { simulator: simulator.native_data_weak(), seed_uri, caps: RwLock::new(HashMap::new()), event_queue: Mutex::new(None), seed_task: Mutex::new(None), events: Arc::new(CapabilitiesEventState::default()), seed_request_completed: AtomicBool::new(false), disconnected: AtomicBool::new(false), }); start_seed_task(&inner)?; Ok(inner) } pub(crate) fn native_from_inner(inner: Arc, simulator: Simulator) -> Self { Self { simulator: Box::new(simulator.native_clone_without_caps()), inner, } } pub fn subscribe_capabilities_received( &self, handler: Option>, ) -> Subscription { let Some(handler) = handler else { return Subscription::detached(); }; let id = self.inner.events.next_id.fetch_add(1, Ordering::Relaxed); mutex(&self.inner.events.handlers).push((id, handler)); let state = Arc::downgrade(&self.inner.events); Subscription::new(move || { if let Some(state) = state.upgrade() { mutex(&state.handlers).retain(|(candidate, _)| *candidate != id); } }) } #[must_use] /// Returns the fixed capability names requested from a seed capability. /// /// # Panics /// /// Panics only if constructing an in-memory LLSD array fails. pub fn all_capabilities() -> OSDArray { OSDArray::new_with_list( ALL_CAPABILITIES .iter() .map(|name| OSD::String((*name).to_owned())) .collect(), ) .expect("static capability list is valid") } pub fn capabilities(&self) -> Result, Error> { let mut names: Vec<_> = read(&self.inner.caps).keys().cloned().collect(); names.sort_unstable(); Ok(names) } pub fn capability_name_from_uri(&self, cap: Uri) -> Result { read(&self.inner.caps) .iter() .find_map(|(name, uri)| (uri == &cap).then(|| name.clone())) .ok_or(Error::InvalidOperation) } pub fn capability_uri(&self, capability: String) -> Result, Error> { Ok(read(&self.inner.caps).get(&capability).cloned()) } pub(crate) fn seed_request_finished(&self) -> bool { self.inner.seed_request_completed.load(Ordering::Acquire) } pub fn disconnect(&self, immediate: bool) -> Result<(), Error> { self.inner.disconnect(immediate); Ok(()) } pub fn get_mesh_cap_uri(&self) -> Result, Error> { let caps = read(&self.inner.caps); Ok(caps .get("ViewerAsset") .or_else(|| caps.get("GetMesh2")) .or_else(|| caps.get("GetMesh")) .cloned()) } pub fn get_texture_cap_uri(&self) -> Result, Error> { let caps = read(&self.inner.caps); Ok(caps .get("ViewerAsset") .or_else(|| caps.get("GetTexture")) .cloned()) } #[must_use] pub fn event_queue(&self) -> Option { mutex(&self.inner.event_queue).clone() } #[must_use] pub fn is_event_queue_running(&self) -> bool { mutex(&self.inner.event_queue) .as_ref() .is_some_and(EventQueueClient::running) } #[must_use] pub fn seed_caps_uri(&self) -> Uri { self.inner.seed_uri.clone() } } impl fmt::Debug for Caps { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { formatter .debug_struct("Caps") .field("capability_count", &read(&self.inner.caps).len()) .field("event_queue_running", &self.is_event_queue_running()) .field( "disconnected", &self.inner.disconnected.load(Ordering::Acquire), ) .finish_non_exhaustive() } } fn start_seed_task(inner: &Arc) -> Result<(), Error> { let cancellation = CancellationTokenSource::new(); let token = cancellation.token(); let weak = Arc::downgrade(inner); let (ready_sender, ready_receiver) = std::sync::mpsc::sync_channel(1); let handle = thread::Builder::new() .name("libremetaverse-seed-caps".to_owned()) .spawn(move || { let runtime = tokio::runtime::Builder::new_current_thread() .enable_all() .build(); let Ok(runtime) = runtime else { let _ = ready_sender.send(false); return; }; let _ = ready_sender.send(true); runtime.block_on(run_seed_requests(weak, token)); }) .map_err(|_| Error::InvalidOperation)?; if ready_receiver.recv_timeout(Duration::from_secs(2)) != Ok(true) { cancellation.cancel(); let _ = handle.join(); return Err(Error::InvalidOperation); } *mutex(&inner.seed_task) = Some(SeedTask { cancellation, handle, }); Ok(()) } async fn run_seed_requests(weak: Weak, cancellation: CancellationToken) { let mut retry = 0_u32; loop { let Some(inner) = weak.upgrade() else { return; }; if cancellation.is_cancellation_requested() || inner.disconnected.load(Ordering::Acquire) { return; } let Some(simulator) = inner.simulator() else { return; }; // Seed discovery is HTTP-only. It is valid for composed/fake clients // to install a simulator and query capabilities before a UDP circuit // is started, and tying this request to UDP connectivity creates an // unnecessary race for every capability-backed manager. let http = simulator.client.native_http_caps_client(); let response = http .post_with_uri_osd_format_osd_cancellation_token_i_progress( inner.seed_uri.clone(), OSDFormat::Xml, OSD::Array(Caps::all_capabilities().snapshot()), cancellation.clone(), None, ) .await; match response { Ok((response, data)) if response.is_success_status_code() => { if install_seed_response(&inner, &simulator, data, cancellation.clone()) .await .is_ok() { finish_seed_request(&inner, simulator); return; } } Ok((response, _)) if response.status_code == 404 => { finish_seed_request(&inner, simulator); return; } Err(Error::Cancelled) if cancellation.is_cancellation_requested() => return, Ok(_) | Err(_) => {} } retry = retry.saturating_add(1); let delay = seed_retry_delay(&inner.seed_uri, retry); tokio::select! { () = tokio::time::sleep(delay) => {} () = cancellation.cancelled() => return, } } } async fn install_seed_response( inner: &Arc, simulator: &Simulator, data: Vec, cancellation: CancellationToken, ) -> Result<(), Error> { let OSD::Map(map) = OSDParser::deserialize_with_bytes(data)? else { return Err(Error::Parse { position: 0, context: "seed capability response is not a map", }); }; let mut discovered = HashMap::with_capacity(map.len()); for (name, value) in map { if let Some(uri) = value.as_uri()? { if !is_http_uri(&uri) { continue; } simulator .client .native_caps_rate_limiter() .register_cap_uri(name.clone(), uri.clone())?; discovered.insert(name, uri); } } *write(&inner.caps) = discovered; if let Some(uri) = read(&inner.caps).get("EventQueueGet").cloned() { let mut queue = EventQueueClient::new_for_caps(uri, simulator.native_clone_without_caps())?; let connected_inner = Arc::downgrade(inner); queue.on_connected = Some(EventQueueClientConnectedCallback::from_handler(move || { if let Some(inner) = connected_inner.upgrade() && let Some(simulator) = inner.simulator() { simulator.native_raise_event_queue_running(); } })); let event_inner = Arc::downgrade(inner); queue.on_event = Some(EventQueueClientEventCallback::from_handler( move |event_name, body| { if let Some(inner) = event_inner.upgrade() { inner.route_event(event_name, body); } }, )); queue.start()?; *mutex(&inner.event_queue) = Some(queue); } // SimulatorFeatures is a GET capability, not merely an event-queue // message. Load its initial snapshot before announcing that capability // discovery is complete, while still allowing grids that omit or reject // the optional endpoint to finish seed discovery normally. let simulator_features_uri = read(&inner.caps).get("SimulatorFeatures").cloned(); if let Some(uri) = simulator_features_uri && let Ok((response, bytes)) = simulator .client .native_http_caps_client() .get(uri, cancellation.clone(), None) .await && response.is_success_status_code() && bytes.len() <= 8 * 1024 * 1024 { let _ = simulator.features.set_features(None, Some(bytes), None); } cancellation.throw_if_cancellation_requested()?; Ok(()) } fn finish_seed_request(inner: &CapsInner, simulator: Simulator) { if !inner.seed_request_completed.swap(true, Ordering::AcqRel) { inner.emit_capabilities_received(simulator); } } fn seed_retry_delay(uri: &Uri, retry: u32) -> Duration { let multiplier = 1_u32 .checked_shl(retry.saturating_sub(1).min(30)) .unwrap_or(u32::MAX); let base = INITIAL_SEED_RETRY .saturating_mul(multiplier) .min(MAXIMUM_SEED_RETRY); let mut hash = 0xcbf2_9ce4_8422_2325_u64 ^ u64::from(retry); for byte in uri.0.as_bytes() { hash ^= u64::from(*byte); hash = hash.wrapping_mul(0x100_0000_01b3); } let max_jitter = base / 8; if max_jitter.is_zero() { return base; } let nanos = u64::try_from(max_jitter.as_nanos()) .unwrap_or(u64::MAX) .saturating_add(1); base.saturating_add(Duration::from_nanos(hash % nanos)) } fn is_http_uri(uri: &Uri) -> bool { reqwest::Url::parse(&uri.0) .ok() .is_some_and(|uri| matches!(uri.scheme(), "http" | "https") && uri.host().is_some()) } const ALL_CAPABILITIES: &[&str] = &[ "AbuseCategories", "AcceptFriendship", "AcceptGroupInvite", "AgentPreferences", "AgentProfile", "AgentState", "AttachmentResources", "AvatarPickerSearch", "AvatarRenderInfo", "CharacterProperties", "ChatSessionRequest", "CopyInventoryFromNotecard", "CreateInventoryCategory", "DeclineFriendship", "DeclineGroupInvite", "DispatchRegionInfo", "DirectDelivery", "EnvironmentSettings", "EstateAccess", "EstateChangeInfo", "EventQueueGet", "ExtEnvironment", "FetchLib2", "FetchLibDescendents2", "FetchInventory2", "FetchInventoryDescendents2", "IncrementCOFVersion", "RequestTaskInventory", "InterestList", "InventoryThumbnailUpload", "GetDisplayNames", "GetExperiences", "AgentExperiences", "FindExperienceByName", "GetExperienceInfo", "GetAdminExperiences", "GetCreatorExperiences", "ExperiencePreferences", "GroupExperiences", "UpdateExperience", "IsExperienceAdmin", "IsExperienceContributor", "RegionExperiences", "ExperienceQuery", "GetMesh", "GetMesh2", "GetMetadata", "GetObjectCost", "GetObjectPhysicsData", "GetTexture", "GroupAPIv1", "GroupMemberData", "GroupProposalBallot", "HomeLocation", "LandResources", "LSLSyntax", "MapLayer", "MapLayerGod", "MeshUploadFlag", "ModifyMaterialParams", "ModifyRegion", "NavMeshGenerationStatus", "NewFileAgentInventory", "NewFileAgentInventoryVariablePrice", "ObjectAnimation", "ObjectMedia", "ObjectMediaNavigate", "ObjectNavMeshProperties", "ParcelPropertiesUpdate", "ParcelVoiceInfoRequest", "ProductInfoRequest", "ProvisionVoiceAccountRequest", "ReadOfflineMsgs", "RegionObjects", "RegionSchedule", "RemoteParcelRequest", "RenderMaterials", "RequestTextureDownload", "ResourceCostSelected", "RetrieveNavMeshSrc", "SearchStatRequest", "SearchStatTracking", "SendPostcard", "SendUserReport", "SendUserReportWithScreenshot", "ServerReleaseNotes", "SetDisplayName", "SimConsoleAsync", "SimulatorFeatures", "StartGroupProposal", "TerrainNavMeshProperties", "TextureStats", "UntrustedSimulatorMessage", "UpdateAgentInformation", "UpdateAgentLanguage", "UpdateAvatarAppearance", "UpdateGestureAgentInventory", "UpdateGestureTaskInventory", "UpdateNotecardAgentInventory", "UpdateNotecardTaskInventory", "UpdateScriptAgent", "UpdateScriptTask", "UpdateSettingsAgentInventory", "UpdateSettingsTaskInventory", "UploadAgentProfileImage", "UpdateMaterialAgentInventory", "UpdateMaterialTaskInventory", "UploadBakedTexture", "UserInfo", "ViewerAsset", "ViewerBenefits", "ViewerMetrics", "ViewerStartAuction", "ViewerStats", "VoiceSignalingRequest", "InventoryAPIv3", "LibraryAPIv3", ]; #[cfg(test)] #[allow(clippy::too_many_lines)] // The seed exchange is one contiguous protocol scenario. mod tests { use super::*; use crate::{GridClient, HttpCapsClient}; use libremetaverse_types::compat::{HttpMessageHandler, HttpRequest, HttpResponse}; use std::collections::BTreeMap; use std::sync::mpsc; use std::time::Instant; fn lock(value: &Mutex) -> std::sync::MutexGuard<'_, T> { value .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) } fn response(status_code: u16, body: Vec) -> HttpResponse { HttpResponse { status_code, headers: BTreeMap::new(), content_type: Some(HttpCapsClient::LLSD_XML.to_owned()), body, } } fn seed_body() -> Vec { OSDParser::serialize_llsd_xml_bytes(OSD::Map(HashMap::from([ ( "EventQueueGet".to_owned(), OSD::Uri(Uri( "https://caps.example.test/queue?token=eq-secret".to_owned() )), ), ( "ViewerAsset".to_owned(), OSD::Uri(Uri("https://asset.example.test/viewer".to_owned())), ), ( "GetMesh2".to_owned(), OSD::Uri(Uri("https://asset.example.test/mesh2".to_owned())), ), ( "GetTexture".to_owned(), OSD::Uri(Uri("https://asset.example.test/texture".to_owned())), ), ( "SimulatorFeatures".to_owned(), OSD::Uri(Uri("https://caps.example.test/features".to_owned())), ), ( "InvalidScheme".to_owned(), OSD::Uri(Uri("file:///private/capability".to_owned())), ), ]))) .unwrap() } #[test] fn seed_discovery_preserves_request_shape_registers_caps_and_starts_queue() { let requests = Arc::new(Mutex::new(Vec::::new())); let recorded = Arc::clone(&requests); let release_seed = Arc::new(AtomicBool::new(false)); let release = Arc::clone(&release_seed); let handler = HttpMessageHandler::new(move |request, _cancellation| { let recorded = Arc::clone(&recorded); let release = Arc::clone(&release); async move { let is_seed = request.uri.0.contains("/seed"); let is_features = request.uri.0.contains("/features"); lock(&recorded).push(request); if is_seed { while !release.load(Ordering::Acquire) { tokio::time::sleep(Duration::from_millis(1)).await; } response(200, seed_body()) } else if is_features { response( 200, OSDParser::serialize_llsd_xml_bytes(OSD::Map(HashMap::from([( "MeshRezEnabled".to_owned(), OSD::Boolean(true), )]))) .unwrap(), ) } else { response(404, Vec::new()) } } }); let mut client = GridClient::new().unwrap(); client.set_http_caps_client(HttpCapsClient::new(handler).unwrap()); let simulator = Simulator::new(client, "127.0.0.1:13001".parse().unwrap(), 2, None, None).unwrap(); simulator.native_set_connected_for_tests(true); simulator .native_set_seed_caps( Some(Uri( "https://caps.example.test/seed?token=seed-secret".to_owned() )), Some(false), ) .unwrap(); let caps = simulator.clone().caps.expect("caps snapshot"); let (sender, receiver) = mpsc::channel(); let _subscription = caps.subscribe_capabilities_received(Some(Arc::new(move |args| { sender.send(args.simulator().handle).unwrap(); }))); release_seed.store(true, Ordering::Release); assert_eq!(receiver.recv_timeout(Duration::from_secs(2)).unwrap(), 2); assert_eq!( caps.seed_caps_uri(), Uri("https://caps.example.test/seed?token=seed-secret".to_owned()) ); assert_eq!( caps.get_texture_cap_uri().unwrap(), Some(Uri("https://asset.example.test/viewer".to_owned())) ); assert_eq!( caps.get_mesh_cap_uri().unwrap(), Some(Uri("https://asset.example.test/viewer".to_owned())) ); assert_eq!( caps.capability_name_from_uri(Uri("https://asset.example.test/mesh2".to_owned())) .unwrap(), "GetMesh2" ); assert_eq!( caps.capability_uri("InvalidScheme".to_owned()).unwrap(), None ); assert!(caps.event_queue().is_some()); assert_eq!( simulator.features.get("MeshRezEnabled".to_owned()).unwrap(), Some(OSD::Boolean(true)) ); let deadline = Instant::now() + Duration::from_secs(2); while lock(&requests).len() < 2 && Instant::now() < deadline { thread::sleep(Duration::from_millis(5)); } let requests = lock(&requests); let OSD::Array(requested) = OSDParser::deserialize_llsd_xml_with_bytes(requests[0].body.clone()).unwrap() else { panic!("capability request array") }; let requested: Vec<_> = requested .iter() .map(|value| value.as_string().unwrap()) .collect(); assert_eq!( requested, ALL_CAPABILITIES .iter() .map(|value| (*value).to_owned()) .collect::>() ); assert_eq!( requests[0].content_type.as_deref(), Some(HttpCapsClient::LLSD_XML) ); drop(requests); caps.disconnect(true).unwrap(); } #[test] fn all_capabilities_matches_reference_and_debug_redacts_seed_uri() { let values: Vec<_> = Caps::all_capabilities() .snapshot() .into_iter() .map(|value| value.as_string().unwrap()) .collect(); assert_eq!( values, ALL_CAPABILITIES .iter() .map(|value| (*value).to_owned()) .collect::>() ); assert!(values.contains(&"EventQueueGet".to_owned())); assert!(values.contains(&"InventoryAPIv3".to_owned())); let handler = HttpMessageHandler::new(|_request, cancellation| async move { cancellation.cancelled().await; response(499, Vec::new()) }); let mut client = GridClient::new().unwrap(); client.set_http_caps_client(HttpCapsClient::new(handler).unwrap()); let simulator = Simulator::new(client, "127.0.0.1:13002".parse().unwrap(), 3, None, None).unwrap(); simulator.native_set_connected_for_tests(true); simulator .native_set_seed_caps( Some(Uri( "https://caps.example.test/seed?token=do-not-log".to_owned() )), Some(false), ) .unwrap(); let caps = simulator.clone().caps.unwrap(); let debug = format!("{caps:?}"); assert!(!debug.contains("do-not-log")); caps.disconnect(true).unwrap(); } }