kynos 0.1.0

An idiomatic, performance-focused REST API framework with OpenAPI 3.1 and 3.2 support.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
//! Every limit Kynos ships, held to the response its `Short` type names.
//!
//! One reason: the whole interceptor design says the declaration and the
//! behaviour are the same text, and a `ShortCircuit` type is exactly where that
//! could be false without the compiler noticing — `STATUSES` is a constant, and
//! nothing checks it against what `into_response` writes.
//!
//! Real durations in the 10–50 ms band rather than `tokio::time::pause()`,
//! which needs `tokio/test-util` — not a feature this workspace enables. The
//! nextest profile's `slow-timeout` is 30 s, so the band has three orders of
//! magnitude of headroom.

#![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};

/// A concurrency limit of one, which most of the cases below want.
///
/// `Concurrency::new` takes a `NonZeroUsize` because zero would be a service
/// that refuses everything; naming the conversion once keeps the cases about
/// what they are testing.
fn one() -> NonZeroUsize {
    NonZeroUsize::new(1).expect("one is not zero")
}

// --- BodySize: 413 -------------------------------------------------------

/// A body past the limit is refused, and the refusal is the status the type
/// declares.
#[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);
    // The limit is in the detail, so the client is told what it exceeded rather
    // than only that it exceeded something.
    assert!(reply.text().contains("16"), "{}", reply.text());
}

/// The control: the same request under a limit it fits inside.
#[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);
}

/// A declared length is refused before a byte is read, which is the branch a
/// streaming upload depends on and the one a length-less body cannot take.
#[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);
}

/// Every covered operation declares the 413, because configuring a limit and
/// documenting it are the same action.
#[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"
            );
        }
    }
}

// --- Timeout: 408 --------------------------------------------------------

/// A handler that outlives the limit.
#[kynos::get("/slow")]
async fn slow() -> NoContent {
    tokio::time::sleep(Duration::from_millis(400)).await;
    NoContent
}

/// One that does not, differing in exactly that.
#[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);
}

// --- Concurrency: 503 ----------------------------------------------------

/// Two requests overlap, so the second meets a full table.
///
/// Driven with `join!` on two futures rather than two spawned tasks: the
/// nextest profile fails a test that leaks a task, and a spawned request could
/// outlive the body of this one.
#[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 {
        // Long enough for the first request to have taken the only slot, and
        // far inside the 400 ms it holds it for.
        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);

    // No `Retry-After`: how long a slot takes to free is a property of the
    // requests already running, and a number invented here is one the service
    // cannot honour. The header is described because a *deployment* may know;
    // this one does not.
    assert!(refused.field(header::RETRY_AFTER.as_str()).is_none());
}

/// The control: the slot is released when the first request finishes, so the
/// same second request succeeds once it is free.
#[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"
        );
    }
}

/// A limit that short-circuits still leaves the rest of the router alone: a
/// request to a path no operation declares is still a 404 rather than the
/// limit's own status.
#[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);
}

// --- What applies when nothing is mounted --------------------------------

/// A service with no `BodySize` accepts a body of any size.
///
/// Recorded rather than fixed. `docs/nfr.md` read "body size, header count and
/// header size limits are enforced by default", and only the second and third
/// are: they are hyper's, set on the connection. A body cap is an interceptor
/// and `Router::build` mounts none.
///
/// Making one default was considered and rejected, and any one of three reasons
/// is sufficient. It would add 413 to every operation of every application that
/// never asked for one. It would make a user's own `BodySize` a `const` compile
/// error, since `statuses_disjoint` is what stops two interceptors claiming a
/// status. And it would buffer a body that declares no length, which is exactly
/// the streaming upload the limit is supposed to leave alone.
///
/// The framework's own rule — configuring a limit and documenting it are one
/// action — has a converse, and this is it: a limit nobody configured must not
/// be documented either.
#[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"
    );
}

/// And says so in the description: no operation declares a 413.
///
/// The other half. A service that accepted any body while *claiming* a 413
/// would be the defect `tests/matrix.rs` found in `BodyRejection`, which is
/// recorded in `docs/testing.md`.
#[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"
            );
        }
    }
}

/// A timeout and a body limit stacked together each declare their own status.
///
/// The router below mounts the arrangement the slow-body rule asks for —
/// `Timeout` outside `BodySize`, because `BodySize` reads a length-less body
/// frame by frame and only a timeout wrapping it ends the exchange. **This
/// test does not check that.** It sends no request, and what it asserts is
/// true in either mounting order.
///
/// The name used to say otherwise. Pinning the read needs a client that dribbles
/// a chunked body over a real socket, which this harness cannot express;
/// `docs/middleware.md` is where the rule is stated and says nothing checks it.
#[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");

    // The chain runs outermost-first, so the timeout is written first. What
    // follows reads the description, not the exchange.
    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"));
}

// --- Concurrency scope ----------------------------------------------------

/// Two endpoints, one `Concurrency` each: the caps are separate.
///
/// "Maximum concurrent requests per endpoint" needs no new API. An
/// `EndpointBuilder` has its own interceptor list, and `Router::build`
/// composes it with the router's, so one instance per endpoint *is* a
/// per-endpoint cap. Recorded because the alternative — a `per_route()` mode
/// keyed on the matched path — would cost a lock and a lookup on the request
/// path to express what the mount site already says.
#[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 {
        // Long enough for the first request to have taken `/slow`'s only slot.
        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"
    );
}

/// A bounded queue absorbs a burst instead of shedding it.
///
/// The wait is not a declaration: the answer when it expires is the same 503,
/// and a delay is not a response. What it changes is which of the two a client
/// gets, and only where the deployment asked.
#[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"
    );
}

/// A queue that expires still sheds, with the status it always had.
#[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);
}

/// A deployment that knows how long to wait can say so.
///
/// `AtCapacity` has always *described* a `Retry-After` and nothing could
/// produce one — the shape `assert_declared_responses_covered` exists to catch.
#[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")
    );
}

/// A slot is released when the chain's future is dropped, not only when it
/// finishes.
///
/// The reason the permit is a guard rather than a counter pair: a client that
/// disconnects mid-request drops the future at an await point, and a slot that
/// leaked there would shrink the limit until the process restarted.
#[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");

    // Abandon a request that has taken the only slot. Boxed rather than
    // `tokio::pin!`ed, because that macro shadows the binding with a
    // `Pin<&mut _>` and dropping *that* leaves the future alive to the end of
    // the scope — which is a test that passes for the wrong reason.
    {
        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"
    );
}