//! Client-owned asset download, cache, and upload coordinator. #![allow( clippy::missing_errors_doc, clippy::too_many_arguments, clippy::too_many_lines, clippy::needless_pass_by_value, clippy::unused_self )] use crate::assets::{Asset, AssetMesh, AssetTexture}; use crate::download_manager::{DownloadManager, DownloadRequest}; use crate::network_manager::RawPacketReceivedEventArgs; use crate::packet_catalog::GeneratedPacket; use crate::packets::{ AssetUploadCompletePacket, AssetUploadRequestPacket, ConfirmXferPacketPacket, InitiateDownloadPacket, PacketType, RequestXferPacket, SendXferPacketPacket, TransferInfoPacket, TransferPacketPacket, TransferRequestPacket, }; use crate::{ AssetCache, AssetDownload, AssetUploadEventArgs, Error, EstateAssetType, GridClient, ImageReceiveProgressEventArgs, ImageType, InitiateDownloadEventArgs, InventoryItem, SourceType, StatusCode, XferDownload, XferReceivedEventArgs, }; use futures_channel::oneshot; use futures_util::{FutureExt, pin_mut, select_biased}; use libremetaverse_types::compat::{ CancellationToken, CancellationTokenSource, EventHandler, Subscription, Uri, }; use libremetaverse_types::{AssetType, UUID, Utils}; use serde_json::Value; use std::collections::{BTreeMap, HashMap}; use std::panic::{AssertUnwindSafe, catch_unwind}; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; fn block_on(future: F) -> F::Output { use std::task::{Context, Poll, Wake, Waker}; struct ThreadWake(std::thread::Thread); impl Wake for ThreadWake { fn wake(self: Arc) { self.0.unpark(); } } let waker = Waker::from(Arc::new(ThreadWake(std::thread::current()))); let mut context = Context::from_waker(&waker); let mut future = std::pin::pin!(future); loop { match future.as_mut().poll(&mut context) { Poll::Ready(value) => return value, Poll::Pending => std::thread::park(), } } } fn mutex(value: &Mutex) -> std::sync::MutexGuard<'_, T> { value .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) } struct EventSlot { next: AtomicU64, handlers: Mutex>>, } impl Default for EventSlot { fn default() -> Self { Self { next: AtomicU64::new(1), handlers: Mutex::new(HashMap::new()), } } } impl EventSlot { fn subscribe(self: &Arc, handler: EventHandler) -> Subscription { let id = self.next.fetch_add(1, Ordering::Relaxed); mutex(&self.handlers).insert(id, handler); let weak = Arc::downgrade(self); Subscription::new(move || { if let Some(slot) = weak.upgrade() { mutex(&slot.handlers).remove(&id); } }) } } impl EventSlot { fn emit(&self, value: T) { let handlers: Vec<_> = mutex(&self.handlers).values().cloned().collect(); for handler in handlers { let argument = value.clone(); let _ = catch_unwind(AssertUnwindSafe(|| handler(argument))); } } } #[cfg(test)] mod event_slot_tests { use super::*; #[test] fn dispatch_isolates_panics_and_does_not_retain_dropped_handlers() { let slot = Arc::new(EventSlot::::default()); let released = Arc::new(()); let released_probe = Arc::downgrade(&released); let panicking = slot.subscribe(Arc::new(|_| panic!("subscriber failure"))); let delivered = Arc::new(AtomicU64::new(0)); let delivered_handler = Arc::clone(&delivered); let releasing = Arc::clone(&released); let healthy = slot.subscribe(Arc::new(move |value| { std::hint::black_box(&releasing); delivered_handler.fetch_add(value as u64, Ordering::AcqRel); })); drop(released); slot.emit(1); assert_eq!(delivered.load(Ordering::Acquire), 1); assert!(released_probe.upgrade().is_some()); drop(panicking); drop(healthy); assert!(released_probe.upgrade().is_none()); slot.emit(1); assert_eq!(delivered.load(Ordering::Acquire), 1); } } struct UploadState { asset_id: UUID, data: Vec, packet_num: u32, transaction_id: UUID, transferred: usize, type_: AssetType, xfer_id: u64, confirmation: Option>, } struct XferState { download: XferDownload, chunks: HashMap>, expected_size: Option, final_packet: Option, } #[derive(Default)] struct AssetEvents { uploaded: Arc>, upload_progress: Arc>, image_progress: Arc>, initiate: Arc>, xfer: Arc>, } pub(crate) struct AssetManagerInner { client: crate::client_core::ClientWeakHandle, cache: AssetCache, downloads: DownloadManager, image_cancellations: Mutex>, pending_upload: Mutex>>>, xfer_uploads: Mutex>>>, xfer_downloads: Mutex>, raw_packet_subscription: Mutex>, events: AssetEvents, } pub struct AssetManager { pub cache: AssetCache, inner: Arc, } impl Clone for AssetManager { fn clone(&self) -> Self { Self { cache: self.cache.clone(), inner: Arc::clone(&self.inner), } } } impl std::fmt::Debug for AssetManager { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("AssetManager").finish_non_exhaustive() } } impl AssetManager { pub(crate) fn native_stage_estate_upload(&self, data: Vec) -> Result { if data.is_empty() || data.len() > crate::asset_models::MAX_ASSET_BYTES { return Err(Error::Argument); } let id = UUID::random()?; let upload = Arc::new(Mutex::new(UploadState { asset_id: id, data, packet_num: 0, transaction_id: id, transferred: 0, type_: AssetType::Unknown, xfer_id: 0, confirmation: None, })); let mut pending = mutex(&self.inner.pending_upload); if pending.is_some() { return Err(Error::InvalidOperation); } *pending = Some(upload); Ok(id) } pub(crate) fn native_cancel_staged_estate_upload(&self, id: UUID) { let mut pending = mutex(&self.inner.pending_upload); if pending .as_ref() .is_some_and(|upload| mutex(upload).transaction_id == id) { pending.take(); } } pub(crate) fn native_from_inner(inner: Arc) -> Result { inner.client.upgrade().ok_or(Error::InvalidOperation)?; Ok(Self { cache: inner.cache.clone(), inner, }) } pub(crate) fn native_inner(&self) -> Arc { Arc::clone(&self.inner) } pub fn new(client: Option) -> Result { let client = client.ok_or(Error::ArgumentNull)?; let cache = AssetCache::new(client.clone())?; let downloads = DownloadManager::new(client.clone())?; let inner_cache = cache.clone(); let manager = Self { cache, inner: Arc::new(AssetManagerInner { client: client.native_weak_handle(), cache: inner_cache, downloads, image_cancellations: Mutex::new(HashMap::new()), pending_upload: Mutex::new(None), xfer_uploads: Mutex::new(HashMap::new()), xfer_downloads: Mutex::new(HashMap::new()), raw_packet_subscription: Mutex::new(None), events: AssetEvents::default(), }), }; manager.install_packet_handler()?; Ok(manager) } fn install_packet_handler(&self) -> Result<(), Error> { let network = self.client()?.native_network()?; let weak = Arc::downgrade(&self.inner); let subscription = network.subscribe_raw_packet(Arc::new(move |event| { if let Some(inner) = weak.upgrade() { let _ = AssetManager::native_from_inner(inner) .and_then(|manager| manager.handle_raw_packet(event)); } })); *mutex(&self.inner.raw_packet_subscription) = Some(subscription); Ok(()) } fn decode_packet(bytes: &[u8]) -> Result { let mut packet = T::new_generated(); let mut position = 0_i32; let mut end = i32::try_from(bytes.len()).map_err(|_| Error::Argument)? - 1; let mut zero = vec![0; crate::udp_transport::UdpTransportConfig::default().max_decoded_packet_size]; packet.decode_from_bytes(bytes, &mut position, &mut end, Some(&mut zero))?; Ok(packet) } fn send_packet( &self, simulator: &crate::Simulator, packet: &T, packet_type: PacketType, ) -> Result<(), Error> { let bytes = packet.encode_packet()?; simulator.native_send_packet_data( bytes.clone(), i32::try_from(bytes.len()).map_err(|_| Error::Argument)?, packet_type, T::new_generated().generated_header().zerocoded, ) } fn handle_raw_packet(&self, event: RawPacketReceivedEventArgs) -> Result<(), Error> { match event.packet_type { PacketType::RequestXfer => { self.handle_request_xfer(Self::decode_packet(&event.data)?, event.simulator) } PacketType::ConfirmXferPacket => { self.handle_confirm_xfer(Self::decode_packet(&event.data)?, event.simulator) } PacketType::AssetUploadComplete => { self.handle_upload_complete(Self::decode_packet(&event.data)?) } PacketType::SendXferPacket => { self.handle_send_xfer(Self::decode_packet(&event.data)?, event.simulator) } PacketType::InitiateDownload => { self.handle_initiate_download(Self::decode_packet(&event.data)?) } _ => Ok(()), } } fn handle_initiate_download(&self, packet: InitiateDownloadPacket) -> Result<(), Error> { let event = InitiateDownloadEventArgs::new( Utils::bytes_to_string_with_bytes(packet.file_data.sim_filename)?, Utils::bytes_to_string_with_bytes(packet.file_data.viewer_filename)?, )?; self.inner.events.initiate.emit(event); Ok(()) } fn handle_send_xfer( &self, packet: SendXferPacketPacket, simulator: crate::Simulator, ) -> Result<(), Error> { let id = packet.xfer_id.id; let packet_number = packet.xfer_id.packet & 0x7fff_ffff; let is_final = packet.xfer_id.packet & 0x8000_0000 != 0; let mut data = packet.data_packet.data; let completed = { let mut downloads = mutex(&self.inner.xfer_downloads); let Some(state) = downloads.get_mut(&id) else { return Ok(()); }; if packet_number == 0 { if data.len() < 4 { return Err(Error::Argument); } let size = usize::try_from(u32::from_le_bytes( data[..4].try_into().map_err(|_| Error::Argument)?, )) .map_err(|_| Error::Argument)?; if size > crate::asset_models::MAX_ASSET_BYTES { return Err(Error::Argument); } state.expected_size = Some(size); data.drain(..4); } state.chunks.entry(packet_number).or_insert(data); if is_final { state.final_packet = Some(packet_number); } if let (Some(size), Some(final_packet)) = (state.expected_size, state.final_packet) { if (0..=final_packet).all(|number| state.chunks.contains_key(&number)) { let mut assembled = Vec::with_capacity(size); for number in 0..=final_packet { let chunk = state.chunks.get(&number).ok_or(Error::InvalidOperation)?; let remaining = size.saturating_sub(assembled.len()); assembled.extend_from_slice(&chunk[..chunk.len().min(remaining)]); } if assembled.len() != size { return Err(Error::Argument); } state.download.packet_num = final_packet; state.download.base.size = i32::try_from(size).map_err(|_| Error::Argument)?; state.download.base.transferred = state.download.base.size; state.download.base.asset_data = assembled; state.download.base.success = true; Some(state.download.clone()) } else { None } } else { None } }; let mut confirm = ConfirmXferPacketPacket::new_with_constructor()?; confirm.xfer_id.id = id; confirm.xfer_id.packet = packet.xfer_id.packet; self.send_packet(&simulator, &confirm, PacketType::ConfirmXferPacket)?; if let Some(download) = completed { mutex(&self.inner.xfer_downloads).remove(&id); self.inner .events .xfer .emit(XferReceivedEventArgs::new(download)?); } Ok(()) } fn handle_request_xfer( &self, packet: RequestXferPacket, simulator: crate::Simulator, ) -> Result<(), Error> { let Some(upload) = mutex(&self.inner.pending_upload).take() else { return Ok(()); }; { let mut state = mutex(&upload); state.xfer_id = packet.xfer_id.id; } mutex(&self.inner.xfer_uploads).insert(packet.xfer_id.id, Arc::clone(&upload)); self.send_next_upload(&upload, &simulator)?; Ok(()) } fn handle_confirm_xfer( &self, packet: ConfirmXferPacketPacket, simulator: crate::Simulator, ) -> Result<(), Error> { let Some(upload) = mutex(&self.inner.xfer_uploads) .get(&packet.xfer_id.id) .cloned() else { return Ok(()); }; let snapshot = { let state = mutex(&upload); if packet.xfer_id.packet & 0x7fff_ffff != state.packet_num.saturating_sub(1) { return Ok(()); } crate::AssetUpload { asset_id: state.asset_id, packet_num: state.packet_num, type_: state.type_, xfer_id: state.xfer_id, base: crate::Transfer { asset_data: Vec::new(), asset_type: state.type_, id: state.transaction_id, size: i32::try_from(state.data.len()).unwrap_or(i32::MAX), success: false, transferred: i32::try_from(state.transferred).unwrap_or(i32::MAX), ..crate::Transfer::new()? }, } }; self.inner .events .upload_progress .emit(crate::AssetUploadEventArgs::new(snapshot)?); if mutex(&upload).transferred < mutex(&upload).data.len() { self.send_next_upload(&upload, &simulator)?; } Ok(()) } fn handle_upload_complete(&self, packet: AssetUploadCompletePacket) -> Result<(), Error> { let mut found = mutex(&self.inner.pending_upload) .take() .filter(|upload| mutex(upload).asset_id == packet.asset_block.uuid); if found.is_none() { let id = mutex(&self.inner.xfer_uploads) .iter() .find_map(|(id, upload)| { (mutex(upload).asset_id == packet.asset_block.uuid).then_some(*id) }); if let Some(id) = id { found = mutex(&self.inner.xfer_uploads).remove(&id); } } if let Some(upload) = found { let mut state = mutex(&upload); let transaction = state.transaction_id; if let Some(sender) = state.confirmation.take() { let _ = sender.send(if packet.asset_block.success { transaction } else { UUID::zero() }); } let event = crate::AssetUpload { asset_id: state.asset_id, packet_num: state.packet_num, type_: state.type_, xfer_id: state.xfer_id, base: crate::Transfer { asset_data: Vec::new(), asset_type: state.type_, id: transaction, size: i32::try_from(state.data.len()).unwrap_or(i32::MAX), success: packet.asset_block.success, transferred: i32::try_from(state.transferred).unwrap_or(i32::MAX), ..crate::Transfer::new()? }, }; drop(state); self.inner .events .uploaded .emit(crate::AssetUploadEventArgs::new(event)?); } Ok(()) } fn send_next_upload( &self, upload: &Arc>, simulator: &crate::Simulator, ) -> Result<(), Error> { let mut state = mutex(upload); let packet_number = state.packet_num; let start = state.transferred; let payload_capacity = if packet_number == 0 { 996 } else { 1000 }; let end = start.saturating_add(payload_capacity).min(state.data.len()); if start >= end { return Ok(()); } let mut data = Vec::with_capacity((end - start) + usize::from(packet_number == 0)); if packet_number == 0 { data.extend_from_slice( &u32::try_from(state.data.len()) .map_err(|_| Error::Argument)? .to_le_bytes(), ); } data.extend_from_slice(&state.data[start..end]); let final_packet = end == state.data.len(); let xfer_id = state.xfer_id; state.packet_num = state.packet_num.saturating_add(1); state.transferred = end; drop(state); let mut packet = SendXferPacketPacket::new_with_constructor()?; packet.xfer_id.id = xfer_id; packet.xfer_id.packet = packet_number | if final_packet { 0x8000_0000 } else { 0 }; packet.data_packet.data = data; self.send_packet(simulator, &packet, PacketType::SendXferPacket) } fn client(&self) -> Result { self.inner.client.upgrade().ok_or(Error::InvalidOperation) } fn capability(&self, name: &str) -> Result, Error> { let Some(sim) = self.client()?.native_network()?.current_sim() else { return Ok(None); }; let Some(caps) = sim.native_caps() else { return Ok(None); }; let deadline = Instant::now() + Duration::from_secs(2); loop { if let Some(uri) = caps.capability_uri(name.to_owned())? { return Ok(Some(uri)); } if caps.seed_request_finished() || Instant::now() >= deadline { return Ok(None); } std::thread::sleep(Duration::from_millis(1)); } } fn with_query(base: &Uri, name: &str, id: UUID) -> Uri { Uri(format!( "{}{}{}={}", base.0, if base.0.contains('?') { '&' } else { '?' }, name, id )) } fn token(value: Option) -> CancellationToken { value.unwrap_or_default() } async fn fetch( &self, uri: Uri, content_type: Option, cancellation: CancellationToken, ) -> Result>, Error> { let mut request = DownloadRequest::new(uri, content_type, None)?; request.cancellation_token = cancellation; let (response, bytes) = self .inner .downloads .queue_download_with_download_request_ba727387(request) .await?; if (200..300).contains(&response.status_code) { if bytes.len() > crate::asset_models::MAX_ASSET_BYTES { return Err(Error::Argument); } Ok(Some(bytes)) } else { Ok(None) } } fn abort_transfer(&self, simulator: &crate::Simulator, transfer_id: UUID) { if let Ok(mut packet) = crate::packets::TransferAbortPacket::new_with_constructor() { packet.transfer_info.channel_type = crate::ChannelType::Asset as i32; packet.transfer_info.transfer_id = transfer_id; let _ = self.send_packet(simulator, &packet, PacketType::TransferAbort); } } async fn request_asset_udp( &self, asset_id: UUID, item_id: UUID, task_id: UUID, owner_id: UUID, asset_type: AssetType, priority: bool, source_type: SourceType, transfer_id: UUID, cancellation: CancellationToken, ) -> Result>, Error> { let simulator = self .client()? .native_network()? .current_sim() .ok_or(Error::InvalidOperation)?; let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel(); let subscription = self .client()? .native_network()? .subscribe_raw_packet(Arc::new(move |event| { if matches!( event.packet_type, PacketType::TransferInfo | PacketType::TransferPacket ) { let _ = sender.send(event); } })); let mut request = TransferRequestPacket::new_with_constructor()?; request.transfer_info.channel_type = crate::ChannelType::Asset as i32; request.transfer_info.priority = if priority { 101.0 } else { 100.0 }; request.transfer_info.source_type = source_type as i32; request.transfer_info.transfer_id = transfer_id; if source_type == SourceType::SimInventoryItem { let mut client = self.client()?; request .transfer_info .params .extend_from_slice(&client.native_self().agent_id().get_bytes()?); request .transfer_info .params .extend_from_slice(&client.native_self().session_id().get_bytes()?); request .transfer_info .params .extend_from_slice(&owner_id.get_bytes()?); request .transfer_info .params .extend_from_slice(&task_id.get_bytes()?); request .transfer_info .params .extend_from_slice(&item_id.get_bytes()?); } request .transfer_info .params .extend_from_slice(&asset_id.get_bytes()?); request .transfer_info .params .extend_from_slice(&(asset_type as i32).to_le_bytes()); self.send_packet(&simulator, &request, PacketType::TransferRequest)?; let mut expected_size = None; let mut chunks = BTreeMap::>::new(); let deadline = tokio::time::sleep(Duration::from_secs(30)); tokio::pin!(deadline); let result = loop { let event = tokio::select! { event = receiver.recv() => event.ok_or(Error::InvalidOperation)?, () = cancellation.cancelled() => break Err(Error::Cancelled), () = &mut deadline => break Err(Error::InvalidOperation), }; match event.packet_type { PacketType::TransferInfo => { let packet: TransferInfoPacket = Self::decode_packet(&event.data)?; if packet.transfer_info.transfer_id != transfer_id { continue; } if packet.transfer_info.status != StatusCode::OK as i32 { break Ok(None); } let size = usize::try_from(packet.transfer_info.size).map_err(|_| Error::Argument)?; if size > crate::asset_models::MAX_ASSET_BYTES { break Err(Error::Argument); } expected_size = Some(size); } PacketType::TransferPacket => { let packet: TransferPacketPacket = Self::decode_packet(&event.data)?; if packet.transfer_data.transfer_id != transfer_id || packet.transfer_data.packet < 0 { continue; } if packet.transfer_data.status < StatusCode::OK as i32 { break Ok(None); } chunks .entry(packet.transfer_data.packet) .or_insert(packet.transfer_data.data); } _ => continue, } let Some(size) = expected_size else { continue; }; let mut assembled = Vec::with_capacity(size); for packet_number in 0_i32.. { let Some(chunk) = chunks.get(&packet_number) else { break; }; let remaining = size.saturating_sub(assembled.len()); assembled.extend_from_slice(&chunk[..chunk.len().min(remaining)]); if assembled.len() == size { break; } } if assembled.len() == size { break Ok(Some(assembled)); } }; drop(subscription); if result.is_err() { self.abort_transfer(&simulator, transfer_id); } result } pub fn asset_type_to_string(&self, asset_type: AssetType) -> Result { Ok(match asset_type { AssetType::Texture => "texture", AssetType::Sound => "sound", AssetType::CallingCard => "callcard", AssetType::Landmark => "landmark", AssetType::Script => "script", AssetType::Clothing => "clothing", AssetType::Object => "object", AssetType::Notecard => "notecard", AssetType::Folder => "category", AssetType::LSLText => "lsltext", AssetType::LSLBytecode => "lslbyte", AssetType::TextureTGA => "txtr_tga", AssetType::Bodypart => "bodypart", AssetType::SoundWAV => "snd_wav", AssetType::ImageTGA => "img_tga", AssetType::ImageJPEG => "jpeg", AssetType::Animation => "animatn", AssetType::Gesture => "gesture", AssetType::Simstate => "simstate", AssetType::Link => "link", AssetType::LinkFolder => "link_f", AssetType::Mesh => "mesh", AssetType::Widget => "widget", AssetType::Person => "person", AssetType::Material => "material", AssetType::Settings => "settings", AssetType::Unknown => "invalid", } .to_owned()) } pub fn create_asset_wrapper(&self, type_: AssetType) -> Result { Asset::native_new(type_, UUID::zero(), Vec::new()) } pub async fn request_asset_with_uuid_asset_type_boolean_cancellation_token( &self, asset_id: UUID, type_: AssetType, priority: bool, cancellation_token: Option, ) -> Result, Error> { self.request_asset_with_uuid_asset_type_boolean_source_type_cancellation_token( asset_id, type_, priority, SourceType::Asset, cancellation_token, ) .await } pub async fn request_asset_with_uuid_asset_type_boolean_source_type_cancellation_token( &self, asset_id: UUID, type_: AssetType, priority: bool, source_type: SourceType, cancellation_token: Option, ) -> Result, Error> { if asset_id == UUID::zero() { return Ok(None); } if let Some(bytes) = self.cache.get_cached_asset_bytes_with_uuid(asset_id)? { return Asset::native_new(type_, asset_id, bytes).map(Some); } let token = Self::token(cancellation_token); let bytes = if source_type == SourceType::Asset { if let Some(cap) = self.capability("ViewerAsset")? { let uri = Self::with_query( &cap, &format!("{}_id", self.asset_type_to_string(type_)?), asset_id, ); self.fetch(uri, None, token.clone()).await? } else { self.request_asset_udp( asset_id, UUID::zero(), UUID::zero(), UUID::zero(), type_, priority, source_type, UUID::random()?, token.clone(), ) .await? } } else { self.request_asset_udp( asset_id, UUID::zero(), UUID::zero(), UUID::zero(), type_, priority, source_type, UUID::random()?, token.clone(), ) .await? }; let Some(bytes) = bytes else { return Ok(None); }; self.cache .save_asset_to_cache_with_uuid_bytes_cancellation_token( asset_id, bytes.clone(), Some(token), ) .await?; Asset::native_new(type_, asset_id, bytes).map(Some) } pub async fn request_asset_with_uuid_uuid_uuid_asset_type_boolean_source_type_uuid_cancellation_token( &self, asset_id: UUID, item_id: UUID, task_id: UUID, asset_type: AssetType, priority: bool, source_type: SourceType, transaction_id: UUID, cancellation_token: Option, ) -> Result, Error> { if asset_id == UUID::zero() { return Ok(None); } if let Some(bytes) = self.cache.get_cached_asset_bytes_with_uuid(asset_id)? { return Asset::native_new(asset_type, asset_id, bytes).map(Some); } let token = Self::token(cancellation_token); let transfer_id = if transaction_id == UUID::zero() { UUID::random()? } else { transaction_id }; let Some(bytes) = self .request_asset_udp( asset_id, item_id, task_id, UUID::zero(), asset_type, priority, source_type, transfer_id, token.clone(), ) .await? else { return Ok(None); }; self.cache .save_asset_to_cache_with_uuid_bytes_cancellation_token( asset_id, bytes.clone(), Some(token), ) .await?; Asset::native_new(asset_type, asset_id, bytes).map(Some) } pub async fn request_inventory_asset_with_inventory_item_boolean_uuid_cancellation_token( &self, item: InventoryItem, priority: bool, transfer_id: UUID, cancellation_token: Option, ) -> Result, Error> { self.request_inventory_asset_with_uuid_uuid_uuid_uuid_asset_type_boolean_uuid_cancellation_token(item.asset_uuid(), item.base.uuid(), UUID::zero(), item.base.owner_id(), item.asset_type(), priority, transfer_id, cancellation_token).await } pub async fn request_inventory_asset_with_uuid_uuid_uuid_uuid_asset_type_boolean_uuid_cancellation_token( &self, asset_id: UUID, item_id: UUID, task_id: UUID, owner_id: UUID, asset_type: AssetType, priority: bool, transfer_id: UUID, cancellation_token: Option, ) -> Result, Error> { if asset_id == UUID::zero() { return Ok(None); } if let Some(bytes) = self.cache.get_cached_asset_bytes_with_uuid(asset_id)? { return Asset::native_new(asset_type, asset_id, bytes).map(Some); } let token = Self::token(cancellation_token); let id = if transfer_id == UUID::zero() { UUID::random()? } else { transfer_id }; let Some(bytes) = self .request_asset_udp( asset_id, item_id, task_id, owner_id, asset_type, priority, SourceType::SimInventoryItem, id, token.clone(), ) .await? else { return Ok(None); }; self.cache .save_asset_to_cache_with_uuid_bytes_cancellation_token( asset_id, bytes.clone(), Some(token), ) .await?; Asset::native_new(asset_type, asset_id, bytes).map(Some) } pub async fn request_mesh( &self, mesh_id: UUID, cancellation_token: Option, ) -> Result, Error> { if mesh_id == UUID::zero() { return Ok(None); } if let Some(bytes) = self.cache.get_cached_asset_bytes_with_uuid(mesh_id)? { return AssetMesh::new_with_uuid_bytes(mesh_id, bytes).map(Some); } let Some(cap) = self.capability("GetMesh2")?.or(self.capability("GetMesh")?) else { return Ok(None); }; let token = Self::token(cancellation_token); let Some(bytes) = self .fetch( Self::with_query(&cap, "mesh_id", mesh_id), None, token.clone(), ) .await? else { return Ok(None); }; let mesh = AssetMesh::new_with_uuid_bytes(mesh_id, bytes.clone())?; self.cache .save_asset_to_cache_with_uuid_bytes_cancellation_token(mesh_id, bytes, Some(token)) .await?; Ok(Some(mesh)) } pub async fn request_image( &self, texture_id: UUID, _image_type: Option, cancellation_token: Option, ) -> Result, Error> { if texture_id == UUID::zero() { return Ok(None); } if let Some(bytes) = self.cache.get_cached_asset_bytes_with_uuid(texture_id)? { return AssetTexture::new_with_uuid_bytes(texture_id, bytes).map(Some); } let Some(cap) = self.capability("GetTexture")? else { return Ok(None); }; let parents: Vec<_> = cancellation_token.into_iter().collect(); let source = CancellationTokenSource::new_linked(&parents); mutex(&self.inner.image_cancellations).insert(texture_id, source.clone()); let events = Arc::clone(&self.inner.events.image_progress); let progress = move |report: crate::HttpCapsClientProgressReport| { let received = i32::try_from(report.bytes_transferred()).unwrap_or(i32::MAX); let total = report .total_bytes() .flatten() .and_then(|value| i32::try_from(value).ok()) .unwrap_or(received) .max(received); if let Ok(event) = ImageReceiveProgressEventArgs::new(texture_id, received, total) { events.emit(event); } }; let mut request = DownloadRequest::new( Self::with_query(&cap, "texture_id", texture_id), Some("image/x-j2c".into()), Some(Box::new(progress)), )?; request.cancellation_token = source.token(); let result = self .inner .downloads .queue_download_with_download_request_ba727387(request) .await .and_then(|(response, bytes)| { if bytes.len() > crate::asset_models::MAX_ASSET_BYTES { Err(Error::Argument) } else if (200..300).contains(&response.status_code) { Ok(Some(bytes)) } else { Ok(None) } }); mutex(&self.inner.image_cancellations).remove(&texture_id); let Some(bytes) = result? else { return Ok(None); }; self.cache .save_asset_to_cache_with_uuid_bytes_cancellation_token( texture_id, bytes.clone(), Some(source.token()), ) .await?; AssetTexture::new_with_uuid_bytes(texture_id, bytes).map(Some) } pub fn request_image_cancel(&self, texture_id: UUID) -> Result<(), Error> { if let Some(source) = mutex(&self.inner.image_cancellations).remove(&texture_id) { source.cancel(); } Ok(()) } pub(crate) fn native_image_transfer_count(&self) -> i32 { i32::try_from(mutex(&self.inner.image_cancellations).len()).unwrap_or(i32::MAX) } pub(crate) fn native_cancel_image_requests(&self) { let requests = mutex(&self.inner.image_cancellations) .drain() .map(|(_, source)| source) .collect::>(); for source in requests { source.cancel(); } } pub async fn request_server_baked_image( &self, avatar_id: UUID, texture_id: UUID, bake_name: String, cancellation_token: Option, ) -> Result, Error> { if avatar_id == UUID::zero() || texture_id == UUID::zero() || bake_name.is_empty() || bake_name.contains('/') { return Ok(None); } if let Some(bytes) = self.cache.get_cached_asset_bytes_with_uuid(texture_id)? { return AssetTexture::new_with_uuid_bytes(texture_id, bytes).map(Some); } let Some(base) = self .client()? .native_network()? .agent_appearance_service_url() else { return Ok(None); }; let uri = Uri(format!( "{}/texture/{}/{}/{}", base.trim_end_matches('/'), avatar_id, bake_name, texture_id )); let token = Self::token(cancellation_token); let Some(bytes) = self .fetch(uri, Some("image/x-j2c".into()), token.clone()) .await? else { return Ok(None); }; self.cache .save_asset_to_cache_with_uuid_bytes_cancellation_token( texture_id, bytes.clone(), Some(token), ) .await?; AssetTexture::new_with_uuid_bytes(texture_id, bytes).map(Some) } async fn uploader_asset( &self, cap_name: &str, data: Vec, cancellation_token: Option, ) -> Result { if data.is_empty() || data.len() > crate::asset_models::MAX_ASSET_BYTES { return Err(Error::Argument); } let Some(cap) = self.capability(cap_name)? else { return Ok(UUID::zero()); }; let token = Self::token(cancellation_token); let client = self.client()?; let (_, metadata) = client .native_http_caps_client() .post_with_uri_string_bytes_cancellation_token_i_progress( cap, "application/llsd+xml".into(), b"".to_vec(), token.clone(), None, ) .await?; let meta: Value = serde_json::from_slice(&metadata).map_err(|_| Error::Argument)?; let uploader = meta .get("uploader") .and_then(Value::as_str) .ok_or(Error::Argument)?; let (_, response) = client .native_http_caps_client() .post_with_uri_string_bytes_cancellation_token_i_progress( Uri(uploader.to_owned()), "application/octet-stream".into(), data, token, None, ) .await?; let result: Value = serde_json::from_slice(&response).map_err(|_| Error::Argument)?; if result.get("state").and_then(Value::as_str) != Some("complete") { return Ok(UUID::zero()); } result .get("new_asset") .or_else(|| result.get("new_asset_id")) .and_then(Value::as_str) .map(|id| UUID::new_with_string(id.to_owned())) .transpose() .map(|id| id.unwrap_or_else(UUID::zero)) } pub async fn request_upload_baked_texture( &self, texture_data: Vec, cancellation_token: Option, ) -> Result { self.uploader_asset("UploadBakedTexture", texture_data, cancellation_token) .await } pub async fn request_upload_large_texture( &self, data: Vec, width_pixels: i32, cancellation_token: Option, ) -> Result { if width_pixels <= 0 { return Err(Error::Argument); } self.request_upload_with_asset_type_bytes_boolean_uuid_cancellation_token( AssetType::Texture, data, false, UUID::random()?, cancellation_token, ) .await } pub fn get_texture_upload_cost(&self, width_pixels: i32) -> Result { if width_pixels <= 0 { return Err(Error::Argument); } Ok(self.client()?.settings_ref().upload_cost()) } pub async fn request_upload_with_asset_type_bytes_boolean_uuid_cancellation_token( &self, type_: AssetType, data: Vec, store_local: bool, transaction_id: UUID, cancellation_token: Option, ) -> Result { let token = Self::token(cancellation_token); token.throw_if_cancellation_requested()?; if data.is_empty() || data.len() > crate::asset_models::MAX_ASSET_BYTES { return Err(Error::Argument); } let mut client = self.client()?; let network = client.native_network()?; let simulator = network.current_sim().ok_or(Error::InvalidOperation)?; let secure_session = client.native_self().secure_session_id(); let asset_id = UUID::combine(transaction_id, secure_session)?; let (sender, receiver) = oneshot::channel(); let upload = Arc::new(Mutex::new(UploadState { asset_id, data: data.clone(), packet_num: 0, transaction_id, transferred: if data.len() <= 1000 { data.len() } else { 0 }, type_, xfer_id: 0, confirmation: Some(sender), })); { let mut pending = mutex(&self.inner.pending_upload); if pending.is_some() { return Err(Error::InvalidOperation); } *pending = Some(Arc::clone(&upload)); } let mut packet = AssetUploadRequestPacket::new_with_constructor()?; packet.asset_block.transaction_id = transaction_id; packet.asset_block.type_ = type_ as i8; packet.asset_block.store_local = store_local; packet.asset_block.tempfile = false; if data.len() <= 1000 { packet.asset_block.asset_data = data; } if let Err(error) = self.send_packet(&simulator, &packet, PacketType::AssetUploadRequest) { mutex(&self.inner.pending_upload).take(); return Err(error); } let (timeout_sender, timeout_receiver) = oneshot::channel(); std::thread::Builder::new() .name("asset-upload-timeout".into()) .spawn(move || { std::thread::sleep(Duration::from_secs(20)); let _ = timeout_sender.send(()); }) .map_err(|_| Error::InvalidOperation)?; let completion = receiver.fuse(); let cancelled = token.cancelled().fuse(); let timeout = timeout_receiver.fuse(); pin_mut!(completion, cancelled, timeout); let result = select_biased! { value = completion => value.map_err(|_| Error::InvalidOperation).and_then(|value| if value == UUID::zero() { Err(Error::InvalidOperation) } else { Ok(value) }), () = cancelled => Err(Error::Cancelled), _ = timeout => Err(Error::InvalidOperation), }; if result.is_err() { let mut pending = mutex(&self.inner.pending_upload); if pending .as_ref() .is_some_and(|value| Arc::ptr_eq(value, &upload)) { pending.take(); } } result } pub fn request_upload_with_uuid_asset_type_bytes_boolean_uuid( &self, asset_id: &mut UUID, type_: AssetType, data: Vec, store_local: bool, transaction_id: UUID, ) -> Result { let mut client = self.client()?; *asset_id = UUID::combine(transaction_id, client.native_self().secure_session_id())?; block_on( self.request_upload_with_asset_type_bytes_boolean_uuid_cancellation_token( type_, data, store_local, transaction_id, None, ), ) } pub fn request_asset_xfer( &self, filename: String, delete_on_completion: bool, use_big_packets: bool, v_file_id: UUID, v_file_type: AssetType, from_cache: bool, ) -> Result { if filename.as_bytes().contains(&0) || filename.len() > 1024 { return Err(Error::Argument); } let simulator = self .client()? .native_network()? .current_sim() .ok_or(Error::InvalidOperation)?; let id = UUID::random()?.get_u_long()?; let mut download = XferDownload::new()?; download.filename.clone_from(&filename); download.v_file_id = v_file_id; download.xfer_id = id; download.base.asset_type = v_file_type; download.base.id = v_file_id; mutex(&self.inner.xfer_downloads).insert( id, XferState { download, chunks: HashMap::new(), expected_size: None, final_packet: None, }, ); let mut packet = RequestXferPacket::new_with_constructor()?; packet.xfer_id.id = id; packet.xfer_id.filename = filename.into_bytes(); packet.xfer_id.filename.push(0); packet.xfer_id.delete_on_completion = delete_on_completion; packet.xfer_id.use_big_packets = use_big_packets; packet.xfer_id.v_file_id = v_file_id; packet.xfer_id.v_file_type = v_file_type as i16; packet.xfer_id.file_path = u8::from(from_cache); if let Err(error) = self.send_packet(&simulator, &packet, PacketType::RequestXfer) { mutex(&self.inner.xfer_downloads).remove(&id); return Err(error); } Ok(id) } pub fn request_estate_asset( &self, transfer: AssetDownload, eat: EstateAssetType, ) -> Result<(), Error> { let simulator = transfer .simulator .clone() .or(self.client()?.native_network()?.current_sim()) .ok_or(Error::InvalidOperation)?; let mut client = self.client()?; let agent_id = client.native_self().agent_id(); let session_id = client.native_self().session_id(); let mut params = Vec::with_capacity(36); params.extend_from_slice(&agent_id.get_bytes()?); params.extend_from_slice(&session_id.get_bytes()?); params.extend_from_slice(&(eat as i32).to_le_bytes()); let mut packet = TransferRequestPacket::new_with_constructor()?; packet.transfer_info.channel_type = transfer.channel as i32; packet.transfer_info.priority = transfer.priority; packet.transfer_info.source_type = transfer.source as i32; packet.transfer_info.transfer_id = transfer.base.id; packet.transfer_info.params = params; self.send_packet(&simulator, &packet, PacketType::TransferRequest) } #[must_use] pub fn http_download_manager(&self) -> DownloadManager { self.inner.downloads.clone() } pub fn subscribe_asset_uploaded( &self, handler: EventHandler, ) -> Subscription { self.inner.events.uploaded.subscribe(handler) } pub fn subscribe_upload_progress( &self, handler: EventHandler, ) -> Subscription { self.inner.events.upload_progress.subscribe(handler) } pub fn subscribe_image_receive_progress( &self, handler: EventHandler, ) -> Subscription { self.inner.events.image_progress.subscribe(handler) } pub fn subscribe_initiate_download( &self, handler: EventHandler, ) -> Subscription { self.inner.events.initiate.subscribe(handler) } pub fn subscribe_xfer_received( &self, handler: EventHandler, ) -> Subscription { self.inner.events.xfer.subscribe(handler) } }