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(value: &Mutex) -> std::sync::MutexGuard<'_, T> { value .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) } fn response(status_code: u16, body: impl Into>) -> 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>>, 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::>(), ["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"); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn download_queue_applies_backpressure_and_cancellation_drains_every_job() { let active = Arc::new(AtomicUsize::new(0)); let handler_active = Arc::clone(&active); let handler = HttpMessageHandler::new(move |_request, cancellation| { let active = Arc::clone(&handler_active); async move { active.fetch_add(1, Ordering::AcqRel); cancellation.cancelled().await; active.fetch_sub(1, Ordering::AcqRel); response(200, Vec::new()) } }); 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(1); downloads .queue_download_with_download_request( DownloadRequest::new( Uri("http://example.test/backpressure/initial".into()), None, None, ) .expect("request"), ) .expect("initial request"); tokio::time::timeout(Duration::from_secs(2), async { while active.load(Ordering::Acquire) == 0 { tokio::task::yield_now().await; } }) .await .expect("initial request started"); let mut rejected = 0; for index in 0..300 { let result = downloads.queue_download_with_download_request( DownloadRequest::new( Uri(format!("http://example.test/backpressure/{index}")), None, None, ) .expect("request"), ); match result { Ok(()) => {} Err(Error::InvalidOperation) => rejected += 1, Err(error) => panic!("unexpected queue error: {error:?}"), } } assert!(rejected > 0, "the bounded queue must reject excess work"); assert!(downloads.active_download_count() <= 258); downloads.dispose().expect("dispose"); assert_eq!(active.load(Ordering::Acquire), 0); assert_eq!(downloads.active_download_count(), 0); assert!(!downloads.dispatcher_running()); assert!(downloads.is_disposed()); } async fn read_http_request(stream: &mut tokio::net::TcpStream) -> Vec { 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; } } } #[test] fn native_transport_rejects_invalid_resource_policy_before_building() { assert!(matches!( HttpCapsClient::with_native_transport( "MetaCrate-caps-test/0.0.1", Duration::ZERO, 1, None, CapsHttpLimits::default(), ), Err(Error::Argument) )); assert!(matches!( HttpCapsClient::with_native_transport( "MetaCrate-caps-test/0.0.1", Duration::from_secs(1), 0, None, CapsHttpLimits::default(), ), Err(Error::Argument) )); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn native_http_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 = (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 client = HttpCapsClient::with_native_transport( "MetaCrate-caps-test/0.0.1", Duration::from_secs(5), 4, 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"); }