Files
MetaCrate/crates/libremetaverse/tests/caps_http.rs
Chili Palmer 52f62d8c39
Some checks failed
Native code generation / deterministic (push) Failing after 2m18s
Imaging and meshing gate / native (push) Failing after 1m30s
JPEG 2000 feature / linux (push) Successful in 2m40s
Native Rust workspace compile / compile (push) Failing after 57s
Skia feature / linux (push) Successful in 31m8s
Implement native asset pipeline and cache (#64)
2026-08-10 07:51:17 +00:00

823 lines
29 KiB
Rust

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_is_independent() {
let attempts = Arc::new(AtomicUsize::new(0));
let handler_attempts = Arc::clone(&attempts);
let release = Arc::new(tokio::sync::Notify::new());
let handler_release = Arc::clone(&release);
let handler = HttpMessageHandler::new(move |_request, _cancellation| {
handler_attempts.fetch_add(1, Ordering::AcqRel);
let release = Arc::clone(&handler_release);
async move {
release.notified().await;
response(200, b"shared".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!(second.await.expect("second task"), Err(Error::Cancelled));
release.notify_waiters();
assert_eq!(
first.await.expect("first task").expect("first download").1,
b"shared"
);
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");
}