Implement capability HTTP and downloads (#54)
All checks were successful
Native code generation / deterministic (push) Successful in 11m59s
Imaging and meshing gate / native (push) Successful in 3m55s
JPEG 2000 feature / linux (push) Successful in 2m26s
Native Rust workspace compile / compile (push) Successful in 3m58s
Skia feature / linux (push) Successful in 31m44s

This commit is contained in:
2026-08-09 15:44:42 +00:00
parent d298b4f4c4
commit 3032b16e7b
15 changed files with 5833 additions and 668 deletions

1207
Cargo.lock generated

File diff suppressed because it is too large Load Diff

View File

@@ -282,3 +282,15 @@ snapshotted before invocation, and manager/transport workers retain no dropped
client owner. The ownership, client owner. The ownership,
dispatch, lifecycle, and fake-server verification contracts are documented in dispatch, lifecycle, and fake-server verification contracts are documented in
[`docs/network-manager.md`](docs/network-manager.md). [`docs/network-manager.md`](docs/network-manager.md).
The capability transport is now native Rust as well. `HttpCapsClient` supports
the C# GET/POST/PUT/PATCH/DELETE and LLSD overloads through an injectable
`reqwest`/fake-handler boundary, with HTTPS, bounded redirects, gzip/deflate,
streaming progress, cancellation, explicit request/response/decompression
limits, and rejection of outbound non-HTTP capability URIs. The per-category
oldest-first token buckets retain the reference defaults and cap-name mapping.
The download manager adds a bounded queue and concurrency gate, canonical-URI
deduplication, subscriber progress fanout, shared cancellation, and transient
retry handling while preserving permanent 401/403/404/410 failures. Capability
URLs and tokens are absent from errors and diagnostics. Ownership, limits,
injection, and offline fake-server coverage are documented in
[`docs/caps-http.md`](docs/caps-http.md).

View File

@@ -4,7 +4,7 @@ Generated by `python3 tools/generate_api_shims.py`; do not edit by hand.
| Assembly | Types | Members | Status | | Assembly | Types | Members | Status |
|---|---:|---:|---| |---|---:|---:|---|
| `LibreMetaverse` | 2,711 | 27,281 | native implementation: 47 types / 13,615 members; remaining surface is callable failure-only shims | | `LibreMetaverse` | 2,711 | 27,281 | native implementation: 53 types / 13,673 members; remaining surface is callable failure-only shims |
| `LibreMetaverse.Imaging.Abstractions` | 3 | 20 | native implementation: 3 types / 20 members; no generated shims remain | | `LibreMetaverse.Imaging.Abstractions` | 3 | 20 | native implementation: 3 types / 20 members; no generated shims remain |
| `LibreMetaverse.Imaging.Skia` | 1 | 3 | native implementation: 1 type / 3 members; no generated shims remain | | `LibreMetaverse.Imaging.Skia` | 1 | 3 | native implementation: 1 type / 3 members; no generated shims remain |
| `LibreMetaverse.LslTools` | 164 | 768 | callable failure-only shim | | `LibreMetaverse.LslTools` | 164 | 768 | callable failure-only shim |

View File

@@ -802,7 +802,7 @@ fn sorted_entries(values: &HashMap<String, OSD>) -> Vec<(&String, &OSD)> {
} }
fn decode_base16(value: &str) -> Result<Vec<u8>, &'static str> { fn decode_base16(value: &str) -> Result<Vec<u8>, &'static str> {
if value.len() % 2 != 0 { if !value.len().is_multiple_of(2) {
return Err("invalid notation LLSD base16 length"); return Err("invalid notation LLSD base16 length");
} }
value value

View File

@@ -356,6 +356,45 @@ impl CancellationToken {
} }
} }
/// Registers an executor-neutral callback and removes it when the returned
/// guard is dropped.
#[must_use = "dropping the guard immediately unregisters the callback"]
pub fn register_callback(&self, callback: Arc<dyn Fn() + Send + Sync>) -> Subscription {
if self.is_cancellation_requested() {
callback();
return Subscription::detached();
}
let id = self.0.next_id.fetch_add(1, Ordering::Relaxed);
self.0
.callbacks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(id, callback);
if self.is_cancellation_requested() {
if let Some(callback) = self
.0
.callbacks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&id)
{
callback();
}
return Subscription::detached();
}
let state = Arc::downgrade(&self.0);
Subscription::new(move || {
let Some(state) = state.upgrade() else {
return;
};
state
.callbacks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&id);
})
}
fn register(&self, callback: Arc<dyn Fn() + Send + Sync>) -> CancellationRegistration { fn register(&self, callback: Arc<dyn Fn() + Send + Sync>) -> CancellationRegistration {
if self.is_cancellation_requested() { if self.is_cancellation_requested() {
callback(); callback();
@@ -528,7 +567,31 @@ pub struct DictionaryKeys<TKey, TValue>(pub Vec<TKey>, pub PhantomData<fn(TValue
pub struct DictionaryEntry(pub Object, pub Object); pub struct DictionaryEntry(pub Object, pub Object);
pub struct RateLimitLease; /// Result of a mapped rate-limit acquisition.
///
/// Token-bucket leases do not return a token when dropped; the value records
/// whether a token was acquired or the bounded waiting queue was already full.
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
pub struct RateLimitLease {
acquired: bool,
}
impl RateLimitLease {
#[must_use]
pub const fn acquired() -> Self {
Self { acquired: true }
}
#[must_use]
pub const fn rejected() -> Self {
Self { acquired: false }
}
#[must_use]
pub const fn is_acquired(self) -> bool {
self.acquired
}
}
struct CancellationSourceState { struct CancellationSourceState {
token: CancellationToken, token: CancellationToken,
@@ -641,7 +704,160 @@ pub struct Thread;
pub struct Task<T>(pub PhantomData<fn() -> T>); pub struct Task<T>(pub PhantomData<fn() -> T>);
pub struct TaskCompletionSource<T>(pub PhantomData<fn(T)>); struct TaskCompletionState<T> {
next_id: AtomicU64,
result: Mutex<Option<Result<T, crate::Error>>>,
waiters: Mutex<BTreeMap<u64, Waker>>,
}
/// Cloneable completion source used by mapped task-returning APIs.
#[derive(Clone)]
pub struct TaskCompletionSource<T>(Arc<TaskCompletionState<T>>);
impl<T> TaskCompletionSource<T> {
#[must_use]
pub fn new() -> Self {
Self(Arc::new(TaskCompletionState {
next_id: AtomicU64::new(1),
result: Mutex::new(None),
waiters: Mutex::new(BTreeMap::new()),
}))
}
/// Completes the task once and wakes every waiter.
#[must_use]
pub fn try_set_result(&self, value: T) -> bool {
self.try_complete(Ok(value))
}
/// Completes the task with a typed failure once.
#[must_use]
pub fn try_set_error(&self, error: crate::Error) -> bool {
self.try_complete(Err(error))
}
/// Completes the task with mapped cancellation once.
#[must_use]
pub fn try_set_cancelled(&self) -> bool {
self.try_set_error(crate::Error::Cancelled)
}
fn try_complete(&self, value: Result<T, crate::Error>) -> bool {
let mut result = self
.0
.result
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if result.is_some() {
return false;
}
*result = Some(value);
drop(result);
for (_, waiter) in std::mem::take(
&mut *self
.0
.waiters
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
) {
waiter.wake();
}
true
}
#[must_use]
pub fn future(&self) -> TaskCompletionFuture<T> {
TaskCompletionFuture {
state: Arc::clone(&self.0),
waiter_id: None,
}
}
#[must_use]
pub fn is_completed(&self) -> bool {
self.0
.result
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.is_some()
}
}
impl<T> Default for TaskCompletionSource<T> {
fn default() -> Self {
Self::new()
}
}
impl<T> fmt::Debug for TaskCompletionSource<T> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("TaskCompletionSource")
.field("completed", &self.is_completed())
.finish_non_exhaustive()
}
}
/// Future associated with a [`TaskCompletionSource`].
pub struct TaskCompletionFuture<T> {
state: Arc<TaskCompletionState<T>>,
waiter_id: Option<u64>,
}
impl<T: Clone> Future for TaskCompletionFuture<T> {
type Output = Result<T, crate::Error>;
fn poll(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
if let Some(result) = self
.state
.result
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
{
return Poll::Ready(result);
}
let id = self.waiter_id.unwrap_or_else(|| {
let id = self.state.next_id.fetch_add(1, Ordering::Relaxed);
self.waiter_id = Some(id);
id
});
self.state
.waiters
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(id, context.waker().clone());
let ready = self
.state
.result
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(result) = ready {
self.state
.waiters
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&id);
self.waiter_id = None;
Poll::Ready(result)
} else {
Poll::Pending
}
}
}
impl<T> Drop for TaskCompletionFuture<T> {
fn drop(&mut self) {
if let Some(id) = self.waiter_id {
self.state
.waiters
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&id);
}
}
}
/// Cloneable clock boundary used by runtime-neutral client code. /// Cloneable clock boundary used by runtime-neutral client code.
#[derive(Clone)] #[derive(Clone)]
@@ -741,7 +957,19 @@ pub struct HttpMethod(pub String);
pub struct AuthenticationHeaderValue(pub String); pub struct AuthenticationHeaderValue(pub String);
pub trait IProgress<T>: Send + Sync {} /// Executor-neutral progress sink used by mapped upload/download APIs.
pub trait IProgress<T>: Send + Sync {
fn report(&self, value: T);
}
impl<T, F> IProgress<T> for F
where
F: Fn(T) + Send + Sync,
{
fn report(&self, value: T) {
self(value);
}
}
pub struct IEnumerator; pub struct IEnumerator;

View File

@@ -14,9 +14,12 @@ jpeg2000 = ["libremetaverse-imaging/jpeg2000"]
[dependencies] [dependencies]
bcdec_rs = { version = "0.2.0", optional = true } bcdec_rs = { version = "0.2.0", optional = true }
flate2 = "1.1.2"
futures-util = "0.3.31"
libremetaverse-imaging = { path = "../libremetaverse-imaging" } libremetaverse-imaging = { path = "../libremetaverse-imaging" }
libremetaverse-structured-data = { path = "../libremetaverse-structured-data" } libremetaverse-structured-data = { path = "../libremetaverse-structured-data" }
libremetaverse-types = { path = "../libremetaverse-types" } libremetaverse-types = { path = "../libremetaverse-types" }
reqwest = { version = "0.13.4", default-features = false, features = ["rustls", "stream"] }
roxmltree = "0.21.1" roxmltree = "0.21.1"
tokio = { version = "1.47.1", features = ["macros", "net", "rt", "sync", "time"] } tokio = { version = "1.47.1", features = ["macros", "net", "rt", "sync", "time"] }

File diff suppressed because it is too large Load Diff

View File

@@ -114,12 +114,18 @@ struct ClientRuntime {
cancellation: CancellationTokenSource, cancellation: CancellationTokenSource,
services: Mutex<Vec<ServiceRegistration>>, services: Mutex<Vec<ServiceRegistration>>,
agent_throttle_sender: Mutex<Option<Arc<dyn crate::udp_transport::AgentThrottleSender>>>, agent_throttle_sender: Mutex<Option<Arc<dyn crate::udp_transport::AgentThrottleSender>>>,
http_caps_client: Mutex<crate::caps_http::HttpCapsClient>,
caps_rate_limiter: Mutex<crate::caps_http::CapsRateLimiter>,
shutdown_complete: Condvar, shutdown_complete: Condvar,
shutdown_wait: Mutex<()>, shutdown_wait: Mutex<()>,
} }
impl ClientRuntime { impl ClientRuntime {
fn new(services: Vec<Arc<dyn ClientService>>) -> Self { fn new(
services: Vec<Arc<dyn ClientService>>,
http_caps_client: crate::caps_http::HttpCapsClient,
caps_rate_limiter: crate::caps_http::CapsRateLimiter,
) -> Self {
Self { Self {
state: AtomicU8::new(ClientLifecycleState::Active as u8), state: AtomicU8::new(ClientLifecycleState::Active as u8),
cancellation: CancellationTokenSource::new(), cancellation: CancellationTokenSource::new(),
@@ -134,6 +140,8 @@ impl ClientRuntime {
.collect(), .collect(),
), ),
agent_throttle_sender: Mutex::new(None), agent_throttle_sender: Mutex::new(None),
http_caps_client: Mutex::new(http_caps_client),
caps_rate_limiter: Mutex::new(caps_rate_limiter),
shutdown_complete: Condvar::new(), shutdown_complete: Condvar::new(),
shutdown_wait: Mutex::new(()), shutdown_wait: Mutex::new(()),
} }
@@ -206,6 +214,15 @@ impl ClientRuntime {
.lock() .lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) .unwrap_or_else(std::sync::PoisonError::into_inner)
.take(); .take();
self.http_caps_client
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.shutdown();
let _ = self
.caps_rate_limiter
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.dispose();
let shutdown_wait = self let shutdown_wait = self
.shutdown_wait .shutdown_wait
.lock() .lock()
@@ -309,6 +326,43 @@ impl GridClient {
.clone() .clone()
} }
pub(crate) fn native_http_caps_client(&self) -> crate::caps_http::HttpCapsClient {
self.runtime
.http_caps_client
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub(crate) fn native_set_http_caps_client(&mut self, value: crate::caps_http::HttpCapsClient) {
*self
.runtime
.http_caps_client
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = value;
}
pub(crate) fn native_caps_rate_limiter(&self) -> crate::caps_http::CapsRateLimiter {
self.runtime
.caps_rate_limiter
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub(crate) fn native_set_caps_rate_limiter(
&mut self,
value: crate::caps_http::CapsRateLimiter,
) {
self.native_http_caps_client()
.set_rate_limiter(Some(value.clone()));
*self
.runtime
.caps_rate_limiter
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = value;
}
pub(crate) fn shutdown(&self) -> Result<(), ClientCoreError> { pub(crate) fn shutdown(&self) -> Result<(), ClientCoreError> {
self.runtime.shutdown() self.runtime.shutdown()
} }
@@ -348,6 +402,10 @@ impl Drop for GridClient {
} }
impl crate::IGridClient for GridClient { impl crate::IGridClient for GridClient {
fn http_caps_client(&self) -> crate::caps_http::HttpCapsClient {
self.native_http_caps_client()
}
fn settings(&self) -> Settings { fn settings(&self) -> Settings {
self.settings.clone() self.settings.clone()
} }
@@ -404,10 +462,44 @@ impl GridClientBuilder {
/// Returns the first invalid configuration field. /// Returns the first invalid configuration field.
pub fn build(self) -> Result<GridClient, ClientCoreError> { pub fn build(self) -> Result<GridClient, ClientCoreError> {
self.settings.validate()?; self.settings.validate()?;
let caps_rate_limiter = crate::caps_http::CapsRateLimiter::new_with_clock_and_overrides(
self.time_provider.clone(),
None,
)
.map_err(|_| ClientCoreError::InvalidConfiguration {
field: "CapsRateLimiter",
reason: "could not construct the configured token buckets",
})?;
let max_connections = usize::try_from(Settings::max_http_connections()).map_err(|_| {
ClientCoreError::InvalidConfiguration {
field: "Settings.MaxHttpConnections",
reason: "must be positive",
}
})?;
let caps_timeout = u64::try_from(self.settings.timing.caps_timeout).map_err(|_| {
ClientCoreError::InvalidConfiguration {
field: "Timing.CapsTimeout",
reason: "must be positive",
}
})?;
let http_caps_client = crate::caps_http::HttpCapsClient::production(
&Settings::user_agent(),
std::time::Duration::from_millis(caps_timeout),
max_connections,
caps_rate_limiter.clone(),
)
.map_err(|_| ClientCoreError::InvalidConfiguration {
field: "HttpCapsClient",
reason: "could not construct the HTTP transport",
})?;
Ok(GridClient { Ok(GridClient {
settings: self.settings, settings: self.settings,
time_provider: self.time_provider, time_provider: self.time_provider,
runtime: Arc::new(ClientRuntime::new(self.services)), runtime: Arc::new(ClientRuntime::new(
self.services,
http_caps_client,
caps_rate_limiter,
)),
}) })
} }
} }

View File

@@ -0,0 +1,609 @@
//! Bounded, deduplicating HTTP download dispatcher.
#![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::{Error, GridClient, HttpCapsClientProgressReport};
use libremetaverse_types::compat::{
CancellationToken, CancellationTokenSource, HttpResponse, IProgress, Subscription,
TaskCompletionSource, Uri,
};
use std::collections::HashMap;
use std::fmt;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, Weak};
use std::thread::{self, JoinHandle};
use std::time::Duration;
use tokio::sync::{Notify, mpsc, oneshot};
use tokio::task::JoinSet;
const DEFAULT_PARALLEL_DOWNLOADS: usize = 8;
const MAX_PARALLEL_DOWNLOADS: usize = 32;
const DOWNLOAD_QUEUE_CAPACITY: usize = 256;
type DownloadValue = (HttpResponse, Vec<u8>);
type DownloadCompletion = TaskCompletionSource<DownloadValue>;
type DownloadResult = Result<DownloadValue, Error>;
fn mutex<T>(value: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
/// Mapped request state with the same defaults as the C# constructor.
pub struct DownloadRequest {
pub address: Uri,
pub attempt: i32,
pub cancellation_token: CancellationToken,
pub completion_tcs: Option<DownloadCompletion>,
pub content_type: Option<String>,
pub download_progress_callback: Option<Box<dyn IProgress<HttpCapsClientProgressReport>>>,
pub retries: i32,
}
impl DownloadRequest {
pub fn new(
address: Uri,
content_type: Option<String>,
download_progress_callback: Option<Box<dyn IProgress<HttpCapsClientProgressReport>>>,
) -> Result<Self, Error> {
validate_http_uri(&address)?;
Ok(Self {
address,
attempt: 0,
cancellation_token: CancellationToken::default(),
completion_tcs: None,
content_type,
download_progress_callback,
retries: 5,
})
}
}
impl fmt::Debug for DownloadRequest {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("DownloadRequest")
.field("address", &"<redacted capability URI>")
.field("attempt", &self.attempt)
.field(
"cancellation_requested",
&self.cancellation_token.is_cancellation_requested(),
)
.field("has_completion", &self.completion_tcs.is_some())
.field("has_content_type", &self.content_type.is_some())
.field("has_progress", &self.download_progress_callback.is_some())
.field("retries", &self.retries)
.finish()
}
}
struct ActiveDownload {
cancellation: CancellationTokenSource,
cancellation_guards: Mutex<Vec<Subscription>>,
progress: Mutex<Vec<Arc<dyn IProgress<HttpCapsClientProgressReport>>>>,
completion_sources: Mutex<Vec<DownloadCompletion>>,
waiters: Mutex<Vec<oneshot::Sender<DownloadResult>>>,
completed: AtomicBool,
}
impl ActiveDownload {
fn new() -> Self {
Self {
cancellation: CancellationTokenSource::new(),
cancellation_guards: Mutex::new(Vec::new()),
progress: Mutex::new(Vec::new()),
completion_sources: Mutex::new(Vec::new()),
waiters: Mutex::new(Vec::new()),
completed: AtomicBool::new(false),
}
}
fn attach_cancellation(&self, token: &CancellationToken) {
let cancellation = self.cancellation.clone();
let guard = token.register_callback(Arc::new(move || cancellation.cancel()));
mutex(&self.cancellation_guards).push(guard);
}
fn attach_progress(&self, progress: Option<Box<dyn IProgress<HttpCapsClientProgressReport>>>) {
if let Some(progress) = progress {
mutex(&self.progress).push(Arc::from(progress));
}
}
fn add_waiter(&self) -> oneshot::Receiver<DownloadResult> {
let (sender, receiver) = oneshot::channel();
mutex(&self.waiters).push(sender);
receiver
}
fn attach_completion_source(&self, completion_source: Option<DownloadCompletion>) {
if let Some(completion_source) = completion_source {
mutex(&self.completion_sources).push(completion_source);
}
}
fn complete(&self, result: DownloadResult) {
if self.completed.swap(true, Ordering::AcqRel) {
return;
}
for waiter in std::mem::take(&mut *mutex(&self.waiters)) {
let _ = waiter.send(result.clone());
}
for completion_source in std::mem::take(&mut *mutex(&self.completion_sources)) {
match result.clone() {
Ok(value) => {
let _ = completion_source.try_set_result(value);
}
Err(Error::Cancelled) => {
let _ = completion_source.try_set_cancelled();
}
Err(error) => {
let _ = completion_source.try_set_error(error);
}
}
}
mutex(&self.progress).clear();
mutex(&self.cancellation_guards).clear();
}
fn report(&self, report: HttpCapsClientProgressReport) {
let handlers = mutex(&self.progress).clone();
for handler in handlers {
let _ = catch_unwind(AssertUnwindSafe(|| handler.report(report)));
}
}
}
struct ProgressFanout(Weak<ActiveDownload>);
impl IProgress<HttpCapsClientProgressReport> for ProgressFanout {
fn report(&self, report: HttpCapsClientProgressReport) {
if let Some(active) = self.0.upgrade() {
active.report(report);
}
}
}
struct DownloadJob {
key: String,
address: Uri,
attempt: i32,
retries: i32,
active: Arc<ActiveDownload>,
client: GridClient,
}
struct GateState {
active: usize,
}
struct DynamicGate {
limit: AtomicUsize,
state: Mutex<GateState>,
notify: Notify,
}
impl DynamicGate {
fn new(limit: usize) -> Self {
Self {
limit: AtomicUsize::new(limit),
state: Mutex::new(GateState { active: 0 }),
notify: Notify::new(),
}
}
fn set_limit(&self, value: usize) {
self.limit.store(value, Ordering::Release);
self.notify.notify_waiters();
}
async fn acquire(
self: &Arc<Self>,
cancellation: CancellationToken,
) -> Result<GatePermit, Error> {
loop {
cancellation.throw_if_cancellation_requested()?;
let notified = self.notify.notified();
{
let mut state = mutex(&self.state);
if state.active < self.limit.load(Ordering::Acquire) {
state.active += 1;
return Ok(GatePermit(Arc::clone(self)));
}
}
tokio::select! {
() = notified => {}
() = cancellation.cancelled() => return Err(Error::Cancelled),
}
}
}
}
struct GatePermit(Arc<DynamicGate>);
impl Drop for GatePermit {
fn drop(&mut self) {
let mut state = mutex(&self.0.state);
state.active = state.active.saturating_sub(1);
drop(state);
self.0.notify.notify_one();
}
}
struct DownloadManagerInner {
client: GridClient,
active: Mutex<HashMap<String, Arc<ActiveDownload>>>,
sender: Mutex<Option<mpsc::Sender<DownloadJob>>>,
dispatcher: Mutex<Option<JoinHandle<()>>>,
shutdown: CancellationTokenSource,
gate: Arc<DynamicGate>,
parallel_downloads: AtomicUsize,
disposed: AtomicBool,
}
impl DownloadManagerInner {
fn remove_active(&self, key: &str, active: &Arc<ActiveDownload>) {
let mut downloads = mutex(&self.active);
if downloads
.get(key)
.is_some_and(|candidate| Arc::ptr_eq(candidate, active))
{
downloads.remove(key);
}
}
fn shutdown(&self) {
if self.disposed.swap(true, Ordering::AcqRel) {
return;
}
self.shutdown.cancel();
for active in mutex(&self.active).values() {
active.cancellation.cancel();
}
mutex(&self.sender).take();
if let Some(dispatcher) = mutex(&self.dispatcher).take()
&& dispatcher.thread().id() != thread::current().id()
{
let _ = dispatcher.join();
}
let remaining = std::mem::take(&mut *mutex(&self.active));
for active in remaining.into_values() {
active.complete(Err(Error::Cancelled));
}
}
}
impl Drop for DownloadManagerInner {
fn drop(&mut self) {
self.shutdown();
}
}
/// Native, bounded download manager with URI-level request deduplication.
#[derive(Clone)]
pub struct DownloadManager(Arc<DownloadManagerInner>);
impl DownloadManager {
pub fn new(client: GridClient) -> Result<Self, Error> {
let (sender, receiver) = mpsc::channel(DOWNLOAD_QUEUE_CAPACITY);
let gate = Arc::new(DynamicGate::new(DEFAULT_PARALLEL_DOWNLOADS));
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|_| Error::InvalidOperation)?;
let inner = Arc::new(DownloadManagerInner {
client,
active: Mutex::new(HashMap::new()),
sender: Mutex::new(Some(sender)),
dispatcher: Mutex::new(None),
shutdown: CancellationTokenSource::new(),
gate: Arc::clone(&gate),
parallel_downloads: AtomicUsize::new(DEFAULT_PARALLEL_DOWNLOADS),
disposed: AtomicBool::new(false),
});
let weak = Arc::downgrade(&inner);
let shutdown = inner.shutdown.token();
let dispatcher = thread::Builder::new()
.name("libremetaverse-download-dispatcher".to_owned())
.spawn(move || download_dispatcher(runtime, receiver, weak, gate, shutdown))
.map_err(|_| Error::InvalidOperation)?;
*mutex(&inner.dispatcher) = Some(dispatcher);
Ok(Self(inner))
}
pub fn dispose(&self) -> Result<(), Error> {
self.0.shutdown();
Ok(())
}
#[must_use]
pub fn parallel_downloads(&self) -> i32 {
i32::try_from(self.0.parallel_downloads.load(Ordering::Acquire)).unwrap_or(i32::MAX)
}
pub fn set_parallel_downloads(&mut self, value: i32) {
let value = usize::try_from(value)
.unwrap_or(1)
.clamp(1, MAX_PARALLEL_DOWNLOADS);
self.0.parallel_downloads.store(value, Ordering::Release);
self.0.gate.set_limit(value);
}
pub fn queue_download_with_download_request(
&self,
request: DownloadRequest,
) -> Result<(), Error> {
self.enqueue(request, false).map(|_| ())
}
pub fn queue_download_with_download_request_cancellation_token(
&self,
mut request: DownloadRequest,
cancellation_token: CancellationToken,
) -> Result<(), Error> {
request.cancellation_token = cancellation_token;
self.queue_download_with_download_request(request)
}
pub async fn queue_download_with_download_request_ba727387(
&self,
request: DownloadRequest,
) -> DownloadResult {
let cancellation = request.cancellation_token.clone();
let (active, receiver) = self.enqueue(request, true)?;
await_download(
active,
receiver.ok_or(Error::InvalidOperation)?,
cancellation,
)
.await
}
pub async fn queue_download_with_uri_string_i_progress_cancellation_token_int32(
&self,
address: Uri,
content_type: Option<String>,
progress_callback: Option<Box<dyn IProgress<HttpCapsClientProgressReport>>>,
cancellation_token: Option<CancellationToken>,
retries: Option<i32>,
) -> DownloadResult {
let cancellation = cancellation_token.unwrap_or_default();
let mut request = DownloadRequest::new(address, content_type, progress_callback)?;
request.cancellation_token = cancellation.clone();
request.retries = retries.unwrap_or(5).max(0);
let (active, receiver) = self.enqueue(request, true)?;
await_download(
active,
receiver.ok_or(Error::InvalidOperation)?,
cancellation,
)
.await
}
pub async fn download_with_uri_string_i_progress_cancellation_token(
&self,
address: Uri,
content_type: String,
progress: Box<dyn IProgress<HttpCapsClientProgressReport>>,
cancellation_token: CancellationToken,
) -> DownloadResult {
self.queue_download_with_uri_string_i_progress_cancellation_token_int32(
address,
Some(content_type),
Some(progress),
Some(cancellation_token),
Some(5),
)
.await
}
pub async fn download_with_uri_i_progress_cancellation_token(
&self,
address: Uri,
progress: Box<dyn IProgress<HttpCapsClientProgressReport>>,
cancellation_token: CancellationToken,
) -> DownloadResult {
self.download_with_uri_string_i_progress_cancellation_token(
address,
String::new(),
progress,
cancellation_token,
)
.await
}
fn enqueue(
&self,
mut request: DownloadRequest,
with_waiter: bool,
) -> Result<
(
Arc<ActiveDownload>,
Option<oneshot::Receiver<DownloadResult>>,
),
Error,
> {
if self.0.disposed.load(Ordering::Acquire) {
return Err(Error::InvalidOperation);
}
let key = validate_http_uri(&request.address)?;
let selected = {
let mut active = mutex(&self.0.active);
if let Some(existing) = active.get(&key) {
Some((Arc::clone(existing), false))
} else if request.cancellation_token.is_cancellation_requested() {
None
} else {
let download = Arc::new(ActiveDownload::new());
active.insert(key.clone(), Arc::clone(&download));
Some((download, true))
}
};
let Some((active, is_new)) = selected else {
if let Some(completion_source) = request.completion_tcs.take() {
let _ = completion_source.try_set_cancelled();
}
return Err(Error::Cancelled);
};
active.attach_cancellation(&request.cancellation_token);
active.attach_progress(request.download_progress_callback.take());
active.attach_completion_source(request.completion_tcs.take());
let receiver = with_waiter.then(|| active.add_waiter());
if !is_new {
return Ok((active, receiver));
}
let job = DownloadJob {
key: key.clone(),
address: request.address,
attempt: request.attempt.max(0),
retries: request.retries.max(0),
active: Arc::clone(&active),
client: self.0.client.clone(),
};
let sender = mutex(&self.0.sender)
.clone()
.ok_or(Error::InvalidOperation)?;
if sender.try_send(job).is_err() {
self.0.remove_active(&key, &active);
active.complete(Err(Error::InvalidOperation));
return Err(Error::InvalidOperation);
}
Ok((active, receiver))
}
}
impl fmt::Debug for DownloadManager {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("DownloadManager")
.field("active_downloads", &mutex(&self.0.active).len())
.field("parallel_downloads", &self.parallel_downloads())
.field("queue_capacity", &DOWNLOAD_QUEUE_CAPACITY)
.field("disposed", &self.0.disposed.load(Ordering::Acquire))
.finish()
}
}
async fn await_download(
active: Arc<ActiveDownload>,
receiver: oneshot::Receiver<DownloadResult>,
cancellation: CancellationToken,
) -> DownloadResult {
tokio::select! {
result = receiver => result.map_err(|_| Error::Cancelled)?,
() = cancellation.cancelled() => {
active.cancellation.cancel();
Err(Error::Cancelled)
}
}
}
fn download_dispatcher(
runtime: tokio::runtime::Runtime,
mut receiver: mpsc::Receiver<DownloadJob>,
manager: Weak<DownloadManagerInner>,
gate: Arc<DynamicGate>,
shutdown: CancellationToken,
) {
runtime.block_on(async move {
let mut jobs = JoinSet::new();
loop {
tokio::select! {
biased;
() = shutdown.cancelled() => break,
Some(job) = receiver.recv() => {
let manager = manager.clone();
let gate = Arc::clone(&gate);
jobs.spawn(async move {
run_download_job(job, manager, gate).await;
});
}
Some(_) = jobs.join_next(), if !jobs.is_empty() => {}
else => break,
}
}
while let Ok(job) = receiver.try_recv() {
job.active.complete(Err(Error::Cancelled));
if let Some(manager) = manager.upgrade() {
manager.remove_active(&job.key, &job.active);
}
}
while jobs.join_next().await.is_some() {}
});
}
async fn run_download_job(
job: DownloadJob,
manager: Weak<DownloadManagerInner>,
gate: Arc<DynamicGate>,
) {
let cancellation = job.active.cancellation.token();
let permit = gate.acquire(cancellation.clone()).await;
let result = match permit {
Ok(_permit) => perform_download(&job, cancellation).await,
Err(error) => Err(error),
};
if let Some(manager) = manager.upgrade() {
manager.remove_active(&job.key, &job.active);
}
job.active.complete(result);
}
async fn perform_download(job: &DownloadJob, cancellation: CancellationToken) -> DownloadResult {
let mut attempt = job.attempt;
loop {
cancellation.throw_if_cancellation_requested()?;
let progress: Box<dyn IProgress<HttpCapsClientProgressReport>> =
Box::new(ProgressFanout(Arc::downgrade(&job.active)));
let result = job
.client
.native_http_caps_client()
.get(job.address.clone(), cancellation.clone(), Some(progress))
.await;
match result {
Ok((response, data)) if response.is_success_status_code() => {
return Ok((response, data));
}
Ok((response, _)) => {
if is_permanent_status(response.status_code) || attempt >= job.retries {
return Err(Error::HttpRequest);
}
}
Err(Error::Cancelled) => return Err(Error::Cancelled),
Err(error) if attempt >= job.retries => return Err(error),
Err(_) => {}
}
attempt = attempt.saturating_add(1);
let base_delay = u64::try_from(200_i32.saturating_mul(attempt).min(2_000)).unwrap_or(2_000);
let delay =
Duration::from_millis(base_delay.saturating_add(retry_jitter(&job.key, attempt)));
tokio::select! {
() = tokio::time::sleep(delay) => {}
() = cancellation.cancelled() => return Err(Error::Cancelled),
}
}
}
fn is_permanent_status(status: u16) -> bool {
matches!(status, 401 | 403 | 404 | 410)
}
fn retry_jitter(key: &str, attempt: i32) -> u64 {
let hash = key.bytes().fold(
0xcbf2_9ce4_8422_2325_u64 ^ u64::try_from(attempt).unwrap_or_default(),
|hash, byte| hash.wrapping_mul(0x0000_0100_0000_01b3) ^ u64::from(byte),
);
hash % 200
}
fn validate_http_uri(uri: &Uri) -> Result<String, Error> {
let parsed = reqwest::Url::parse(&uri.0).map_err(|_| Error::Argument)?;
if !matches!(parsed.scheme(), "http" | "https") || parsed.host().is_none() {
return Err(Error::Argument);
}
Ok(parsed.to_string())
}

View File

@@ -6900,107 +6900,20 @@ pub enum CapsCategory {
pub use crate::network_manager::CapsEventDictionary; pub use crate::network_manager::CapsEventDictionary;
/// C# type: `T:LibreMetaverse.CapsRateLimiter`. /// C# type: `T:LibreMetaverse.CapsRateLimiter`.
pub struct CapsRateLimiter;
impl CapsRateLimiter {
/// C# member: `M:LibreMetaverse.CapsRateLimiter.#ctor`. /// C# member: `M:LibreMetaverse.CapsRateLimiter.#ctor`.
pub fn new_with_constructor() -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented("M:LibreMetaverse.CapsRateLimiter.#ctor")
}
/// C# member: `M:LibreMetaverse.CapsRateLimiter.#ctor(System.Collections.Generic.IReadOnlyDictionary{LibreMetaverse.CapsCategory,LibreMetaverse.CapsRateLimiterOptions})`. /// C# member: `M:LibreMetaverse.CapsRateLimiter.#ctor(System.Collections.Generic.IReadOnlyDictionary{LibreMetaverse.CapsCategory,LibreMetaverse.CapsRateLimiterOptions})`.
pub fn new_with_i_read_only_dictionary(
overrides: Option<
std::collections::HashMap<
libremetaverse::CapsCategory,
libremetaverse::CapsRateLimiterOptions,
>,
>,
) -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.CapsRateLimiter.#ctor(System.Collections.Generic.IReadOnlyDictionary{LibreMetaverse.CapsCategory,LibreMetaverse.CapsRateLimiterOptions})",
)
}
/// C# member: `M:LibreMetaverse.CapsRateLimiter.AcquireAsync(System.Uri,System.Threading.CancellationToken)`. /// C# member: `M:LibreMetaverse.CapsRateLimiter.AcquireAsync(System.Uri,System.Threading.CancellationToken)`.
pub async fn acquire(
&self,
uri: libremetaverse_types::compat::Uri,
cancellation_token: Option<libremetaverse_types::compat::CancellationToken>,
) -> Result<libremetaverse_types::compat::RateLimitLease, crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.CapsRateLimiter.AcquireAsync(System.Uri,System.Threading.CancellationToken)",
)
}
/// C# member: `M:LibreMetaverse.CapsRateLimiter.Dispose`. /// C# member: `M:LibreMetaverse.CapsRateLimiter.Dispose`.
pub fn dispose(&self) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented("M:LibreMetaverse.CapsRateLimiter.Dispose")
}
/// C# member: `M:LibreMetaverse.CapsRateLimiter.RegisterCapUri(System.String,System.Uri)`. /// C# member: `M:LibreMetaverse.CapsRateLimiter.RegisterCapUri(System.String,System.Uri)`.
pub fn register_cap_uri( pub use crate::caps_http::CapsRateLimiter;
&self,
cap_name: String,
uri: libremetaverse_types::compat::Uri,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.CapsRateLimiter.RegisterCapUri(System.String,System.Uri)",
)
}
}
/// C# type: `T:LibreMetaverse.CapsRateLimiterOptions`. /// C# type: `T:LibreMetaverse.CapsRateLimiterOptions`.
pub struct CapsRateLimiterOptions;
impl CapsRateLimiterOptions {
/// C# member: `M:LibreMetaverse.CapsRateLimiterOptions.#ctor`. /// C# member: `M:LibreMetaverse.CapsRateLimiterOptions.#ctor`.
pub fn new() -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented("M:LibreMetaverse.CapsRateLimiterOptions.#ctor")
}
/// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.QueueLimit`. /// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.QueueLimit`.
pub fn queue_limit(&self) -> i32 {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.QueueLimit"
)
}
/// Setter for C# member: `P:LibreMetaverse.CapsRateLimiterOptions.QueueLimit`.
pub fn set_queue_limit(&mut self, value: i32) {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.QueueLimit"
)
}
/// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.ReplenishmentPeriod`. /// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.ReplenishmentPeriod`.
pub fn replenishment_period(&self) -> std::time::Duration {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.ReplenishmentPeriod"
)
}
/// Setter for C# member: `P:LibreMetaverse.CapsRateLimiterOptions.ReplenishmentPeriod`.
pub fn set_replenishment_period(&mut self, value: std::time::Duration) {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.ReplenishmentPeriod"
)
}
/// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.TokenLimit`. /// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.TokenLimit`.
pub fn token_limit(&self) -> i32 {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.TokenLimit"
)
}
/// Setter for C# member: `P:LibreMetaverse.CapsRateLimiterOptions.TokenLimit`.
pub fn set_token_limit(&mut self, value: i32) {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.TokenLimit"
)
}
/// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.TokensPerPeriod`. /// C# member: `P:LibreMetaverse.CapsRateLimiterOptions.TokensPerPeriod`.
pub fn tokens_per_period(&self) -> i32 { pub use crate::caps_http::CapsRateLimiterOptions;
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.TokensPerPeriod"
)
}
/// Setter for C# member: `P:LibreMetaverse.CapsRateLimiterOptions.TokensPerPeriod`.
pub fn set_tokens_per_period(&mut self, value: i32) {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.CapsRateLimiterOptions.TokensPerPeriod"
)
}
}
/// C# type: `T:LibreMetaverse.ChannelType`. /// C# type: `T:LibreMetaverse.ChannelType`.
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
@@ -10839,11 +10752,13 @@ impl GridClient {
} }
/// C# member: `P:LibreMetaverse.GridClient.CapsRateLimiter`. /// C# member: `P:LibreMetaverse.GridClient.CapsRateLimiter`.
pub fn caps_rate_limiter(&self) -> libremetaverse::CapsRateLimiter { pub fn caps_rate_limiter(&self) -> libremetaverse::CapsRateLimiter {
libremetaverse_types::unimplemented_api!("P:LibreMetaverse.GridClient.CapsRateLimiter") /* native client-core implementation */
self.native_caps_rate_limiter()
} }
/// Setter for C# member: `P:LibreMetaverse.GridClient.CapsRateLimiter`. /// Setter for C# member: `P:LibreMetaverse.GridClient.CapsRateLimiter`.
pub fn set_caps_rate_limiter(&mut self, value: libremetaverse::CapsRateLimiter) { pub fn set_caps_rate_limiter(&mut self, value: libremetaverse::CapsRateLimiter) {
libremetaverse_types::unimplemented_api!("P:LibreMetaverse.GridClient.CapsRateLimiter") /* native client-core implementation */
self.native_set_caps_rate_limiter(value)
} }
/// C# member: `P:LibreMetaverse.GridClient.Directory`. /// C# member: `P:LibreMetaverse.GridClient.Directory`.
pub fn directory(&self) -> libremetaverse::DirectoryManager { pub fn directory(&self) -> libremetaverse::DirectoryManager {
@@ -10895,11 +10810,13 @@ impl GridClient {
} }
/// C# member: `P:LibreMetaverse.GridClient.HttpCapsClient`. /// C# member: `P:LibreMetaverse.GridClient.HttpCapsClient`.
pub fn http_caps_client(&self) -> libremetaverse::HttpCapsClient { pub fn http_caps_client(&self) -> libremetaverse::HttpCapsClient {
libremetaverse_types::unimplemented_api!("P:LibreMetaverse.GridClient.HttpCapsClient") /* native client-core implementation */
self.native_http_caps_client()
} }
/// Setter for C# member: `P:LibreMetaverse.GridClient.HttpCapsClient`. /// Setter for C# member: `P:LibreMetaverse.GridClient.HttpCapsClient`.
pub fn set_http_caps_client(&mut self, value: libremetaverse::HttpCapsClient) { pub fn set_http_caps_client(&mut self, value: libremetaverse::HttpCapsClient) {
libremetaverse_types::unimplemented_api!("P:LibreMetaverse.GridClient.HttpCapsClient") /* native client-core implementation */
self.native_set_http_caps_client(value)
} }
/// C# member: `P:LibreMetaverse.GridClient.InterestList`. /// C# member: `P:LibreMetaverse.GridClient.InterestList`.
pub fn interest_list(&self) -> libremetaverse::InterestListManager { pub fn interest_list(&self) -> libremetaverse::InterestListManager {
@@ -12830,406 +12747,39 @@ impl HomeInfo {
} }
/// C# type: `T:LibreMetaverse.HttpCapsClient`. /// C# type: `T:LibreMetaverse.HttpCapsClient`.
pub struct HttpCapsClient;
impl HttpCapsClient {
/// C# member: `F:LibreMetaverse.HttpCapsClient.HDR_LLSD_BINARY`. /// C# member: `F:LibreMetaverse.HttpCapsClient.HDR_LLSD_BINARY`.
pub fn hdr_llsd_binary() -> libremetaverse_types::compat::MediaTypeHeaderValue {
libremetaverse_types::unimplemented_api!("F:LibreMetaverse.HttpCapsClient.HDR_LLSD_BINARY")
}
/// C# member: `F:LibreMetaverse.HttpCapsClient.HDR_LLSD_JSON`. /// C# member: `F:LibreMetaverse.HttpCapsClient.HDR_LLSD_JSON`.
pub fn hdr_llsd_json() -> libremetaverse_types::compat::MediaTypeHeaderValue {
libremetaverse_types::unimplemented_api!("F:LibreMetaverse.HttpCapsClient.HDR_LLSD_JSON")
}
/// C# member: `F:LibreMetaverse.HttpCapsClient.HDR_LLSD_XML`. /// C# member: `F:LibreMetaverse.HttpCapsClient.HDR_LLSD_XML`.
pub fn hdr_llsd_xml() -> libremetaverse_types::compat::MediaTypeHeaderValue {
libremetaverse_types::unimplemented_api!("F:LibreMetaverse.HttpCapsClient.HDR_LLSD_XML")
}
/// C# member: `F:LibreMetaverse.HttpCapsClient.LLSD_BINARY`. /// C# member: `F:LibreMetaverse.HttpCapsClient.LLSD_BINARY`.
pub const LLSD_BINARY: &'static str = "application/llsd+binary";
/// C# member: `F:LibreMetaverse.HttpCapsClient.LLSD_JSON`. /// C# member: `F:LibreMetaverse.HttpCapsClient.LLSD_JSON`.
pub const LLSD_JSON: &'static str = "application/llsd+json";
/// C# member: `F:LibreMetaverse.HttpCapsClient.LLSD_XML`. /// C# member: `F:LibreMetaverse.HttpCapsClient.LLSD_XML`.
pub const LLSD_XML: &'static str = "application/llsd+xml";
/// C# member: `M:LibreMetaverse.HttpCapsClient.#ctor(System.Net.Http.HttpMessageHandler)`. /// C# member: `M:LibreMetaverse.HttpCapsClient.#ctor(System.Net.Http.HttpMessageHandler)`.
pub fn new(
handler: libremetaverse_types::compat::HttpMessageHandler,
) -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.#ctor(System.Net.Http.HttpMessageHandler)",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn delete_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.DeleteAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn delete_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.DeleteAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn delete_request_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.DeleteRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.DeleteRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn delete_request_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.DeleteRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.GetAsync(System.Uri,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.GetAsync(System.Uri,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn get(
&self,
uri: libremetaverse_types::compat::Uri,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.GetAsync(System.Uri,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.GetRequestAsync(System.Uri,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.GetRequestAsync(System.Uri,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn get_request(
&self,
uri: libremetaverse_types::compat::Uri,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.GetRequestAsync(System.Uri,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PatchAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PatchAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn patch_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PatchAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PatchAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PatchAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn patch_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PatchAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PatchRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PatchRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn patch_request_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PatchRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PatchRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PatchRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn patch_request_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PatchRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PostAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PostAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn post_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PostAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PostAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PostAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn post_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PostAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PostRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PostRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn post_request_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PostRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PostRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PostRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn post_request_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PostRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PutAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PutAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn put_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PutAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PutAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PutAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn put_with_uri_string_bytes_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PutAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PutRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PutRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn put_request_with_uri_osd_format_osd_cancellation_token_i_progress(
&self,
uri: libremetaverse_types::compat::Uri,
format: libremetaverse_structured_data::OSDFormat,
payload: libremetaverse_structured_data::OSD,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PutRequestAsync(System.Uri,LibreMetaverse.StructuredData.OSDFormat,LibreMetaverse.StructuredData.OSD,System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
/// C# member: `M:LibreMetaverse.HttpCapsClient.PutRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.PutRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub async fn put_request_with_uri_string_bytes_cancellation_token_i_progress( pub use crate::caps_http::HttpCapsClient;
&self,
uri: libremetaverse_types::compat::Uri,
content_type: String,
payload: Vec<u8>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
progress: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.PutRequestAsync(System.Uri,System.String,System.Byte[],System.Threading.CancellationToken,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
}
/// C# type: `T:LibreMetaverse.HttpCapsClient.ProgressReport`. /// C# type: `T:LibreMetaverse.HttpCapsClient.ProgressReport`.
pub struct HttpCapsClientProgressReport;
impl HttpCapsClientProgressReport {
/// C# member: `M:LibreMetaverse.HttpCapsClient.ProgressReport.#ctor(System.Nullable{System.Int64},System.Int64,System.Nullable{System.Double})`. /// C# member: `M:LibreMetaverse.HttpCapsClient.ProgressReport.#ctor(System.Nullable{System.Int64},System.Int64,System.Nullable{System.Double})`.
pub fn new(
total_bytes: Option<Option<i64>>,
bytes_transferred: i64,
percent: Option<Option<f64>>,
) -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.HttpCapsClient.ProgressReport.#ctor(System.Nullable{System.Int64},System.Int64,System.Nullable{System.Double})",
)
}
/// C# member: `P:LibreMetaverse.HttpCapsClient.ProgressReport.BytesTransferred`. /// C# member: `P:LibreMetaverse.HttpCapsClient.ProgressReport.BytesTransferred`.
pub fn bytes_transferred(&self) -> i64 {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.HttpCapsClient.ProgressReport.BytesTransferred"
)
}
/// C# member: `P:LibreMetaverse.HttpCapsClient.ProgressReport.Percent`. /// C# member: `P:LibreMetaverse.HttpCapsClient.ProgressReport.Percent`.
pub fn percent(&self) -> Option<Option<f64>> {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.HttpCapsClient.ProgressReport.Percent"
)
}
/// C# member: `P:LibreMetaverse.HttpCapsClient.ProgressReport.TotalBytes`. /// C# member: `P:LibreMetaverse.HttpCapsClient.ProgressReport.TotalBytes`.
pub fn total_bytes(&self) -> Option<Option<i64>> { pub use crate::caps_http::HttpCapsClientProgressReport;
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.HttpCapsClient.ProgressReport.TotalBytes"
)
}
}
/// C# type: `T:LibreMetaverse.IBakingTextureProvider`. /// C# type: `T:LibreMetaverse.IBakingTextureProvider`.
pub trait IBakingTextureProvider: std::any::Any { pub trait IBakingTextureProvider: std::any::Any {
@@ -32665,156 +32215,27 @@ pub mod formatters {
pub mod http { pub mod http {
/// C# type: `T:LibreMetaverse.Http.DownloadManager`. /// C# type: `T:LibreMetaverse.Http.DownloadManager`.
pub struct DownloadManager;
impl DownloadManager {
/// C# member: `M:LibreMetaverse.Http.DownloadManager.#ctor(LibreMetaverse.GridClient)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.#ctor(LibreMetaverse.GridClient)`.
pub fn new(client: libremetaverse::GridClient) -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.#ctor(LibreMetaverse.GridClient)",
)
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.Dispose`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.Dispose`.
pub fn dispose(&self) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented("M:LibreMetaverse.Http.DownloadManager.Dispose")
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.DownloadAsync(System.Uri,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.DownloadAsync(System.Uri,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken)`.
pub async fn download_with_uri_i_progress_cancellation_token(
&self,
address: libremetaverse_types::compat::Uri,
progress: Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.DownloadAsync(System.Uri,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken)",
)
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.DownloadAsync(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.DownloadAsync(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken)`.
pub async fn download_with_uri_string_i_progress_cancellation_token(
&self,
address: libremetaverse_types::compat::Uri,
content_type: String,
progress: Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
cancellation_token: libremetaverse_types::compat::CancellationToken,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.DownloadAsync(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken)",
)
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownload(LibreMetaverse.Http.DownloadRequest)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownload(LibreMetaverse.Http.DownloadRequest)`.
pub fn queue_download_with_download_request(
&self,
req: libremetaverse::http::DownloadRequest,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.QueueDownload(LibreMetaverse.Http.DownloadRequest)",
)
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownload(LibreMetaverse.Http.DownloadRequest,System.Threading.CancellationToken)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownload(LibreMetaverse.Http.DownloadRequest,System.Threading.CancellationToken)`.
pub fn queue_download_with_download_request_cancellation_token(
&self,
req: libremetaverse::http::DownloadRequest,
cancellation_token: libremetaverse_types::compat::CancellationToken,
) -> Result<(), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.QueueDownload(LibreMetaverse.Http.DownloadRequest,System.Threading.CancellationToken)",
)
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownloadAsync(LibreMetaverse.Http.DownloadRequest)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownloadAsync(LibreMetaverse.Http.DownloadRequest)`.
pub async fn queue_download_with_download_request_ba727387(
&self,
req: libremetaverse::http::DownloadRequest,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.QueueDownloadAsync(LibreMetaverse.Http.DownloadRequest)",
)
}
/// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownloadAsync(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken,System.Int32)`. /// C# member: `M:LibreMetaverse.Http.DownloadManager.QueueDownloadAsync(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken,System.Int32)`.
pub async fn queue_download_with_uri_string_i_progress_cancellation_token_int32(
&self,
address: libremetaverse_types::compat::Uri,
content_type: Option<String>,
progress_callback: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
cancellation_token: Option<libremetaverse_types::compat::CancellationToken>,
retries: Option<i32>,
) -> Result<(libremetaverse_types::compat::HttpResponse, Vec<u8>), crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadManager.QueueDownloadAsync(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport},System.Threading.CancellationToken,System.Int32)",
)
}
/// C# member: `P:LibreMetaverse.Http.DownloadManager.ParallelDownloads`. /// C# member: `P:LibreMetaverse.Http.DownloadManager.ParallelDownloads`.
pub fn parallel_downloads(&self) -> i32 { pub use crate::download_manager::DownloadManager;
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.Http.DownloadManager.ParallelDownloads"
)
}
/// Setter for C# member: `P:LibreMetaverse.Http.DownloadManager.ParallelDownloads`.
pub fn set_parallel_downloads(&mut self, value: i32) {
libremetaverse_types::unimplemented_api!(
"P:LibreMetaverse.Http.DownloadManager.ParallelDownloads"
)
}
}
/// C# type: `T:LibreMetaverse.Http.DownloadRequest`. /// C# type: `T:LibreMetaverse.Http.DownloadRequest`.
pub struct DownloadRequest {
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.Address`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.Address`.
pub address: libremetaverse_types::compat::Uri,
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.Attempt`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.Attempt`.
pub attempt: i32,
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.CancellationToken`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.CancellationToken`.
pub cancellation_token: libremetaverse_types::compat::CancellationToken,
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.CompletionTcs`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.CompletionTcs`.
pub completion_tcs: Option<
libremetaverse_types::compat::TaskCompletionSource<(
libremetaverse_types::compat::HttpResponse,
Vec<u8>,
)>,
>,
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.ContentType`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.ContentType`.
pub content_type: Option<String>,
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.DownloadProgressCallback`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.DownloadProgressCallback`.
pub download_progress_callback: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
/// C# member: `F:LibreMetaverse.Http.DownloadRequest.Retries`. /// C# member: `F:LibreMetaverse.Http.DownloadRequest.Retries`.
pub retries: i32,
}
impl DownloadRequest {
/// C# member: `M:LibreMetaverse.Http.DownloadRequest.#ctor(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`. /// C# member: `M:LibreMetaverse.Http.DownloadRequest.#ctor(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})`.
pub fn new( pub use crate::download_manager::DownloadRequest;
address: libremetaverse_types::compat::Uri,
content_type: Option<String>,
download_progress_callback: Option<
Box<
dyn libremetaverse_types::compat::IProgress<
libremetaverse::HttpCapsClientProgressReport,
>,
>,
>,
) -> Result<Self, crate::Error> {
libremetaverse_types::not_implemented(
"M:LibreMetaverse.Http.DownloadRequest.#ctor(System.Uri,System.String,System.IProgress{LibreMetaverse.HttpCapsClient.ProgressReport})",
)
}
}
/// C# type: `T:LibreMetaverse.Http.EventQueueClient`. /// C# type: `T:LibreMetaverse.Http.EventQueueClient`.
pub struct EventQueueClient { pub struct EventQueueClient {

View File

@@ -5,7 +5,9 @@ extern crate self as libremetaverse;
#[rustfmt::skip] // Deterministic machine output is formatted by the pinned generator. #[rustfmt::skip] // Deterministic machine output is formatted by the pinned generator.
mod attention_catalog; mod attention_catalog;
mod bit_pack; mod bit_pack;
mod caps_http;
mod client_core; mod client_core;
mod download_manager;
#[rustfmt::skip] // Deterministic machine output is formatted by the pinned generator. #[rustfmt::skip] // Deterministic machine output is formatted by the pinned generator.
mod foliage_catalog; mod foliage_catalog;
mod generated; mod generated;
@@ -198,6 +200,7 @@ impl imaging::Baker {
} }
} }
pub use caps_http::CapsHttpLimits;
pub use client_core::{ pub use client_core::{
ClientCoreError, ClientLifecycleState, ClientService, GridClientBuilder, ShutdownPhase, ClientCoreError, ClientLifecycleState, ClientService, GridClientBuilder, ShutdownPhase,
}; };

View File

@@ -0,0 +1,815 @@
use flate2::Compression;
use flate2::write::{GzEncoder, ZlibEncoder};
use libremetaverse::http::{DownloadManager, DownloadRequest};
use libremetaverse::{
CapsCategory, CapsHttpLimits, CapsRateLimiter, CapsRateLimiterOptions, GridClient,
HttpCapsClient, HttpCapsClientProgressReport,
};
use libremetaverse_structured_data::{OSD, OSDFormat};
use libremetaverse_types::Error;
use libremetaverse_types::compat::{
CancellationToken, CancellationTokenSource, HttpMessageHandler, HttpRequest, HttpResponse,
TaskCompletionSource, TimeProvider, Uri,
};
use std::collections::{BTreeMap, HashMap};
use std::io::Write as _;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime};
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
fn mutex<T>(value: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
value
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn response(status_code: u16, body: impl Into<Vec<u8>>) -> HttpResponse {
let body = body.into();
HttpResponse {
status_code,
headers: BTreeMap::from([("Content-Length".into(), body.len().to_string())]),
content_type: None,
body,
}
}
fn handler_returning(reply: HttpResponse) -> (Arc<Mutex<Vec<HttpRequest>>>, HttpMessageHandler) {
let requests = Arc::new(Mutex::new(Vec::new()));
let recorded = Arc::clone(&requests);
let handler = HttpMessageHandler::new(move |request, _cancellation| {
let recorded = Arc::clone(&recorded);
let reply = reply.clone();
async move {
mutex(&recorded).push(request);
reply
}
});
(requests, handler)
}
fn strict_limits(raw: usize, decompressed: usize) -> CapsHttpLimits {
CapsHttpLimits {
max_request_bytes: raw,
max_response_bytes: raw,
max_decompressed_bytes: decompressed,
max_redirects: 3,
}
}
fn one_token_options(queue_limit: i32) -> CapsRateLimiterOptions {
let mut options = CapsRateLimiterOptions::new().expect("options");
options.set_token_limit(1);
options.set_tokens_per_period(1);
options.set_replenishment_period(Duration::from_secs(1));
options.set_queue_limit(queue_limit);
options
}
#[tokio::test]
async fn byte_and_osd_requests_preserve_methods_content_types_and_payloads() {
let (requests, handler) = handler_returning(response(200, b"reply".to_vec()));
let client = HttpCapsClient::new(handler).expect("client");
let uri = Uri("https://caps.example.test/secret?token=do-not-log".into());
let (_, get_body) = client
.get(uri.clone(), CancellationToken::default(), None)
.await
.expect("GET");
assert_eq!(get_body, b"reply");
client
.post_with_uri_string_bytes_cancellation_token_i_progress(
uri.clone(),
"application/octet-stream".into(),
b"post".to_vec(),
CancellationToken::default(),
None,
)
.await
.expect("POST");
client
.put_with_uri_string_bytes_cancellation_token_i_progress(
uri.clone(),
"text/plain".into(),
b"put".to_vec(),
CancellationToken::default(),
None,
)
.await
.expect("PUT");
client
.patch_with_uri_string_bytes_cancellation_token_i_progress(
uri.clone(),
"text/plain".into(),
b"patch".to_vec(),
CancellationToken::default(),
None,
)
.await
.expect("PATCH");
client
.delete_with_uri_string_bytes_cancellation_token_i_progress(
uri.clone(),
"text/plain".into(),
b"delete".to_vec(),
CancellationToken::default(),
None,
)
.await
.expect("DELETE");
for (format, expected_content_type) in [
(OSDFormat::Xml, HttpCapsClient::LLSD_XML),
(OSDFormat::Binary, HttpCapsClient::LLSD_BINARY),
(OSDFormat::Json, HttpCapsClient::LLSD_JSON),
] {
client
.post_with_uri_osd_format_osd_cancellation_token_i_progress(
uri.clone(),
format,
OSD::String("payload".into()),
CancellationToken::default(),
None,
)
.await
.expect("OSD POST");
assert_eq!(
mutex(&requests)
.last()
.and_then(|request| request.content_type.as_deref()),
Some(expected_content_type)
);
assert!(!mutex(&requests).last().expect("request").body.is_empty());
}
let requests = mutex(&requests);
assert_eq!(
requests
.iter()
.take(5)
.map(|request| request.method.as_str())
.collect::<Vec<_>>(),
["GET", "POST", "PUT", "PATCH", "DELETE"]
);
assert_eq!(requests[1].body, b"post");
assert_eq!(
requests[1].content_type.as_deref(),
Some("application/octet-stream")
);
}
#[tokio::test]
async fn outbound_non_http_and_oversize_requests_are_rejected_before_the_handler() {
let (requests, handler) = handler_returning(response(200, Vec::new()));
let client =
HttpCapsClient::with_handler_and_limits(handler, strict_limits(4, 8)).expect("client");
let non_http = client
.get(
Uri("slcaps://opaque/secret".into()),
CancellationToken::default(),
None,
)
.await;
assert_eq!(non_http, Err(Error::Argument));
let oversized = client
.post_with_uri_string_bytes_cancellation_token_i_progress(
Uri("http://example.test/upload".into()),
"application/octet-stream".into(),
vec![0; 5],
CancellationToken::default(),
None,
)
.await;
assert_eq!(oversized, Err(Error::HttpRequest));
assert!(mutex(&requests).is_empty());
assert!(!format!("{client:?}").contains("opaque"));
assert!(!format!("{client:?}").contains("secret"));
let download = DownloadRequest::new(
Uri("https://caps.example.test/asset?token=download-secret".into()),
None,
None,
)
.expect("download request");
let debug = format!("{download:?}");
assert!(!debug.contains("caps.example.test"));
assert!(!debug.contains("download-secret"));
}
#[tokio::test]
async fn gzip_deflate_and_response_limits_are_enforced() {
let payload = b"compressed capability response".repeat(4);
let mut gzip = GzEncoder::new(Vec::new(), Compression::default());
gzip.write_all(&payload).expect("gzip write");
let gzip = gzip.finish().expect("gzip finish");
let mut gzip_reply = response(200, gzip);
gzip_reply
.headers
.insert("Content-Encoding".into(), "gzip".into());
let (_, gzip_handler) = handler_returning(gzip_reply);
let gzip_client =
HttpCapsClient::with_handler_and_limits(gzip_handler, strict_limits(256, 256))
.expect("gzip client");
let (_, decoded) = gzip_client
.get(
Uri("http://example.test/gzip".into()),
CancellationToken::default(),
None,
)
.await
.expect("gzip response");
assert_eq!(decoded, payload);
let mut deflate = ZlibEncoder::new(Vec::new(), Compression::default());
deflate.write_all(&payload).expect("deflate write");
let deflate = deflate.finish().expect("deflate finish");
let mut deflate_reply = response(200, deflate);
deflate_reply
.headers
.insert("Content-Encoding".into(), "deflate".into());
let (_, deflate_handler) = handler_returning(deflate_reply);
let deflate_client =
HttpCapsClient::with_handler_and_limits(deflate_handler, strict_limits(256, 256))
.expect("deflate client");
let (_, decoded) = deflate_client
.get(
Uri("http://example.test/deflate".into()),
CancellationToken::default(),
None,
)
.await
.expect("deflate response");
assert_eq!(decoded, payload);
let bomb = b"x".repeat(2_048);
let mut gzip = GzEncoder::new(Vec::new(), Compression::best());
gzip.write_all(&bomb).expect("bomb write");
let mut bomb_reply = response(200, gzip.finish().expect("bomb finish"));
bomb_reply
.headers
.insert("Content-Encoding".into(), "gzip".into());
let (_, bomb_handler) = handler_returning(bomb_reply);
let bomb_client =
HttpCapsClient::with_handler_and_limits(bomb_handler, strict_limits(128, 128))
.expect("bomb client");
assert_eq!(
bomb_client
.get(
Uri("http://example.test/bomb".into()),
CancellationToken::default(),
None,
)
.await,
Err(Error::HttpRequest)
);
let mut declared_too_large = response(200, b"short".to_vec());
declared_too_large
.headers
.insert("Content-Length".into(), "129".into());
let (_, length_handler) = handler_returning(declared_too_large);
let length_client =
HttpCapsClient::with_handler_and_limits(length_handler, strict_limits(128, 128))
.expect("length client");
assert_eq!(
length_client
.get(
Uri("http://example.test/declared".into()),
CancellationToken::default(),
None,
)
.await,
Err(Error::HttpRequest)
);
let mut unsupported_reply = response(200, b"opaque".to_vec());
unsupported_reply
.headers
.insert("Content-Encoding".into(), "br".into());
let (_, unsupported_handler) = handler_returning(unsupported_reply);
let unsupported_client =
HttpCapsClient::with_handler_and_limits(unsupported_handler, strict_limits(128, 128))
.expect("unsupported-encoding client");
assert_eq!(
unsupported_client
.get(
Uri("http://example.test/unsupported".into()),
CancellationToken::default(),
None,
)
.await,
Err(Error::HttpRequest)
);
}
#[tokio::test]
async fn upload_and_download_progress_reaches_exact_final_values() {
let body = vec![7_u8; 100_000];
let (_, handler) = handler_returning(response(200, body.clone()));
let reports = Arc::new(Mutex::new(Vec::new()));
let recorded = Arc::clone(&reports);
let client = HttpCapsClient::new(handler).expect("client");
client
.post_with_uri_string_bytes_cancellation_token_i_progress(
Uri("http://example.test/progress".into()),
"application/octet-stream".into(),
body.clone(),
CancellationToken::default(),
Some(Box::new(move |report: HttpCapsClientProgressReport| {
mutex(&recorded).push(report);
})),
)
.await
.expect("request");
let reports = mutex(&reports);
assert!(
reports.len() >= 4,
"upload and download should each report chunks"
);
let final_report = reports.last().expect("final report");
assert_eq!(final_report.bytes_transferred(), 100_000);
assert_eq!(final_report.total_bytes(), Some(Some(100_000)));
assert_eq!(final_report.percent(), Some(Some(100.0)));
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn http_policy_bounds_concurrency_and_grid_shutdown_cancels_inflight_requests() {
let active = Arc::new(AtomicUsize::new(0));
let maximum = Arc::new(AtomicUsize::new(0));
let handler_active = Arc::clone(&active);
let handler_maximum = Arc::clone(&maximum);
let handler = HttpMessageHandler::new(move |_request, _cancellation| {
let active = Arc::clone(&handler_active);
let maximum = Arc::clone(&handler_maximum);
async move {
let current = active.fetch_add(1, Ordering::AcqRel) + 1;
maximum.fetch_max(current, Ordering::AcqRel);
tokio::time::sleep(Duration::from_millis(25)).await;
active.fetch_sub(1, Ordering::AcqRel);
response(200, Vec::new())
}
});
let client = HttpCapsClient::with_handler_policy(handler, None, CapsHttpLimits::default(), 2)
.expect("bounded client");
let mut requests = Vec::new();
for index in 0..6 {
let client = client.clone();
requests.push(tokio::spawn(async move {
client
.get(
Uri(format!("http://example.test/bounded/{index}")),
CancellationToken::default(),
None,
)
.await
}));
}
for request in requests {
request.await.expect("request task").expect("request");
}
assert_eq!(maximum.load(Ordering::Acquire), 2);
let started = Arc::new(AtomicUsize::new(0));
let handler_started = Arc::clone(&started);
let handler = HttpMessageHandler::new(move |_request, cancellation| {
handler_started.fetch_add(1, Ordering::AcqRel);
async move {
cancellation.cancelled().await;
response(200, Vec::new())
}
});
let mut grid = GridClient::new().expect("grid client");
grid.set_http_caps_client(HttpCapsClient::new(handler).expect("HTTP client"));
let inflight_client = grid.http_caps_client();
let inflight = tokio::spawn(async move {
inflight_client
.get(
Uri("http://example.test/inflight".into()),
CancellationToken::default(),
None,
)
.await
});
while started.load(Ordering::Acquire) == 0 {
tokio::task::yield_now().await;
}
grid.dispose_with_method().expect("grid shutdown");
assert_eq!(
tokio::time::timeout(Duration::from_millis(100), inflight)
.await
.expect("shutdown wake")
.expect("request task"),
Err(Error::Cancelled)
);
}
#[tokio::test(start_paused = true)]
async fn rate_limiter_is_bounded_fifo_categorized_and_clock_injected() {
let seconds = Arc::new(AtomicU64::new(0));
let clock_seconds = Arc::clone(&seconds);
let clock = TimeProvider::from_fn(move || {
SystemTime::UNIX_EPOCH + Duration::from_secs(clock_seconds.load(Ordering::Acquire))
});
let options = one_token_options(2);
let limiter = CapsRateLimiter::new_with_clock_and_overrides(
clock,
Some(HashMap::from([
(CapsCategory::Default, options.clone()),
(CapsCategory::AssetFetch, options),
])),
)
.expect("limiter");
let default_uri = Uri("https://caps.example.test/default".into());
let asset_uri = Uri("https://caps.example.test/asset".into());
limiter
.register_cap_uri("GetTexture".into(), asset_uri.clone())
.expect("register category");
assert!(
limiter
.acquire(default_uri.clone(), None)
.await
.expect("default token")
.is_acquired()
);
assert!(
limiter
.acquire(asset_uri, None)
.await
.expect("isolated asset token")
.is_acquired()
);
let first_limiter = limiter.clone();
let first_uri = default_uri.clone();
let first = tokio::spawn(async move { first_limiter.acquire(first_uri, None).await });
tokio::task::yield_now().await;
let second_limiter = limiter.clone();
let second = tokio::spawn(async move { second_limiter.acquire(default_uri, None).await });
tokio::task::yield_now().await;
seconds.store(1, Ordering::Release);
tokio::time::advance(Duration::from_secs(1)).await;
tokio::task::yield_now().await;
assert!(
first.is_finished(),
"oldest waiter must receive the first token"
);
assert!(!second.is_finished(), "newer waiter must remain queued");
assert!(
first
.await
.expect("first task")
.expect("first lease")
.is_acquired()
);
seconds.store(2, Ordering::Release);
tokio::time::advance(Duration::from_secs(1)).await;
assert!(
second
.await
.expect("second task")
.expect("second lease")
.is_acquired()
);
}
#[tokio::test]
async fn rate_limiter_rejects_full_queue_and_wakes_cancelled_or_disposed_waiters() {
let no_queue = one_token_options(0);
let limiter = CapsRateLimiter::new_with_i_read_only_dictionary(Some(HashMap::from([(
CapsCategory::Default,
no_queue,
)])))
.expect("limiter");
let uri = Uri("http://example.test/limited".into());
assert!(
limiter
.acquire(uri.clone(), None)
.await
.expect("first")
.is_acquired()
);
assert!(
!limiter
.acquire(uri, None)
.await
.expect("rejected")
.is_acquired()
);
let queued = one_token_options(2);
let limiter = CapsRateLimiter::new_with_i_read_only_dictionary(Some(HashMap::from([(
CapsCategory::Default,
queued,
)])))
.expect("queued limiter");
let uri = Uri("http://example.test/queued".into());
limiter
.acquire(uri.clone(), None)
.await
.expect("consume token");
let cancellation = CancellationTokenSource::new();
let cancelled_limiter = limiter.clone();
let cancelled_uri = uri.clone();
let cancelled_token = cancellation.token();
let cancelled = tokio::spawn(async move {
cancelled_limiter
.acquire(cancelled_uri, Some(cancelled_token))
.await
});
tokio::task::yield_now().await;
cancellation.cancel();
assert_eq!(
tokio::time::timeout(Duration::from_millis(100), cancelled)
.await
.expect("cancel wake")
.expect("task"),
Err(Error::Cancelled)
);
let disposed_limiter = limiter.clone();
let disposed = tokio::spawn(async move { disposed_limiter.acquire(uri, None).await });
tokio::task::yield_now().await;
limiter.dispose().expect("dispose");
assert_eq!(
tokio::time::timeout(Duration::from_millis(100), disposed)
.await
.expect("dispose wake")
.expect("task"),
Err(Error::InvalidOperation)
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn download_manager_retries_transient_but_not_permanent_statuses() {
let attempts = Arc::new(AtomicUsize::new(0));
let handler_attempts = Arc::clone(&attempts);
let handler = HttpMessageHandler::new(move |_request, _cancellation| {
let attempt = handler_attempts.fetch_add(1, Ordering::AcqRel);
async move {
if attempt < 2 {
response(503, b"retry".to_vec())
} else {
response(200, b"done".to_vec())
}
}
});
let mut client = GridClient::new().expect("client");
client.set_http_caps_client(HttpCapsClient::new(handler).expect("HTTP client"));
let downloads = DownloadManager::new(client).expect("downloads");
let (_, body) = downloads
.queue_download_with_uri_string_i_progress_cancellation_token_int32(
Uri("http://example.test/transient".into()),
None,
None,
None,
Some(2),
)
.await
.expect("retried download");
assert_eq!(body, b"done");
assert_eq!(attempts.load(Ordering::Acquire), 3);
downloads.dispose().expect("dispose");
let attempts = Arc::new(AtomicUsize::new(0));
let handler_attempts = Arc::clone(&attempts);
let handler = HttpMessageHandler::new(move |_request, _cancellation| {
handler_attempts.fetch_add(1, Ordering::AcqRel);
async move { response(404, b"missing".to_vec()) }
});
let mut client = GridClient::new().expect("client");
client.set_http_caps_client(HttpCapsClient::new(handler).expect("HTTP client"));
let downloads = DownloadManager::new(client).expect("downloads");
assert_eq!(
downloads
.queue_download_with_uri_string_i_progress_cancellation_token_int32(
Uri("http://example.test/missing".into()),
None,
None,
None,
Some(5),
)
.await,
Err(Error::HttpRequest)
);
assert_eq!(attempts.load(Ordering::Acquire), 1);
downloads.dispose().expect("dispose");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn deduplicated_subscriber_cancellation_cancels_the_shared_download() {
let attempts = Arc::new(AtomicUsize::new(0));
let handler_attempts = Arc::clone(&attempts);
let handler = HttpMessageHandler::new(move |_request, cancellation| {
handler_attempts.fetch_add(1, Ordering::AcqRel);
async move {
cancellation.cancelled().await;
response(200, b"too late".to_vec())
}
});
let mut client = GridClient::new().expect("client");
client.set_http_caps_client(HttpCapsClient::new(handler).expect("HTTP client"));
let downloads = DownloadManager::new(client).expect("downloads");
let uri = Uri("http://example.test/shared".into());
let first_downloads = downloads.clone();
let first_uri = uri.clone();
let first = tokio::spawn(async move {
first_downloads
.queue_download_with_uri_string_i_progress_cancellation_token_int32(
first_uri,
None,
None,
None,
Some(0),
)
.await
});
let cancellation = CancellationTokenSource::new();
let second_downloads = downloads.clone();
let second_token = cancellation.token();
let second = tokio::spawn(async move {
second_downloads
.queue_download_with_uri_string_i_progress_cancellation_token_int32(
uri,
None,
None,
Some(second_token),
Some(0),
)
.await
});
while attempts.load(Ordering::Acquire) == 0 {
tokio::task::yield_now().await;
}
cancellation.cancel();
assert_eq!(first.await.expect("first task"), Err(Error::Cancelled));
assert_eq!(second.await.expect("second task"), Err(Error::Cancelled));
assert_eq!(attempts.load(Ordering::Acquire), 1);
downloads.dispose().expect("dispose");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn downloads_honor_parallel_limit_progress_and_external_completion_source() {
let active = Arc::new(AtomicUsize::new(0));
let maximum = Arc::new(AtomicUsize::new(0));
let handler_active = Arc::clone(&active);
let handler_maximum = Arc::clone(&maximum);
let handler = HttpMessageHandler::new(move |_request, _cancellation| {
let active = Arc::clone(&handler_active);
let maximum = Arc::clone(&handler_maximum);
async move {
let current = active.fetch_add(1, Ordering::AcqRel) + 1;
maximum.fetch_max(current, Ordering::AcqRel);
tokio::time::sleep(Duration::from_millis(30)).await;
active.fetch_sub(1, Ordering::AcqRel);
response(200, vec![1; 100_000])
}
});
let mut client = GridClient::new().expect("client");
client.set_http_caps_client(HttpCapsClient::new(handler).expect("HTTP client"));
let mut downloads = DownloadManager::new(client).expect("downloads");
downloads.set_parallel_downloads(2);
let reports = Arc::new(AtomicUsize::new(0));
let completion = TaskCompletionSource::new();
let mut request = DownloadRequest::new(
Uri("http://example.test/tcs".into()),
None,
Some(Box::new({
let reports = Arc::clone(&reports);
move |_report: HttpCapsClientProgressReport| {
reports.fetch_add(1, Ordering::AcqRel);
}
})),
)
.expect("request");
request.completion_tcs = Some(completion.clone());
downloads
.queue_download_with_download_request(request)
.expect("queue TCS request");
let mut tasks = Vec::new();
for index in 0..3 {
let downloads = downloads.clone();
tasks.push(tokio::spawn(async move {
downloads
.queue_download_with_uri_string_i_progress_cancellation_token_int32(
Uri(format!("http://example.test/{index}")),
None,
None,
None,
Some(0),
)
.await
}));
}
let (_, body) = completion.future().await.expect("TCS completion");
assert_eq!(body.len(), 100_000);
for task in tasks {
assert!(task.await.expect("download task").is_ok());
}
assert!(reports.load(Ordering::Acquire) >= 2);
assert_eq!(maximum.load(Ordering::Acquire), 2);
downloads.dispose().expect("dispose");
}
async fn read_http_request(stream: &mut tokio::net::TcpStream) -> Vec<u8> {
let mut request = Vec::new();
let mut scratch = [0_u8; 1024];
loop {
let count = stream.read(&mut scratch).await.expect("read request");
assert_ne!(count, 0, "client closed before completing request headers");
request.extend_from_slice(&scratch[..count]);
if request.windows(4).any(|window| window == b"\r\n\r\n") {
return request;
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn injected_reqwest_backend_follows_redirects_streams_and_decompresses() {
let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
.await
.expect("bind fake server");
let address = listener.local_addr().expect("server address");
let mut state = 0x1234_5678_u32;
let payload: Vec<u8> = (0..180_000)
.map(|_| {
state = state.wrapping_mul(1_664_525).wrapping_add(1_013_904_223);
state.to_be_bytes()[0]
})
.collect();
let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
encoder.write_all(&payload).expect("gzip write");
let compressed = encoder.finish().expect("gzip finish");
let server = tokio::spawn(async move {
let (mut first, _) = listener.accept().await.expect("first connection");
let first_request = read_http_request(&mut first).await;
assert!(first_request.starts_with(b"GET /start HTTP/1.1"));
first
.write_all(
b"HTTP/1.1 302 Found\r\nLocation: /final\r\nContent-Length: 0\r\nConnection: close\r\n\r\n",
)
.await
.expect("redirect response");
let (mut second, _) = listener.accept().await.expect("second connection");
let second_request = read_http_request(&mut second).await;
assert!(second_request.starts_with(b"GET /final HTTP/1.1"));
assert!(
String::from_utf8_lossy(&second_request)
.to_ascii_lowercase()
.contains("accept-encoding: gzip, deflate")
);
let headers = format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/llsd+binary\r\nContent-Encoding: gzip\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
compressed.len()
);
second.write_all(headers.as_bytes()).await.expect("headers");
for chunk in compressed.chunks(777) {
second.write_all(chunk).await.expect("body chunk");
tokio::task::yield_now().await;
}
});
let reqwest = reqwest::Client::builder()
.redirect(reqwest::redirect::Policy::limited(3))
.build()
.expect("reqwest client");
let client = HttpCapsClient::with_reqwest_client(
reqwest,
None,
CapsHttpLimits {
max_request_bytes: 1_024,
max_response_bytes: 256 * 1_024,
max_decompressed_bytes: 512 * 1_024,
max_redirects: 3,
},
)
.expect("caps client");
let reports = Arc::new(AtomicUsize::new(0));
let report_count = Arc::clone(&reports);
let (reply, body) = client
.get(
Uri(format!("http://{address}/start")),
CancellationToken::default(),
Some(Box::new(move |_report: HttpCapsClientProgressReport| {
report_count.fetch_add(1, Ordering::AcqRel);
})),
)
.await
.expect("redirected request");
assert_eq!(reply.status_code, 200);
assert_eq!(
reply.content_type.as_deref(),
Some("application/llsd+binary")
);
assert_eq!(body, payload);
assert!(
reports.load(Ordering::Acquire) > 1,
"streaming progress reports"
);
server.await.expect("server task");
}

85
docs/caps-http.md Normal file
View File

@@ -0,0 +1,85 @@
# Capability HTTP, rate limiting, and downloads
`HttpCapsClient`, `CapsRateLimiter`, and `Http::DownloadManager` are native Rust
implementations of the corresponding LibreMetaverse classes. They do not call
.NET code or use a helper process. A normal `GridClient` owns a `reqwest`
client using rustls, its category limiter, and their shutdown order. Tests and
embedders can inject either a `reqwest::Client` or an executor-neutral
`HttpMessageHandler`, plus a `TimeProvider` for deterministic token-bucket time.
## HTTP policy
Only absolute `http` and `https` request URIs are accepted. Response headers
are stored as text without interpreting `Location` or `Content-Location`, so a
successful response containing an opaque `slcaps://` value remains usable.
Production redirects are followed only while the next URI is HTTP(S) and the
configured redirect count has not been exhausted.
The default policy bounds request bodies at 32 MiB, wire response bodies at
64 MiB, decompressed bodies at 128 MiB, redirects at 10, and concurrent HTTP
requests at `Settings::max_http_connections()` (32). Declared and observed
response lengths are checked. Gzip and zlib/raw-deflate decoding reads through
a limiting adapter, preventing a small compressed response from expanding
beyond policy. Unsupported content encodings fail closed. Upload and download
progress follows the mapped C# `ProgressReport`, including two-decimal percent
values when a total length is known.
`CapsHttpLimits`, `HttpCapsClient::with_reqwest_client`, and
`HttpCapsClient::with_handler_policy` are the explicit injection seams. Fake
handlers receive the exact method, URI, content type, and request bytes and can
return deterministic status, headers, and bytes without external network
access. Grid-client shutdown cancels in-flight operations before disposing the
rate limiter.
## Rate categories
Each category owns a bounded, oldest-first token bucket. A full waiting queue
returns a non-acquired lease; the HTTP pipeline then proceeds, matching the C#
handler's overload behavior rather than dropping a request.
| Category | Burst | Tokens/second | Queue |
| --- | ---: | ---: | ---: |
| Default | 20 | 10 | 30 |
| RenderMaterials | 4 | 2 | 20 |
| AssetFetch | 24 | 12 | 60 |
| AssetUpload | 4 | 2 | 10 |
| Inventory | 6 | 3 | 20 |
| EventQueue | 3 | 2 | 3 |
| DisplayName | 5 | 2 | 15 |
| Voice | 10 | 5 | 20 |
`RegisterCapUri` applies the complete case-insensitive cap-name mapping from the
reference. Unregistered URIs use `Default`; texture/mesh, upload, inventory,
event queue, display-name/profile, material, and voice endpoints remain
isolated from one another. Cancellation removes queued tickets, and disposal
wakes all waiters immediately.
## Download behavior
The manager owns a 256-entry dispatch queue and permits 1 through 32 active
downloads (8 by default). Absolute URIs are canonicalized before deduplication.
All subscribers to the same active URI share one HTTP request, receive progress,
and complete with cloned response bytes. Cancellation by any subscriber cancels
that shared operation, as in the reference. Caller-provided
`TaskCompletionSource` values are completed by the fire-and-forget API.
HTTP 401, 403, 404, and 410 responses are permanent. Other unsuccessful
statuses and transport failures retry up to the request's configured retry
count with bounded backoff. Disposal cancels active and queued work, wakes every
waiter, joins the dispatcher, and prevents new requests.
Errors expose only typed categories and `Debug` output substitutes a redacted
marker for capability URIs. No request URL, query token, body, or response body
is logged.
Run the deterministic issue gate with:
```sh
cargo test -p libremetaverse --test caps_http
cargo test -p libremetaverse-compat-tests --test network_semantics queue_download
cargo test -p libremetaverse-compat-tests --test network_semantics put_response_with_non_http_location
```
The first suite uses injected handlers for all policy and concurrency cases and
a loopback-only fake HTTP server for the actual reqwest redirect/stream path. It
requires no live grid or Internet service.

File diff suppressed because it is too large Load Diff

View File

@@ -38,6 +38,12 @@ TARGETS = {
# implementations. The generated module keeps catalog markers and re-exports # implementations. The generated module keeps catalog markers and re-exports
# the hand-written type so coverage remains deterministic. # the hand-written type so coverage remains deterministic.
NATIVE_TYPES = { NATIVE_TYPES = {
"T:LibreMetaverse.CapsRateLimiter": "crate::caps_http::CapsRateLimiter",
"T:LibreMetaverse.CapsRateLimiterOptions": "crate::caps_http::CapsRateLimiterOptions",
"T:LibreMetaverse.Http.DownloadManager": "crate::download_manager::DownloadManager",
"T:LibreMetaverse.Http.DownloadRequest": "crate::download_manager::DownloadRequest",
"T:LibreMetaverse.HttpCapsClient": "crate::caps_http::HttpCapsClient",
"T:LibreMetaverse.HttpCapsClient.ProgressReport": "crate::caps_http::HttpCapsClientProgressReport",
"T:LibreMetaverse.Caps.EventQueueCallback": "crate::network_manager::CapsEventQueueCallback", "T:LibreMetaverse.Caps.EventQueueCallback": "crate::network_manager::CapsEventQueueCallback",
"T:LibreMetaverse.CapsEventDictionary": "crate::network_manager::CapsEventDictionary", "T:LibreMetaverse.CapsEventDictionary": "crate::network_manager::CapsEventDictionary",
"T:LibreMetaverse.DisconnectedEventArgs": "crate::network_manager::DisconnectedEventArgs", "T:LibreMetaverse.DisconnectedEventArgs": "crate::network_manager::DisconnectedEventArgs",
@@ -178,6 +184,14 @@ NATIVE_MEMBER_BODIES = {
"self.shutdown().map_err(Into::into)", "self.shutdown().map_err(Into::into)",
"M:LibreMetaverse.GridClient.ToString": "M:LibreMetaverse.GridClient.ToString":
"String::new()", "String::new()",
"P:LibreMetaverse.GridClient.CapsRateLimiter":
"self.native_caps_rate_limiter()",
"P:LibreMetaverse.GridClient.CapsRateLimiter#set":
"self.native_set_caps_rate_limiter(value)",
"P:LibreMetaverse.GridClient.HttpCapsClient":
"self.native_http_caps_client()",
"P:LibreMetaverse.GridClient.HttpCapsClient#set":
"self.native_set_http_caps_client(value)",
"E:LibreMetaverse.NetworkManager.Disconnected": "E:LibreMetaverse.NetworkManager.Disconnected":
"self.native_subscribe_disconnected(handler)", "self.native_subscribe_disconnected(handler)",
"E:LibreMetaverse.NetworkManager.GenericStreamingMessage": "E:LibreMetaverse.NetworkManager.GenericStreamingMessage":