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
816 lines
28 KiB
Rust
816 lines
28 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_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");
|
|
}
|