#![cfg(all(feature = "macros", feature = "json"))]
use std::{num::NonZeroUsize, time::Duration};
use kynos::{
Router,
error::problem::ProblemType,
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"
);
}
struct TookTooLong {
after: Duration,
}
impl From<Duration> for TookTooLong {
fn from(after: Duration) -> Self {
Self { after }
}
}
impl kynos::response::IntoResponse for TookTooLong {
fn into_response(self) -> kynos::http::Response {
let mut response = kynos::http::Response::new(kynos::http::body::Body::from_bytes(
bytes::Bytes::from(format!("gave up after {}ms", self.after.as_millis())),
));
*response.status_mut() = StatusCode::GATEWAY_TIMEOUT;
response
}
}
impl kynos::response::Responses for TookTooLong {
fn responses(registry: &mut kynos::schema::registry::Registry) -> kynos::openapi::Responses {
let _ = registry;
kynos::openapi::Responses::new().with(
504,
kynos::openapi::Response::new("the handler was abandoned"),
)
}
}
impl kynos::response::ShortCircuit for TookTooLong {
const STATUSES: &'static [u16] = &[504];
}
#[tokio::test]
async fn a_substituted_timeout_response_reaches_the_client_and_the_document() {
let service = Router::<()>::new()
.mount(kynos::routes![slow, prompt])
.intercept(Timeout::new(Duration::from_millis(20)).answer_with::<TookTooLong>())
.build(())
.expect("a describable router");
let timed_out = get(&service, "/slow").call().await;
assert_eq!(timed_out.status, StatusCode::GATEWAY_TIMEOUT);
assert!(timed_out.text().contains("20"), "{}", timed_out.text());
assert_eq!(
get(&service, "/prompt").call().await.status,
StatusCode::NO_CONTENT
);
let operation = service.openapi().paths.items["/slow"]
.get
.as_ref()
.expect("the operation exists");
assert!(
operation.responses.responses.contains_key("504"),
"the substituted type's status is what the operation describes"
);
assert!(
!operation.responses.responses.contains_key("408"),
"the default's status is described even though nothing can send it"
);
}
#[tokio::test]
async fn an_unsubstituted_timeout_still_answers_408() {
let service = Router::<()>::new()
.mount(kynos::routes![slow])
.intercept(Timeout::new(Duration::from_millis(20)))
.build(())
.expect("a describable router");
assert_eq!(
get(&service, "/slow").call().await.status,
StatusCode::REQUEST_TIMEOUT
);
}
#[cfg(feature = "openapi32")]
struct Trickle {
remaining: usize,
gap: Duration,
timer: std::pin::Pin<Box<tokio::time::Sleep>>,
}
#[cfg(feature = "openapi32")]
impl Trickle {
fn new(chunks: usize, gap: Duration) -> Self {
Self {
remaining: chunks,
gap,
timer: Box::pin(tokio::time::sleep(gap)),
}
}
}
#[cfg(feature = "openapi32")]
impl futures_core::Stream for Trickle {
type Item = Result<bytes::Bytes, std::convert::Infallible>;
fn poll_next(
self: std::pin::Pin<&mut Self>,
context: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
let this = self.get_mut();
if this.remaining == 0 {
return std::task::Poll::Ready(None);
}
std::task::ready!(this.timer.as_mut().poll(context));
this.remaining -= 1;
let next = tokio::time::Instant::now() + this.gap;
this.timer.as_mut().reset(next);
std::task::Poll::Ready(Some(Ok(bytes::Bytes::from_static(b"chunk"))))
}
}
#[cfg(feature = "openapi32")]
async fn read_to_end(body: kynos::http::body::Body) -> Result<bytes::Bytes, String> {
use http_body_util::BodyExt;
body.collect()
.await
.map(http_body_util::Collected::to_bytes)
.map_err(|error| error.to_string())
}
#[cfg(feature = "openapi32")]
async fn body_of<C: Send + Sync + 'static>(
service: &kynos::router::service::Service<C>,
target: &str,
) -> kynos::http::body::Body {
let mut request = kynos::http::Request::new(kynos::http::body::Body::empty());
*request.uri_mut() = target.parse().expect("a usable request target");
service.call(request).await.into_body()
}
#[cfg(feature = "openapi32")]
#[kynos::get("/steady")]
async fn steady()
-> kynos::response::stream::binary::BinaryStream<Trickle, kynos::extract::media::OctetStream> {
kynos::response::stream::binary::BinaryStream::new(Trickle::new(10, Duration::from_millis(5)))
}
#[cfg(feature = "openapi32")]
#[kynos::get("/stalled")]
async fn stalled()
-> kynos::response::stream::binary::BinaryStream<Trickle, kynos::extract::media::OctetStream> {
kynos::response::stream::binary::BinaryStream::new(Trickle::new(1, Duration::from_secs(30)))
}
#[cfg(feature = "openapi32")]
#[tokio::test]
async fn an_idle_body_timeout_ends_a_stalled_stream() {
let service = Router::<()>::new()
.mount(kynos::routes![stalled])
.intercept(kynos::middleware::limits::BodyTimeout::idle(
Duration::from_millis(100),
))
.build(())
.expect("a describable router");
let failure = read_to_end(body_of(&service, "/stalled").await)
.await
.expect_err("a stalled body outlived its idle limit");
assert!(failure.contains("did not finish"), "{failure}");
}
#[cfg(feature = "openapi32")]
#[tokio::test]
async fn an_idle_body_timeout_leaves_a_steady_stream_alone() {
let service = Router::<()>::new()
.mount(kynos::routes![steady])
.intercept(kynos::middleware::limits::BodyTimeout::idle(
Duration::from_millis(100),
))
.build(())
.expect("a describable router");
let delivered = read_to_end(body_of(&service, "/steady").await)
.await
.expect("a steady stream is inside its idle limit");
assert_eq!(delivered.len(), "chunk".len() * 10);
}
#[cfg(feature = "openapi32")]
#[tokio::test]
async fn a_deadline_ends_a_stream_that_is_still_producing() {
let service = Router::<()>::new()
.mount(kynos::routes![steady])
.intercept(kynos::middleware::limits::BodyTimeout::deadline(
Duration::from_millis(20),
))
.build(())
.expect("a describable router");
let failure = read_to_end(body_of(&service, "/steady").await)
.await
.expect_err("a deadline ends a body however steadily it arrives");
assert!(failure.contains("did not finish"), "{failure}");
}
#[cfg(feature = "openapi32")]
#[tokio::test]
async fn a_body_timeout_declares_no_status() {
let bounded = Router::<()>::new()
.mount(kynos::routes![steady])
.intercept(kynos::middleware::limits::BodyTimeout::idle(
Duration::from_millis(100),
))
.build(())
.expect("a describable router");
let plain = Router::<()>::new()
.mount(kynos::routes![steady])
.build(())
.expect("a describable router");
let described = |service: &kynos::router::service::Service<()>| {
let mut statuses: Vec<String> = service.openapi().paths.items["/steady"]
.get
.as_ref()
.expect("the operation exists")
.responses
.responses
.keys()
.cloned()
.collect();
statuses.sort();
statuses
};
assert_eq!(described(&bounded), described(&plain));
}
#[cfg(all(feature = "openapi32", feature = "json"))]
struct Silent {
timer: std::pin::Pin<Box<tokio::time::Sleep>>,
}
#[cfg(all(feature = "openapi32", feature = "json"))]
impl futures_core::Stream for Silent {
type Item = Result<kynos::response::stream::sse::Event<String>, std::convert::Infallible>;
fn poll_next(
self: std::pin::Pin<&mut Self>,
context: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
std::task::ready!(self.get_mut().timer.as_mut().poll(context));
std::task::Poll::Ready(None)
}
}
#[cfg(all(feature = "openapi32", feature = "json"))]
#[kynos::get("/heartbeat")]
async fn heartbeat() -> kynos::response::stream::sse::Sse<Silent> {
kynos::response::stream::sse::Sse::new(Silent {
timer: Box::pin(tokio::time::sleep(Duration::from_secs(30))),
})
.keep_alive(kynos::response::stream::sse::KeepAlive::new().interval(Duration::from_millis(5)))
}
#[cfg(all(feature = "openapi32", feature = "json"))]
#[tokio::test]
async fn a_keep_alive_frame_resets_an_idle_body_timeout() {
let service = Router::<()>::new()
.mount(kynos::routes![heartbeat])
.intercept(kynos::middleware::limits::BodyTimeout::idle(
Duration::from_millis(100),
))
.build(())
.expect("a describable router");
let outcome = tokio::time::timeout(
Duration::from_millis(400),
read_to_end(body_of(&service, "/heartbeat").await),
)
.await;
assert!(
outcome.is_err(),
"an event stream sending keep-alives every 5 ms tripped a 30 ms idle \
limit: {outcome:?}"
);
}
#[cfg(feature = "openapi32")]
struct EndCounts {
responses: std::sync::atomic::AtomicUsize,
disconnects: std::sync::atomic::AtomicUsize,
}
#[cfg(feature = "openapi32")]
struct CountingEnds(std::sync::Arc<EndCounts>);
#[cfg(feature = "openapi32")]
impl kynos::middleware::Observer<()> for CountingEnds {
fn on_request(
&self,
_: &kynos::http::Request,
_: Option<kynos::router::operation::Route<'_>>,
(): &(),
) {
}
fn on_response(
&self,
_: &kynos::http::Response,
_: Option<kynos::router::operation::Route<'_>>,
_: Duration,
) {
self.0
.responses
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
fn on_disconnect(&self, _: Option<kynos::router::operation::Route<'_>>, _: Duration) {
self.0
.disconnects
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
}
#[cfg(feature = "openapi32")]
#[tokio::test]
async fn a_body_the_timer_ended_is_reported_as_interrupted() {
let counts = std::sync::Arc::new(EndCounts {
responses: std::sync::atomic::AtomicUsize::new(0),
disconnects: std::sync::atomic::AtomicUsize::new(0),
});
let service = Router::<()>::new()
.mount(kynos::routes![stalled])
.observe(CountingEnds(std::sync::Arc::clone(&counts)))
.intercept(kynos::middleware::limits::BodyTimeout::idle(
Duration::from_millis(100),
))
.build(())
.expect("a describable router");
let body = body_of(&service, "/stalled").await;
read_to_end(body)
.await
.expect_err("a stalled body outlived its idle limit");
assert_eq!(
counts.responses.load(std::sync::atomic::Ordering::SeqCst),
1
);
assert_eq!(
counts.disconnects.load(std::sync::atomic::Ordering::SeqCst),
1,
"a response the body timer destroyed was reported as delivered"
);
}
#[cfg(feature = "openapi32")]
#[tokio::test]
async fn a_body_that_finished_is_not_reported_as_interrupted() {
let counts = std::sync::Arc::new(EndCounts {
responses: std::sync::atomic::AtomicUsize::new(0),
disconnects: std::sync::atomic::AtomicUsize::new(0),
});
let service = Router::<()>::new()
.mount(kynos::routes![steady])
.observe(CountingEnds(std::sync::Arc::clone(&counts)))
.intercept(kynos::middleware::limits::BodyTimeout::idle(
Duration::from_millis(200),
))
.build(())
.expect("a describable router");
let body = body_of(&service, "/steady").await;
read_to_end(body).await.expect("a steady stream completes");
assert_eq!(
counts.responses.load(std::sync::atomic::Ordering::SeqCst),
1
);
assert_eq!(
counts.disconnects.load(std::sync::atomic::Ordering::SeqCst),
0
);
}
struct TooLarge;
impl ProblemType for TooLarge {
const TYPE_URI: Option<&'static str> = Some("https://errors.example.com/payload-too-large");
}
#[tokio::test]
async fn a_named_body_limit_publishes_one_type_on_both_halves() {
const URI: &str = "https://errors.example.com/payload-too-large";
let service = support::router()
.intercept(BodySize::new(16).problem_type::<TooLarge>())
.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_eq!(reply.json()["type"], URI);
let declared = serde_json::to_value(service.openapi()).expect("a serializable document");
assert_eq!(
narrowed_type(&declared, "/users", "post", 413),
URI,
"the declared 413 does not narrow to the type the wire sent: {declared}"
);
}
#[tokio::test]
async fn an_unnamed_body_limit_publishes_about_blank_on_both_halves() {
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.json()["type"], "about:blank");
let declared = serde_json::to_value(service.openapi()).expect("a serializable document");
assert_eq!(
narrowed_type(&declared, "/users", "post", 413),
"about:blank"
);
}
fn narrowed_type<'a>(
document: &'a serde_json::Value,
path: &str,
method: &str,
status: u16,
) -> &'a str {
let schema = &document["paths"][path][method]["responses"][status.to_string()]["content"]["application/problem+json"]
["schema"];
schema["allOf"][1]["properties"]["type"]["const"]
.as_str()
.unwrap_or_else(|| panic!("the declared {status} does not narrow `type`: {schema}"))
}