Files
MetaCrate/crates/libremetaverse/src/asset_manager.rs
Chili Palmer c9a1170a27
Some checks failed
API and SemVer surface / api-surface (push) Failing after 1m34s
Native code generation / deterministic (push) Has been cancelled
Concurrency and resource soak audit / soak (push) Has been cancelled
Documentation / documentation (push) Has been cancelled
performance evidence / audit (push) Has been cancelled
First release candidate / non-fuzz-release-gate (push) Has been cancelled
Release platform and feature matrix / audit (push) Has been cancelled
Release platform and feature matrix / matrix (false, linux-stable-minimal, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (false, macos-stable-portable, x86_64-apple-darwin, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (false, windows-stable-portable, x86_64-pc-windows-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-msrv-portable, x86_64-unknown-linux-gnu, 1.96.0) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-stable-default, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-stable-features, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-stable-release-surface, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Dependency and supply-chain audit / audit (push) Has been cancelled
Native release artifact audit / audit (push) Has been cancelled
Imaging and meshing gate / native (push) Failing after 53s
JPEG 2000 feature / linux (push) Successful in 2m49s
Native Rust workspace compile / compile (push) Failing after 59s
Skia feature / linux (push) Has been cancelled
Complete first release candidate audit (#107)
2026-08-12 14:44:28 +00:00

1399 lines
50 KiB
Rust

//! 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<F: std::future::Future>(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>) {
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<T>(value: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
struct EventSlot<T> {
next: AtomicU64,
handlers: Mutex<HashMap<u64, EventHandler<T>>>,
}
impl<T> Default for EventSlot<T> {
fn default() -> Self {
Self {
next: AtomicU64::new(1),
handlers: Mutex::new(HashMap::new()),
}
}
}
impl<T: 'static> EventSlot<T> {
fn subscribe(self: &Arc<Self>, handler: EventHandler<T>) -> 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<T: Clone + 'static> EventSlot<T> {
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::<usize>::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<u8>,
packet_num: u32,
transaction_id: UUID,
transferred: usize,
type_: AssetType,
xfer_id: u64,
confirmation: Option<oneshot::Sender<UUID>>,
}
struct XferState {
download: XferDownload,
chunks: HashMap<u32, Vec<u8>>,
expected_size: Option<usize>,
final_packet: Option<u32>,
}
#[derive(Default)]
struct AssetEvents {
uploaded: Arc<EventSlot<AssetUploadEventArgs>>,
upload_progress: Arc<EventSlot<AssetUploadEventArgs>>,
image_progress: Arc<EventSlot<ImageReceiveProgressEventArgs>>,
initiate: Arc<EventSlot<InitiateDownloadEventArgs>>,
xfer: Arc<EventSlot<XferReceivedEventArgs>>,
}
pub(crate) struct AssetManagerInner {
client: crate::client_core::ClientWeakHandle,
cache: AssetCache,
downloads: DownloadManager,
image_cancellations: Mutex<HashMap<UUID, CancellationTokenSource>>,
pending_upload: Mutex<Option<Arc<Mutex<UploadState>>>>,
xfer_uploads: Mutex<HashMap<u64, Arc<Mutex<UploadState>>>>,
xfer_downloads: Mutex<HashMap<u64, XferState>>,
raw_packet_subscription: Mutex<Option<Subscription>>,
events: AssetEvents,
}
pub struct AssetManager {
pub cache: AssetCache,
inner: Arc<AssetManagerInner>,
}
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<u8>) -> Result<UUID, Error> {
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<AssetManagerInner>) -> Result<Self, Error> {
inner.client.upgrade().ok_or(Error::InvalidOperation)?;
Ok(Self {
cache: inner.cache.clone(),
inner,
})
}
pub(crate) fn native_inner(&self) -> Arc<AssetManagerInner> {
Arc::clone(&self.inner)
}
pub fn new(client: Option<GridClient>) -> Result<Self, Error> {
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<T: GeneratedPacket>(bytes: &[u8]) -> Result<T, Error> {
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<T: GeneratedPacket>(
&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<Mutex<UploadState>>,
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<GridClient, Error> {
self.inner.client.upgrade().ok_or(Error::InvalidOperation)
}
fn capability(&self, name: &str) -> Result<Option<Uri>, 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>) -> CancellationToken {
value.unwrap_or_default()
}
async fn fetch(
&self,
uri: Uri,
content_type: Option<String>,
cancellation: CancellationToken,
) -> Result<Option<Vec<u8>>, 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<Option<Vec<u8>>, 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::<i32, Vec<u8>>::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<String, Error> {
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, Error> {
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<CancellationToken>,
) -> Result<Option<Asset>, 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<CancellationToken>,
) -> Result<Option<Asset>, 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<CancellationToken>,
) -> Result<Option<Asset>, 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<CancellationToken>,
) -> Result<Option<Asset>, 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<CancellationToken>,
) -> Result<Option<Asset>, 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<CancellationToken>,
) -> Result<Option<AssetMesh>, 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<ImageType>,
cancellation_token: Option<CancellationToken>,
) -> Result<Option<AssetTexture>, 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::<Vec<_>>();
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<CancellationToken>,
) -> Result<Option<AssetTexture>, 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<u8>,
cancellation_token: Option<CancellationToken>,
) -> Result<UUID, Error> {
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"<llsd><undef /></llsd>".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<u8>,
cancellation_token: Option<CancellationToken>,
) -> Result<UUID, Error> {
self.uploader_asset("UploadBakedTexture", texture_data, cancellation_token)
.await
}
pub async fn request_upload_large_texture(
&self,
data: Vec<u8>,
width_pixels: i32,
cancellation_token: Option<CancellationToken>,
) -> Result<UUID, Error> {
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<i32, Error> {
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<u8>,
store_local: bool,
transaction_id: UUID,
cancellation_token: Option<CancellationToken>,
) -> Result<UUID, Error> {
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<u8>,
store_local: bool,
transaction_id: UUID,
) -> Result<UUID, Error> {
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<u64, Error> {
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<AssetUploadEventArgs>,
) -> Subscription {
self.inner.events.uploaded.subscribe(handler)
}
pub fn subscribe_upload_progress(
&self,
handler: EventHandler<AssetUploadEventArgs>,
) -> Subscription {
self.inner.events.upload_progress.subscribe(handler)
}
pub fn subscribe_image_receive_progress(
&self,
handler: EventHandler<ImageReceiveProgressEventArgs>,
) -> Subscription {
self.inner.events.image_progress.subscribe(handler)
}
pub fn subscribe_initiate_download(
&self,
handler: EventHandler<InitiateDownloadEventArgs>,
) -> Subscription {
self.inner.events.initiate.subscribe(handler)
}
pub fn subscribe_xfer_received(
&self,
handler: EventHandler<XferReceivedEventArgs>,
) -> Subscription {
self.inner.events.xfer.subscribe(handler)
}
}