use std::sync::{Arc, Mutex as StdMutex};
use std::time::{Duration, Instant};
use futures_util::StreamExt;
use reqwest::Method;
use serde::Deserialize;
use serde_json::{json, Value};
use wiremock::matchers::{body_json, method, path, query_param};
use wiremock::{Mock, MockServer, Request, ResponseTemplate};
use crate::client::{CallOptions, Client};
use crate::error::{Error, Result};
use crate::pagination::{paginate, ListParams, PageFetcher};
use crate::query::QueryBuilder;
use crate::test_support::{client, header, requests};
async fn only_request(server: &MockServer) -> Request {
requests(server)
.await
.into_iter()
.next()
.expect("one request was received")
}
#[tokio::test]
async fn sends_a_raw_jwt_with_no_bearer_prefix() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/v2/rooms"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.mount(&server)
.await;
let _: Value = client(&server)
.json(Method::GET, "/v2/rooms", CallOptions::new())
.await
.unwrap();
let request = only_request(&server).await;
let authorization = header(&request, "authorization").expect("authorization header");
assert!(
!authorization.to_lowercase().starts_with("bearer"),
"the API expects a bare JWT, got {authorization:?}"
);
assert!(authorization.starts_with("eyJ"), "{authorization:?}");
assert_eq!(authorization.split('.').count(), 3);
assert_eq!(
header(&request, "accept").as_deref(),
Some("application/json")
);
assert!(header(&request, "user-agent")
.unwrap()
.starts_with("videosdk-rs/"));
}
#[tokio::test]
async fn a_static_token_is_sent_verbatim() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.mount(&server)
.await;
let client = Client::builder()
.token("a-static-token")
.base_url(server.uri())
.build()
.unwrap();
let _: Value = client
.json(Method::GET, "/x", CallOptions::new())
.await
.unwrap();
let request = only_request(&server).await;
assert_eq!(
header(&request, "authorization").as_deref(),
Some("a-static-token")
);
}
#[tokio::test]
async fn a_client_header_replaces_the_sdk_default_rather_than_duplicating_it() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.mount(&server)
.await;
let client = Client::builder()
.token("sdk-token")
.base_url(server.uri())
.header("user-agent", "my-app/1.0")
.header("accept", "application/vnd.custom+json")
.header("authorization", "Bearer mine")
.build()
.unwrap();
let _: Value = client
.json(Method::GET, "/x", CallOptions::new())
.await
.unwrap();
let request = only_request(&server).await;
for (name, expected) in [
("user-agent", "my-app/1.0"),
("accept", "application/vnd.custom+json"),
("authorization", "Bearer mine"),
] {
let values: Vec<_> = request.headers.get_all(name).iter().collect();
assert_eq!(values.len(), 1, "{name} was sent {} times", values.len());
assert_eq!(values[0], expected, "{name}");
}
}
#[tokio::test]
async fn a_per_call_header_wins_over_a_client_header() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.mount(&server)
.await;
let client = Client::builder()
.token("t")
.base_url(server.uri())
.header("x-tenant", "client-level")
.build()
.unwrap();
let _ = client
.api()
.get(
"/x",
crate::RawRequest::new().header("x-tenant", "call-level"),
)
.await
.unwrap();
let request = only_request(&server).await;
let values: Vec<_> = request.headers.get_all("x-tenant").iter().collect();
assert_eq!(values.len(), 1);
assert_eq!(values[0], "call-level");
}
#[tokio::test]
async fn an_invalid_header_is_a_validation_error() {
let server = MockServer::start().await;
let err = client(&server)
.api()
.get("/x", crate::RawRequest::new().header("bad header", "v"))
.await
.unwrap_err();
assert!(err.is_validation(), "got {err:?}");
assert!(requests(&server).await.is_empty());
}
#[tokio::test]
async fn a_text_response_asks_for_any_content_type() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_string("a,b\n1,2\n"))
.mount(&server)
.await;
let csv = client(&server)
.text(Method::GET, "/export", CallOptions::new())
.await
.unwrap();
assert_eq!(csv, "a,b\n1,2\n");
assert_eq!(
header(&only_request(&server).await, "accept").as_deref(),
Some("*/*")
);
}
#[tokio::test]
async fn a_raw_body_is_sent_verbatim_with_the_caller_s_content_type() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200).set_body_string("v=0\r\n"))
.mount(&server)
.await;
let answer = client(&server)
.api()
.post(
"/v2/whip",
crate::RawRequest::new()
.raw_body(b"v=0\r\no=- 0 0 IN IP4 0.0.0.0\r\n".to_vec())
.header("content-type", "application/sdp")
.expect(crate::Expect::Text),
)
.await
.unwrap();
assert_eq!(answer, json!("v=0\r\n"));
let request = only_request(&server).await;
assert_eq!(
header(&request, "content-type").as_deref(),
Some("application/sdp")
);
assert_eq!(request.body, b"v=0\r\no=- 0 0 IN IP4 0.0.0.0\r\n");
}
#[tokio::test]
async fn builder_requires_credentials() {
static ENV_LOCK: StdMutex<()> = StdMutex::new(());
let _guard = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner());
let saved: Vec<_> = ["VIDEOSDK_API_KEY", "VIDEOSDK_SECRET"]
.iter()
.map(|name| (*name, std::env::var(name).ok()))
.collect();
for (name, _) in &saved {
std::env::remove_var(name);
}
let err = Client::builder().build().unwrap_err();
assert!(matches!(err, Error::Configuration(_)));
assert!(err.to_string().contains("VIDEOSDK_API_KEY"));
assert!(Client::builder().api_key("k").build().is_err());
assert!(Client::builder().token("t").build().is_ok());
std::env::set_var("VIDEOSDK_API_KEY", "from-env");
std::env::set_var("VIDEOSDK_SECRET", "from-env");
assert!(
Client::new().is_ok(),
"credentials come from the environment"
);
for (name, value) in saved {
match value {
Some(value) => std::env::set_var(name, value),
None => std::env::remove_var(name),
}
}
}
#[tokio::test]
async fn base_url_trailing_slashes_are_trimmed() {
let client = Client::builder()
.token("t")
.base_url("https://example.com///")
.build()
.unwrap();
assert_eq!(client.base_url(), "https://example.com");
}
#[tokio::test]
async fn retries_a_rate_limit_for_a_non_idempotent_method() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(429))
.up_to_n_times(1)
.expect(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"ok": true})))
.expect(1)
.mount(&server)
.await;
let value: Value = client(&server)
.json(Method::POST, "/x", CallOptions::new())
.await
.unwrap();
assert_eq!(value["ok"], json!(true));
}
#[tokio::test]
async fn retries_a_server_error_only_for_idempotent_methods() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(503))
.up_to_n_times(2)
.expect(2)
.mount(&server)
.await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.expect(1)
.mount(&server)
.await;
let _: Value = client(&server)
.json(Method::GET, "/x", CallOptions::new())
.await
.unwrap();
drop(server);
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(503))
.expect(1)
.mount(&server)
.await;
let err = client(&server)
.json::<Value>(Method::POST, "/x", CallOptions::new())
.await
.unwrap_err();
assert_eq!(err.status(), Some(503));
}
#[tokio::test]
async fn stops_after_max_retries_and_honors_a_per_call_override() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(500))
.expect(4)
.mount(&server)
.await;
let client = client(&server);
assert!(client
.json::<Value>(Method::GET, "/x", CallOptions::new())
.await
.is_err());
assert!(client
.with_max_retries(0)
.json::<Value>(Method::GET, "/x", CallOptions::new())
.await
.is_err());
}
#[tokio::test]
async fn never_retries_client_errors() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(404).set_body_json(json!({"message": "gone"})))
.expect(1)
.mount(&server)
.await;
let err = client(&server)
.json::<Value>(Method::GET, "/x", CallOptions::new())
.await
.unwrap_err();
assert!(err.is_not_found());
}
#[tokio::test]
async fn honors_retry_after_on_a_rate_limit() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(429).insert_header("retry-after", "1"))
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.mount(&server)
.await;
let started = Instant::now();
let _: Value = client(&server)
.json(Method::GET, "/x", CallOptions::new())
.await
.unwrap();
assert!(
started.elapsed() >= Duration::from_millis(900),
"expected to wait for Retry-After, waited {:?}",
started.elapsed()
);
}
#[tokio::test]
async fn normalizes_an_error_response() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(
ResponseTemplate::new(409)
.insert_header("x-request-id", "req-123")
.set_body_json(json!({"message": "already exists", "code": "duplicate"})),
)
.mount(&server)
.await;
let err = client(&server)
.json::<Value>(Method::GET, "/v2/rooms", CallOptions::new())
.await
.unwrap_err();
assert!(err.is_conflict());
assert_eq!(err.status(), Some(409));
assert_eq!(err.code(), Some("duplicate"));
assert_eq!(err.request_id(), Some("req-123"));
let api = err.as_api().unwrap();
assert_eq!(api.message, "already exists");
assert_eq!(api.method, "GET");
assert_eq!(api.path, "/v2/rooms");
assert!(api.details.is_some());
}
#[tokio::test]
async fn a_payment_required_response_is_a_not_found() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(402))
.mount(&server)
.await;
let err = client(&server)
.json::<Value>(Method::GET, "/x", CallOptions::new())
.await
.unwrap_err();
assert!(
err.is_not_found(),
"the API overloads 402 to mean not-found"
);
}
#[tokio::test]
async fn a_non_json_error_body_still_yields_a_message() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(502).set_body_string("<html>bad gateway</html>"))
.mount(&server)
.await;
let err = client(&server)
.with_max_retries(0)
.json::<Value>(Method::GET, "/x", CallOptions::new())
.await
.unwrap_err();
assert_eq!(err.as_api().unwrap().message, "<html>bad gateway</html>");
}
#[tokio::test]
async fn a_client_timeout_is_reported_as_a_timeout() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(ResponseTemplate::new(200).set_delay(Duration::from_secs(5)))
.mount(&server)
.await;
let err = client(&server)
.with_request_timeout(Duration::from_millis(50))
.json::<Value>(Method::POST, "/x", CallOptions::new())
.await
.unwrap_err();
assert!(err.is_timeout(), "got {err:?}");
}
#[derive(Debug, Deserialize, PartialEq)]
struct Room {
id: String,
}
#[tokio::test]
async fn json_decodes_a_flat_body() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"id": "room-1"})))
.mount(&server)
.await;
let room: Room = client(&server)
.json(Method::GET, "/x", CallOptions::new())
.await
.unwrap();
assert_eq!(room.id, "room-1");
}
#[tokio::test]
async fn data_unwraps_a_data_envelope() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"data": {"id": "room-1"}})))
.mount(&server)
.await;
let room: Room = client(&server)
.data(Method::GET, "/x", CallOptions::new())
.await
.unwrap();
assert_eq!(room.id, "room-1");
}
#[tokio::test]
async fn wrapped_unwraps_a_named_envelope_and_tolerates_a_missing_key() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(path("/present"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"alertRule": {"id": "r-1"}})))
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/absent"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"other": 1})))
.mount(&server)
.await;
let client = client(&server);
let room: Room = client
.wrapped(Method::GET, "/present", "alertRule", CallOptions::new())
.await
.unwrap();
assert_eq!(room.id, "r-1");
let missing: Option<Room> = client
.wrapped(Method::GET, "/absent", "alertRule", CallOptions::new())
.await
.unwrap();
assert_eq!(missing, None);
}
#[tokio::test]
async fn message_unquotes_a_json_string_and_passes_through_text() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/quoted"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!("Room ended")))
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/plain"))
.respond_with(ResponseTemplate::new(200).set_body_string(" Room ended "))
.mount(&server)
.await;
let client = client(&server);
assert_eq!(
client
.message(Method::POST, "/quoted", CallOptions::new())
.await
.unwrap(),
"Room ended"
);
assert_eq!(
client
.message(Method::POST, "/plain", CallOptions::new())
.await
.unwrap(),
"Room ended"
);
}
#[tokio::test]
async fn text_returns_the_raw_body() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_string("a,b\n1,2\n"))
.mount(&server)
.await;
let csv = client(&server)
.text(Method::GET, "/export", CallOptions::new())
.await
.unwrap();
assert_eq!(csv, "a,b\n1,2\n");
}
#[tokio::test]
async fn none_discards_the_body_and_a_204_decodes_as_null() {
let server = MockServer::start().await;
Mock::given(method("DELETE"))
.respond_with(ResponseTemplate::new(204))
.mount(&server)
.await;
let client = client(&server);
client
.none(Method::DELETE, "/x", CallOptions::new())
.await
.unwrap();
let nothing: Option<Room> = client
.json(Method::DELETE, "/x", CallOptions::new())
.await
.unwrap();
assert_eq!(nothing, None);
let err = client
.json::<Room>(Method::DELETE, "/x", CallOptions::new())
.await
.unwrap_err();
assert!(matches!(err, Error::Decode { .. }));
}
#[tokio::test]
async fn a_malformed_body_is_a_decode_error() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_string("not json"))
.mount(&server)
.await;
let err = client(&server)
.json::<Room>(Method::GET, "/v2/rooms", CallOptions::new())
.await
.unwrap_err();
assert!(matches!(err, Error::Decode { .. }));
assert!(err.to_string().contains("GET /v2/rooms"));
}
#[tokio::test]
async fn sends_a_json_body_with_a_content_type() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(body_json(json!({"name": "standup"})))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.expect(1)
.mount(&server)
.await;
let options = CallOptions::json(json!({"name": "standup"})).unwrap();
let _: Value = client(&server)
.json(Method::POST, "/v2/rooms", options)
.await
.unwrap();
let request = only_request(&server).await;
assert_eq!(
header(&request, "content-type").as_deref(),
Some("application/json")
);
}
#[tokio::test]
async fn sends_query_parameters_omitting_absent_ones() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(query_param("page", "2"))
.and(query_param("kinds", "audio,video"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({})))
.expect(1)
.mount(&server)
.await;
let query = QueryBuilder::new()
.opt("page", Some(2))
.opt_str("search", None::<&str>)
.csv("kinds", &["audio", "video"])
.into_pairs();
let _: Value = client(&server)
.json(Method::GET, "/v2/rooms", CallOptions::new().query(query))
.await
.unwrap();
let request = only_request(&server).await;
assert!(!request.url.query().unwrap().contains("search"));
}
#[tokio::test]
async fn put_binary_omits_auth_and_retries() {
let server = MockServer::start().await;
Mock::given(method("PUT"))
.respond_with(ResponseTemplate::new(500))
.up_to_n_times(1)
.expect(1)
.mount(&server)
.await;
Mock::given(method("PUT"))
.respond_with(ResponseTemplate::new(200))
.expect(1)
.mount(&server)
.await;
let url = format!("{}/bucket/key?X-Amz-Signature=secret", server.uri());
client(&server)
.put_binary(&url, b"a,b\n", "text/csv")
.await
.unwrap();
let request = only_request(&server).await;
assert!(
header(&request, "authorization").is_none(),
"sending auth would invalidate the presigned signature"
);
assert_eq!(
header(&request, "content-type").as_deref(),
Some("text/csv")
);
}
#[tokio::test]
async fn put_binary_keeps_the_signature_out_of_errors() {
let server = MockServer::start().await;
Mock::given(method("PUT"))
.respond_with(ResponseTemplate::new(403))
.mount(&server)
.await;
let url = format!("{}/bucket/key?X-Amz-Signature=secret", server.uri());
let err = client(&server)
.put_binary(&url, b"x", "text/csv")
.await
.unwrap_err();
assert_eq!(err.code(), Some("upload_failed"));
let rendered = err.to_string();
assert!(
!rendered.contains("secret"),
"leaked the signature: {rendered}"
);
assert!(!err.as_api().unwrap().path.contains('?'));
}
#[tokio::test]
async fn the_raw_api_facade_applies_auth() {
let server = MockServer::start().await;
Mock::given(method("PATCH"))
.and(path("/v2/anything"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({"ok": 1})))
.mount(&server)
.await;
let client = client(&server);
let value = client
.api()
.patch(
"/v2/anything",
crate::RawRequest::new().body(json!({"a": 1})),
)
.await
.unwrap();
assert_eq!(value["ok"], json!(1));
assert!(header(&only_request(&server).await, "authorization").is_some());
}
#[derive(Debug, Deserialize)]
struct Item {
id: u32,
}
fn item_fetcher(client: Client) -> PageFetcher {
Arc::new(move |page, per_page| {
let client = client.clone();
Box::pin(async move {
let query = QueryBuilder::new()
.opt("page", page)
.opt("perPage", per_page)
.into_pairs();
client
.json::<Value>(Method::GET, "/items", CallOptions::new().query(query))
.await
})
})
}
#[tokio::test]
async fn a_page_walks_to_the_next_page() {
let server = MockServer::start().await;
let body = |page: u32, id: u32| {
json!({
"pageInfo": {"currentPage": page, "perPage": 1, "lastPage": 2, "total": 2},
"data": [{"id": id}],
})
};
Mock::given(method("GET"))
.and(query_param("page", "2"))
.respond_with(ResponseTemplate::new(200).set_body_json(body(2, 2)))
.mount(&server)
.await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(body(1, 1)))
.mount(&server)
.await;
let first: crate::Page<Item> = paginate(
item_fetcher(client(&server)),
&ListParams::default(),
"data",
None,
)
.await
.unwrap();
assert_eq!(first.data[0].id, 1);
assert_eq!(first.total(), 2);
assert!(first.has_next_page());
let second = first.next_page().await.unwrap().expect("a second page");
assert_eq!(second.data[0].id, 2);
assert!(!second.has_next_page());
assert!(second.next_page().await.unwrap().is_none());
}
#[tokio::test]
async fn a_page_streams_every_item_across_pages() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(query_param("page", "2"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"pageInfo": {"currentPage": 2, "perPage": 2, "lastPage": 2, "total": 3},
"data": [{"id": 3}],
})))
.mount(&server)
.await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"pageInfo": {"currentPage": 1, "perPage": 2, "lastPage": 2, "total": 3},
"data": [{"id": 1}, {"id": 2}],
})))
.mount(&server)
.await;
let page: crate::Page<Item> = paginate(
item_fetcher(client(&server)),
&ListParams::default(),
"data",
None,
)
.await
.unwrap();
let ids: Vec<u32> = page
.into_stream()
.map(|item| item.unwrap().id)
.collect()
.await;
assert_eq!(ids, [1, 2, 3]);
}
#[tokio::test]
async fn a_stream_surfaces_a_fetch_error_once_and_stops() {
let server = MockServer::start().await;
Mock::given(method("GET"))
.and(query_param("page", "2"))
.respond_with(ResponseTemplate::new(500))
.mount(&server)
.await;
Mock::given(method("GET"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({
"pageInfo": {"currentPage": 1, "perPage": 1, "lastPage": 2, "total": 2},
"data": [{"id": 1}],
})))
.mount(&server)
.await;
let page: crate::Page<Item> = paginate(
item_fetcher(client(&server).with_max_retries(0)),
&ListParams::default(),
"data",
None,
)
.await
.unwrap();
let results: Vec<Result<Item>> = page.into_stream().collect().await;
assert_eq!(results.len(), 2);
assert!(results[0].is_ok());
assert_eq!(results[1].as_ref().unwrap_err().status(), Some(500));
}