use std::{convert::Infallible, net::Ipv4Addr, num::NonZeroUsize, time::Duration};
use kynos::{
http,
middleware::{
Continued, Interceptor, Next,
compression::{Compression, levels::GzipLevel},
cors::Cors,
limits::{BodySize, BodyTimeout, Concurrency, Timeout},
rate_limit::{
RateLimit,
decision::{Decision, QuotaPolicy, QuotaUnit, RateLimitPolicy, ServiceLimit},
},
request_id::{CorrelationHeaders, Counter, RequestId, RequestIdSource},
trace::Trace,
},
openapi::Method,
prelude::*,
response::status::NoContent,
router::operation::Route,
server::Server,
};
use serde::{Deserialize, Serialize};
#[derive(Schema, Serialize, Deserialize)]
struct User {
id: u64,
name: String,
}
#[derive(HeaderParams)]
struct CorrelationId {
#[header(rename = "X-Correlation-Id")]
correlation_id: String,
}
impl CorrelationHeaders for CorrelationId {
fn from_id(id: http::HeaderValue) -> Self {
Self {
correlation_id: id.to_str().unwrap_or_default().to_owned(),
}
}
}
struct Monotonic;
impl RequestIdSource for Monotonic {
fn next_id(&self) -> http::HeaderValue {
http::HeaderValue::from_static("00000000-0000-4000-8000-000000000000")
}
}
#[derive(HeaderParams)]
struct AdminSecret {
#[header(rename = "X-Admin-Secret")]
secret: Option<String>,
}
#[derive(Debug, thiserror::Error, ApiError)]
#[error("the shared secret was absent or wrong")]
#[problem(status = 401)]
struct SecretRejected;
struct RequireSecret {
secret: &'static str,
}
impl<C: Sync + 'static> Interceptor<C> for RequireSecret {
type Reads = AdminSecret;
type Adds = ();
type Short = SecretRejected;
async fn intercept(
&self,
request: http::Request,
reads: AdminSecret,
context: &C,
next: Next<'_, C>,
) -> Result<Continued<()>, SecretRejected> {
let _ = context;
if reads.secret.as_deref() != Some(self.secret) {
return Err(SecretRejected);
}
Ok(next.run(request).await)
}
}
#[derive(HeaderParams)]
struct TenantHeader {
#[header(rename = "X-Tenant")]
tenant: String,
}
#[derive(Clone)]
struct Tenant;
impl<C: Sync + 'static> Interceptor<C> for Tenant {
type Reads = TenantHeader;
type Adds = ();
type Short = Infallible;
async fn intercept(
&self,
request: http::Request,
reads: TenantHeader,
context: &C,
next: Next<'_, C>,
) -> Result<Continued<()>, Infallible> {
let _ = (context, reads.tenant);
Ok(next.run(request).await)
}
}
#[kynos::get("/users")]
async fn list_users() -> Json<Vec<User>> {
Json(vec![User {
id: 1,
name: "Ada Lovelace".to_owned(),
}])
}
#[kynos::post("/users/avatar")]
async fn upload_avatar(Json(user): Json<User>) -> NoContent {
println!("avatar for {}", user.name);
NoContent
}
#[kynos::get("/reports")]
async fn reports() -> NoContent {
NoContent
}
#[derive(Clone, Debug)]
struct PerProcess {
served: std::sync::Arc<std::sync::atomic::AtomicU32>,
ceiling: u32,
policies: Vec<QuotaPolicy>,
}
impl PerProcess {
fn new(ceiling: u32) -> Self {
Self {
served: std::sync::Arc::default(),
ceiling,
policies: vec![QuotaPolicy {
name: "per-process".into(),
quota: u64::from(ceiling),
window: Some(Duration::from_secs(60)),
unit: QuotaUnit::Requests,
}],
}
}
}
impl RateLimitPolicy<()> for PerProcess {
fn advertised(&self) -> &[QuotaPolicy] {
&self.policies
}
async fn check(&self, _: &http::Request, _: Route<'_>, (): &()) -> Decision {
let served = self
.served
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
if served >= self.ceiling {
return Decision::deny(
Duration::from_secs(60),
ServiceLimit {
name: "per-process".into(),
quota: u64::from(self.ceiling),
remaining: 0,
reset: Duration::from_secs(60),
},
);
}
Decision::allow(ServiceLimit {
name: "per-process".into(),
quota: u64::from(self.ceiling),
remaining: u64::from(self.ceiling - served - 1),
reset: Duration::from_secs(60),
})
}
}
#[tokio::main]
async fn main() -> kynos::Result<()> {
let router = Router::<()>::new()
.intercept(
Cors::new()
.allow_origins(["https://app.example.com"])
.allow_methods([Method::Get, Method::Post])
.allow_headers(["content-type", "x-tenant"])
.expose_headers(["x-correlation-id"])
.max_age(Duration::from_secs(600))
.allow_credentials()
.document_response_headers(),
)
.intercept(
RequestId::new()
.header::<CorrelationId>()
.trust_client(false)
.source(Monotonic),
)
.observe(
Trace::new()
.level(tracing::Level::INFO)
.record_headers(&["x-correlation-id", "x-tenant"]),
)
.intercept(BodyTimeout::idle(Duration::from_secs(15)))
.intercept(
Compression::new()
.min_size(1_024)
.gzip_level(GzipLevel::new(5).expect("DEFLATE defines level 5")),
)
.intercept(Timeout::new(Duration::from_secs(30)))
.intercept(BodySize::new(1_048_576))
.intercept(Concurrency::new(
NonZeroUsize::new(256).expect("a nonzero concurrency limit"),
))
.intercept(RateLimit::new(PerProcess::new(100)))
.group(
Group::new("/tenanted")
.intercept(Tenant)
.mount(kynos::routes![upload_avatar]),
)
.group(
Group::new("/admin")
.intercept(RequireSecret {
secret: "opensesame",
})
.mount(kynos::routes![reports]),
)
.mount(kynos::routes![list_users]);
let document = router.openapi()?;
println!("{}", document.to_json()?);
let _default_source = Counter::default();
Server::new(router.build(())?)
.bind((Ipv4Addr::UNSPECIFIED, 3000))
.serve()
.await
}