#![cfg(all(feature = "macros", feature = "json"))]
use std::{num::NonZeroUsize, time::Duration};
use kynos::{
Router,
http::{Method, StatusCode, header},
middleware::limits::{BodySize, Concurrency, Timeout},
response::status::NoContent,
};
#[path = "support/mod.rs"]
mod support;
use support::{App, User, get, send};
fn one() -> NonZeroUsize {
NonZeroUsize::new(1).expect("one is not zero")
}
#[tokio::test]
async fn a_body_past_the_limit_is_refused_with_the_status_its_type_declares() {
let service = support::router()
.intercept(BodySize::new(16))
.build(App::new())
.expect("a describable router");
let reply = support::post(&service, "/users")
.json(&User {
id: 1,
name: "a name comfortably longer than sixteen bytes".to_owned(),
})
.call()
.await;
assert_eq!(reply.status, StatusCode::PAYLOAD_TOO_LARGE);
assert!(reply.text().contains("16"), "{}", reply.text());
}
#[tokio::test]
async fn a_body_within_the_limit_reaches_its_operation() {
let service = support::router()
.intercept(BodySize::new(4096))
.build(App::new())
.expect("a describable router");
let reply = support::post(&service, "/users")
.json(&User {
id: 1,
name: "fresh".to_owned(),
})
.call()
.await;
assert_eq!(reply.status, StatusCode::CREATED);
}
#[tokio::test]
async fn a_declared_length_past_the_limit_is_refused_without_reading_the_body() {
let service = support::router()
.intercept(BodySize::new(8))
.build(App::new())
.expect("a describable router");
let reply = support::post(&service, "/users")
.header("content-type", "application/json")
.header("content-length", "4096")
.body(&b"{}"[..])
.call()
.await;
assert_eq!(reply.status, StatusCode::PAYLOAD_TOO_LARGE);
}
#[test]
fn a_body_limit_declares_its_status_on_every_operation_it_covers() {
let document = support::router()
.intercept(BodySize::new(4096))
.openapi()
.expect("a describable router");
for (path, item) in &document.paths.items {
for (method, operation) in item.operations() {
assert!(
operation.responses.responses.contains_key("413"),
"{method} {path} is covered by a body limit and does not declare its 413"
);
}
}
}
#[kynos::get("/slow")]
async fn slow() -> NoContent {
tokio::time::sleep(Duration::from_millis(400)).await;
NoContent
}
#[kynos::get("/prompt")]
async fn prompt() -> NoContent {
NoContent
}
#[tokio::test]
async fn a_handler_past_the_limit_is_answered_with_the_status_its_type_declares() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Timeout::new(Duration::from_millis(20)))
.build(())
.expect("a describable router");
let timed_out = get(&service, "/slow").call().await;
assert_eq!(timed_out.status, StatusCode::REQUEST_TIMEOUT);
let in_time = get(&service, "/prompt").call().await;
assert_eq!(in_time.status, StatusCode::NO_CONTENT);
}
#[tokio::test]
async fn a_request_past_the_concurrency_limit_is_refused_while_the_first_runs() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Concurrency::new(one()))
.build(())
.expect("a describable router");
let (held, refused) = tokio::join!(get(&service, "/slow").call(), async {
tokio::time::sleep(Duration::from_millis(50)).await;
get(&service, "/prompt").call().await
});
assert_eq!(held.status, StatusCode::NO_CONTENT);
assert_eq!(refused.status, StatusCode::SERVICE_UNAVAILABLE);
assert!(refused.field(header::RETRY_AFTER.as_str()).is_none());
}
#[tokio::test]
async fn a_released_slot_is_available_to_the_next_request() {
let service = Router::<()>::new()
.mount(kynos::routes![prompt])
.intercept(Concurrency::new(one()))
.build(())
.expect("a describable router");
for _ in 0..3 {
assert_eq!(
get(&service, "/prompt").call().await.status,
StatusCode::NO_CONTENT,
"a slot was not released when its request finished"
);
}
}
#[tokio::test]
async fn a_limit_does_not_answer_for_a_route_that_does_not_exist() {
let service = support::router()
.intercept(Timeout::new(Duration::from_secs(30)))
.intercept(BodySize::new(4096))
.build(App::new())
.expect("a describable router");
let reply = send(&service, Method::GET, "/nothing-here").call().await;
assert_eq!(reply.status, StatusCode::NOT_FOUND);
}
#[tokio::test]
async fn a_service_with_no_body_limit_accepts_a_body_of_any_size() {
let service = support::router()
.build(App::new())
.expect("a describable router");
let reply = support::post(&service, "/users")
.json(&User {
id: 1,
name: "n".repeat(64 * 1024),
})
.call()
.await;
assert_ne!(
reply.status,
StatusCode::PAYLOAD_TOO_LARGE,
"no limit was mounted, so nothing may refuse for size"
);
}
#[test]
fn a_service_with_no_body_limit_declares_no_413() {
let document = support::router().openapi().expect("a describable router");
for (path, item) in &document.paths.items {
for (method, operation) in item.operations() {
assert!(
!operation.responses.responses.contains_key("413"),
"{method:?} {path} declares a 413 that nothing can produce"
);
}
}
}
#[tokio::test]
async fn a_timeout_over_a_body_limit_declares_both_statuses() {
let service = support::router()
.intercept(Timeout::new(Duration::from_millis(30)))
.intercept(BodySize::new(4096))
.build(App::new())
.expect("a describable router");
let document = service.openapi();
let operation = document.paths.items["/users"]
.post
.as_ref()
.expect("the operation exists");
assert!(
operation.responses.responses.contains_key("408"),
"the timeout contributes its status to the operation it covers"
);
assert!(operation.responses.responses.contains_key("413"));
}
#[tokio::test]
async fn one_limit_per_endpoint_caps_each_endpoint_separately() {
let service = Router::<()>::new()
.mount((
kynos::routes![slow].0.intercept(Concurrency::new(one())),
kynos::routes![prompt].0.intercept(Concurrency::new(one())),
))
.build(())
.expect("a describable router");
let (held, other) = tokio::join!(get(&service, "/slow").call(), async {
tokio::time::sleep(Duration::from_millis(50)).await;
get(&service, "/prompt").call().await
});
assert_eq!(held.status, StatusCode::NO_CONTENT);
assert_eq!(
other.status,
StatusCode::NO_CONTENT,
"a cap on one endpoint refused a request to another"
);
}
#[tokio::test]
async fn a_queued_request_waits_for_a_slot_rather_than_being_shed() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Concurrency::new(one()).queue_for(Duration::from_secs(2)))
.build(())
.expect("a describable router");
let (held, queued) = tokio::join!(get(&service, "/slow").call(), async {
tokio::time::sleep(Duration::from_millis(50)).await;
get(&service, "/prompt").call().await
});
assert_eq!(held.status, StatusCode::NO_CONTENT);
assert_eq!(
queued.status,
StatusCode::NO_CONTENT,
"the second request had two seconds to wait for a slot that frees in well under one"
);
}
#[tokio::test]
async fn a_queue_that_expires_sheds_the_same_status() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Concurrency::new(one()).queue_for(Duration::from_millis(20)))
.build(())
.expect("a describable router");
let (_held, shed) = tokio::join!(get(&service, "/slow").call(), async {
tokio::time::sleep(Duration::from_millis(50)).await;
get(&service, "/prompt").call().await
});
assert_eq!(shed.status, StatusCode::SERVICE_UNAVAILABLE);
}
#[tokio::test]
async fn a_configured_retry_after_reaches_a_shed_response() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Concurrency::new(one()).retry_after(Duration::from_secs(5)))
.build(())
.expect("a describable router");
let (_held, shed) = tokio::join!(get(&service, "/slow").call(), async {
tokio::time::sleep(Duration::from_millis(50)).await;
get(&service, "/prompt").call().await
});
assert_eq!(shed.status, StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
shed.field(header::RETRY_AFTER.as_str()).as_deref(),
Some("5")
);
}
#[tokio::test]
async fn a_slot_is_released_when_a_request_is_abandoned() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Concurrency::new(one()))
.build(())
.expect("a describable router");
{
let mut abandoned = Box::pin(get(&service, "/slow").call());
let started = tokio::time::timeout(Duration::from_millis(30), &mut abandoned).await;
assert!(started.is_err(), "the request must still be in flight");
}
assert_eq!(
get(&service, "/prompt").call().await.status,
StatusCode::NO_CONTENT,
"the abandoned request's slot was never released"
);
}