# Tako > Tako is a multi-transport Rust web framework: HTTP/1.1, HTTP/2, HTTP/3, WebTransport, WebSocket, SSE, gRPC, TCP, UDP and Unix sockets behind one router, on Tokio or Compio. Key facts for writing Tako code: - The crates.io package is `tako-rs`; the library is imported as `tako` (`use tako::...`). The crate named `tako` on crates.io is an unrelated project, so never depend on it. - Install with `cargo add tako-rs` (`tako-rs = "2.4"`). Tako requires Rust 1.95 or newer and edition 2024. - Default features give Tokio with HTTP/1.1, raw TCP, Unix sockets, static files, and the core extractors and middleware. HTTP/2, HTTP/3, TLS, WebSocket (`ws`), SSE (`sse`), UDP, gRPC, and the bundled plugins are opt-in Cargo features. - Start servers with `Server::builder().build().try_spawn_http(listener, router)?`; the `serve_*` free functions are deprecated. The `compio` feature switches to `CompioServer`; every transport runs on both runtimes. - Tako has no outbound HTTP client (`tako::client` was removed in 2.3.0); use `reqwest` or `ureq`. - Every docs page is available as Markdown by appending `.mdx` to its URL. --- # Introduction > A multi-transport Rust framework for HTTP, WebSocket, SSE, gRPC, HTTP/3, and raw sockets behind one routing, middleware, and observability model. **Tako** (*"octopus"* in Japanese) is a pragmatic, ergonomic, and extensible Rust framework for services that go beyond plain HTTP. One application surface covers HTTP/1.1, HTTP/2, HTTP/3, WebSocket, SSE, gRPC, TCP, UDP, Unix sockets, and WebTransport — with a shared routing, middleware, and observability model, running on **Tokio** or **Compio**. This site is the long-form companion to the in-source rustdoc. The rustdoc is the reference; this site is the *guide*. When the two disagree, the rustdoc wins for a single API; the site wins for the intent and the recommended pattern. ## What's in here The framework is built around eight workspace crates, re-exported through the `tako-rs` umbrella under the `tako::*` path: ``` tako-rs umbrella re-export (tako::*) tako-rs-core routing, handlers, middleware, body/request types, state, signals, queue, GraphQL / gRPC / OpenAPI helpers tako-rs-extractors concrete request extractors (cookies, form, query, path, JWT, multipart, simd-json, …) tako-rs-server HTTP/1.1, HTTP/2, HTTP/3, WebTransport, TLS, raw TCP/UDP/Unix, PROXY protocol, plus compio variants tako-rs-streams WebSocket, SSE, file streaming, static files, raw QUIC tako-rs-plugins bundled middleware (auth, CSRF, sessions, …) and plugins (CORS, compression, rate limiting, idempotency, metrics) tako-rs-macros the #[tako::route] / #[tako::get] attribute family tako-rs-server-pt optional thread-per-core entry point ``` ## Where to start * **New here?** [Quickstart](/docs/getting-started/quickstart) — install, write a handler, and serve it, end-to-end. * **Need a specific transport?** Start at the [Transports overview](/docs/transports) and pick the page for the protocol you need. * **Wiring request handling?** [Routing](/docs/routing), [Extractors](/docs/extractors), and [Middleware](/docs/middleware) walk the request lifecycle. ## Navigation The left sidebar is grouped into: 1. **Start here** — getting-started + concepts (architecture, the routing model, runtimes, feature flags). Read this once before diving into the catalog. 2. **Transports** — one page per protocol: HTTP, HTTP/3, WebSocket, WebTransport, SSE, gRPC, raw sockets, PROXY protocol. 3. **Request handling** — routing, the extractor catalog, the middleware catalog, and bundled plugins. 4. **Primitives** — streams, queue, signals, observability. 5. **Guides** — long-form tutorials, deployment, and the benchmark page. 6. **Reference** — migration, API stability policy, the cargo feature graph, and the docs.rs API link. ## Who this is for * **Service teams** building APIs that need more than plain REST — WebSockets, SSE, gRPC, raw TCP/UDP, or QUIC in the same binary. * **Platform teams** that want one Rust framework story across protocols, runtimes (Tokio + Compio), and deployment shapes. * **Operators** who want first-class signals, queues, and graceful shutdown without composing them from scratch. If you are building a plain JSON-over-HTTP service and never need the realtime / multi-transport surface, you can use Tako as a straight HTTP framework and ignore the rest. Source: https://tako.rust-dd.com/docs --- # Installation > Add tako-rs to a Cargo project, pick the Tokio or Compio runtime, and opt into the feature flags your service needs. Tako is published on crates.io as **`tako-rs`**. The package name is `tako-rs`; the crate you import is `tako`. Everything ships through this one umbrella — you never depend on the sub-crates directly. ## Add the dependency ```toml [dependencies] tako-rs = "2" tokio = { version = "1", features = ["full"] } anyhow = "1" ``` `tokio` and `anyhow` are not required by Tako itself, but the default runtime is Tokio and the examples use `anyhow::Result` for `main`. With the default feature set you get HTTP/1.1, WebSocket, SSE, raw TCP/UDP, Unix sockets, and PROXY protocol on Tokio — no flags needed. ## Toolchain Tako targets **Rust edition 2024** with an **MSRV of 1.95**. It relies on let-chains, async-fn-in-trait HRTB inference, and other edition-2024 features, so an older toolchain will not compile the crate. ```toml # rust-toolchain.toml [toolchain] channel = "1.95" ``` Tracking the latest stable Rust is recommended. ## Choosing a runtime Tako runs on two async runtimes, selected at build time: * **Tokio** (default) — multi-threaded, work-stealing scheduler with `Send` futures end-to-end. This is what the standalone `Server` builder targets. * **Compio** (opt-in via the `compio` feature) — single-threaded thread-per-core over io\_uring / IOCP / kqueue, with `!Send` futures pinned to their runtime thread. You cannot enable both runtimes in one binary. `cargo build --all-features` does not produce a runnable binary. Pick one runtime per binary; if a deployment needs both, build separate binaries. See [Runtimes](/docs/concepts/runtimes) for the full trade-off. To build on Compio instead of Tokio: ```toml [dependencies] tako-rs = { version = "2", default-features = false, features = ["compio"] } ``` ## Key feature flags The umbrella crate is the contractual surface: a feature like `tako-rs/http2` turns on the matching feature in every sub-crate at once. The flags you reach for most often: | Flag | Turns on | | ------------- | ------------------------------------------------------------------------------ | | `http2` | HTTP/2 (h2c cleartext and ALPN-negotiated h2 over TLS) | | `http3` | HTTP/3 over QUIC | | `tls` | Server-side HTTPS via rustls | | `grpc` | gRPC unary and streaming RPCs with protobuf | | `simd` | SIMD JSON parsing (`sonic-rs` + `simd-json`) | | `multipart` | `Multipart` / typed multipart body extractors | | `plugins` | Bundled middleware and plugins (CORS, compression, rate limiting, idempotency) | | `signals` | In-process pub/sub signal bus and lifecycle signals | | `file-stream` | File streaming, range requests, conditional GET | | `jemalloc` | Re-export `tako::Jemalloc`; your binary sets it as the global allocator | Compio has its own TLS and WebSocket flags, since those features compile differently on the two runtimes: | Flag | Turns on | | ------------ | ----------------------------------------------------------- | | `compio` | The Compio runtime (mutually exclusive with the Tokio path) | | `compio-tls` | TLS on Compio (implies `compio`) | | `compio-ws` | WebSocket on Compio (implies `compio`) | Enable several at once: ```toml [dependencies] tako-rs = { version = "2", features = ["http2", "tls", "simd", "plugins"] } ``` For the complete list — observability backends, OpenAPI integrations, validators, zero-copy extractors, and the combinations that do not compile — see the [feature reference](/docs/reference/features). ## Next steps * [Quickstart](/docs/getting-started/quickstart) — write and serve your first handler. * [Project layout](/docs/getting-started/project-layout) — what each workspace crate owns. * [Runtimes](/docs/concepts/runtimes) — Tokio vs Compio in depth. Source: https://tako.rust-dd.com/docs/getting-started/installation --- # Quickstart > Add the tako-rs dependency, write a handler, register a GET route on Router, and start it with Server::builder and a TcpListener. This walks through installing Tako, writing your first handler, and running it on a local TCP listener. It mirrors the [`examples/hello-world`](https://github.com/rust-dd/tako/tree/main/examples/hello-world) crate in the repository. ## 1. Add the dependency ```toml [dependencies] tako-rs = "2" tokio = { version = "1", features = ["full"] } anyhow = "1" ``` The package is published as `tako-rs`; everything is re-exported under the `tako::*` path. The default feature set covers HTTP/1.1, raw TCP, and Unix sockets on Tokio; WebSocket (`ws`), SSE (`sse`), UDP (`udp`), PROXY protocol (`proxy-protocol`), HTTP/2, HTTP/3, and TLS are opt-in features. See [Installation](/docs/getting-started/installation) for the runtime choice and opt-in feature flags. ## 2. Write the server Source: [examples/hello-world/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/hello-world/src/main.rs) ```rust use anyhow::Result; use tako::Method; use tako::responder::Responder; use tako::router::Router; use tokio::net::TcpListener; async fn hello_world() -> impl Responder { "Hello, World!".into_response() } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello_world); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` ## 3. Run it ```bash cargo run ``` Then, in another terminal: ```bash curl http://127.0.0.1:8080/ # Hello, World! ``` ## What the code does A handler is any `async fn` whose arguments implement `FromRequest` / `FromRequestParts` and whose return type implements `Responder`. The hello-world handler takes zero arguments, so it extracts nothing from the request: ```rust async fn hello_world() -> impl Responder { "Hello, World!".into_response() } ``` `Router::new()` produces an empty router. `route(Method::GET, "/", handler)` registers a handler against a `(method, path)` pair. Convenience shorthands exist too — `router.get(path, handler)`, plus `.post`, `.put`, `.patch`, `.delete`, `.head`, and `.options`: ```rust let mut router = Router::new(); router.route(Method::GET, "/", hello_world); // equivalently: router.get("/", hello_world); ``` `Server::builder().build().spawn_http(listener, router)` returns a server handle. Await `result()` to observe startup or serving errors, or call `shutdown(timeout)` to stop acceptance and drain active requests. ## Adding extractors Handlers compose by taking more arguments. Tako runs the extractor for each argument before invoking the handler: ```rust use serde::Deserialize; use tako::Method; use tako::extractors::json::Json; use tako::extractors::path::Path; use tako::responder::Responder; use tako::router::Router; #[derive(Deserialize)] struct UserPath { id: u64 } #[derive(Deserialize)] struct CreateUser { name: String } async fn get_user(Path(p): Path) -> impl Responder { format!("user_id={}", p.id) } async fn create_user(Json(body): Json) -> impl Responder { format!("created: {}", body.name) } let mut router = Router::new(); router.route(Method::GET, "/users/{id}", get_user); router.route(Method::POST, "/users", create_user); ``` ## Return types Anything implementing `Responder` can be returned from a handler. The blanket impls cover the common shapes: * `&'static str` / `String` — `text/plain` body * `Bytes` / `Vec` — raw body * `Json` — JSON body with `Content-Type: application/json` * `(StatusCode, T)` — status plus body * `(StatusCode, HeaderMap, T)` — status, headers, and body * `StatusCode` alone — empty body, just the status line * `Result` where both arms implement `Responder` — the preferred shape for fallible handlers ## Next steps * [Routing](/docs/routing) — path parameters, nesting, scopes, typed slots. * [Extractors](/docs/extractors) — the full catalog of request extractors. * [State](/docs/state) — share configuration and dependencies with handlers. * [Middleware](/docs/middleware) — auth, CORS, compression, metrics, rate limiting. Source: https://tako.rust-dd.com/docs/getting-started/quickstart --- # Coming from Axum > Map Axum routers, extractors, state, middleware, serving, and tests to their Tako equivalents, with the differences that matter when porting. Tako's handler model is close to Axum's: handlers are async functions, their arguments are typed extractors, and their return values implement a response trait. Path syntax is the same `{id}` form Axum 0.8 uses. Most ports are mechanical; this page maps each Axum building block to Tako and calls out the places where the two differ. ## At a glance | Axum 0.8 | Tako 2.x | | ---------------------------------- | ----------------------------------------------- | | `Router::new().route("/", get(h))` | `router.get("/", h)` | | `get(show).post(update)` | `router.get(p, show)`, `router.post(p, update)` | | `Path`, `Query`, `Json`, `State` | Same names in `tako::extractors` | | `http::HeaderMap` | `header_map::HeaderMap` wrapper | | `IntoResponse` | `Responder` | | `Router::with_state` | `router.with_state(value)` per type | | `middleware::from_fn` | `router.middleware(f)` | | `route_layer` | `route.middleware(f)` | | `tower-http` layers | Bundled middleware and plugins | | `nest`, `fallback` | `router.nest`, `router.fallback` | | `axum::serve` | `Server::builder()` | | `with_graceful_shutdown` | `handle.shutdown_on_signal()` | | `WebSocketUpgrade` | `TakoWs::new(req, handler)` | | `sse::Sse`, `Event` | `Sse::events`, `SseEvent` | | `ServiceExt::oneshot` | `router.dispatch(request)` | ## Dependencies ```toml [dependencies] serde = { version = "1", features = ["derive"] } tako-rs = { version = "2.4", features = ["plugins", "ws", "sse"] } tokio = { version = "1", features = ["macros", "net", "rt-multi-thread"] } ``` The package is `tako-rs` and the Rust import is `tako`. The crate named `tako` on crates.io is an unrelated project. Transports and integrations are Cargo features; the [feature reference](/docs/reference/features) lists them all. ## Routes and handlers Axum chains routes onto a router value: ```rust let app = Router::new() .route("/users", get(list_users).post(create_user)) .route("/users/{id}", get(show_user)); ``` Tako registers each method on a mutable router: ```rust use tako::router::Router; let mut router = Router::new(); router.get("/users", list_users); router.post("/users", create_user); router.get("/users/{id}", show_user); ``` Handlers keep their shape. As in Axum, only the last argument may consume the request body. ## Extractors ```rust use serde::Deserialize; use tako::extractors::{header_map::HeaderMap, json::Json, path::Path, query::Query, state::State}; use tako::{responder::Responder, StatusCode}; #[derive(Deserialize)] struct Pagination { page: u32, } async fn list_users(Query(p): Query, State(state): State) -> String { format!("page {} from {}", p.page, state.db_url) } async fn show_user(Path(id): Path) -> String { format!("user {id}") } async fn user_agent(HeaderMap(headers): HeaderMap) -> String { format!("{:?}", headers.get("user-agent")) } async fn create_user(Json(body): Json) -> impl Responder { (StatusCode::CREATED, Json(body)) } ``` Three differences matter when porting: * `State` holds an `Arc`, and state is keyed by type. Call `with_state` once per type and extract any subset in each handler; there is no `FromRef`. * `HeaderMap` is a tuple wrapper around `http::HeaderMap`, so destructure it. * Buffered body extractors reject bodies over 2 MiB by default, much like Axum's `DefaultBodyLimit`. Change it with `router.body_limit(bytes)`. The [extractor catalog](/docs/extractors) covers the rest, including cookies, forms, multipart, JWT claims, and protobuf. ## Responses and errors Return `impl Responder` where Axum returns `impl IntoResponse`. Strings, JSON, status codes, and `(StatusCode, R)` tuples all implement it. A handler can return `Result` when both `T` and `E` implement `Responder`, so a custom error type implements `Responder` the way it would implement `IntoResponse` in Axum. See [Building a REST API](/docs/tutorials/rest-api) for a full error type. ## State ```rust #[derive(Clone)] struct AppState { db_url: String, } let mut router = Router::new(); router.with_state(AppState { db_url: "postgres://localhost/app".into() }); ``` Routers nested under this one see the same state. Because state is looked up at runtime instead of being part of the router's type, a handler that extracts a type nobody registered compiles, and responds with `500 Internal Server Error`. ## Middleware An Axum `middleware::from_fn` function ports with the same signature: ```rust use tako::middleware::Next; use tako::types::{Request, Response}; async fn require_api_key(req: Request, next: Next) -> Response { if req.headers().contains_key("x-api-key") { next.run(req).await } else { StatusCode::UNAUTHORIZED.into_response() } } router.middleware(require_api_key); router.get("/admin", admin).middleware(require_api_key); ``` `router.middleware` wraps every route; `.middleware` on a route wraps only that route, like Axum's `route_layer`. Tako does not implement Tower's `Service` and `Layer` traits, so `tower-http` layers do not plug in. The common ones have bundled equivalents: | Axum ecosystem | Tako | | ------------------------------------------- | ------------------------------------------------------------- | | `CorsLayer` | `tako::plugins::cors::CorsBuilder` (`plugins`) | | `CompressionLayer` | `tako::plugins::compression::CompressionBuilder` (`plugins`) | | `DefaultBodyLimit`, `RequestBodyLimitLayer` | `router.body_limit(bytes)` | | `TimeoutLayer` | `router.timeout(duration)` or `route.timeout(duration)` | | `SetRequestIdLayer` | `tako::middleware::request_id::RequestId` | | `tower-sessions` | `tako::middleware::session::SessionMiddleware` | | `tower_governor` | `tako::plugins::rate_limiter::RateLimiterBuilder` (`plugins`) | Plugins register with `router.plugin(...)`, and middleware structs convert with `.into_middleware()`: ```rust use tako::middleware::{request_id::RequestId, IntoMiddleware}; use tako::plugins::{compression::CompressionBuilder, cors::CorsBuilder}; router.plugin(CorsBuilder::new().build()); router.plugin(CompressionBuilder::new().build()); router.middleware(RequestId::new().into_middleware()); ``` Authentication (`BasicAuth`, `BearerAuth`, `ApiKeyAuth`, `JwtAuth`), CSRF, and security headers are covered in [Middleware](/docs/middleware). ## Nesting and fallbacks ```rust let mut api = Router::new(); api.get("/users/{id}", show_user); let mut app = Router::new(); app.nest("/api", api); app.fallback(|_req: Request| async { (StatusCode::NOT_FOUND, "nothing here") }); ``` ## Serving and shutdown Axum hands the router to `axum::serve`: ```rust let listener = tokio::net::TcpListener::bind("0.0.0.0:8080").await?; axum::serve(listener, app) .with_graceful_shutdown(shutdown_signal()) .await?; ``` Tako spawns the server and returns a handle: ```rust use tako::Server; use tokio::net::TcpListener; let listener = TcpListener::bind("0.0.0.0:8080").await?; let handle = Server::builder().build().try_spawn_http(listener, app)?; handle.shutdown_on_signal().await?; ``` `shutdown_on_signal` waits for Ctrl+C or SIGTERM, then drains in-flight requests within the drain timeout. Use `handle.result().await?` to wait without handling signals. The same builder serves TLS, h2c, HTTP/3, and Unix sockets; see [Transports](/docs/transports). Axum needs `axum-server` or a rustls acceptor for TLS and has no HTTP/3 server. ## WebSocket and SSE A Tako WebSocket handler receives the whole request instead of a `WebSocketUpgrade` extractor, and the socket is a `tokio-tungstenite` `WebSocketStream`: ```rust use futures_util::{SinkExt, StreamExt}; use tako::ws::TakoWs; async fn echo(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { while let Some(Ok(msg)) = ws.next().await { if (msg.is_text() || msg.is_binary()) && ws.send(msg).await.is_err() { break; } } }) } ``` Axum's `Sse` takes a stream of `Result`; Tako's `Sse::events` takes a stream of `SseEvent`: ```rust use futures_util::{stream, StreamExt}; use tako::sse::{Sse, SseEvent}; async fn ticks(_: Request) -> impl Responder { Sse::events(stream::iter(0..3).map(|i| SseEvent::data(format!("tick {i}")))) } ``` See [WebSocket](/docs/transports/websocket) and [Server-Sent Events](/docs/transports/sse) for keep-alive, limits, and Compio. ## Testing Axum tests drive the router through `tower::ServiceExt::oneshot`. A Tako router dispatches a request directly, with no server or socket: ```rust use tako::body::TakoBody; let mut request = Request::new(TakoBody::empty()); *request.uri_mut() = "/api/users/7".parse()?; let response = app.dispatch(request).await; assert_eq!(response.status(), StatusCode::OK); ``` ## No direct equivalent * **Tower compatibility.** Tako has its own middleware and plugin traits. * **`#[debug_handler]`.** Handler errors surface as ordinary trait-bound errors. * **Typed router state.** State lives in a per-router, type-keyed store instead of `Router`. Tako also has pieces Axum does not ship: typed route macros such as `#[tako::get("/users/{id: u64}")]` (see [Routing](/docs/routing)), background [queues](/docs/queue), in-process [signals](/docs/signals), and the [thread-per-core server](/docs/deployment#thread-per-core). Source: https://tako.rust-dd.com/docs/getting-started/coming-from-axum --- # Using Tako with AI assistants > Point coding agents at Tako's llms.txt, Markdown docs, and Context7 index, and give them project rules so they write current Tako 2.x code. Tako is younger than the Rust frameworks most language models learned from, so coding assistants tend to write Axum code, call APIs that do not exist, or add the wrong crate. Give them the current docs and a few project rules and they write working Tako 2.x code. ## Machine-readable docs | URL | Contents | | ------------------------------------------------------- | -------------------------------------------------------------------------------------------------- | | [`/llms.txt`](/llms.txt) | Index of every docs page with a one-line summary, in the [llms.txt](https://llmstxt.org) format | | [`/llms-full.txt`](/llms-full.txt) | The whole guide as one Markdown file, with the runnable examples inlined | | `/docs/.mdx` | Any page as Markdown, for example [`/docs/routing.mdx`](https://tako.rust-dd.com/docs/routing.mdx) | | [docs.rs/tako-rs](https://docs.rs/tako-rs/latest/tako/) | The rustdoc API reference | Docs pages also return Markdown when a request prefers it with `Accept: text/markdown`, and every page has a **Copy Markdown** button plus links that open the page in ChatGPT, Claude, or Cursor. ## Context7 [Context7](https://context7.com/rust-dd/tako) indexes the Tako docs and examples under the library ID `/rust-dd/tako`. With the Context7 MCP server installed, add "use library /rust-dd/tako" to a prompt to pull those docs in. The crate named `tako` on crates.io belongs to a different project, so pick the `rust-dd` entry. ## Project rules Paste these rules into `AGENTS.md`, `CLAUDE.md`, or a Cursor rule file in any project that uses Tako: ```md ## Tako (Rust web framework) - Depend on `tako-rs` (`cargo add tako-rs`) and import it as `tako`. Never add the crate named `tako`; it is an unrelated project. - Target Tako 2.x. Register routes on a mutable router (`let mut router = Router::new(); router.get("/", handler);`) and start servers with `Server::builder().build().try_spawn_http(listener, router)?`. Do not use the deprecated `serve_*` functions. - Handlers are async functions. Extractors live in `tako::extractors::*` (`Path`, `Query`, `Json`, `State`, `HeaderMap`), and responses implement `tako::responder::Responder`, not `IntoResponse`. Only the last handler argument may consume the body. - Transports and integrations are Cargo features: `http2`, `http3`, `tls`, `ws`, `sse`, `grpc`, `plugins`, `per-thread`, `compio`. Enable the feature before using its module. - Tako is not Tower-compatible. Use `router.middleware(...)` and the bundled plugins instead of `tower-http` layers. - Docs: https://tako.rust-dd.com/llms.txt and https://tako.rust-dd.com/llms-full.txt ``` ## Checking generated code * Run `cargo check` with the features the code uses; a missing feature shows up as an unresolved import. * Compare against the runnable crates in [`examples/`](https://github.com/rust-dd/tako/tree/main/examples), one per transport and integration. * Porting from Axum? [Coming from Axum](/docs/getting-started/coming-from-axum) maps each building block. Source: https://tako.rust-dd.com/docs/getting-started/ai-assistants --- # Project layout > The eight-crate tako workspace — what the tako-rs umbrella re-exports and which sub-crate owns each part of the framework. Tako is a Cargo workspace of eight crates. Application code depends on a single one — the `tako-rs` umbrella — and reaches everything else through the `tako::*` path. The remaining seven crates are the internal building blocks; you rarely name them directly. ## The umbrella crate `tako-rs` (imported as `tako`) re-exports every public type from the sub-crates at its original path. A handler imports `tako::router::Router`, `tako::responder::Responder`, `tako::extractors::json::Json`, and so on — regardless of which sub-crate actually defines each item. The umbrella also owns the cargo features: a feature like `tako-rs/multipart` turns on the matching feature in `tako-rs-extractors` and `tako-rs-core` together. ```toml [dependencies] tako-rs = "2" ``` A `tako::prelude` module bundles the types most handlers reach for — `Router`, `Responder`, `Request`, `Response`, `Method`, `StatusCode`, the common body / query / form / path extractors, and the middleware traits — so a single `use tako::prelude::*;` covers everyday code. ## The eight crates | Crate | Owns | Reach for it when | | -------------------- | --------------------------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------- | | `tako-rs` | Umbrella re-export (`tako::*`), cargo feature surface, prelude, allocator hook | Always — this is your only direct dependency | | `tako-rs-core` | Router, handlers, middleware and plugin traits, body / request / response types, state, signals, queue, plus GraphQL / gRPC / OpenAPI helpers | Foundational — pulled in by everything | | `tako-rs-extractors` | Concrete request extractors: cookies, form, query, path, headers, JWT, multipart, simd-json, and more | Reading typed data out of requests | | `tako-rs-server` | HTTP/1.1, HTTP/2, HTTP/3, WebTransport, TLS, raw TCP / UDP / Unix, PROXY protocol, plus compio variants | Binding listeners and serving transports | | `tako-rs-streams` | WebSocket, SSE, file streaming, static file serving, raw QUIC sessions | Realtime and streaming responses | | `tako-rs-plugins` | Bundled middleware (auth, CSRF, sessions, …) and plugins (CORS, compression, rate limiting, idempotency, metrics) | Production middleware off the shelf | | `tako-rs-macros` | The `#[tako::route]` / `#[tako::get]` / `.post` / … attribute family | Declarative route registration with typed path slots | | `tako-rs-server-pt` | Optional thread-per-core entry point (`serve_per_thread`) | Thread-per-core deployments via `SO_REUSEPORT` | ## How the pieces fit `tako-rs-core` defines the framework primitives every other crate depends on: the `Router` that maps `(method, path)` to handlers, the `Request` / `Response` type aliases, the `Responder` trait for return values, and the `IntoMiddleware` trait for the middleware chain. The extractor, server, streams, and plugins crates each layer concrete implementations on top of those traits, and the macros crate generates route registrations that plug into the same router. The `tako-rs-server` crate is what actually drives I/O: it accepts connections on a listener and dispatches each request through the core router. The thread-per-core variant lives separately in `tako-rs-server-pt` because it hosts the same thread-safe `Router` across N current-thread workers. For the architecture-level view of how a request flows through these layers, see [Architecture](/docs/concepts/architecture) and the [request lifecycle](/docs/concepts/request-lifecycle). For the cargo feature that each crate exposes, see the [feature reference](/docs/reference/features). Source: https://tako.rust-dd.com/docs/getting-started/project-layout --- # Architecture > The big picture of tako — Router, handlers, extractors, middleware, and responders in the core, with the server crate driving transports. Tako separates *what* a service does from *how* it is served. The framework primitives — routing, extraction, middleware, responses, state — live in `tako-rs-core` and are independent of any transport. The `tako-rs-server` crate drives concrete transports and feeds requests into that core. The result is one application model that works the same whether the bytes arrive over HTTP/1.1, HTTP/2, HTTP/3, a raw TCP socket, or a Unix socket. ## Layers | Layer | Crate | Responsibility | | ---------- | --------------------------------------------------- | -------------------------------------------------------------------------------------- | | Transports | `tako-rs-server`, `tako-rs-streams` | Accept connections, parse protocol frames, hand each request to the router | | Routing | `tako-rs-core` (`router`) | Match `(Method, path)` to a handler; run the middleware chain; answer `405` / fallback | | Middleware | `tako-rs-core` (`middleware`), `tako-rs-plugins` | Wrap the handler — auth, CORS, compression, rate limiting, metrics | | Extraction | `tako-rs-core` (`extractors`), `tako-rs-extractors` | Turn raw request parts into typed handler arguments | | Handler | your code | `async fn` from extracted arguments to a `Responder` | | Response | `tako-rs-core` (`responder`) | Convert the handler's return value into a `Response` | ## The core types A handful of types from `tako-rs-core` define the whole surface: * **`Router`** is the dispatch type. It maps `(Method, path)` pairs to handlers, runs global and per-route middleware, and produces a `405 Method Not Allowed` (with a populated `Allow` header) when the path matches but the method does not. Path syntax is `matchit`-compatible. * **`Request`** and **`Response`** are thin aliases — `http::Request` and `http::Response` — so the whole `http` crate ecosystem composes directly. `TakoBody` is Tako's streaming body type. * A **handler** is any `async fn` whose arguments implement `FromRequest` / `FromRequestParts` and whose return type implements `Responder`. The `Handler` trait is implemented automatically; you never write it by hand. * **`Responder`** converts a return value into a `Response`. Blanket impls cover `&str`, `String`, `Bytes`, `Json`, `(StatusCode, T)`, `StatusCode`, and `Result`. * **`IntoMiddleware`** plus the **`Next`** chain compose request processing. Middleware takes `(Request, Next)` and returns a `Response`; calling `next.run(req)` advances to the next middleware and eventually the handler. ```rust use tako::middleware::Next; use tako::types::{Request, Response}; async fn logging(req: Request, next: Next) -> Response { println!("{} {}", req.method(), req.uri()); next.run(req).await } ``` ## State model Application state is shared with handlers in two scopes. **Global state** is keyed by type and read with the `State` extractor. **Per-router state** is instance-local: `Router::with_state` attaches a typed container to one router instance, replacing the v1 process-global slot. The [State](/docs/state) page covers both; per-router state is what lets two independently nested routers carry different configuration. ## How transports plug in The server crate is the only layer that touches sockets. `Server::spawn_http` builds a default `Server` and accepts connections on a `TcpListener`, dispatching each request through the same `Router`. Other transports follow the identical shape through dedicated builder methods — `spawn_tls`, `spawn_h2c`, `spawn_h3`, `server_tcp`, `server_udp`, `server_unix`, and the compio variants — all of which drive the unchanged core router. WebTransport sessions start as CONNECT requests on the HTTP/3 server. Streaming and upgrade transports (WebSocket, SSE, file streaming) live in `tako-rs-streams` and hook into the same request flow. This is the payoff of the layering: routing, middleware, extraction, and responses are written once and reused across every transport. To follow a single request through all of these layers in order, read the [request lifecycle](/docs/concepts/request-lifecycle). For the crate-by-crate breakdown, see [project layout](/docs/getting-started/project-layout). Source: https://tako.rust-dd.com/docs/concepts/architecture --- # Request lifecycle > Trace a request through tako end to end — listener, server, Router match, middleware chain, extractors, handler, and Responder. This page traces a single request through Tako from the moment bytes arrive on a socket to the moment a response goes back out. Every stage uses the real type names from `tako-rs-core`. ## The path of a request 1. **Listener** — a `TcpListener` (or TLS / HTTP/2 / HTTP/3 / Unix listener) accepts a connection. 2. **Server** — `Server::spawn_http(listener, router)` drives the accept loop and parses each request into a `Request` (`http::Request`). 3. **Router match** — `Router` looks up the `(Method, path)` pair against its `matchit`-backed table. 4. **Middleware chain** — matched route and global middleware wrap the handler, each advancing the chain with `Next::run`. 5. **Extractors** — each handler argument is produced from the request via `FromRequest` / `FromRequestParts`. 6. **Handler** — your `async fn` runs against the extracted arguments and returns a value. 7. **Responder** — the return value is converted to a `Response` via `into_response`. 8. **Response** — the `Response` (`http::Response`) travels back out the same connection. ## 1–2. Listener and server Start the listener through the server builder: ```rust let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello_world); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); ``` `serve` builds a default `Server`, accepts connections, and turns each one into a `Request`. `Request` is an alias for `http::Request`, where `TakoBody` is Tako's streaming body type — so the body can be read incrementally rather than fully buffered. The matching `Response` alias is `http::Response`. ## 3. Router match `Router` maps each `(Method, path)` pair to a handler. Routes are registered with the explicit form or a method shorthand: ```rust router.route(Method::GET, "/users/{id}", show_user); router.get("/health", health); ``` If the path matches but the method does not, the router returns `405 Method Not Allowed` with an `Allow` header listing the supported methods. If no route matches at all, the router's `fallback` runs. Path syntax is `matchit`-compatible: `{name}` for a free segment and `{*rest}` for a catch-all. ## 4. Middleware chain Matched middleware wrap the handler. Each middleware takes the `Request` and a `Next`, does its work, and calls `next.run(req)` to advance the chain — the last link is the handler itself: ```rust use tako::middleware::Next; use tako::types::{Request, Response}; async fn auth(req: Request, next: Next) -> Response { // inspect or reject the request before it reaches the handler next.run(req).await } ``` Global middleware (added on the router) runs on every request; per-route middleware is attached by chaining `.middleware(mw)` on the `Arc` that `Router::route` returns. A middleware can short-circuit by returning a `Response` without calling `next.run`. ## 5. Extractors A handler's arguments are produced from the request before the handler body runs. Extractors implement one of two traits: * **`FromRequestParts`** — reads from the request head only (method, URI, headers, path/query parameters, state). Cheap and body-preserving. * **`FromRequest`** — consumes the body (e.g. `Json`, `Form`, `Bytes`, `Multipart`). At most one body-consuming extractor per handler. ```rust async fn show_user(Path(p): Path, State(db): State) -> impl Responder { // Path reads parameters from the head; State reads shared application state format!("user_id={}", p.id) } ``` If an extractor fails — malformed JSON, a path parameter that does not deserialise — the error is shaped by the router's error handlers (`client_error_handler` / `use_problem_json()`) instead of reaching the handler. ## 6–7. Handler and Responder The handler is an ordinary `async fn`. Its return type must implement `Responder`, which converts the value into a `Response`: ```rust async fn hello_world() -> impl Responder { "Hello, World!".into_response() } ``` Blanket `Responder` impls cover `&str`, `String`, `Bytes` / `Vec`, `Json`, `(StatusCode, T)`, `(StatusCode, HeaderMap, T)`, `StatusCode` alone, and `Result` where both arms are responders — the preferred shape for fallible handlers. ## 8. Response The resulting `Response` unwinds back through any middleware that is still awaiting `next.run`, giving each a chance to modify the response (add headers, compress the body, record metrics) before it is written to the connection and sent to the client. ## Where to go next * [Routing](/docs/routing) — path parameters, nesting, scopes, fallbacks. * [Extractors](/docs/extractors) — the full extractor catalog and ordering. * [Middleware](/docs/middleware) — the bundled middleware and writing your own. * [Architecture](/docs/concepts/architecture) — the layered view of these same pieces. Source: https://tako.rust-dd.com/docs/concepts/request-lifecycle --- # Runtimes > Tako's dual-runtime model — Tokio vs Compio, which transports each supports, TLS and HTTP/2 coverage, and how to switch. Tako runs on two async runtimes, chosen at build time. The application model — `Router`, handlers, extractors, middleware, responders — is identical on both. Transport availability and some middleware adapters differ. Router and route timeouts work on both; `middleware::timeout::Timeout` currently has a Tokio adapter only. * **Tokio** (default) — multi-threaded, work-stealing scheduler with `Send` futures end-to-end. This is what the standalone `Server` builder targets. * **Compio** (opt-in via the `compio` feature) — single-threaded thread-per-core over io\_uring / IOCP / kqueue. Futures are `!Send` and pinned to the runtime thread that produced them. ## One runtime selection per build Without `compio`, Tako selects Tokio. Enabling `compio` switches the HTTP server, router timers, middleware, queues, and per-thread workers to Compio. `cargo build --all-features` compiles and selects Compio; the Tokio-only raw QUIC helper (`RawQuicSession`) is excluded from that build. To use it, choose a Tokio feature set without `compio`. ## Switching runtimes By default you get Tokio. To build on Compio, disable defaults and enable the `compio` feature; add `compio-tls` and `compio-ws` for TLS and WebSocket support on that runtime. ```toml [dependencies] tako-rs = "2" ``` ```toml [dependencies] tako-rs = { version = "2", default-features = false, features = [ "compio", "compio-tls", "compio-ws", ] } ``` The umbrella crate swaps the server type for you: with `compio` on, the re-exports surface `CompioServer` / `CompioServerBuilder` instead of `Server` / `ServerBuilder`. The core router API is shared by both. The HTTP `serve_*` helpers are deprecated; use the selected builder and its `try_spawn_*` methods. ## Transport coverage Every transport runs on both runtimes: | Protocol | Tokio | Compio | Feature flag | | -------------------------- | ----- | ------ | --------------------- | | HTTP/1.1 | yes | yes | *default* | | HTTP/2 | yes | yes | `http2` | | HTTP/3 (QUIC) | yes | yes | `http3` | | TLS (rustls) | yes | yes | `tls` / `compio-tls` | | WebSocket | yes | yes | `ws` / `compio-ws` | | WebTransport | yes | yes | `webtransport` | | SSE | yes | yes | `sse` | | gRPC (unary and streaming) | yes | yes | `grpc` | | Raw TCP | yes | yes | *default* | | Raw UDP | yes | yes | `udp` | | Unix sockets | yes | yes | *default* (unix only) | | PROXY protocol v1/v2 | yes | yes | `proxy-protocol` | HTTP/2 over TLS, cleartext h2c, HTTP/3, WebTransport, Unix sockets, and PROXY protocol are supported on both runtimes. Only the raw QUIC helper (`RawQuicSession`) is Tokio-only. Raw TCP and UDP use each runtime’s native I/O API. ## Crossing the runtime boundary internally Where the framework itself must cross the boundary — for example, hyper's HTTP/2 service running over Compio TLS — it uses `send_wrapper::SendWrapper` to satisfy hyper's `Send` bound at the type level. The wrapper panics on cross-thread access, so a misuse becomes a loud panic rather than undefined behaviour. The soundness contract is per-runtime, not global: every wrapped future is constructed and polled on the same Compio runtime thread and never handed back to a multi-threaded Tokio executor. ## Thread-per-core The `per-thread` and `per-thread-compio` features host the same thread-safe `Router` on N current-thread workers, fanned out with `SO_REUSEPORT`; Tokio workers also hand new connections to the least busy worker. `serve_per_thread(addr, router, config)` blocks until Ctrl+C or SIGTERM and then drains the workers. `spawn_per_thread(addr, router, config)` returns the worker thread handles and a `PerThreadShutdown` for explicit control. Enabling `per-thread-compio` also selects Compio for all shared runtime-dependent code. These entry points live in `tako-rs-server-pt`. ## Related * [Installation](/docs/getting-started/installation) — selecting the runtime in `Cargo.toml`. * [Feature reference](/docs/reference/features) — the full runtime-selection feature table. * [Transports](/docs/transports) — per-protocol guides. Source: https://tako.rust-dd.com/docs/concepts/runtimes --- # Design philosophy > The principles behind tako — one service many transports, one model two runtimes, primitives included, and performance knobs when they matter. Tako is built for teams that want fewer moving parts in production. Rather than gluing several partially overlapping libraries together, it offers one coherent framework that spans protocols, runtimes, and the operational primitives real services need. Four principles drive the design. ## One service, many transports A single application surface covers HTTP/1.1, HTTP/2, HTTP/3, WebSocket, SSE, gRPC, raw TCP, UDP, Unix sockets, and WebTransport — with one routing, middleware, and observability model. You serve REST, realtime, and protocol-gateway workloads from the same binary without switching frameworks. The mechanism behind this is the layering: routing, extraction, middleware, and responses live in `tako-rs-core` and are transport-agnostic, while `tako-rs-server` drives the concrete protocols and feeds requests into that unchanged core. Add a WebSocket or a gRPC endpoint next to your HTTP routes and the request model does not change. ## One mental model, two runtimes The same framework style runs on **Tokio** or **Compio**, selected at build time to match the deployment constraint. Tokio gives you a multi-threaded work-stealing scheduler; Compio gives you single-threaded thread-per-core over io\_uring / IOCP / kqueue. Handlers, extractors, and middleware are written once and work on either — only the I/O machinery underneath differs. See [Runtimes](/docs/concepts/runtimes) for the trade-offs and the transport coverage on each. ## Application primitives included Middleware, auth, metrics, signals, queues, graceful shutdown, and streaming are part of the framework story, not an afterthought you assemble yourself. The bundled middleware covers JWT / Basic / Bearer / API-key auth, CSRF, sessions, body limits, request IDs, security headers, rate limiting, CORS, idempotency, and compression. Beyond request handling, an in-process signal bus, a background job queue, and first-class graceful shutdown ship in the box. This is what makes Tako a strong fit for protocol-heavy infrastructure — gateways, internal platforms, telemetry collectors, and control planes — where coordinating these concerns by hand is the bulk of the work. ## Performance knobs when they matter The fast paths are available without fragmenting the API. SIMD JSON (`sonic-rs` / `simd-json`), optional zero-copy extractors, brotli / gzip / deflate / zstd compression, jemalloc, and HTTP/3 are all opt-in cargo features layered on the same types — you turn them on where they pay off and ignore them elsewhere. The default build stays lean; the knobs are there when a workload needs them. See the [feature reference](/docs/reference/features) for the full set. ## Who this is for Tako is a strong fit when your service needs one or more of: * **More than REST** — HTTP APIs alongside WebSockets, SSE, gRPC, TCP, UDP, or QUIC in the same application. * **Realtime coordination** — built-in signals, queues, and streaming primitives instead of composing everything manually. * **Framework consolidation** — one coherent crate rather than several overlapping libraries. * **Protocol-heavy infrastructure** — gateways, internal platforms, telemetry collectors, control planes, and edge services. If you are building a plain JSON-over-HTTP service and never need the realtime or multi-transport surface, you can use Tako as a straight HTTP framework and ignore the rest. Nothing about the multi-transport design taxes the simple case. ## Related * [Architecture](/docs/concepts/architecture) — how the layers realise these principles. * [Project layout](/docs/getting-started/project-layout) — the crates that make up the framework. * [Quickstart](/docs/getting-started/quickstart) — the simple case, end to end. Source: https://tako.rust-dd.com/docs/concepts/design-philosophy --- # Tako vs Axum vs Actix Web > How Tako compares with Axum and Actix Web on transports, runtimes, bundled middleware, ecosystem, and throughput, and when to pick each one. Axum and Actix Web are the two most widely used Rust web frameworks. Tako shares their handler model: async functions, typed extractors, and responses built from a trait. The three differ mainly in how much ships in the box and how many transports one application can serve. The Tako maintainers wrote this page; it aims to be fair to all three, and the throughput numbers link to a reproducible methodology. ## Summary * **Axum** keeps a small core and builds on the Tower ecosystem. Pick it when you want Tower middleware, the largest community, and mostly HTTP/1.1 and HTTP/2 APIs. * **Actix Web** has been around since 2017 and has the fastest default setup in hello-world benchmarks; Tako's opt-in per-thread server is 3% ahead of it on loopback and 7% ahead when the server is the bottleneck. Pick Actix Web when a long track record and fast defaults matter most. * **Tako** puts several transports behind one router and middleware stack, and bundles auth, sessions, rate limiting, and metrics. Pick it when one service speaks HTTP/3, WebSocket, SSE, gRPC, or raw TCP/UDP alongside plain HTTP, or when you want Compio or a thread-per-core server inside the same framework. ## Feature comparison Versions compared: Tako 2.4, Axum 0.8, and Actix Web 4. | | Tako | Axum | Actix Web | | ------------------------------------------------ | ----------------------------------------------------------- | ------------------------------------------------------------------------ | --------------------------------------------------------------------------- | | HTTP/1.1 | Yes | Yes | Yes | | HTTP/2 | Yes (`http2`, including h2c) | Yes (`http2` feature) | Yes | | HTTP/3 (QUIC) | Yes (`http3`) | No built-in support | No built-in support | | WebTransport | W3C sessions over HTTP/3 (`webtransport`) | No built-in support | No built-in support | | TLS | Built in: rustls, mTLS, SNI, hot reload | Via `axum-server` or a rustls acceptor | Built in: rustls or OpenSSL | | WebSocket | Built in (`ws`) | Built in (`ws` feature) | `actix-ws` crate | | Server-Sent Events | Built in (`sse`) | Built in | `actix-web-lab` crate | | gRPC | Unary, server, client, and bidirectional streaming (`grpc`) | Through `tonic`, which shares hyper and Tower | No first-party support | | Unix sockets | Yes | Yes, `axum::serve` accepts a `UnixListener` | Yes (`bind_uds`) | | Raw TCP and UDP servers | Built in | Not in scope; use Tokio directly | Not part of the web framework | | Runtime | Tokio, or Compio (io\_uring on Linux, IOCP on Windows) | Tokio | actix-rt, a Tokio runtime per worker thread | | Thread-per-core server | `per-thread` feature | Not built in | Worker-per-core by design | | Middleware model | Tako middleware and plugins; no Tower compatibility | Tower layers and `middleware::from_fn` | Actix `Transform` and `middleware::from_fn` | | Auth, sessions, CSRF, rate limiting, idempotency | Bundled, feature-gated | Separate crates such as `tower-http`, `tower-sessions`, `tower_governor` | Separate crates such as `actix-session`, `actix-identity`, `actix-governor` | | OpenAPI | `utoipa` and `vespera` integrations | Community crates (`utoipa-axum`, `aide`) | Community crates (`utoipa-actix-web`, `apistos`) | | GraphQL | `async-graphql` integration | `async-graphql-axum` | `async-graphql-actix-web` | | Metrics | Prometheus and OpenTelemetry plugins | Community crates | Community crates | | First release | 2025 | 2021 | 2017 | ## Throughput Hello-world requests per second, the median of five 15-second `wrk` runs in a 24 vCPU Linux container, measured with tako-rs 2.4.0 in October 2026: | Framework | Loopback, 1,000 connections | Server-bound | Pipelined | | --------------- | --------------------------: | -----------: | ---------: | | Tako per-thread | 1,842,063 | 612,091 | 15,033,668 | | Actix Web | 1,786,640 | 570,865 | 12,806,477 | | Tako | 1,401,625 | 483,526 | 11,512,467 | | Tako + jemalloc | 1,386,954 | 469,807 | 10,853,271 | | Axum | 1,132,259 | 389,844 | n/a | Tako's default server runs on Tokio's work-stealing runtime, like Axum, and serves about 24% more requests at 1,000 connections and when the server is the bottleneck. The thread-per-core `per-thread` server runs one runtime per worker, like Actix Web, and serves 3% more requests than Actix Web on loopback, 7% more server-bound, and 17% more pipelined. Axum has no pipelined number because `axum::serve` leaves `TCP_NODELAY` off. A hello-world route measures framework overhead, not application performance; read the [benchmark methodology](/docs/benchmarks) before drawing conclusions. ## When to choose which Choose **Axum** when: * you want Tower middleware and the widest set of third-party integrations; * the service is mostly REST or JSON over HTTP/1.1 and HTTP/2; * a small core that you compose yourself is the goal. Choose **Actix Web** when: * you want the longest production track record among Rust web frameworks; * top hello-world throughput on the default setup matters most. Choose **Tako** when: * one process serves HTTP next to HTTP/3, WebSocket, SSE, gRPC, or raw TCP/UDP, and those transports should share routing, middleware, and signals; * you want auth, sessions, CSRF, rate limiting, idempotency, and metrics without assembling them from separate crates; * you want Compio or a thread-per-core server without leaving the framework. Know the trade-offs before choosing Tako: * The community is smaller and there are fewer third-party integrations. Tower layers do not plug in. * The API still moves. Versions 2.1, 2.2, and 2.3 had breaking changes, each with a [migration guide](/docs/reference/migration-2-3). * The minimum supported Rust version is 1.95. ## Porting an Axum service Handlers, extractors, and path syntax carry over almost unchanged. [Coming from Axum](/docs/getting-started/coming-from-axum) maps each Axum building block to its Tako equivalent. Source: https://tako.rust-dd.com/docs/concepts/comparison --- # Transports > Every Tako transport — HTTP, HTTP/3, WebSocket, WebTransport, SSE, gRPC, raw sockets, and PROXY protocol — with runtime and feature-flag support. Tako runs the same `Router` across many transports. A single `ServerConfig` flows into every one of them, so header read timeouts, keep-alive, drain timeout, max connections, H2 caps, H3 caps, and the PROXY read timeout are all sourced from one struct. You pick which transport a listener serves at spawn time — the routing, middleware, and observability model does not change. ## Transport matrix | Transport | Tokio | Compio | Feature flag | | -------------------------- | ----- | ------ | --------------------- | | HTTP/1.1 | yes | yes | *default* | | HTTP/2 | yes | yes | `http2` | | HTTP/3 (QUIC) | yes | yes | `http3` | | TLS (rustls) | yes | yes | `tls` / `compio-tls` | | WebSocket | yes | yes | `ws` / `compio-ws` | | WebTransport | yes | yes | `webtransport` | | SSE | yes | yes | `sse` | | gRPC (unary and streaming) | yes | yes | `grpc` | | Raw TCP | yes | yes | *default* | | Raw UDP | yes | yes | `udp` | | Unix sockets | yes | yes | *default* (unix only) | | PROXY protocol v1/v2 | yes | yes | `proxy-protocol` | Every transport runs on both runtimes: HTTP/2 over TLS, cleartext h2c, HTTP/3, WebTransport, raw TCP/UDP, Unix sockets, PROXY protocol, gRPC handlers, WebSocket, and SSE. Only the raw QUIC session helper behind `webtransport` is Tokio-only. See [Runtimes](/docs/concepts/runtimes) and the [feature reference](/docs/reference/features). ## How transports attach to a server Use `Server::builder()` on Tokio or `CompioServer::builder()` on Compio. The `try_spawn_*` methods validate configuration before returning a handle; `handle.result().await` reports asynchronous startup and serving errors. `shutdown(timeout).await` stops acceptance and bounds the drain. The HTTP `serve_*` helpers are deprecated; raw socket and per-thread helpers remain. ```rust use std::time::Duration; use tako::{Server, ServerConfig}; use tako::router::Router; use tokio::net::TcpListener; #[tokio::main] async fn main() -> Result<(), tako::types::BoxError> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "hi" }); let server = Server::builder() .config(ServerConfig { drain_timeout: Duration::from_secs(30), max_connections: Some(10_000), ..ServerConfig::default() }) .build(); let handle = server.try_spawn_http(listener, router)?; tako::shutdown_signal().await?; handle.shutdown(Duration::from_secs(30)).await; Ok(()) } ``` The handle is runtime-agnostic — both `Server` (tokio) and `CompioServer` (compio) return the same `ServerHandle` type. ## Per-transport pages * [HTTP](/docs/transports/http) — HTTP/1.1, HTTP/2 cleartext (h2c), and HTTP/2 over TLS via ALPN. * [HTTP/3](/docs/transports/http3) — HTTP/3 over QUIC on either runtime; TLS required. * [WebSocket](/docs/transports/websocket) — RFC-6455 upgrade inside an HTTP handler, on both runtimes. * [WebTransport](/docs/transports/webtransport) — W3C WebTransport sessions on the HTTP/3 server, reachable from browsers; requires the `webtransport` feature. * [SSE](/docs/streams) — Server-Sent Events; see the Streams primitives. * [gRPC](/docs/transports/grpc) — unary and streaming RPCs with protobuf. * [TCP / UDP / Unix](/docs/transports/tcp-udp-unix) — raw byte transports and Unix domain sockets. * [PROXY protocol](/docs/transports/proxy-protocol) — v1/v2 header parsing behind an L4 load balancer. ## WebSocket, SSE, and the HTTP handler WebSocket and SSE live inside HTTP handlers. `TakoWs::new(req, fut)` performs an HTTP/1.1 Upgrade; it works over plain HTTP/1.1 and TLS negotiated as HTTP/1.1. HTTP/2 Extended CONNECT is not implemented. `tako::sse::Sse` wraps a stream in a `text/event-stream` response and can use any supported HTTP transport. ## Raw transports `spawn_tcp_raw` / `tako::server_tcp::serve_tcp` and `spawn_udp_raw` / `tako::server_udp::serve_udp` skip the HTTP layer entirely. The handler closure receives each accepted stream (TCP) or datagram (UDP), so you can implement a custom line protocol, relay, or proxy. These do not use a `Router`. See [TCP / UDP / Unix](/docs/transports/tcp-udp-unix). Source: https://tako.rust-dd.com/docs/transports --- # HTTP > Serve HTTP/1.1, HTTP/2 cleartext (h2c), and HTTP/2 over TLS with ALPN negotiation on either runtime. HTTP/1.1 is the default transport — no feature flag required. HTTP/2 is added with the `http2` feature, and TLS with `tls`. The same `Router` serves all three; you choose the wire protocol at spawn time. HTTP/1.1, HTTP/2, and TLS all run on both the Tokio and Compio runtimes. ## HTTP/1.1 Bind a listener and pass it to `Server::spawn_http` on Tokio, or `CompioServer::spawn_http` on Compio. The returned handle reports failures and controls graceful shutdown. ```rust use anyhow::Result; use tako::Method; use tako::responder::Responder; use tako::router::Router; use tokio::net::TcpListener; async fn hello() -> impl Responder { "Hello, World!".into_response() } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` ```rust use anyhow::Result; use compio::net::TcpListener; use tako::Method; use tako::responder::Responder; use tako::router::Router; async fn hello() -> impl Responder { "Hello, World!".into_response() } #[compio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello); tako::CompioServer::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` The listener type, server builder, and `main` attribute select the runtime. The `Router` API is shared. See [Runtimes](/docs/concepts/runtimes) for how to pick. Source: [examples/hello-world/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/hello-world/src/main.rs) ```rust use anyhow::Result; use tako::Method; use tako::responder::Responder; use tako::router::Router; use tokio::net::TcpListener; async fn hello_world() -> impl Responder { "Hello, World!".into_response() } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello_world); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` ## The Server builder For graceful drain, a max-connection ceiling, or an owned shutdown trigger, build a `Server` explicitly. `spawn_http` returns a `ServerHandle`. ```rust use std::time::Duration; use tako::{Server, ServerConfig}; use tako::router::Router; use tokio::net::TcpListener; #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "hi" }); let server = Server::builder() .config(ServerConfig { drain_timeout: Duration::from_secs(30), max_connections: Some(10_000), ..ServerConfig::default() }) .build(); let handle = server.spawn_http(listener, router); tokio::signal::ctrl_c().await?; handle.shutdown(Duration::from_secs(30)).await; Ok(()) } ``` `ServerConfig` carries the production knobs shared by every transport: header read timeout, keep-alive and keep-alive timeout, HTTP/2 caps, HTTP/3 caps, `max_connections`, the PROXY read timeout, and the TLS handshake timeout. Its `Default` mirrors the historical hardcoded values (30 s drain, 30 s header read, 100 H2 streams), so anything you do not set keeps a safe default. ## HTTP/2 cleartext (h2c) Enable the `http2` feature. Prior-knowledge h2c is typically used behind an L7 proxy (Envoy, Nginx) that terminates TLS and forwards cleartext h2c upstream. Use `Server::spawn_h2c` on Tokio or `CompioServer::spawn_h2c` on Compio; both honor the HTTP/2 caps and keep-alive settings in `ServerConfig`. ```rust use tako::{Server, ServerConfig}; use tako::router::Router; use tokio::net::TcpListener; #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "h2c!" }); let server = Server::builder().config(ServerConfig::default()).build(); let handle = server.spawn_h2c(listener, router); handle.join().await; Ok(()) } ``` ```rust use compio::net::TcpListener; use tako::CompioServer; use tako::router::Router; use tako::types::BoxError; #[compio::main] async fn main() -> Result<(), BoxError> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.get("/", || async { "h2c!" }); let handle = CompioServer::builder().build().try_spawn_h2c(listener, router)?; handle.result().await?; Ok(()) } ``` h2c has no ALPN negotiation — the client must use prior knowledge that the endpoint speaks HTTP/2. Expose h2c only behind a proxy you control, never directly to browsers. ## HTTP/2 over TLS (ALPN h2) For browser-facing HTTP/2, combine `http2` with `tls` and use `spawn_tls`. When both features are on, Tako advertises `h2` and `http/1.1` via ALPN and the client negotiates the highest it supports — HTTP/2 if available, HTTP/1.1 otherwise. Attach a `TlsCert` to the builder. ```rust use tako::{Server, TlsCert}; use tako::router::Router; use tokio::net::TcpListener; #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = TcpListener::bind("127.0.0.1:8443").await?; let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "h2!" }); let server = Server::builder() .tls(TlsCert::pem_paths("cert.pem", "key.pem")) .build(); let handle = server.spawn_tls(listener, router); handle.join().await; Ok(()) } ``` H2 caps — `h2_max_concurrent_streams`, `h2_max_header_list_size`, `h2_max_send_buf_size`, the `h2_max_pending_accept_reset_streams` CVE-2023-44487 mitigation, and `h2_keep_alive_interval` — all live on `ServerConfig`. ### TLS without HTTP/2 The `tls` feature works on its own; without `http2`, ALPN advertises only `http/1.1`. The same `spawn_tls` builder path applies. ## TLS material `TlsCert` is the certificate source handed to the builder via `.tls(...)`: * `TlsCert::pem_paths(cert, key)` — load PEM cert and key from disk. * `TlsCert::der(certs, key)` — pre-loaded DER chain and key (for certs sourced from secret storage rather than the filesystem). * `TlsCert::resolver(resolver)` — a user-supplied `rustls::server::ResolvesServerCert` for SNI multi-cert serving or hot-reloadable certificates. `ReloadableResolver` backs the latter with an atomic, lock-free swap per handshake. Mutual TLS is configured by attaching a `ClientAuth` policy (`Optional` or `Required`) to any `TlsCert` variant via `.with_client_auth(...)`. ## Compio TLS On the Compio runtime, TLS lives behind the `compio-tls` feature. `CompioServer::builder().tls(cert).build().spawn_tls(listener, router)` mirrors the tokio path; the listener is a `compio::net::TcpListener`. HTTP/3 runs on Compio too, through `CompioServer::spawn_h3` — see [HTTP/3](/docs/transports/http3). ## Related * [Routing](/docs/routing) — declaring routes and handlers. * [HTTP/3](/docs/transports/http3) — the QUIC transport. * [Feature flags](/docs/reference/features) — `http2`, `tls`, `compio-tls`. Source: https://tako.rust-dd.com/docs/transports/http --- # HTTP/3 > Serve HTTP/3 over QUIC on Tokio (quinn) or Compio (compio-quic) with mandatory TLS, congestion control, and stream caps. HTTP/3 runs over QUIC (UDP) with the `h3` crate for the protocol: `quinn` drives QUIC on Tokio and `compio-quic` on Compio. It is gated behind the `http3` feature. The same `Router` you use for HTTP/1.1 and HTTP/2 drives HTTP/3 unchanged, and handlers see a `ConnInfo` whose transport is `Http3`. TLS is **mandatory** for QUIC, so a `TlsCert` is always required. On Compio a request body is bound to the runtime thread that received it, like every compio task; moving it to another thread panics. Enable `webtransport` to also accept [WebTransport sessions](/docs/transports/webtransport) on the same server. ## Quickstart QUIC binds a UDP socket, not a `TcpListener`, so the address is passed at spawn time. `spawn_h3` returns a handle for startup errors and shutdown; `try_spawn_h3` binds before returning, so certificate and bind errors surface immediately. ```rust use anyhow::Result; use tako::Method; use tako::responder::Responder; use tako::router::Router; async fn hello() -> impl Responder { "Hello, HTTP/3 World!".into_response() } #[tokio::main] async fn main() -> Result<()> { let mut router = Router::new(); router.route(Method::GET, "/", hello); tako::Server::builder() .tls(tako::TlsCert::pem_paths("cert.pem", "key.pem")) .build().spawn_h3("[::]:4433", router).result().await?; Ok(()) } ``` ```rust use tako::router::Router; use tako::types::BoxError; use tako::{CompioServer, TlsCert}; #[compio::main] async fn main() -> Result<(), BoxError> { let mut router = Router::new(); router.get("/", || async { "Hello, HTTP/3 World!" }); let handle = CompioServer::builder() .tls(TlsCert::pem_paths("cert.pem", "key.pem")) .build() .try_spawn_h3("[::]:4433", router)?; handle.result().await?; Ok(()) } ``` Test with a recent curl: ``` curl --http3 -k https://localhost:4433/ ``` Source: [examples/hello-world-http3/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/hello-world-http3/src/main.rs) ```rust use anyhow::Result; use tako::Method; use tako::responder::Responder; use tako::router::Router; async fn hello_world() -> impl Responder { "Hello, HTTP/3 World!".into_response() } async fn json_example() -> impl Responder { r#"{"message": "Hello from HTTP/3!", "protocol": "h3"}"#.into_response() } #[tokio::main] async fn main() -> Result<()> { // HTTP/3 requires TLS certificates // You can generate self-signed certificates for testing with: // openssl req -x509 -newkey rsa:4096 -keyout key.pem -out cert.pem -days 365 -nodes -subj "/CN=localhost" let mut router = Router::new(); router.route(Method::GET, "/", hello_world); router.route(Method::GET, "/json", json_example); println!("Starting HTTP/3 server on [::]:4433"); println!("Test with: curl --http3 -k https://localhost:4433/"); tako::Server::builder() .tls(tako::TlsCert::pem_paths("cert.pem", "key.pem")) .build() .spawn_h3("[::]:4433", router) .result() .await .expect("HTTP/3 server failed"); Ok(()) } ``` ## The Server builder For graceful drain and an owned shutdown handle, build a `Server`, attach a `TlsCert`, and call `spawn_h3` with the bind address. It returns the same `ServerHandle` as every other transport. ```rust use tako::{Server, TlsCert}; use tako::router::Router; #[tokio::main] async fn main() -> anyhow::Result<()> { let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "Hello over HTTP/3" }); let server = Server::builder() .tls(TlsCert::pem_paths("cert.pem", "key.pem")) .build(); let handle = server.spawn_h3("[::]:4433", router); handle.join().await; Ok(()) } ``` Without a `TlsCert`, `spawn_h3` reports the error through `result()` and `try_spawn_h3` returns it immediately — QUIC cannot run without TLS. Every `TlsCert` variant (PEM paths, DER, `Resolver`, mTLS) flows through the same rustls-config helper as HTTPS. ## Tuning QUIC HTTP/3 knobs live on `ServerConfig` alongside the HTTP/1 and HTTP/2 settings and apply on both runtimes: | Field | Purpose | | -------------------------------- | ---------------------------------------------------------------------------- | | `h3_max_concurrent_bidi_streams` | Cap on client-initiated bidirectional streams (default 100). | | `h3_max_concurrent_uni_streams` | Cap on client-initiated unidirectional streams (default 8). | | `h3_max_idle_timeout` | Idle timeout with no QUIC packets either direction (default 30 s). | | `h3_congestion` | Congestion controller: `Cubic` (default), `NewReno`, or `Bbr`. | | `h3_enable_datagrams` | Enable QUIC datagrams (RFC 9221) on the connection. | | `h3_use_retry` | Issue a QUIC Retry packet to validate source addresses (anti-amplification). | | `h3_goaway_grace` | Per-connection grace for in-flight streams after GOAWAY. | ```rust use tako::{Server, ServerConfig, TlsCert}; use tako::router::Router; use tako_rs_server::H3Congestion; #[tokio::main] async fn main() -> anyhow::Result<()> { let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "h3 tuned" }); let server = Server::builder() .tls(TlsCert::pem_paths("cert.pem", "key.pem")) .config(ServerConfig { h3_congestion: H3Congestion::Bbr, h3_use_retry: true, h3_enable_datagrams: true, ..ServerConfig::default() }) .build(); let handle = server.spawn_h3("[::]:4433", router); handle.join().await; Ok(()) } ``` `ServerConfig` is re-exported as `tako::ServerConfig`; the `H3Congestion` enum lives in `tako-rs-server`, so import it from there. The effective per-connection GOAWAY grace at runtime is `min(h3_goaway_grace, drain_timeout)` — a long per-connection grace cannot push the total shutdown past the global drain budget. `h3_goaway_grace` larger than `drain_timeout` is a no-op beyond the global ceiling. Raise `drain_timeout` too if you need a longer overall drain window. ## Certificate provider note With `http3`, rustls is compiled with two crypto providers: `aws-lc-rs` (its default) and `ring` (pulled in by the QUIC stack). Tako installs `aws-lc-rs` as the process-wide `CryptoProvider` when none is installed yet. To use a different one, install it at startup — for example `rustls::crypto::ring::default_provider().install_default()` — **before** constructing any server; once a provider is installed, Tako uses it as-is. ## Related * [HTTP](/docs/transports/http) — HTTP/1.1 and HTTP/2. * [WebTransport](/docs/transports/webtransport) — browser WebTransport sessions on the HTTP/3 server, behind `webtransport`. * [Feature flags](/docs/reference/features) — the `http3` feature and what it pulls in. Source: https://tako.rust-dd.com/docs/transports/http3 --- # WebSocket > Upgrade an HTTP request to a WebSocket inside a handler, then run the send/receive message loop on either runtime. WebSocket is not a separate listener — it lives inside an HTTP handler on an HTTP/1.1 listener, optionally over TLS. `TakoWs::new(req, fut)` performs the RFC-6455 server-side handshake and hands the upgraded stream to your closure. The handler owns the stream and drives the message loop. WebSocket is available on both runtimes: enable `ws` for Tokio (`tako::ws`), and behind the `compio-ws` feature for Compio (`tako::ws_compio`). ## Tokio On Tokio, `TakoWs` returns a `tokio_tungstenite` `WebSocketStream`. It is a `Sink` + `Stream` of `Message`, so you use `SinkExt::send` and `StreamExt::next`. ```rust use futures_util::SinkExt; use futures_util::StreamExt; use tako::Method; use tako::responder::Responder; use tako::types::Request; use tako::ws::TakoWs; use tokio_tungstenite::tungstenite::Message; use tokio_tungstenite::tungstenite::Utf8Bytes; async fn ws_echo(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { let _ = ws.send(Message::Text("Welcome to Tako WS!".into())).await; while let Some(Ok(msg)) = ws.next().await { match msg { Message::Text(txt) => { let _ = ws .send(Message::Text(Utf8Bytes::from(format!("Echo: {txt}")))) .await; } Message::Binary(bin) => { let _ = ws.send(Message::Binary(bin)).await; } Message::Ping(p) => { let _ = ws.send(Message::Pong(p)).await; } Message::Close(_) => { let _ = ws.send(Message::Close(None)).await; break; } _ => {} } } }) } #[tokio::main] async fn main() { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080") .await .unwrap(); let mut router = tako::router::Router::new(); router.route(Method::GET, "/ws/echo", ws_echo); tako::Server::builder().build().spawn_http(listener, router).result().await.unwrap(); } ``` The handler is just another route — register it with `router.route` and serve it on an HTTP/1.1 transport. Because the closure returns a `Responder`, the upgrade response (`101 Switching Protocols`) is produced for you; the upgraded socket is driven on a spawned task once the handshake completes. Source: [examples/websocket/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/websocket/src/main.rs) ```rust use std::time::Duration; use futures_util::SinkExt; use futures_util::StreamExt; use tako::Method; use tako::responder::Responder; use tako::types::Request; use tako::ws::TakoWs; use tokio_stream::wrappers::IntervalStream; use tokio_tungstenite::tungstenite::Message; use tokio_tungstenite::tungstenite::Utf8Bytes; pub async fn ws_echo(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { let _ = ws.send(Message::Text("Welcome to Tako WS!".into())).await; while let Some(Ok(msg)) = ws.next().await { match msg { Message::Text(txt) => { let _ = ws .send(Message::Text(Utf8Bytes::from(format!("Echo: {txt}")))) .await; } Message::Binary(bin) => { let _ = ws.send(Message::Binary(bin)).await; } Message::Ping(p) => { let _ = ws.send(Message::Pong(p)).await; } Message::Close(_) => { let _ = ws.send(Message::Close(None)).await; break; } _ => {} } } }) } pub async fn ws_tick(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { let mut ticker = IntervalStream::new(tokio::time::interval(Duration::from_secs(1))).enumerate(); loop { tokio::select! { msg = ws.next() => { match msg { Some(Ok(Message::Close(_))) | None => break, _ => {} } } Some((i, _)) = ticker.next() => { let _ = ws.send(Message::Text(Utf8Bytes::from(format!("tick #{i}")))).await; } } } }) } #[tokio::main] async fn main() { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080") .await .unwrap(); let mut router = tako::router::Router::new(); router.route(Method::GET, "/ws/echo", ws_echo); router.route(Method::GET, "/ws/tick", ws_tick); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); } ``` ## Compio On Compio, enable the `compio-ws` feature and use `TakoWsCompio`. The upgraded type is `CompioWebSocket`, which exposes an explicit `read()` method rather than the `Stream` interface. ```rust use compio::ws::tungstenite::Message; use tako::Method; use tako::responder::Responder; use tako::types::Request; use tako::ws_compio::CompioWebSocket; use tako::ws_compio::TakoWsCompio; use tako::ws_compio::UpgradedStream; async fn ws_echo(req: Request) -> impl Responder { TakoWsCompio::new(req, |mut ws: CompioWebSocket| async move { let _ = ws.send(Message::Text("Welcome to Tako WS (compio)!".into())).await; loop { match ws.read().await { Ok(Message::Text(txt)) => { let _ = ws.send(Message::Text(format!("Echo: {txt}").into())).await; } Ok(Message::Binary(bin)) => { let _ = ws.send(Message::Binary(bin)).await; } Ok(Message::Ping(p)) => { let _ = ws.send(Message::Pong(p)).await; } Ok(Message::Close(_)) => { let _ = ws.close(None).await; break; } Ok(_) => {} Err(_) => break, } } }) } #[compio::main] async fn main() { let listener = compio::net::TcpListener::bind("127.0.0.1:8080") .await .unwrap(); let mut router = tako::router::Router::new(); router.route(Method::GET, "/ws/echo", ws_echo); tako::CompioServer::builder().build().spawn_http(listener, router).result().await.unwrap(); } ``` Source: [examples/websocket-compio/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/websocket-compio/src/main.rs) ```rust //! WebSocket example using compio runtime. //! //! Run with: cargo run --example websocket-compio --features compio-ws use std::time::Duration; use compio::ws::tungstenite::Message; use tako::Method; use tako::responder::Responder; use tako::types::Request; use tako::ws_compio::CompioWebSocket; use tako::ws_compio::TakoWsCompio; use tako::ws_compio::UpgradedStream; pub async fn ws_echo(req: Request) -> impl Responder { TakoWsCompio::new(req, |mut ws: CompioWebSocket| async move { if let Err(e) = ws .send(Message::Text("Welcome to Tako WS (compio)!".into())) .await { tracing::error!("Failed to send welcome message: {}", e); return; } loop { match ws.read().await { Ok(Message::Text(txt)) => { let response = format!("Echo: {}", txt); if let Err(e) = ws.send(Message::Text(response.into())).await { tracing::error!("Failed to send echo: {}", e); break; } } Ok(Message::Binary(bin)) => { if let Err(e) = ws.send(Message::Binary(bin)).await { tracing::error!("Failed to send binary: {}", e); break; } } Ok(Message::Ping(p)) => { if let Err(e) = ws.send(Message::Pong(p)).await { tracing::error!("Failed to send pong: {}", e); break; } } Ok(Message::Close(_)) => { let _ = ws.close(None).await; break; } Ok(_) => {} Err(e) => { tracing::error!("WebSocket error: {}", e); break; } } } }) } pub async fn ws_count(req: Request) -> impl Responder { TakoWsCompio::new(req, |mut ws: CompioWebSocket| async move { let mut count: u64 = 0; loop { // Send current count let msg = format!("count: {}", count); if let Err(e) = ws.send(Message::Text(msg.into())).await { tracing::error!("Failed to send count: {}", e); break; } count += 1; // Wait 1 second compio::time::sleep(Duration::from_secs(1)).await; // Check for close message (non-blocking read attempt) // For simplicity, we just keep counting until connection drops } }) } #[compio::main] async fn main() { let listener = compio::net::TcpListener::bind("127.0.0.1:8080") .await .unwrap(); let mut router = tako::router::Router::new(); router.route(Method::GET, "/ws/echo", ws_echo); router.route(Method::GET, "/ws/count", ws_count); println!("WebSocket server running at:"); println!(" - ws://127.0.0.1:8080/ws/echo"); println!(" - ws://127.0.0.1:8080/ws/count"); tako::CompioServer::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); } ``` ## Configuration On Tokio, `TakoWs` is a builder. Each option is a chained method before the handler runs: | Method | Effect | | --------------------------------- | ------------------------------------------------------------------------------------------ | | `.protocols(list)` | Subprotocol allow-list; the first server-preferred match the client offers is echoed back. | | `.max_frame_size(n)` | Cap a single WebSocket frame, in bytes. | | `.max_message_size(n)` | Cap a reassembled message, in bytes. | | `.allowed_origins(list)` | Reject upgrades whose `Origin` is not on the list with `403`. | | `.upgrade_timeout(d)` | Drop the task if the client never completes the upgrade. | | `.max_lifetime(d)` | Hard cap on total conversation lifetime after upgrade. | | `.keep_alive(WsKeepAlive { .. })` | Deprecated; has no effect. Implement keep-alive in the handler. | ```rust use std::time::Duration; use tako::ws::TakoWs; # fn _doc(req: tako::types::Request) -> impl tako::responder::Responder { TakoWs::new(req, |ws| async move { /* ... */ }) .protocols(["chat", "superchat"]) .max_message_size(1 << 20) .allowed_origins(["https://app.example.com"]) .upgrade_timeout(Duration::from_secs(10)) .max_lifetime(Duration::from_secs(3600)) # } ``` The handler owns the socket and must implement ping intervals and pong deadlines. `max_lifetime` caps total conversation lifetime; it does not detect idle peers. ## TLS and HTTP/2 listeners The `examples/websocket-http2` TLS listener supports HTTP/2 for ordinary HTTP requests. WebSocket upgrades on that listener still negotiate HTTP/1.1. HTTP/2 Extended CONNECT is not supported. Invalid upgrade methods return 405, unsupported WebSocket versions return 426 with `Sec-WebSocket-Version: 13`, and malformed handshake headers or keys return 400. ## Related * [Routing](/docs/routing) — registering the upgrade handler as a route. * [HTTP](/docs/transports/http) — the listeners WebSocket rides on. * [Streams](/docs/streams) — SSE and other streaming responses. * [Feature flags](/docs/reference/features) — the `compio-ws` feature. Source: https://tako.rust-dd.com/docs/transports/websocket --- # WebTransport > Browser-compatible W3C WebTransport over HTTP/3 on Tokio or Compio, routed through the Router, plus a raw QUIC session helper on Tokio. Tako serves W3C WebTransport over HTTP/3, the protocol behind the browser's `new WebTransport(url)`. A session carries bidirectional and unidirectional streams plus unreliable datagrams over one QUIC connection. It runs on the HTTP/3 server on either runtime, behind the `webtransport` feature, which also enables `http3`. ```toml [dependencies] tako-rs = { version = "2", features = ["webtransport"] } ``` ## Accepting sessions A session starts as an HTTP/3 extended CONNECT request. It goes through the router like any other request, so path matching, extractors, and middleware such as authentication run first. Register a `CONNECT` route whose handler takes the `WebTransport` extractor and returns `on_session`: ```rust use tako::Method; use tako::TlsCert; use tako::router::Router; use tako::types::{BoxError, Response}; use tako::webtransport::{WebTransport, WebTransportSession}; use tokio::io::AsyncWriteExt; async fn echo(wt: WebTransport) -> Response { wt.on_session(|session: WebTransportSession| async move { while let Ok(Some(stream)) = session.accept_bi().await { tokio::spawn(async move { let (mut send, mut recv) = stream.split(); let _ = tokio::io::copy(&mut recv, &mut send).await; let _ = send.shutdown().await; }); } }) } #[tokio::main] async fn main() -> Result<(), BoxError> { let mut router = Router::new(); router.route(Method::CONNECT, "/echo", echo); tako::Server::builder() .tls(TlsCert::pem_paths("cert.pem", "key.pem")) .build() .try_spawn_h3("0.0.0.0:4433", router)? .result() .await?; Ok(()) } ``` `on_session` answers `200 OK`, and the callback runs once the server has sent the response. To refuse a session, return any non-2xx response instead. The `WebTransport` extractor rejects requests that are not a WebTransport CONNECT with `400 Bad Request`, and a path without a CONNECT route answers `404`. In the browser: ```js const transport = new WebTransport("https://example.com:4433/echo"); await transport.ready; const stream = await transport.createBidirectionalStream(); ``` ## Streams and datagrams | Method | Result | | ----------------------------------- | ----------------------------------------------- | | `accept_bi()` | The next bidirectional stream the client opens | | `accept_uni()` | The next unidirectional stream the client opens | | `open_bi()` / `open_uni()` | A stream opened by the server | | `read_datagram()` | The next datagram of the session | | `send_datagram(bytes)` | Sends an unreliable datagram | | `remote_address()` / `session_id()` | The client and the session | `BidiStream`, `SendStream`, and `RecvStream` implement `AsyncRead` and `AsyncWrite` from both `tokio::io` and `futures::io`. `BidiStream::split` returns separate sending and receiving sides; shutting down the sending side finishes the stream. `WebTransportSession` is `Clone`, so one task can read datagrams while another accepts streams. Each QUIC connection carries one session. Other HTTP/3 requests that arrive on the same connection are still served by the router. QUIC datagrams are always on with the `webtransport` feature, whatever `ServerConfig::h3_enable_datagrams` says, because WebTransport requires them. ## Closing a session A session ends when: * the browser calls `transport.close({ closeCode, reason })` or goes away; * the server calls `session.close(code, reason).await`; * the handler returns, which closes the session with code 0; * the server shuts down. After that, `accept_bi`, `accept_uni`, and `read_datagram` return `Ok(None)`, so accept loops end on their own. `session.closed().await` resolves with a `WebTransportClose { code, reason }` from whichever side closed first, and the browser reads the server's code and reason from `transport.closed`. ## Runtimes The session and its streams are `Send`, so the example above hands each stream to `tokio::spawn`. Serve the same router with `CompioServer::builder()` and `try_spawn_h3`. The session and its streams are `!Send` and stay on the connection's thread, so spawn per-stream work with `compio::runtime::spawn` and use the `futures::io` traits: ```rust use futures_util::io::AsyncWriteExt; use tako::types::Response; use tako::webtransport::{WebTransport, WebTransportSession}; async fn echo(wt: WebTransport) -> Response { wt.on_session(|session: WebTransportSession| async move { while let Ok(Some(stream)) = session.accept_bi().await { compio::runtime::spawn(async move { let (mut send, recv) = stream.split(); let _ = futures_util::io::copy(recv, &mut send).await; let _ = send.close().await; }) .detach(); } }) } ``` ## Certificates for local development Browsers only connect to a certificate they trust. During development, skip the certificate store by passing the certificate's SHA-256 hash to `serverCertificateHashes`: ```js const transport = new WebTransport("https://127.0.0.1:4433/echo", { serverCertificateHashes: [{ algorithm: "sha-256", value: certificateHash }], }); ``` The certificate must be an ECDSA P-256 certificate valid for at most 14 days. `examples/webtransport` generates one at startup, serves a test page at `http://127.0.0.1:3000`, and echoes a stream and a datagram: Source: [examples/webtransport/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/webtransport/src/main.rs) ```rust //! W3C WebTransport from a browser, with no certificate to install. //! //! Run it and open in Chrome, Edge, or Firefox. The //! page opens a session to `https://127.0.0.1:4433/echo` with //! `serverCertificateHashes`, which accepts the short-lived self-signed //! certificate generated at startup, then echoes a stream and a datagram. use rcgen::CertificateParams; use rcgen::KeyPair; use sha2::Digest; use sha2::Sha256; use tako::Method; use tako::TlsCert; use tako::body::TakoBody; use tako::router::Router; use tako::types::BoxError; use tako::types::Response; use tako::webtransport::WebTransport; use tako::webtransport::WebTransportSession; use time::Duration; use time::OffsetDateTime; use tokio::io::AsyncWriteExt; async fn echo(wt: WebTransport) -> Response { wt.on_session(|session: WebTransportSession| async move { println!("session from {}", session.remote_address()); let datagrams = tokio::spawn({ let session = session.clone(); async move { while let Ok(Some(datagram)) = session.read_datagram().await { let _ = session.send_datagram(datagram); } } }); while let Ok(Some(stream)) = session.accept_bi().await { tokio::spawn(async move { let (mut send, mut recv) = stream.split(); let _ = tokio::io::copy(&mut recv, &mut send).await; let _ = send.shutdown().await; }); } datagrams.abort(); let close = session.closed().await; println!("session closed: code {} ({:?})", close.code, close.reason); }) } /// A self-signed ECDSA P-256 certificate valid for ten days: the kind /// `serverCertificateHashes` accepts. Returns it with its SHA-256 hash. fn certificate() -> Result<(TlsCert, Vec), BoxError> { let key = KeyPair::generate()?; let mut params = CertificateParams::new(vec!["localhost".into(), "127.0.0.1".into()])?; params.not_before = OffsetDateTime::now_utc() - Duration::hours(1); params.not_after = OffsetDateTime::now_utc() + Duration::days(10); let cert = params.self_signed(&key)?; let hash = Sha256::digest(cert.der()).to_vec(); Ok((TlsCert::der(vec![cert.der().clone()], key.into()), hash)) } #[tokio::main] async fn main() -> Result<(), BoxError> { let (tls, hash) = certificate()?; let hash: Vec = hash.iter().map(u8::to_string).collect(); let page = include_str!("index.html").replace("CERT_HASH", &hash.join(",")); let mut pages = Router::new(); pages.get("/", move || { let mut page = Response::new(TakoBody::from(page.clone())); page.headers_mut().insert( tako::header::CONTENT_TYPE, tako::header::HeaderValue::from_static("text/html; charset=utf-8"), ); async move { page } }); let mut sessions = Router::new(); sessions.route(Method::CONNECT, "/echo", echo); // Browsers load pages over HTTP/3 only from a trusted certificate, so the // page comes over HTTP/1.1; `127.0.0.1` still counts as a secure context. let listener = tokio::net::TcpListener::bind("127.0.0.1:3000").await?; let page_server = tako::Server::builder() .build() .try_spawn_http(listener, pages)?; let webtransport = tako::Server::builder() .tls(tls) .build() .try_spawn_h3("127.0.0.1:4433", sessions)?; println!("open http://127.0.0.1:3000"); tokio::select! { result = page_server.result() => result?, result = webtransport.result() => result?, } Ok(()) } ``` ## Raw QUIC sessions `serve_webtransport(addr, cert_path, key_path, handler)` is a separate, lower-level API. It binds a plain QUIC endpoint and hands each connection to your handler as a `RawQuicSession`, with no HTTP/3 or WebTransport handshake. Browsers cannot connect to it; use it for private QUIC tunnels between peers you control. It is Tokio-only and needs the `webtransport` feature. Source: [examples/raw-quic/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/raw-quic/src/main.rs) ```rust use tako::webtransport::RawQuicSession; use tako::webtransport::serve_webtransport; async fn handle_session(session: RawQuicSession) { let remote = session.remote_address(); println!("New WebTransport session from {remote}"); // Handle bidirectional streams (echo server) loop { tokio::select! { // Accept bidirectional streams result = session.accept_bi() => { match result { Ok((mut send, mut recv)) => { println!("[{remote}] New bidirectional stream"); tokio::spawn(async move { let mut buf = vec![0u8; 4096]; loop { match recv.read(&mut buf).await { Ok(Some(n)) => { println!("[stream] Echoing {n} bytes"); if send.write_all(&buf[..n]).await.is_err() { break; } } _ => break, } } }); } Err(e) => { println!("[{remote}] Connection closed: {e}"); break; } } } // Read unreliable datagrams result = session.read_datagram() => { match result { Ok(data) => { println!("[{remote}] Datagram: {} bytes", data.len()); // Echo datagram back let _ = session.send_datagram(data); } Err(e) => { println!("[{remote}] Datagram error: {e}"); break; } } } } } println!("Session closed: {remote}"); } #[tokio::main] async fn main() { tracing_subscriber::fmt::init(); println!("WebTransport echo server on [::]:4433"); println!("Requires TLS certificates: cert.pem and key.pem"); println!(); println!("Generate self-signed certs for testing:"); println!(" openssl req -x509 -newkey ec -pkeyopt ec_paramgen_curve:prime256v1 \\"); println!(" -keyout key.pem -out cert.pem -days 365 -nodes -subj '/CN=localhost'"); serve_webtransport("[::]:4433", "cert.pem", "key.pem", |session| { Box::pin(handle_session(session)) }) .await; } ``` `RawQuicSession` wraps a `quinn::Connection` and exposes its stream and datagram surface directly: | Method | Purpose | | ---------------------- | --------------------------------------------------------------------- | | `remote_address()` | Peer socket address. | | `accept_bi()` | Accept an incoming bidirectional stream (`SendStream`, `RecvStream`). | | `accept_uni()` | Accept an incoming unidirectional stream (`RecvStream`). | | `open_bi()` | Open a new bidirectional stream. | | `open_uni()` | Open a new unidirectional stream. | | `read_datagram()` | Read an unreliable datagram. | | `send_datagram(bytes)` | Send an unreliable datagram. | | `close(code, reason)` | Close the session with a QUIC error code. | | `connection()` | Borrow the underlying `quinn::Connection` for advanced use. | `serve_webtransport_with_shutdown(addr, cert, key, handler, signal)` adds a graceful-shutdown future: the endpoint stops accepting, drains in-flight sessions for up to 30 s, then aborts the rest. The endpoint installs the `ring` rustls `CryptoProvider` if none is installed yet and advertises the `h3` ALPN protocol. ## Related * [HTTP/3](/docs/transports/http3): the QUIC transport WebTransport runs on. * [Feature flags](/docs/reference/features): the `webtransport` feature. Source: https://tako.rust-dd.com/docs/transports/webtransport --- # Server-Sent Events > Stream text/event-stream responses from any HTTP handler with the Sse responder, structured events, and proxy-safe keep-alive. Server-Sent Events (SSE) push a one-way, long-lived stream of `text/event-stream` frames from the server to the browser's `EventSource`. In Tako an SSE response is just a [`Responder`](/docs/routing) returned from an ordinary HTTP handler, so it flows through the same [`Router`](/docs/routing), middleware, and runtime as the rest of your app. `tako::sse::Sse` lives in `tako-rs-streams`, sits behind the `sse` feature, and works on both the **tokio** and **compio** runtimes. There are two construction paths: * `Sse::new(stream)` — the legacy raw-bytes path. Each `Bytes` item is wrapped as `data: …\n\n`. * `Sse::events(stream)` — the structured path. Each item is an `SseEvent` carrying any of `event:`, `id:`, `retry:`, comment, and `data:` fields. ## A minimal event stream `Sse::new` accepts any `Stream` that is `Send + 'static`. Each yielded chunk becomes a single `data:` event: ```rust use bytes::Bytes; use futures_util::{StreamExt, stream}; use tako::Method; use tako::responder::Responder; use tako::router::Router; use tako::sse::Sse; use tako::types::Request; async fn ticker(_: Request) -> impl Responder { let s = stream::unfold(0u64, |i| async move { tokio::time::sleep(std::time::Duration::from_secs(1)).await; Some((Bytes::from(format!("tick: {i}")), i + 1)) }); Sse::new(s) } #[tokio::main] async fn main() { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await.unwrap(); let mut router = Router::new(); router.route(Method::GET, "/events", ticker); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); } ``` The browser side is the standard `EventSource`: ```js const es = new EventSource("/events"); es.onmessage = (e) => console.log(e.data); ``` ## Response headers Both constructors emit a `200 OK` response with headers tuned for streaming through reverse proxies: * `Content-Type: text/event-stream` * `Cache-Control: no-cache, no-store, must-revalidate` * `Connection: keep-alive` * `X-Accel-Buffering: no` — disables nginx response buffering, which would otherwise hold frames until the connection closes. You do not set these yourself; the `Responder` impl builds the response. ## Structured events For named events, reconnection ids, or retry hints, build a `Stream` and hand it to `Sse::events`. `SseEvent` is a builder: * `SseEvent::data(d)` — a `data:` payload (multi-line strings are split into one `data:` field per line). * `SseEvent::comment(c)` — a `:` comment line, invisible to the `EventSource`. * `SseEvent::retry(duration)` — a `retry:` reconnection-delay hint. * `.event(name)` — sets the `event:` field so the client can `addEventListener(name, …)`. * `.id(value)` — sets the `id:` field, which the browser echoes back as `Last-Event-ID` on reconnect. ```rust use std::time::Duration; use futures_util::stream; use tako::responder::Responder; use tako::sse::{Sse, SseEvent}; use tako::types::Request; async fn feed(_: Request) -> impl Responder { let events = stream::iter([ SseEvent::data("hello"), SseEvent::data("again").event("greeting").id("1"), SseEvent::retry(Duration::from_secs(5)), ]); Sse::events(events) } ``` `Sse::new` does **not** sanitize embedded `\n` / `\r` in the raw bytes. A message containing `\n\nevent:click\n\n` is parsed by the browser as a second, synthetic event — a field-injection bug if the payload comes from untrusted input. The `Sse::events` path rebuilds every line with strict prefixes and strips CR, so it is safe for caller-controlled data. Reserve `Sse::new` for bytes you have already encoded yourself. ## Keep-alive Idle SSE connections can be closed by intermediaries that enforce read timeouts. Call `.keep_alive(period)` on either responder to interleave a `:keepalive\n\n` comment frame whenever the inner stream has been silent for `period`. The timer resets every time a real event is emitted, so active streams pay nothing. ```rust use std::time::Duration; use futures_util::stream; use tako::sse::{Sse, SseEvent}; let events = stream::iter([SseEvent::data("hi")]); let response = Sse::events(events).keep_alive(Duration::from_secs(15)); ``` Keep-alive is backed by `tokio::time::Sleep`. It is available on the structured and raw paths through the same `.keep_alive` method. ## Honoring reconnection When a client reconnects it sends the last `id:` it saw back as the `Last-Event-ID` request header. Two helpers read it so you can resume from the right cursor: * `tako::sse::last_event_id(headers)` — returns the trimmed value as `String` when it is valid UTF-8. * `tako::sse::last_event_id_bytes(headers)` — returns the raw trimmed bytes, preserving non-UTF-8 (e.g. opaque binary cursors) that the UTF-8 helper would drop. ```rust use tako::responder::Responder; use tako::sse::{last_event_id, Sse, SseEvent}; use tako::types::Request; use futures_util::stream; async fn resume(req: Request) -> impl Responder { let cursor = last_event_id(req.headers()).unwrap_or_default(); let events = stream::iter([ SseEvent::data(format!("resuming after {cursor}")).id("42"), ]); Sse::events(events) } ``` ## When to reach for something else SSE is one-way (server → client), text-framed, and rides on a normal HTTP request, which makes it ideal for live tickers, log tails, and progress feeds. For full-duplex messaging use a [WebSocket](/docs/transports/websocket); for binary, datagram-style channels over QUIC use [WebTransport](/docs/transports/webtransport). SSE also works over [HTTP/3](/docs/transports/http3) — see the `http3-sse` example. ## Examples * `examples/streams` — `Sse::new` ticker alongside hand-rolled `text/event-stream` bodies. * `examples/http3-sse` — SSE served over HTTP/3. Source: https://tako.rust-dd.com/docs/transports/sse --- # gRPC > Serve unary, server-streaming, client-streaming, and bidirectional gRPC over HTTP/2 on ordinary Tako routes, on Tokio or Compio. Tako serves gRPC by treating each method as a normal HTTP route. There is no separate gRPC server: a gRPC handler is an ordinary async function registered on the same [`Router`](/docs/routing) as your REST endpoints. All four RPC shapes are supported: unary, server streaming, client streaming, and bidirectional. The helpers live in `tako-rs-core` behind `grpc` and work on either runtime. Use an HTTP/2 listener: h2c or TLS on either runtime (TLS on Compio needs `compio-tls`). ```toml [dependencies] tako-rs = { version = "2", features = ["grpc", "http2"] } prost = "0.14" futures-util = "0.3" ``` The `grpc` feature gates the `tako::grpc` module. Protobuf message types come from [`prost`](https://docs.rs/prost); in production they are generated by `prost-build` from a `.proto` file. `futures-util` provides the `StreamExt` combinators used by streaming handlers. | RPC shape | Handler takes | Handler returns | | ---------------- | ----------------------- | --------------------------- | | Unary | `GrpcRequest` | `GrpcResponse` | | Server streaming | `GrpcRequest` | `GrpcServerStream` | | Client streaming | `GrpcClientStream` | `GrpcResponse` | | Bidirectional | `GrpcBidi` | `GrpcServerStream` | ## Defining messages Request and reply types are `prost::Message` derives. Inline definitions are fine for a quick service; real projects generate them: ```rust use prost::Message; #[derive(Clone, PartialEq, Message)] pub struct HelloRequest { #[prost(string, tag = "1")] pub name: String, } #[derive(Clone, PartialEq, Message)] pub struct HelloReply { #[prost(string, tag = "1")] pub message: String, } #[derive(Clone, PartialEq, Message)] pub struct Number { #[prost(int64, tag = "1")] pub value: i64, } ``` ## Unary RPCs `GrpcRequest` decodes the gRPC-framed protobuf body and exposes it as `req.message`. `GrpcResponse` frames the reply and sets the gRPC headers. Build a success with `GrpcResponse::ok(msg)` or an error with `GrpcResponse::error(code, message)` where `code` is a `GrpcStatusCode`: ```rust use tako::grpc::{GrpcRequest, GrpcResponse, GrpcStatusCode}; async fn say_hello(req: GrpcRequest) -> GrpcResponse { let name = &req.message.name; if name.is_empty() { return GrpcResponse::error(GrpcStatusCode::InvalidArgument, "name must not be empty"); } GrpcResponse::ok(HelloReply { message: format!("Hello, {name}!"), }) } ``` A successful reply sends the message and then `grpc-status: 0` in the HTTP/2 trailers. An error is a trailers-only response: HTTP `200 OK` with `grpc-status` and `grpc-message` in the headers and no body, per the gRPC HTTP/2 mapping. ## Server streaming Return a `GrpcServerStream` built from any `Stream` of `Result`. Each `Ok` item is sent as one message as soon as the stream yields it, and the call ends with `grpc-status: 0` when the stream finishes: ```rust use futures_util::StreamExt; use tako::grpc::{GrpcRequest, GrpcServerStream}; use tako::responder::Responder; async fn count(req: GrpcRequest) -> impl Responder { let numbers = futures_util::stream::iter(1..=req.message.value) .map(|value| Ok(Number { value })); GrpcServerStream::new(numbers) } ``` The first `Err(status)` ends the call with that status, and later items are not polled. Use `with_metadata(headers)` to send initial metadata with the response headers. ## Client streaming `GrpcClientStream` is a `Stream` of `Result` that yields messages while the client is still sending them. The `?` operator converts a `GrpcError` into a `GrpcStatus`, so a handler that returns `Result<_, GrpcStatus>` can bail out on a bad message: ```rust use futures_util::StreamExt; use tako::grpc::{GrpcClientStream, GrpcResponse, GrpcStatus}; async fn sum(mut numbers: GrpcClientStream) -> Result, GrpcStatus> { let mut value = 0; while let Some(number) = numbers.next().await { value += number?.value; } Ok(GrpcResponse::ok(Number { value })) } ``` The stream ends after the first error. A request that stops in the middle of a message yields `GrpcError::InvalidFrame`. ## Bidirectional streaming `GrpcBidi` reads `Req` messages like `GrpcClientStream` and names the reply type. `respond` builds the `GrpcServerStream` from the inbound stream. The reply starts as soon as the handler returns, so both directions are open at the same time and each reply can go out before the client sends its next message: ```rust use futures_util::StreamExt; use tako::grpc::{GrpcBidi, GrpcStatus}; use tako::responder::Responder; async fn running_total(bidi: GrpcBidi) -> impl Responder { bidi.respond(|numbers| { numbers.scan(0, |total, number| { let reply = number .map(|number| { *total += number.value; Number { value: *total } }) .map_err(GrpcStatus::from); futures_util::future::ready(Some(reply)) }) }) } ``` When replies do not map one-to-one onto requests, spawn a task that reads the inbound stream, and return a stream over the receiving end of a channel, such as `tokio_stream::wrappers::ReceiverStream`. ## Status and errors `GrpcStatusCode` is the full canonical set (`Ok`, `InvalidArgument`, `NotFound`, `DeadlineExceeded`, `Unimplemented`, `ResourceExhausted`, …). A `GrpcStatus` pairs a code with an optional message: * in a unary handler, return `GrpcResponse::error(code, message)`; * in a stream, yield `Err(GrpcStatus::error(code, message))` to end the call after the messages already sent; * to reject a call before any message, return `Result<_, GrpcStatus>` from the handler. `GrpcStatus` is itself a responder that sends a trailers-only response. Status messages are percent-encoded as the spec requires, so non-ASCII text survives the trip. ## Attaching the service to the router A gRPC method maps to `POST /./`. Register it like any other route. gRPC requires HTTP/2 on the wire, so enable the `http2` feature and serve over a transport that negotiates h2: `try_spawn_h2c` for cleartext behind a proxy, or TLS with ALPN (see [HTTP](/docs/transports/http)): ```rust use tako::Method; use tako::router::Router; use tako::types::BoxError; #[tokio::main] async fn main() -> Result<(), BoxError> { let mut router = Router::new(); router.route(Method::POST, "/grpc.Greeter/SayHello", say_hello); let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; tako::Server::builder() .build() .try_spawn_h2c(listener, router)? .result() .await?; Ok(()) } ``` ```rust use tako::Method; use tako::router::Router; use tako::types::BoxError; #[compio::main] async fn main() -> Result<(), BoxError> { let mut router = Router::new(); router.route(Method::POST, "/grpc.Greeter/SayHello", say_hello); let listener = compio::net::TcpListener::bind("127.0.0.1:8080").await?; tako::CompioServer::builder() .build() .try_spawn_h2c(listener, router)? .result() .await?; Ok(()) } ``` The path string is the gRPC contract — it must be the exact `/./` a client like `grpcurl`, Envoy, or a generated stub will call. There is no reflection-based auto-registration; you list each method explicitly. ## Deadlines and cancellation An `Option` handler argument reads the client's `grpc-timeout` header. Pass it to `with_deadline` and the stream ends with `DeadlineExceeded` once the deadline passes, even if the stream is waiting for its next item: ```rust use std::time::Duration; use tako::grpc::{GrpcDeadline, GrpcRequest, GrpcServerStream}; use tako::responder::Responder; async fn slow_count(deadline: Option, req: GrpcRequest) -> impl Responder { let last = req.message.value; let numbers = futures_util::stream::unfold(1, move |value| async move { if value > last { return None; } tokio::time::sleep(Duration::from_millis(200)).await; Some((Ok(Number { value }), value + 1)) }); GrpcServerStream::new(numbers).with_deadline(deadline) } ``` In a unary handler, bound the work with `deadline.remaining()`, for example with `tokio::time::timeout`. Parsing is overflow-safe: an absurdly large value is treated as "no deadline" rather than panicking. `read_grpc_deadline(&mut req)` stores the deadline in the request extensions for middleware, and the extractor reuses it. When a client cancels a call or disconnects, the server drops the reply stream together with everything it owns, so a producer task notices through its closed channel. An inbound stream whose client goes away yields `GrpcError::BodyReadError`. ## Message size limits Each inbound message is capped at `MAX_GRPC_MESSAGE_SIZE` (4 MiB), matching the default `grpc-go` and `tonic` server limits, for unary and streaming requests alike. A length prefix that advertises more is rejected with `ResourceExhausted` before the message is buffered, so a hostile client cannot force a large allocation with a few bytes. The extractors read the body message by message and never buffer it whole. Compressed messages are **not** supported. A request whose compressed-flag byte is set is rejected with `Unimplemented`, signalling the client to retry uncompressed. ## Calling it with grpcurl Tako has no server reflection, so pass the `.proto` file to the client: ```bash grpcurl -plaintext -import-path proto -proto numbers.proto \ -d '{"value": 3}' 127.0.0.1:50051 numbers.Numbers/Count grpcurl -plaintext -import-path proto -proto numbers.proto \ -d '{"value": 1} {"value": 2} {"value": 3}' 127.0.0.1:50051 numbers.Numbers/RunningTotal ``` Generated `tonic` clients work too: the streaming calls are tested against a tonic client on both runtimes. ## Other helpers The module also contains a gRPC-Web bridge, a `GrpcInterceptor` chain, and `grpc.health.v1` / reflection scaffolding. These are lower-level building blocks; for browser-facing realtime, prefer [SSE](/docs/transports/sse) or [WebSocket](/docs/transports/websocket). ## Examples * `examples/grpc-unary`: a `Greeter` and an `Echo` service registered as POST routes. * `examples/grpc-streaming`: server, client, and bidirectional streaming with a `proto/numbers.proto` for grpcurl. Source: [examples/grpc-streaming/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/grpc-streaming/src/main.rs) ```rust //! Server-streaming, client-streaming, and bidirectional gRPC on Tako. //! //! Run the server, then call it with grpcurl from this directory: //! //! ```text //! grpcurl -plaintext -import-path proto -proto numbers.proto \ //! -d '{"value": 3}' 127.0.0.1:50051 numbers.Numbers/Count //! grpcurl -plaintext -import-path proto -proto numbers.proto \ //! -d '{"value": 1} {"value": 2} {"value": 3}' 127.0.0.1:50051 numbers.Numbers/Sum //! grpcurl -plaintext -import-path proto -proto numbers.proto \ //! -d '{"value": 1} {"value": 2} {"value": 3}' 127.0.0.1:50051 numbers.Numbers/RunningTotal //! ``` use std::time::Duration; use futures_util::StreamExt; use prost::Message; use tako::Method; use tako::grpc::GrpcBidi; use tako::grpc::GrpcClientStream; use tako::grpc::GrpcDeadline; use tako::grpc::GrpcRequest; use tako::grpc::GrpcResponse; use tako::grpc::GrpcServerStream; use tako::grpc::GrpcStatus; use tako::grpc::GrpcStatusCode; use tako::responder::Responder; use tako::router::Router; use tako::types::BoxError; /// Generated by `prost-build` from `proto/numbers.proto` in a real project. #[derive(Clone, PartialEq, Message)] pub struct Number { #[prost(int64, tag = "1")] pub value: i64, } /// Server streaming: one request, a stream of replies. async fn count( deadline: Option, req: GrpcRequest, ) -> Result { let last = req.message.value; if !(1..=100).contains(&last) { return Err(GrpcStatus::error( GrpcStatusCode::InvalidArgument, "value must be between 1 and 100", )); } let numbers = futures_util::stream::unfold(1, move |value| async move { if value > last { return None; } tokio::time::sleep(Duration::from_millis(200)).await; Some((Ok(Number { value }), value + 1)) }); // Ends the call with DEADLINE_EXCEEDED if the client's `grpc-timeout` passes. Ok(GrpcServerStream::new(numbers).with_deadline(deadline)) } /// Client streaming: a stream of requests, one reply. async fn sum(mut numbers: GrpcClientStream) -> Result, GrpcStatus> { let mut value = 0; while let Some(number) = numbers.next().await { value += number?.value; } Ok(GrpcResponse::ok(Number { value })) } /// Bidirectional streaming: each reply is sent as soon as its request arrives. async fn running_total(bidi: GrpcBidi) -> impl Responder { bidi.respond(|numbers| { numbers.scan(0, |total, number| { let reply = number .map(|number| { *total += number.value; Number { value: *total } }) .map_err(GrpcStatus::from); futures_util::future::ready(Some(reply)) }) }) } #[tokio::main] async fn main() -> Result<(), BoxError> { let mut router = Router::new(); router.route(Method::POST, "/numbers.Numbers/Count", count); router.route(Method::POST, "/numbers.Numbers/Sum", sum); router.route(Method::POST, "/numbers.Numbers/RunningTotal", running_total); let listener = tokio::net::TcpListener::bind("127.0.0.1:50051").await?; println!("gRPC (h2c) on http://127.0.0.1:50051"); tako::Server::builder() .build() .try_spawn_h2c(listener, router)? .result() .await?; Ok(()) } ``` Source: https://tako.rust-dd.com/docs/transports/grpc --- # TCP, UDP & Unix sockets > Drop below HTTP to handle raw TCP streams, UDP datagrams, and Unix-domain-socket connections with a per-connection async handler. Not every service speaks HTTP. Tako exposes three raw listeners in `tako-rs-server` that skip the HTTP layer entirely and hand each accepted connection (or datagram) straight to your async closure: raw **TCP**, raw **UDP**, and **Unix domain sockets**. Use them for custom line protocols, relays, gateways, game servers, telemetry collectors, or local IPC. Raw TCP and Unix sockets are part of the default build; raw UDP needs the `udp` feature. Every listener exists on both runtimes under the same function name: with the `compio` feature the handler receives compio stream types and reads and writes through compio's owned-buffer I/O. Unix domain sockets are **unix-only** and compile out on other platforms. | Listener | Function | Handler receives | | -------- | ------------------------------- | --------------------------------------------------------------------------- | | TCP | `tako::server_tcp::serve_tcp` | `(TcpStream, SocketAddr)` | | UDP | `tako::server_udp::serve_udp` | `(Vec, SocketAddr, Arc)` | | Unix | `tako::server_unix::serve_unix` | `(UnixStream, SocketAddr)` on Tokio, `(UnixStream, UnixPeerAddr)` on Compio | Each handler returns a `Pin>` — wrap the body in `Box::pin(async move { … })`. The address argument is the peer; for UDP you also get a clone of the bound socket so you can reply. ## Raw TCP `serve_tcp` binds the address, sets `TCP_NODELAY` on each connection, and spawns your handler per accepted stream. The handler owns the `TcpStream` and reads/writes it directly with `tokio::io` traits. A common shape is a read loop that echoes: ```rust use tako::server_tcp::serve_tcp; use tokio::io::{AsyncReadExt, AsyncWriteExt}; #[tokio::main] async fn main() -> std::io::Result<()> { serve_tcp("127.0.0.1:9001", |mut stream, addr| { Box::pin(async move { let mut buf = vec![0u8; 4096]; loop { let n = stream.read(&mut buf).await?; if n == 0 { break; } stream.write_all(&buf[..n]).await?; eprintln!("echoed {n} bytes for {addr}"); } Ok(()) }) }) .await?; Ok(()) } ``` ## Raw UDP UDP is datagram-oriented, so there is no per-connection stream. `serve_udp` binds one socket, receives each packet (buffer sized for the 64 KiB max datagram), and calls your handler with the bytes, the sender's address, and an `Arc` of the socket for sending replies: ```rust use tako::server_udp::serve_udp; #[tokio::main] async fn main() -> std::io::Result<()> { serve_udp("127.0.0.1:9000", |data, addr, socket| { Box::pin(async move { let _ = socket.send_to(&data, addr).await; }) }) .await?; Ok(()) } ``` The UDP handler returns `()`, not a `Result` — there is no connection to tear down, so send errors are handled inside the closure (here ignored with `let _ =`). ## Unix domain sockets Unix sockets give you local, file-backed IPC without binding a TCP port — ideal for sidecars, supervisors, and admin endpoints. `serve_unix` is the raw equivalent of `serve_tcp`, handing you a `UnixStream` per connection: ```rust use tako::server_unix::serve_unix; use tokio::io::{AsyncReadExt, AsyncWriteExt}; #[tokio::main] async fn main() -> std::io::Result<()> { serve_unix("/tmp/tako.sock", |mut stream, _addr| { Box::pin(async move { let mut buf = vec![0u8; 4096]; let n = stream.read(&mut buf).await?; stream.write_all(&buf[..n]).await?; Ok(()) }) }) .await?; Ok(()) } ``` ```rust use compio::io::{AsyncRead, AsyncWriteExt}; use tako::server_unix::serve_unix; #[compio::main] async fn main() -> std::io::Result<()> { serve_unix("/tmp/tako.sock", |mut stream, _peer| { Box::pin(async move { let compio::BufResult(result, mut buf) = stream.read(vec![0u8; 4096]).await; buf.truncate(result?); let compio::BufResult(result, _) = stream.write_all(buf).await; result }) }) .await } ``` To serve your full HTTP [`Router`](/docs/routing) over a socket file instead of TCP — the typical "behind nginx/HAProxy on a local socket" deployment — use `Server::spawn_unix_http` on Tokio or `CompioServer::spawn_unix_http` on Compio. The `try_spawn_unix_http` variants bind the socket before returning, so bind errors surface immediately. The peer path is exposed to handlers as `tako::conn_info::UnixPeerAddr` in the request extensions, next to a `ConnInfo` whose transport is `Unix`: ```rust use tako::router::Router; use tako::Server; #[tokio::main] async fn main() { let mut router = Router::new(); router.route(tako::Method::GET, "/", || async { "hello from unix" }); Server::builder().build().spawn_unix_http("/tmp/tako-http.sock", router) .result().await.expect("Unix HTTP server failed"); } ``` ```rust use tako::CompioServer; use tako::router::Router; use tako::types::BoxError; #[compio::main] async fn main() -> Result<(), BoxError> { let mut router = Router::new(); router.get("/", || async { "hello from unix" }); let handle = CompioServer::builder() .build() .try_spawn_unix_http("/tmp/tako-http.sock", router) .await?; handle.result().await?; Ok(()) } ``` ```bash curl --unix-socket /tmp/tako-http.sock http://localhost/ ``` A path starting with `@` is treated as a Linux **abstract-namespace** socket (e.g. `@tako.sock`), which never touches the filesystem. Abstract sockets are Linux-only; on other unix platforms an `@`-path returns an `Unsupported` error. Unix sockets are **unix-only** and are absent on non-unix targets. Both runtimes clean up a stale filesystem socket before binding — it refuses to unlink a path that is not actually an `AF_UNIX` socket (symlink-escalation guard) and removes the file again on graceful shutdown so the next run can re-bind. ## Graceful shutdown Each function has `…_with_shutdown` and `…_with_shutdown_and_drain` variants that take a shutdown-signal future. On signal the listener stops accepting and in-flight work is drained, with a default 30-second bound you can override on the `_and_drain` form. For HTTP-over-unix, `serve_unix_http_with_config` accepts a [`ServerConfig`](/docs/concepts/runtimes) for connection caps and timeouts. ## HTTP alternatives If you want HTTP rather than a raw protocol, you usually do not need these: plain HTTP is covered in [HTTP/1.1 & HTTP/2](/docs/transports/http), and the builder's `spawn_unix_http` mirrors `serve_unix_http`. Reach for the raw listeners only when you control both ends of a non-HTTP protocol. ## Examples * `examples/tcp-echo` — line-echo over raw TCP. * `examples/udp-echo` — datagram echo over raw UDP. * `examples/unix-socket` — HTTP `Router` served over a Unix socket with `UnixPeerAddr`. Source: https://tako.rust-dd.com/docs/transports/tcp-udp-unix --- # PROXY protocol > Recover the real client address behind an L4 load balancer by parsing PROXY protocol v1/v2 headers into request extensions. When Tako sits behind an L4 (TCP) load balancer — HAProxy, AWS NLB, fly.io edges — the kernel only sees the balancer's address, not the real client's. The [PROXY protocol](https://www.haproxy.org/download/2.8/doc/proxy-protocol.txt) solves this: the balancer prepends a small header to every TCP connection carrying the original source and destination. Tako parses that header and rewrites the request so handlers see the true client. Both **v1** (human-readable text) and **v2** (binary, TLV-extensible) formats are supported. This lives in `tako-rs-server` behind the `proxy-protocol` feature and works on both runtimes. ## Serving with PROXY protocol `try_spawn_proxy_protocol` serves your [`Router`](/docs/routing) exactly like [`try_spawn_http`](/docs/transports/http), but reads and strips the PROXY header on each connection before dispatching. A handler that reports what the server recovered: ```rust use tako::Method; use tako::proxy_protocol::ProxyHeader; use tako::responder::Responder; use tako::router::Router; use tako::types::{BoxError, Request}; async fn handler(req: Request) -> impl Responder { let real_addr = req .extensions() .get::() .map(|a| a.to_string()) .unwrap_or_else(|| "unknown".into()); let proxy_info = req .extensions() .get::() .map(|h| format!("version={:?}, transport={:?}, src={:?}", h.version, h.transport, h.source)) .unwrap_or_else(|| "no PROXY header".into()); format!("Real client: {real_addr}\nPROXY header: {proxy_info}\n") } ``` Serve it on either runtime: ```rust #[tokio::main] async fn main() -> Result<(), BoxError> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", handler); let handle = tako::Server::builder() .build() .try_spawn_proxy_protocol(listener, router)?; handle.result().await?; Ok(()) } ``` ```rust #[compio::main] async fn main() -> Result<(), BoxError> { let listener = compio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", handler); let handle = tako::CompioServer::builder() .build() .try_spawn_proxy_protocol(listener, router)?; handle.result().await?; Ok(()) } ``` ## What handlers see After parsing, the server populates the request extensions: * A `std::net::SocketAddr` — the real client address from the PROXY header. This **overrides** the raw TCP peer address, so any extractor or middleware that reads the client IP gets the true value. * A `ConnInfo` carrying the same address. * The full `ProxyHeader` struct, for code that needs the raw fields. The server also rewrites forwarding headers defensively: inbound `Forwarded` and `X-Forwarded-*` are stripped (a client behind the hop must not spoof its address), and a single RFC 7239 `Forwarded: for=…` is re-emitted from the PROXY-supplied source so downstream middleware sees a consistent view. ## The ProxyHeader struct `ProxyHeader` exposes the parsed connection plus PROXY v2 metadata: * `version` — `ProxyVersion::V1` or `V2`. * `transport` — `ProxyTransport::Tcp`, `Udp`, or `Unknown`. * `source` / `destination` — the original `SocketAddr`s (`None` for `UNKNOWN`/`LOCAL`). * `source_unix` / `destination_unix` — `AF_UNIX` paths when the family is Unix. * `authority`, `alpn`, `aws_vpc_endpoint_id`, `tls`, `unique_id` — typed PROXY v2 TLVs surfaced where they map cleanly. * `tlvs` — the raw TLV list, for custom or future-defined types. * `crc32c_verified` — see below. ## Security and robustness The v2 parser is TLV-aware and CRC32C-verified. When a `PP2_TYPE_CRC32C` TLV is present, the checksum is recomputed over the full reconstructed header (with the CRC field zeroed): * `Some(true)` — TLV present and matched. * `Some(false)` — present but **mismatched**; the header may be corrupt or spoofed. * `None` — no CRC TLV (it is optional), or a v1 header. A CRC32C mismatch is logged but does **not** abort the parse — the framework hands you `crc32c_verified == Some(false)` and lets the operator decide. If you require verified headers, reject connections where it is `Some(false)`. The advertised v2 address length is also capped (536 bytes) so a malformed header cannot force a large per-connection allocation. The PROXY header is parsed under a `proxy_read_timeout` (from [`ServerConfig`](/docs/concepts/runtimes)) so a client that connects but never sends the header cannot pin a worker task forever. ## Manual parsing on raw TCP If you are running a custom protocol over [raw TCP](/docs/transports/tcp-udp-unix) rather than HTTP, call `tako::proxy_protocol::read_proxy_protocol(&mut stream)` yourself. It consumes the header and leaves the stream positioned at the start of your protocol data. The same code compiles on Compio, where `serve_tcp` hands the handler a compio stream and `read_proxy_protocol` reads it through compio I/O: ```rust use tako::proxy_protocol::read_proxy_protocol; use tako::server_tcp::serve_tcp; #[tokio::main] async fn main() -> std::io::Result<()> { serve_tcp("0.0.0.0:8080", |mut stream, _addr| { Box::pin(async move { let header = read_proxy_protocol(&mut stream).await?; eprintln!("real client: {:?}", header.source); Ok(()) }) }) .await?; Ok(()) } ``` ## Configuration and shutdown Pass a `ServerConfig` to the builder with `.config(…)` for connection caps, keep-alive, header and drain timeouts, and the PROXY read timeout — all sourced from one struct, as with every other transport. The returned `ServerHandle` drives graceful shutdown. The older `serve_http_with_proxy_protocol*` free functions are deprecated and Tokio-only. ## Testing You can exercise the v1 text format by hand — prepend the header line to a raw HTTP request: ```bash printf 'PROXY TCP4 192.168.1.100 10.0.0.1 56324 8080\r\nGET / HTTP/1.1\r\nHost: localhost\r\n\r\n' | nc 127.0.0.1 8080 ``` ## Example * `examples/proxy-protocol` — HTTP server reading the real client address from the PROXY header, with a `nc`-based v1 test command. Source: https://tako.rust-dd.com/docs/transports/proxy-protocol --- # Routing > Map (method, path) pairs to handlers, run middleware, nest sub-routers, and emit method-aware responses with Tako's Router. `Router` is the core dispatch type. It maps `(method, path)` pairs to handlers, runs middleware, and answers `405 Method Not Allowed` (with a populated `Allow` header) when the path matches but the method does not. Path syntax is `matchit`-compatible: `{name}` for a free segment, `{*rest}` for a catch-all, and `{name: T}` only inside the `#[tako::route]` macro family for compile-time-typed slots. ## Building a router ```rust use tako::Method; use tako::router::Router; let mut router = Router::new(); // Register routes with the explicit method form, … router.route(Method::GET, "/", root); router.route(Method::POST, "/users", create_user); // … or the shorthand methods. router.get("/health", health); router.put("/users/{id}", update_user); router.delete("/users/{id}", delete_user); async fn root() -> &'static str { "ok" } async fn health() -> &'static str { "ok" } async fn create_user() -> &'static str { "created" } async fn update_user() -> &'static str { "updated" } async fn delete_user() -> &'static str { "deleted" } ``` `Router::route` returns `Arc`. Attach per-route middleware by chaining `.middleware(mw)` on the returned route (the [`Router::middleware`](https://docs.rs/tako-rs/latest/tako/router/struct.Router.html) on the router itself adds *global* middleware that runs on every request). ## Dynamic path parameters ```rust use serde::Deserialize; use tako::Method; use tako::extractors::path::Path; use tako::responder::Responder; use tako::router::Router; #[derive(Deserialize)] struct UserPath { id: u64 } async fn show_user(Path(p): Path) -> impl Responder { format!("user_id={}", p.id) } let mut router = Router::new(); router.route(Method::GET, "/users/{id}", show_user); ``` For catch-all suffixes use `{*rest}`: ```rust router.get("/static/{*rest}", serve_static); ``` See [Extractors](/docs/extractors) for the full catalog of types you can bind to handler arguments, including [`Path`](/docs/extractors/request-meta). ### Typed slots via the route macro The `#[tako::get]` / `#[tako::post]` / `#[tako::route]` macros parse `{name: T}` at registration time and reject paths whose parameters don't deserialise: ```rust use tako::get; use tako::extractors::typed_params::TypedParams; use tako::responder::Responder; use tako::router::Router; #[get("/users/{id: u64}/posts/{post_id: u64}")] async fn show_post(TypedParams(p): TypedParams) -> impl Responder { format!("user={} post={}", p.user_id, p.post_id) } #[derive(serde::Deserialize)] struct ShowPost { id: u64, post_id: u64 } #[tokio::main] async fn main() -> anyhow::Result<()> { let mut router = Router::new(); router.mount_all(); // picks up every macro-registered route let listener = tako::bind_with_port_fallback("127.0.0.1:3000").await?; tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` `mount_all()` collects every route registered via the macro through the `linkme` distributed-slice. To mount them under a prefix use `mount_all_into("/api/v1")`. See [`examples/typed-routes`](https://github.com/rust-dd/tako/tree/main/examples/typed-routes). ## Sub-routing: `nest` and `scope` Both attach a prefix to a group of routes. They differ in *when* the grouping happens: * `nest(prefix, child_router)` — mount an entire existing router under a path prefix. Useful when you want a self-contained child router that can be tested independently. * `scope(prefix, |r| { ... })` — declare a temporary prefix and add routes inside the closure. Useful for inline grouping inside one builder. ```rust use tako::router::Router; use tako::responder::Responder; async fn list_users() -> impl Responder { "users" } async fn create_user() -> impl Responder { "created" } async fn dashboard() -> impl Responder { "admin home" } // child router, mounted with `nest` let mut api_v1 = Router::new(); api_v1.get("/users", list_users); api_v1.post("/users", create_user); let mut root = Router::new(); root.nest("/api/v1", api_v1); // inline grouping with `scope` root.scope("/admin", |r| { r.get("/", dashboard); }); ``` `nest` clones each child route via `Route::cloned_with_path`, so re-nesting the same child does not double-stack its middleware. The child router's *global* middleware chain is prepended to each nested route's middleware chain at registration time. Route-level plugins and the child's fallback / error handlers are **not** inherited — they are router-scoped. ## Method-aware responses When a path matches but the method does not, Tako returns `405 Method Not Allowed` with an `Allow` header listing the supported methods. This is different from v1, which returned plain `404 Not Found` for the same case. ```text $ curl -s -X DELETE http://127.0.0.1:8080/users -i HTTP/1.1 405 Method Not Allowed allow: GET, POST content-length: 0 ``` ## Trailing-slash redirects (`route_with_tsr`) `route(Method::GET, "/foo", h)` matches only `/foo`. If you also want `/foo/` to redirect to the canonical form (`308 Permanent Redirect`), register the route with `route_with_tsr`: ```rust router.route_with_tsr(Method::GET, "/api", api_handler); // "/api" -> handler // "/api/" -> 308 -> "/api" ``` The root path (`"/"`) is rejected at registration time because it has no canonical sibling to redirect to. ## Fallback and error handlers ```rust router.fallback(|_req: tako::types::Request| async { (http::StatusCode::NOT_FOUND, "not found") }); router.client_error_handler(|err| async move { (http::StatusCode::BAD_REQUEST, format!("{err}")) }); router.use_problem_json(); // map all errors through RFC 7807 Problem+JSON ``` `fallback` runs when no route matches at all. `error_handler` / `client_error_handler` shape how extractor / handler errors hit the wire. `use_problem_json()` is a convenience preset. ## What changed since 1.x A short summary; the full table lives in the [Migration guide](/docs/reference/migration): | Area | 1.x | 2.0 | | --------------------- | -------------------------------------------------- | ------------------------------------- | | Sub-routing | `Router::merge` (mutates a shared `Arc`) | `nest` / `scope` | | Wrong-method response | `404 Not Found` | `405 + Allow` | | Macro path syntax | `{id: u64}` only, always materialised as `Params` | `{id}` and `{id: u64}` both supported | | Per-router state | `GLOBAL_STATE` (one slot per `TypeId` per process) | `Router::with_state` (instance-local) | Server bootstrap and TLS knobs also changed substantially — see the [Transports overview](/docs/transports) and the migration guide. Source: https://tako.rust-dd.com/docs/routing --- # State > Share database pools, configuration, and handles with handlers via per-router State or process-global state in Tako. Tako gives handlers access to shared application state — database pools, configuration, signal arbiters, queue handles — through two mechanisms: * **Per-router state** via `Router::with_state(value)`. New in 2.0. Each `Router` carries its own typed store, so two routers in the same process can hold different values of the same type. * **Process-global state** via `tako::state::{set_state, get_state}`. Inherited from 1.x; still supported. The store is keyed by `TypeId`, so there is only one slot per `T` per process. The `State` extractor reads from the per-router store first and falls back to the process-global store, so most code only needs to care about one of the two paths at a time. ## Per-router state (recommended) ```rust use std::sync::Arc; use tako::Method; use tako::extractors::state::State; use tako::responder::Responder; use tako::router::Router; #[derive(Clone)] struct AppConfig { greeting: &'static str, } async fn hello(State(cfg): State) -> impl Responder { format!("{}, world!", cfg.greeting) } #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.with_state(AppConfig { greeting: "Hello" }); router.route(Method::GET, "/", hello); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` `with_state` writes into the router's `Arc`. The `State` extractor surfaces the stored value as `Arc`. The hot path is fast-checked with an `AtomicBool::Acquire`: when no caller ever invoked `with_state`, the state lookup short-circuits and adds no measurable overhead. You can call `with_state` multiple times, once per type: ```rust let mut router = Router::new(); router .with_state(db_pool) .with_state(redis_pool) .with_state(metrics); ``` Each handler can extract any subset: ```rust async fn handler( State(db): State, State(metrics): State, ) -> impl Responder { ... } ``` ## Multiple routers, distinct state A single process can host several routers — for example a public API on port `8080` and an internal admin API on `8081` — each with its own state: ```rust let mut public = Router::new(); public.with_state(public_cfg); public.get("/", hello); let mut admin = Router::new(); admin.with_state(admin_cfg); admin.get("/metrics", scrape_metrics); let server = Server::builder().build(); let public_h = server.spawn_http(public_listener, public); let admin_h = server.spawn_http(admin_listener, admin); tokio::join!(public_h.join(), admin_h.join()); ``` Both routers carry an `AppConfig`, but the values differ. With `GLOBAL_STATE` this required a newtype wrapper per router; with `with_state` it just works. ## Process-global state (legacy) The 1.x pattern is still available — useful when you don't have a `Router` in hand (background tasks, signal handlers, queue workers): ```rust use tako::state::{get_state, set_state}; #[derive(Clone)] struct Counter(u64); set_state(Counter(0)); // later, from a background task: if let Some(counter) = get_state::() { println!("count = {}", counter.0); } ``` `set_state` registers a value keyed by its type. `get_state::()` returns `Option>`. Because the store is global and keyed by `TypeId`, two callers cannot store distinct `T` values for the same `T` — last writer wins. Reach for `with_state` whenever the value belongs to a single router. See [`examples/with-state`](https://github.com/rust-dd/tako/tree/main/examples/with-state) for a runnable demonstration, and the [Migration guide](/docs/reference/migration) for the full upgrade path from `GLOBAL_STATE` to `Router::with_state`. Source: https://tako.rust-dd.com/docs/state --- # Extractors > Typed handler arguments that read shape out of a request — JSON, form, query, path, headers, cookies, JWT, and more. Extractors read shape out of a request and surface it as typed handler arguments. Any type implementing [`FromRequest`](https://docs.rs/tako-rs/latest/tako/extractors/trait.FromRequest.html) (consumes the body) or [`FromRequestParts`](https://docs.rs/tako-rs/latest/tako/extractors/trait.FromRequestParts.html) (headers / URL only) can be used as a handler parameter. A handler may take any number of extractors; Tako runs them in order before calling the handler. ```rust use serde::{Deserialize, Serialize}; use tako::Method; use tako::extractors::json::Json; use tako::extractors::path::Path; use tako::extractors::query::Query; use tako::extractors::state::State; use tako::responder::Responder; use tako::router::Router; #[derive(Deserialize)] struct ListQuery { page: u32, per_page: u32 } #[derive(Deserialize)] struct UserPath { id: u64 } #[derive(Deserialize, Serialize)] struct CreatePost { title: String, body: String } #[derive(Clone)] struct Db; // imagine a real pool here async fn list_posts( Path(UserPath { id }): Path, Query(q): Query, State(_db): State, ) -> impl Responder { format!("user={id}, page={}, per_page={}", q.page, q.per_page) } async fn create_post( Path(UserPath { id }): Path, Json(p): Json, ) -> Json { println!("creating post for user={id}: {:?}", p.title); Json(p) } #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.with_state(Db); router.route(Method::GET, "/users/{id}/posts", list_posts); router.route(Method::POST, "/users/{id}/posts", create_post); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` ## How binding works Each extractor is a wrapper type you destructure in the argument position — `Json(payload): Json`, `Path(p): Path`, and so on. The body extractors (`FromRequest`) consume the request body and must come last in the argument list, because the body can only be read once. The parts extractors (`FromRequestParts`) only touch headers and the URI, so any number of them can run before the body extractor. If extraction fails — a malformed JSON body, a missing required query parameter, a bad `Authorization` header — the extractor's error type is turned into a `Responder` (usually a `400` or `401`) and the handler is never called. Route-level error handlers and [`use_problem_json()`](/docs/routing) can reshape those errors before they hit the wire. Source: [examples/extractors-multi/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/extractors-multi/src/main.rs) ```rust use anyhow::Result; use serde::Deserialize; use serde::Serialize; use tako::Method; use tako::extractors::json::Json; use tako::extractors::params::Params; use tako::extractors::query::Query; use tako::router::Router; use tokio::net::TcpListener; #[derive(Deserialize)] struct Pagination { page: u32, per_page: u32, } #[derive(Deserialize)] struct UserPath { id: u64, } #[derive(Deserialize, Serialize, Clone)] struct CreateUser { name: String, email: String, } #[derive(Serialize)] struct Created { id: u64, name: String, email: String, } // GET /users/{id}/posts?per_page=10&page=2 // Demonstrates multiple extractors: Params + Query async fn list_user_posts(Params(user): Params, Query(p): Query) -> String { format!( "user_id={}, page={}, per_page={}", user.id, p.page, p.per_page ) } // POST /users with JSON body {"name":"...","email":"..."} // Demonstrates a body extractor async fn create(Json(user): Json) -> Json { // Normally you'd persist the user; here we just echo back with an id Json(Created { id: 1, name: user.name, email: user.email, }) } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/users/{id}/posts", list_user_posts); router.route(Method::POST, "/users", create); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` ## Catalog Tako ships 22+ bundled extractors. They fall into five groups, each with its own reference page. ### Body — `/docs/extractors/body` | Extractor | Description | | -------------------------------------- | ------------------------------------------- | | [`Json`](/docs/extractors/body) | JSON body (with optional SIMD acceleration) | | [`Form`](/docs/extractors/body) | URL-encoded form body | | [`Bytes`](/docs/extractors/body) | Raw request body stream | | [`Protobuf`](/docs/extractors/body) | Protocol Buffers body (feature `protobuf`) | | [`SimdJson`](/docs/extractors/body) | Force SIMD JSON parsing (feature `simd`) | | [`Multipart`](/docs/extractors/body) | Multipart form data (feature `multipart`) | ### Request metadata — `/docs/extractors/request-meta` | Extractor | Description | | ------------------------------------------------- | --------------------------- | | [`Query`](/docs/extractors/request-meta) | URL query parameters | | [`Path`](/docs/extractors/request-meta) | Typed route path parameters | | [`Params`](/docs/extractors/request-meta) | Dynamic path params map | | [`HeaderMap`](/docs/extractors/request-meta) | Full request headers | | [`Accept`](/docs/extractors/request-meta) | Content negotiation | | [`AcceptLanguage`](/docs/extractors/request-meta) | Language negotiation | | [`Range`](/docs/extractors/request-meta) | HTTP Range header | | [`IpAddr`](/docs/extractors/request-meta) | Client IP address | ### Auth — `/docs/extractors/auth` | Extractor | Description | | ------------------------------------------------- | ------------------------------------- | | [`Basic`](/docs/extractors/auth) | HTTP Basic credentials (`BasicAuth`) | | [`Bearer`](/docs/extractors/auth) | Bearer token (`BearerAuth`) | | [`ApiKey`](/docs/extractors/auth) | API key from header or query | | [`JwtClaimsUnverified`](/docs/extractors/auth) | JWT claims, signature **not** checked | | [`JwtClaimsVerified`](/docs/extractors/auth) | Signature-verified JWT claims | ### Cookies — `/docs/extractors/cookies` | Extractor | Description | | ------------------------------------------- | --------------------------------------- | | [`CookieJar`](/docs/extractors/cookies) | Cookie reading / writing | | [`CookieSigned`](/docs/extractors/cookies) | HMAC-signed cookies (`SignedCookieJar`) | | [`CookiePrivate`](/docs/extractors/cookies) | Encrypted cookies (`PrivateCookieJar`) | ### State — `/docs/state` | Extractor | Description | | ------------------------- | ------------------------------------- | | [`State`](/docs/state) | Shared application state, as `Arc` | ## Beyond the catalog The crate also bundles `QueryMulti` (repeated keys), `RawPath` / `RawQuery`, `MatchedPath`, `OriginalUri`, `Host`, `Scheme`, `TypedHeader`, `Extension`, `ConnectInfo`, `ContentLengthLimit`, and `Validated` (behind the `validator` / `garde` features). The `zero-copy-extractors` feature enables borrowing variants of the body and header extractors for hot-path handlers. See the [crate rustdoc](https://docs.rs/tako-rs/latest/tako/extractors/) for the complete list. ## Writing an extractor The last handler argument must implement `FromRequest`; earlier ones must implement `FromRequestParts`. An extractor that only reads headers, the URI, or extensions implements both, so it fits any position. Set `ENTRIES` to the framework entries it reads; this one reads none: ```rust use http::request::Parts; use tako::StatusCode; use tako::extractors::Entries; use tako::extractors::FromRequest; use tako::extractors::FromRequestParts; use tako::header::HeaderMap; use tako::types::Request; struct RequestId(String); fn request_id(headers: &HeaderMap) -> Result { headers .get("x-request-id") .and_then(|value| value.to_str().ok()) .map(|id| RequestId(id.to_owned())) .ok_or(StatusCode::BAD_REQUEST) } impl<'a> FromRequestParts<'a> for RequestId { type Error = StatusCode; const ENTRIES: Entries = Entries::NONE; async fn from_request_parts(parts: &'a mut Parts) -> Result { request_id(&parts.headers) } } impl<'a> FromRequest<'a> for RequestId { type Error = StatusCode; const ENTRIES: Entries = Entries::NONE; async fn from_request(req: &'a mut Request) -> Result { request_id(req.headers()) } } ``` The router can attach these entries to a request's extensions: | Entry | Attaches | Read by | | ----------------------- | ------------------------------------ | --------------------------------------------------------------------- | | `Entries::CONN` | `ConnInfo` and the peer `SocketAddr` | `ConnectInfo`, `IpAddr` | | `Entries::MATCHED_PATH` | `MatchedPath` | `MatchedPath` | | `Entries::PARAMS` | the matched path parameters | `Path`, `Params`, `TypedParams` | | `Entries::STATE` | router and route state | `State`, `IpAddr` | | `Entries::BODY_LIMIT` | the body size limit | `Json`, `Form`, `Bytes`, `String`, `Protobuf`, `SimdJson` | | `Entries::SIMD_JSON` | the route's SIMD JSON mode | `Json` | A route without middleware or other per-request hooks gets only the entries its handler's extractors list, which saves work on every request; all other requests get every entry. Combine entries with `Entries::CONN.union(Entries::STATE)`. `ENTRIES` defaults to `Entries::ALL`, which is always correct. An extractor that reads an entry it does not list finds it missing on those routes. Source: https://tako.rust-dd.com/docs/extractors --- # Body extractors > Read and deserialize request bodies — Json, Form, raw Bytes, Protobuf, SIMD JSON, and multipart uploads. Body extractors consume the request body and deserialize it into a typed handler argument. They implement [`FromRequest`](https://docs.rs/tako-rs/latest/tako/extractors/trait.FromRequest.html), so a handler may take **at most one** of them, and it must come after any [request-metadata extractors](/docs/extractors/request-meta) — the body is read once. On a content-type mismatch or a deserialization failure the extractor's error type becomes a `400 Bad Request` response and the handler is skipped. ## `Json` The default JSON body extractor. `T` must implement `serde::Deserialize`. `Json` also implements `Responder`, so handlers can return it to emit a JSON response with `Content-Type: application/json`. ```rust use serde::{Deserialize, Serialize}; use tako::extractors::json::Json; #[derive(Deserialize, Serialize)] struct CreateUser { name: String, email: String, } async fn create_user(Json(user): Json) -> Json { println!("creating {}", user.name); Json(user) } ``` With the `simd` feature on, `Json` automatically dispatches to a SIMD-accelerated parser for large payloads; the threshold is configurable per route via `Route::simd_json(SimdJsonMode)`. For an unconditional SIMD path, use [`SimdJson`](#simdjsont) below. ## `Form` Parses an `application/x-www-form-urlencoded` body via `serde_urlencoded`. `T` must implement `serde::Deserialize`. The extractor checks that the `Content-Type` starts with `application/x-www-form-urlencoded` (the `; charset=utf-8` variant is accepted) and otherwise returns `FormError::InvalidContentType`. ```rust use serde::Deserialize; use tako::extractors::form::Form; #[derive(Deserialize)] struct LoginForm { username: String, password: String, } async fn login(Form(form): Form) { println!("login attempt for {}", form.username); } ``` `FormError` covers `InvalidContentType`, `BodyReadError`, `InvalidUtf8`, `ParseError`, and `DeserializationError`; all map to `400 Bad Request`. ## `Bytes` Raw access to the underlying body stream, without buffering. `Bytes<'a>` wraps `&'a mut TakoBody`, so you drive the read yourself (for streaming, custom framing, or hashing). Its error type is `Infallible`. ```rust use http_body_util::BodyExt; use tako::extractors::bytes::Bytes; async fn raw_body(Bytes(body): Bytes<'_>) { let collected = body.collect().await.unwrap().to_bytes(); println!("read {} bytes", collected.len()); } ``` This `Bytes<'a>` is a body **reference**, distinct from the refcounted buffer `bytes::Bytes` from the `bytes` crate. In a handler that needs both, import it aliased: `use tako::extractors::bytes::Bytes as BytesBody;`. ## `Protobuf` Requires the `protobuf` feature. Decodes a Protocol Buffers body into a `prost::Message`. The extractor accepts `application/x-protobuf` or `application/protobuf` (with optional parameters) and rejects everything else with `ProtobufError::InvalidContentType`. `Protobuf` also implements `Responder`, encoding the message back out with `Content-Type: application/x-protobuf`. ```rust use prost::Message; use tako::extractors::protobuf::Protobuf; #[derive(Clone, PartialEq, Message)] struct CreateUserRequest { #[prost(string, tag = "1")] pub name: String, #[prost(string, tag = "2")] pub email: String, } async fn create_user(Protobuf(req): Protobuf) -> String { format!("creating {}", req.name) } ``` ## `SimdJson` Requires the `simd` feature. Forces SIMD-accelerated JSON parsing regardless of payload size, backed by the `simd_json` crate. The sibling `SonicJson` uses the `sonic_rs` backend with the same API. Both validate the JSON content type, read the full body, and implement `Responder` for SIMD-serialized JSON responses. ```rust use serde::{Deserialize, Serialize}; use tako::extractors::simdjson::SimdJson; #[derive(Deserialize, Serialize)] struct User { name: String, email: String, age: u32, } async fn create_user(SimdJson(user): SimdJson) -> SimdJson { println!("creating {}", user.name); SimdJson(user) } ``` `SimdJsonError` distinguishes `InvalidContentType`, `MissingContentType`, `BodyReadError`, and `DeserializationError`, all mapping to `400 Bad Request`. ## `Multipart` Requires the `multipart` feature. Two extractors handle `multipart/form-data` (file uploads and mixed forms): * **`TakoMultipart<'a>`** — raw access. Wraps a `multer::Multipart`; you iterate fields manually with `next_field()`. * **`TakoTypedMultipart<'a, T, F>`** — strongly typed. Deserializes every part into `T`, using the file-field type `F` (one of `UploadedFile`, `InMemoryFile`, or `BufferedUploadedFile`) for parts that carry a filename. ```rust use serde::Deserialize; use tako::extractors::multipart::{TakoTypedMultipart, UploadedFile}; #[derive(Deserialize)] struct FileUploadForm { title: String, description: String, file: UploadedFile, } async fn upload( TakoTypedMultipart { data: form, .. }: TakoTypedMultipart<'_, FileUploadForm, UploadedFile>, ) { println!("uploaded {:?} ({} bytes)", form.file.file_name, form.file.size); println!("saved to {:?}", form.file.path); } ``` `UploadedFile` streams to a temp file whose on-disk name is a fresh UUID — the client-supplied filename is preserved only in `UploadedFile.file_name` and never influences the path, which closes off path-traversal. The temp file is removed on drop (RAII); call `persist(dest)` or `disarm_cleanup()` to keep it. Limits come from a `MultipartConfig` (inserted into request extensions or set as global state): `total_size_limit`, `per_part_size_limit` (default 1 MiB), `max_parts`, `allowed_content_types`, `disk_spill_threshold`, and `field_chunk_timeout`. The typed extractor enforces `max_parts` and the content-type allow-list; the raw extractor leaves part-count enforcement to the caller. Source: [examples/multipart/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/multipart/src/main.rs) ```rust use anyhow::Result; use http::Method; use http::StatusCode; use tako::extractors::FromRequest; use tako::extractors::multipart::InMemoryFile; use tako::extractors::multipart::TakoMultipart; use tako::extractors::multipart::TakoTypedMultipart; use tako::extractors::multipart::UploadedFile; use tako::responder::Responder; use tako::router::Router; use tako::types::Request; use tokio::net::TcpListener; async fn upload_file(mut req: Request) -> impl Responder { #[derive(serde::Deserialize)] struct Form { description: String, file: UploadedFile, } let TakoTypedMultipart:: { data, .. } = TakoTypedMultipart::from_request(&mut req).await.unwrap(); println!( "uploaded {:?} ({} bytes) — description: {}", data.file.file_name, data.file.size, data.description, ); (StatusCode::OK, "File uploaded successfully") } async fn upload_mem(mut req: Request) -> impl Responder { #[derive(serde::Deserialize)] struct ImgForm { title: String, image: InMemoryFile, } let TakoTypedMultipart:: { data, .. } = TakoTypedMultipart::from_request(&mut req).await.unwrap(); println!( "image '{}' received ({} bytes in memory)", data.title, data.image.data.len(), ); (StatusCode::OK, "Image uploaded successfully") } async fn raw_with_file(mut req: Request) -> impl Responder { let TakoMultipart(mut mp) = TakoMultipart::from_request(&mut req).await.unwrap(); let mut total_files = 0usize; while let Some(mut field) = mp.next_field().await.unwrap() { if field.file_name().is_some() { let fname = field .file_name() .map(|s| s.to_owned()) .unwrap_or_else(|| "".into()); total_files += 1; let mut size = 0usize; while let Some(chunk) = field.chunk().await.unwrap() { size += chunk.len(); } println!("received {fname} ({size} bytes)"); } } (StatusCode::OK, format!("processed {total_files} file(s)")) } async fn raw_text(mut req: Request) -> impl Responder { use std::collections::HashMap; use tako::types::BuildHasher; let TakoMultipart(mut mp) = TakoMultipart::from_request(&mut req).await.unwrap(); let mut map: HashMap = HashMap::with_hasher(BuildHasher::default()); while let Some(field) = mp.next_field().await.unwrap() { if field.file_name().is_some() { return (StatusCode::BAD_REQUEST, "file not accepted"); } let name = field.name().unwrap_or("noname").to_owned(); let text = field.text().await.unwrap(); map.insert(name, text); } (StatusCode::OK, "text form processed") } async fn typed_text(mut req: Request) -> impl Responder { #[derive(serde::Deserialize)] struct LoginForm { username: String, password: String, } let TakoTypedMultipart:: { data, .. } = TakoTypedMultipart::from_request(&mut req).await.unwrap(); println!( "login attempt: username='{}' (password length: {})", data.username, data.password.len(), ); (StatusCode::OK, "typed text processed") } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::POST, "/upload_file", upload_file); router.route(Method::POST, "/upload_mem", upload_mem); router.route(Method::POST, "/raw_with_file", raw_with_file); router.route(Method::POST, "/raw_text", raw_text); router.route(Method::POST, "/typed_text", typed_text); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` ## Related * Return-side encoding and the `Responder` trait — see [Routing](/docs/routing). * Bound body sizes per handler with `ContentLengthLimit`, or globally with the [Body Limit middleware](/docs/middleware). * Non-body inputs — [request metadata extractors](/docs/extractors/request-meta). Source: https://tako.rust-dd.com/docs/extractors/body --- # Request metadata extractors > Read URL, path, header, and connection metadata — Query, Path, Params, HeaderMap, Accept, AcceptLanguage, Range, and IpAddr. These extractors read the URL, path captures, headers, and connection info — never the body. They implement [`FromRequestParts`](https://docs.rs/tako-rs/latest/tako/extractors/trait.FromRequestParts.html), so any number of them can be combined in one handler, and they run before a [body extractor](/docs/extractors/body). ## `Query` Deserializes the URL query string into `T` (any `serde::Deserialize`). Optional fields map naturally to `Option<_>`. ```rust use serde::Deserialize; use tako::extractors::query::Query; #[derive(Deserialize)] struct SearchQuery { q: String, page: Option, limit: Option, } async fn search(Query(query): Query) -> String { let page = query.page.unwrap_or(1); format!("searching '{}' (page {page})", query.q) } ``` `QueryError` distinguishes `MissingQueryString`, `ParseError`, and `DeserializationError`. For query strings with repeated keys (`?tag=a&tag=b`) reach for `QueryMulti`. ## `Path` Typed route path parameters (axum parity). `T` may be a single primitive (`Path`), a tuple (`Path<(u64, String)>`), a `Vec<_>` for repeated captures, an `Option<_>` (`None` when nothing matched), or a struct deriving `serde::Deserialize`. ```rust use serde::Deserialize; use tako::extractors::path::Path; use tako::responder::Responder; #[derive(Deserialize)] struct UserPath { id: u64 } async fn show_user(Path(p): Path) -> impl Responder { format!("user_id={}", p.id) } ``` The route must declare the matching capture, e.g. `router.route(Method::GET, "/users/{id}", show_user)` — see [Routing](/docs/routing). For the verbatim request path (no captures, no decoding) use `RawPath`. ## `Params` The dynamic path-params extractor. Like `Path`, it deserializes the captured segments into `T`, and it is the form the `#[tako::route]` macro family materializes for `{name}` slots. ```rust use serde::Deserialize; use tako::extractors::params::Params; use tako::extractors::query::Query; #[derive(Deserialize)] struct UserPath { id: u64 } #[derive(Deserialize)] struct Pagination { page: u32, per_page: u32 } // GET /users/{id}/posts?page=2&per_page=10 async fn list_posts( Params(user): Params, Query(p): Query, ) -> String { format!("user={}, page={}, per_page={}", user.id, p.page, p.per_page) } ``` `ParamsError` is either `MissingPathParams` (an internal routing error) or `DeserializationError`. ## `HeaderMap` Gives a handler the full request header set. `HeaderMap(pub http::HeaderMap)` is owned, so you read values without lifetime juggling. Its error type is `Infallible`. ```rust use tako::extractors::header_map::HeaderMap; async fn inspect(HeaderMap(headers): HeaderMap) { if let Some(ua) = headers.get("user-agent") { println!("user-agent: {ua:?}"); } } ``` For a single strongly-typed header, the `typed-header` feature adds `TypedHeader`. Source: [examples/json-header-map/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/json-header-map/src/main.rs) ```rust use anyhow::Result; use serde::Deserialize; use serde::Serialize; use tako::Method; use tako::extractors::header_map::HeaderMap; use tako::extractors::json::Json; use tako::router::Router; use tokio::net::TcpListener; #[derive(Deserialize)] struct Input { name: String, } #[derive(Serialize)] struct Output { name: String, user_agent: Option, } /// POST /echo /// Body: {"name": "Alice"} /// /// Demonstrates using both `Json` and `HeaderMap` extractors in the handler signature. async fn echo_with_headers( HeaderMap(headers): HeaderMap, Json(payload): Json, ) -> Json { let user_agent = headers .get("user-agent") .and_then(|v| v.to_str().ok()) .map(|s| s.to_string()); Json(Output { name: payload.name, user_agent, }) } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::POST, "/echo", echo_with_headers); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` ## `Accept` Parses the `Accept` header into a preference-sorted media-type list and exposes content-negotiation helpers: `prefers(mt)`, `accepts(mt)`, `preferred()`, and `types()`. Wildcards (`*/*`, `image/*`) are matched correctly — `image/*` matches `image/png` but not `imagezzz`. ```rust use tako::extractors::accept::Accept; use tako::responder::Responder; async fn negotiate(accept: Accept) -> impl Responder { if accept.prefers("application/json") { r#"{"message":"hello"}"#.to_string() } else { "hello".to_string() } } ``` ## `AcceptLanguage` Parses `Accept-Language` into a `Vec` (each a `language` tag plus a `quality` from 0.0 to 1.0), sorted by quality per RFC 7231. ```rust use tako::extractors::acc_lang::AcceptLanguage; async fn localize(langs: AcceptLanguage) -> String { match langs.languages.first() { Some(top) => format!("preferred language: {}", top.language), None => "no preference".to_string(), } } ``` `AcceptLanguageError` distinguishes `MissingHeader`, `InvalidHeader`, and `ParseError`. ## `Range` Parses an RFC 9110 `bytes=` Range header into `Range { specs: Vec }`. Each `RangeSpec` is one of `Inclusive { start, end }`, `From { start }`, or `Suffix { length }`, and `RangeSpec::resolve(total_size)` turns a spec into a concrete inclusive `[start, end]` (or `None` when unsatisfiable). Multi-range requests populate the full `specs` list; single-range responders can take the first entry. ```rust use tako::extractors::range::Range; async fn serve_partial(range: Range) -> String { match range.specs.first().and_then(|s| s.resolve(10_000)) { Some((start, end)) => format!("serving bytes {start}-{end}"), None => "range not satisfiable".to_string(), } } ``` ## `IpAddr` Extracts the client IP. **By default it returns the transport-level peer IP** and ignores forwarded headers (`X-Forwarded-For`, `X-Real-IP`, `Forwarded`, …), because any direct client can forge them. It exposes inspection helpers like `is_private()`, `is_loopback()`, `is_ipv4()`. ```rust use tako::extractors::ipaddr::IpAddr; async fn whoami(ip: IpAddr) -> String { if ip.is_private() { format!("private client: {ip}") } else { format!("client: {ip}") } } ``` To honor forwarded headers behind a known proxy, set an `IpAddrConfig` with `trusted_proxies` listing your load-balancer fleet via `tako_rs_core::state::set_state`. Only when the direct peer is in that list are headers consulted, in priority order (`Forwarded`, `X-Forwarded-For`, `X-Real-IP`, `X-Client-IP`, `CF-Connecting-IP`, `True-Client-IP`). ## Related * Route capture syntax (`{name}`, `{*rest}`, typed `{name: T}`) — [Routing](/docs/routing). * Body inputs — [body extractors](/docs/extractors/body). * Authorization headers and tokens — [auth extractors](/docs/extractors/auth). Source: https://tako.rust-dd.com/docs/extractors/request-meta --- # Auth extractors > Pull credentials and identity from a request — Basic, Bearer, ApiKey, and unverified vs verified JWT claims. These extractors read credentials and identity out of a request. Most are header-only ([`FromRequestParts`](https://docs.rs/tako-rs/latest/tako/extractors/trait.FromRequestParts.html)), so they compose freely with the other extractors. They **extract and shape** credentials — enforcing a policy (rejecting unauthenticated requests for a whole route group, verifying a signature) is the job of the [auth middleware](/docs/middleware/auth). ## `Basic` HTTP Basic credentials (RFC 7617), referred to as `BasicAuth` in the catalog. The Base64 token is decoded and split on the first colon into `username` and `password`; the raw token is preserved. ```rust use tako::extractors::basic::Basic; async fn protected(basic: Basic) -> String { // Validate against your user store — do not hardcode in production. format!("authenticated user: {}", basic.username) } ``` `BasicAuthError` covers the missing / malformed / non-Basic / bad-Base64 / bad-UTF-8 / missing-colon cases. ## `Bearer` A Bearer token (RFC 6750), referred to as `BearerAuth` in the catalog. `token` holds the value without the `Bearer ` prefix; `with_bearer` keeps the full header value. ```rust use tako::extractors::bearer::Bearer; async fn api(bearer: Bearer) -> String { // Verify the token (JWT signature, opaque-token lookup, …) here. format!("token: {}", bearer.token) } ``` `BearerAuthError` distinguishes missing / invalid / non-Bearer header forms. ## `ApiKey` API-key authentication is provided as the `ApiKeyAuth` middleware in `tako-rs-plugins`, not as a standalone extractor. It validates a key pulled from a configurable location — `ApiKeyLocation::Header(name)`, `Query(name)`, or `HeaderOrQuery(header, query)` — against a fixed key, a set of keys, or a custom verify closure. ```rust use tako::middleware::api_key_auth::{ApiKeyAuth, ApiKeyLocation}; // Single key from a custom header. let auth = ApiKeyAuth::new("secret-api-key").header_name("X-Custom-Key"); // Multiple keys, read from a query parameter. let multi = ApiKeyAuth::from_keys(["key1", "key2"]) .location(ApiKeyLocation::Query("api_key")); // Dynamic verification. let dynamic = ApiKeyAuth::with_verify(|key| key.starts_with("live_")); ``` Attach it like any other middleware — see the [auth middleware](/docs/middleware/auth) page. ## `JwtClaimsUnverified` Decodes the claims segment of a JWT into `T` (any `serde::Deserialize`) without checking the signature, `exp`, or `nbf`. `JwtClaimsUnverified` does **NOT** verify the token signature. It only base64-decodes the claims. Treat its output as untrusted unless an upstream middleware has already verified the signature for this request. For trusted claims, use [`JwtClaimsVerified`](#jwtclaimsverifiedt). ```rust use serde::{Deserialize, Serialize}; use tako::extractors::jwt::JwtClaimsUnverified; #[derive(Debug, Deserialize, Serialize)] struct UserClaims { sub: String, exp: u64, email: String, role: String, } // Inspection of an untrusted payload only. async fn inspect(JwtClaimsUnverified(claims): JwtClaimsUnverified) -> String { format!("claimed email (unverified): {}", claims.email) } ``` The lower-level `Jwt` extractor (`token`, `header`) exposes the raw token and the decoded `header()` / `claims()` / `signature()` segments, with the same no-signature-check caveat. `JwtError` maps every failure to a `401`. ## `JwtClaimsVerified` The verifying counterpart, from `tako-rs-plugins`. It does not parse the token itself — it reads the claims that the [`JwtAuth`](/docs/middleware/auth) middleware inserted into request extensions **after** verifying the signature (against a JWKS / shared key) and applying the verifier’s time checks and configured issuer/audience constraints. `T` (here `C`) must be the verifier's `Claims` type. ```rust use serde::{Deserialize, Serialize}; use tako::extractors::jwt::JwtClaimsVerified; #[derive(Clone, Debug, Deserialize, Serialize)] struct Claims { sub: String, role: String, } async fn dashboard(JwtClaimsVerified(claims): JwtClaimsVerified) -> String { format!("trusted user {} ({})", claims.sub, claims.role) } ``` If the `JwtAuth` middleware did not run for the request — a typical wiring mistake — extraction fails with `UnverifiedClaims`, returning `401 Unauthorized` ("request was not authenticated by JwtAuth middleware"). Wire the middleware on the route or group that uses this extractor; see [auth middleware](/docs/middleware/auth). ## Choosing between them | Need | Use | | ---------------------------------- | --------------------------------------------- | | Inspect an untrusted token payload | `JwtClaimsUnverified` | | Trusted, signature-checked claims | `JwtAuth` middleware + `JwtClaimsVerified` | | Raw user/password | `Basic` | | Opaque or custom bearer token | `Bearer` | | Shared API key | `ApiKeyAuth` middleware | See also the [auth middleware](/docs/middleware/auth) for enforcing these across a route group. Source: https://tako.rust-dd.com/docs/extractors/auth --- # Cookie extractors > Read and write request cookies — plain CookieJar, HMAC-signed CookieSigned, and encrypted CookiePrivate, with key rotation. Three cookie jars cover three trust levels. All parse the incoming `Cookie:` header into a jar and emit `Set-Cookie` headers for cookies you add; they differ in whether values are plaintext, signed, or encrypted. | Extractor | Catalog name | Protection | | --------------- | ------------------ | -------------------------------------------- | | `CookieJar` | `CookieJar` | none (plaintext) | | `CookieSigned` | `SignedCookieJar` | HMAC integrity — readable but tamper-evident | | `CookiePrivate` | `PrivateCookieJar` | encryption — unreadable and tamper-proof | ## `CookieJar` Plain cookie reading and writing. Wraps the `cookie` crate's jar and parses the `Cookie:` header on extraction. Its error type is `Infallible`. ```rust use cookie::Cookie; use tako::extractors::cookie_jar::CookieJar; async fn handle(mut jar: CookieJar) { if let Some(session) = jar.get("session_id") { println!("session: {}", session.value()); } jar.add(Cookie::new("visited", "true")); } ``` `remove(name)` copies the existing cookie's `Path` / `Domain` into the removal marker so the browser actually invalidates it (rather than leaving an identically-named cookie on another path or the apex domain in place). ## `CookieSigned` HMAC-signed cookies (the catalog's `SignedCookieJar`). Values stay **readable** by the client but any tampering is detected on read. Construct with a `cookie::Key`; `add` signs and `get` verifies. ```rust use cookie::{Cookie, Key}; use tako::extractors::cookie_signed::CookieSigned; let key = Key::generate(); let mut signed = CookieSigned::new(key); signed.add(Cookie::new("username", "alice")); if let Some(c) = signed.get("username") { assert_eq!(c.value(), "alice"); // verified } ``` `CookieSignedError` covers `MissingKey` / `InvalidKey` (→ `500`) and `VerificationFailed` / `InvalidCookieFormat` / `InvalidSignature` (→ `400`). ## `CookiePrivate` Encrypted cookies (the catalog's `PrivateCookieJar`). Values are **unreadable** by the client and tamper-proof. Same shape as `CookieSigned` — `add` encrypts, `get` decrypts. ```rust use cookie::{Cookie, Key}; use tako::extractors::cookie_private::CookiePrivate; let key = Key::generate(); let mut private = CookiePrivate::new(key); private.add(Cookie::new("secret", "sensitive_data")); if let Some(c) = private.get("secret") { assert_eq!(c.value(), "sensitive_data"); // decrypted } ``` `CookiePrivateError` covers `MissingKey` / `InvalidKey` (→ `500`) and `DecryptionFailed` / `InvalidCookieFormat` (→ `400`). Use `CookiePrivate` whenever the value is sensitive (session tokens, user identifiers); use `CookieSigned` when the value may be visible but must not be forged. ## Key rotation with `KeyRing` Both the signed and private jars support a `KeyRing` so you can rotate keys without invalidating cookies signed under an older key. The ring has one `active` key (used to sign / encrypt new cookies) plus any number of `previous` keys (tried for verification / decryption only); each key carries a string `kid`. ```rust use cookie::Key; use tako::extractors::cookie_signed::{CookieSigned, KeyRing}; let ring = KeyRing::new("v2", Key::generate()) .with_previous("v1", Key::generate()); let signed = CookieSigned::with_ring(ring); ``` `get_with_kid(name)` returns which `kid` admitted a cookie (handy for logging mid-rotation). `KeyRing::revoke(kid)` drops a previous key once it is past its retention window, and `previous_kids()` lists the currently trusted ids. The same `KeyRing` type from `cookie_signed` is reused by `CookiePrivate::with_ring`. When a jar is extracted from a request, a `KeyRing` placed in request extensions takes precedence over a bare single key — the standard way to provide the key material is to insert it via middleware on the route group. ## Related * Enforcing cookie-based identity across a route group — see the [auth middleware](/docs/middleware/auth). * Cookie sessions are also available as a bundled [middleware](/docs/middleware). Source: https://tako.rust-dd.com/docs/extractors/cookies --- # Middleware > How Tako's middleware model works — registering the chain, ordering, and the full catalog of bundled middleware. Middleware in Tako wraps a handler with cross-cutting logic — auth, logging, rate limiting, compression — without changing the handler itself. Anything implementing `IntoMiddleware` (from `tako-rs-core`) can be attached either to a single route or to the whole router. ## The model A middleware is a function `Fn(Request, Next) -> Future`. The `Next` value is the rest of the chain: call `next.run(req).await` to continue, or return a `Response` early to short-circuit (the way every auth middleware returns `401` without ever invoking the handler). The bundled middleware types are *builders*. You configure one, then call `.into_middleware()` to produce the function Tako registers: ```rust use tako::middleware::IntoMiddleware; use tako::middleware::basic_auth::BasicAuth; let mw = BasicAuth::single("admin", "pw") .realm("Admin") .into_middleware(); ``` ## Registering the chain Register globally with `Router::middleware` (runs for every request), or per-route with `Route::middleware` (the `route(...)` / `route_with_tsr(...)` methods return an `Arc` you can chain onto): ```rust use tako::Method; use tako::middleware::IntoMiddleware; use tako::middleware::basic_auth::BasicAuth; use tako::middleware::bearer_auth::BearerAuth; use tako::middleware::request_id::RequestId; use tako::responder::Responder; use tako::router::Router; use tako::types::Request; async fn admin_only(_: Request) -> impl Responder { "secret" } async fn webhook(_: Request) -> impl Responder { "ack" } #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let basic = BasicAuth::single("admin", "pw").realm("Admin").into_middleware(); let bearer = BearerAuth::static_token("hunter2").into_middleware(); let mut router = Router::new(); router.middleware(RequestId::new().into_middleware()); // global: every request gets X-Request-ID router .route(Method::GET, "/admin", admin_only) .middleware(basic); router .route(Method::POST, "/webhook", webhook) .middleware(bearer); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` ## Ordering Middleware wraps from the outside in: a middleware registered earlier sits *outside* one registered later, so it sees the request first and the response last. Global (router-level) middleware wraps route-level middleware. Put request-shaping and observability concerns (request ID, body limit, security headers) on the outside, and authorization closer to the handler. Some entries in the catalog are **plugins**, not raw middleware. Plugins (CORS, compression, rate limiter, idempotency, metrics) implement `TakoPlugin` and are registered with `Router::plugin` / `Route::plugin` rather than `.into_middleware()`. See [Plugins](/docs/plugins) for the distinction. ## State and stores Stateful middleware — sessions, rate limiting, idempotency, JWKS rotation, CSRF — keep their persistence behind the `tako_rs_plugins::stores` traits. The in-process default is backed by `scc::HashMap`; companion crates can implement the same traits to swap in Redis or Postgres without touching handler code. ## Catalog All of the following ship in `tako-rs-plugins` and are re-exported under `tako::middleware::*` (and `tako::plugins::*` for the plugin-style entries). ### Authentication — [details](/docs/middleware/auth) | Middleware | Type | Description | | ------------ | -------------------------- | --------------------------------------------------------------------- | | JWT Auth | `jwt_auth::JwtAuth` | Verify JWT tokens; pluggable `JwtVerifier`, JWKS rotation, revocation | | Basic Auth | `basic_auth::BasicAuth` | RFC 7617 HTTP Basic, constant-time credential compare | | Bearer Auth | `bearer_auth::BearerAuth` | RFC 6750 static/dynamic bearer-token validation | | API Key Auth | `api_key_auth::ApiKeyAuth` | Header- or query-based API keys | ### Security — [details](/docs/middleware/security) | Middleware | Type | Description | | ---------------- | ----------------------------------- | -------------------------------------------------------------- | | CSRF | `csrf::Csrf` | Double-submit cookie, optional session binding | | Sessions | `session::SessionMiddleware` | Cookie sessions over an in-memory store | | Security Headers | `security_headers::SecurityHeaders` | HSTS, X-Frame-Options, CSP, COOP/COEP/CORP, Permissions-Policy | | Body Limit | `body_limit::BodyLimit` | Reject oversized request bodies | ### Traffic — [details](/docs/middleware/traffic) | Middleware | Type | Description | | ------------ | ------------------------------------------- | ---------------------------------------------------------------- | | Rate Limiter | `plugins::rate_limiter::RateLimiterBuilder` | Token-bucket or GCRA, composite keys, `RateLimit-*` headers | | CORS | `plugins::cors::CorsBuilder` | Cross-Origin Resource Sharing with preflight handling | | Compression | `plugins::compression::CompressionBuilder` | gzip / brotli / deflate / zstd, negotiated via `Accept-Encoding` | | Idempotency | `plugins::idempotency::IdempotencyBuilder` | `Idempotency-Key` de-duplication of unsafe methods | ### Metrics & observability — [details](/docs/middleware/metrics) | Middleware | Type | Description | | --------------- | ---------------------------------------------------------------- | ----------------------------------------------------- | | Metrics | `plugins::metrics::{PrometheusMetricsConfig, OtelMetricsConfig}` | Export request metrics to Prometheus or OpenTelemetry | | Request ID | `request_id::RequestId` | Generate / propagate `X-Request-ID` | | Upload Progress | `upload_progress::UploadProgress` | Track upload bytes via callback or extension | ### Also bundled Beyond the grouped catalog, `tako-rs-plugins` ships several more middleware: `access_log::AccessLog`, `traceparent::Traceparent`, `etag::Etag`, `timeout::Timeout`, `tenant::Tenant`, `circuit_breaker::CircuitBreaker`, `problem_json::ProblemJson`, `healthcheck`, plus the feature-gated `ip_filter::IpFilter` (`ip-filter`), `hmac_signature::HmacSignature` (`hmac-signature`), and `json_schema::JsonSchema` (`json-schema`). `timeout::Timeout` currently ships a Tokio-runtime path. A Compio variant is on the follow-up list — on Compio, prefer a per-route timeout instead. Source: https://tako.rust-dd.com/docs/middleware --- # Authentication > JWT, Basic, Bearer, and API-key authentication middleware for Tako routes. Tako ships four authentication middleware in `tako-rs-plugins`, re-exported under `tako::middleware::*`. Each is a builder; configure it, then call `.into_middleware()` and attach it globally or per-route. All four return `401 Unauthorized` with the appropriate `WWW-Authenticate` challenge when credentials are missing or invalid, and the static-credential paths use constant-time comparison to avoid timing leaks. This page covers the **middleware** that *guards* a route. To read already-parsed credentials inside a handler (for example `JwtClaimsVerified` produced by `JwtAuth`), see the [auth extractors](/docs/extractors/auth). ## Basic auth RFC 7617 HTTP Basic. Configure a single user, multiple users, or a custom verify closure that returns `bool`. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::basic_auth::BasicAuth; // Single static user. let single = BasicAuth::single("admin", "pw") .realm("Admin Area") .into_middleware(); // Multiple static users. let multi = BasicAuth::multiple([ ("alice", "secret1"), ("bob", "secret2"), ]).into_middleware(); // Dynamic verification — returns bool, not a user object. let dynamic = BasicAuth::with_verify(|user, pass| { user == "admin" && pass == "pw" }).into_middleware(); ``` `realm(...)` sets the `WWW-Authenticate` realm (default `"Restricted"`). `users_with_verify(...)` combines a static map with a fallback closure. ## Bearer auth RFC 6750 bearer tokens from the `Authorization` header. The scheme name is matched case-insensitively. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::bearer_auth::BearerAuth; // Single static token. let single = BearerAuth::static_token("my-secret-token").into_middleware(); // Multiple valid tokens. let multi = BearerAuth::static_tokens([ "development-key", "staging-key", ]).into_middleware(); // Dynamic verification. let dynamic = BearerAuth::with_verify(|token| { token.starts_with("user_") }).into_middleware(); ``` `static_tokens_with_verify(...)` combines static tokens with a fallback closure. For decoding *claims* out of a JWT bearer token, reach for `JwtAuth` below instead — `BearerAuth` only answers yes/no. ## API-key auth Validate an API key drawn from a header (default `X-API-Key`), a query parameter, or either. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::api_key_auth::{ApiKeyAuth, ApiKeyLocation}; // Single key from the default X-API-Key header. let basic = ApiKeyAuth::new("secret-api-key").into_middleware(); // Multiple keys from a custom header. let multi = ApiKeyAuth::from_keys(["key1", "key2"]) .header_name("X-Custom-Key") .into_middleware(); // From a query parameter. let query = ApiKeyAuth::new("secret") .location(ApiKeyLocation::Query("api_key")) .into_middleware(); // Dynamic verification. let dynamic = ApiKeyAuth::with_verify(|key| { key.len() == 32 && key.chars().all(|c| c.is_ascii_hexdigit()) }).into_middleware(); ``` `ApiKeyLocation` is `Header(name)`, `Query(name)`, or `HeaderOrQuery(header, query)` (header tried first). `header_name(...)` and `query_param(...)` are shorthands. ## JWT auth `JwtAuth` is trait-based: implement `JwtVerifier` with your preferred JWT library, or enable the `jwt-simple` cargo feature for the bundled `MultiKeyVerifier` (HMAC, RSA, RSA-PSS, ECDSA, EdDSA, BLAKE2b). On success the decoded claims are inserted into the request extensions, where a handler can read them. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::jwt_auth::{JwtAuth, VerifyConstraints}; // `verifier` implements JwtVerifier (e.g. the jwt-simple MultiKeyVerifier). let jwt = JwtAuth::new(verifier) .constraints(VerifyConstraints { issuer: Some("https://issuer.example".into()), audience: Some("my-api".into()), leeway_secs: 30, }) .into_middleware(); ``` Capabilities: * **JWKS rotation** — the `jwt-simple` `MultiKeyVerifier` selects keys by `kid` and exposes `rotate_key` / `revoke_kid` for runtime rotation without a restart. * **Shared key provider** — `.store(provider)` on `JwtAuth` calls `JwksProvider::keys_for(kid)` before verification. Providers return `VerificationKey { algorithm, bytes }`; the verifier checks its algorithm allow-list before using raw MAC or DER key bytes. Provider failures return `503`; invalid candidates return `401`. An empty key set permits the static verifier fallback. No built-in network JWKS fetcher is supplied. * **Constraints** — `VerifyConstraints` enforces `iss`, `aud`, and a clock `leeway_secs` uniformly across algorithms. The default verifier **fails closed** if constraints are set but the verifier cannot enforce them. * **Revocation** — `.revocation(list, extractor)` checks a `RevocationList` (the bundled `InMemoryRevocationList` is keyed by `jti`) after signature verification. * **Introspection** — `.introspect(closure)` runs an async callback on every request, the correct hook for opaque tokens or tenant-scoped revocation. A `VerifyConstraints` set on `JwtAuth` only takes effect if the verifier enforces it. The default `validate_constraints` rejects (returns `401`) when any non-default constraint is configured on a verifier that does not override it — this prevents silently dropping `iss` / `aud` / `leeway`. ## Where to register Attach an auth middleware per-route to guard only the protected paths, or globally if the whole router is behind auth: ```rust use tako::Method; router .route(Method::GET, "/admin", admin_only) .middleware(BasicAuth::single("admin", "pw").into_middleware()); ``` See the [middleware model](/docs/middleware) for chain ordering, and the runnable [`examples/auth`](https://github.com/rust-dd/tako/tree/main/examples/auth) for a full Basic + Bearer setup. Source: https://tako.rust-dd.com/docs/middleware/auth --- # Security > CSRF protection, security response headers, body-size limits, and cookie sessions for Tako. `tako-rs-plugins` bundles the hardening middleware most services need: CSRF protection, a curated set of security response headers, request body limits, and cookie sessions. Each is a builder; configure it and call `.into_middleware()`. ## CSRF `csrf::Csrf` defaults to the **double-submit cookie** pattern: a random token is placed in a cookie and must be echoed back in a request header (default `x-csrf-token`). The middleware verifies the two match. When a `Session` extension is present, the token is also bound to the active session id, so a token minted under a previous session (for example before privilege rotation) is rejected. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::csrf::Csrf; use tako::middleware::session::SameSite; let csrf = Csrf::new() .cookie_name("csrf_token") .header_name("x-csrf-token") .secure(true) // set the cookie Secure flag .same_site(SameSite::Strict) // default .exempt("/webhooks") // bypass CSRF for a path prefix .trust_origin("https://app.example.com") .into_middleware(); ``` Key knobs: * `secure(bool)` — toggles the cookie `Secure` flag (required when `same_site = None`). * `same_site(...)` — `Strict` (default), `Lax`, or `None`. * `exempt(prefix)` — skip CSRF entirely for a path prefix (e.g. webhook endpoints that authenticate by signature instead). * `trust_origin(origin)` — strict `Origin` / `Referer` allow-list used as a fallback for legacy clients that send neither cookie nor header. * `bind_to_session(bool)` — defaults to `true`; bind the token to the session. Both CSRF and session cookies default to `Secure`. For local plain HTTP, explicitly call `.secure(false)`. For shared CSRF state use `.store(backend)` with `CsrfTokenStore`, register session middleware first, and configure `.single_use(true)` for atomic token consumption. `.token_ttl(duration)` controls expiry. Backend failures return 503. ## Sessions `session::SessionMiddleware` provides cookie-based sessions over an in-memory `MemorySessionStore` backend by default. Sessions are keyed by a random cookie value and carry arbitrary `serde`-compatible data. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::session::{SameSite, SessionMiddleware, SessionTtl}; let sessions = SessionMiddleware::new() .cookie_name("tako_session") // default .ttl(SessionTtl { idle_secs: 3_600, absolute_secs: Some(86_400) }) .secure(true) .http_only(true) // default .same_site(SameSite::Lax) // default .into_middleware(); ``` `SessionTtl` separates an **idle** timeout (default 1 h) from an **absolute** lifetime cap (default 24 h), so a stolen session id cannot be refreshed forever. A live session re-emits `Set-Cookie` with a refreshed `Max-Age` on every request (rolling refresh). `Session::rotate` swaps the session id while keeping the data — call it after login to defend against fixation. `SessionMiddleware::handle()` returns a handle with a `revoke_all` API for emergency purges. The default store is in-process and not shared across instances. For a clustered deployment, back sessions with a shared store via the `SessionMiddleware::new().store(backend)` builder and `tako::stores::SessionStore`. ## Security headers `security_headers::SecurityHeaders` emits a curated set of response headers following OWASP / MDN guidance. Three are always on; the rest are opt-in. Always emitted: * `X-Content-Type-Options: nosniff` * `X-Frame-Options: DENY` * `Referrer-Policy: strict-origin-when-cross-origin` ```rust use tako::middleware::IntoMiddleware; use tako::middleware::security_headers::SecurityHeaders; let headers = SecurityHeaders::new() .hsts(true) .hsts_max_age(31_536_000) .hsts_include_subdomains(true) .hsts_preload(true) .csp("default-src 'self'") .frame_options("SAMEORIGIN") .coop("same-origin") .coep("require-corp") .corp("same-origin") .permissions_policy("geolocation=(), camera=()") .into_middleware(); ``` For inline scripts, `csp_with_nonce(template)` generates a per-request nonce, exposes it as a `CspNonce` extension for handlers to interpolate, and substitutes it into the emitted header. `csp_report_only(template)` emits a report-only policy instead. `X-XSS-Protection` is intentionally **not** emitted — modern browsers ignore it and OWASP recommends removing it. CSP is the authoritative replacement. ## Body limit Buffered extractors already use a 2 MiB default router limit. Configure `router.body_limit(bytes)` or explicitly call `router.disable_body_limit()`. The middleware below can add route-dependent limits and early rejection. `body_limit::BodyLimit` rejects oversized request bodies before they are read, preventing resource-exhaustion attacks. It fast-rejects using `Content-Length` when present. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::body_limit::BodyLimit; // Static 1 MiB limit. let limit = BodyLimit::new(1024 * 1024).into_middleware(); // Dynamic limit based on the request. let dynamic = BodyLimit::with_dynamic_limit(|req| { if req.uri().path().starts_with("/upload") { 50 * 1024 * 1024 } else { 1024 * 1024 } }).into_middleware(); ``` `new_with_dynamic(static_limit, closure)` combines a static cap with a dynamic override. See the [middleware model](/docs/middleware) for chain ordering — register body limit and security headers on the outside of the chain. Source: https://tako.rust-dd.com/docs/middleware/security --- # Traffic > Rate limiting, CORS, response compression, and idempotency-key de-duplication for Tako. The traffic-shaping entries — rate limiter, CORS, compression, and idempotency — are **plugins**, not raw middleware. They implement `TakoPlugin` and are registered with `Router::plugin` (global) or `Route::plugin` (per-route) instead of `.into_middleware()`. All four require the `plugins` feature; compression's zstd path additionally needs `zstd`. ```toml [dependencies] tako-rs = { version = "2", features = ["plugins"] } ``` ## Rate limiter `plugins::rate_limiter::RateLimiterBuilder` builds a token-bucket (default) or GCRA limiter. The default key is the peer IP; `key_fn` composes per-route, per-tenant, or per-user buckets. It emits the IETF `RateLimit-Limit`, `RateLimit-Remaining`, `RateLimit-Reset`, and `Retry-After` headers. ```rust use tako::plugins::rate_limiter::{Algorithm, RateLimiterBuilder, UnkeyedBehavior}; let limiter = RateLimiterBuilder::new() .requests_per_second(100) .algorithm(Algorithm::TokenBucket) // or Algorithm::Gcra .on_unkeyed(UnkeyedBehavior::Reject) // requests with no discoverable IP (or Allow) .build(); router.plugin(limiter); ``` `requests_per_second(n)` and `requests_per_minute(n)` are shorthands over the lower-level `max_requests` / `refill_rate` / `refill_interval_ms` knobs. `build()` panics on a zero `max_requests`, `refill_rate`, or `refill_interval_ms` — a zero rate or cap would silently deny every request or poison the GCRA arithmetic, so it is rejected at startup. For a hard throttle use a deliberately tiny rate with a long interval. ## CORS `plugins::cors::CorsBuilder` handles preflight `OPTIONS` requests, validates origins, and adds the CORS response headers. Apply it router-wide or to a specific route for a tighter policy. ```rust use http::Method; use tako::plugins::cors::CorsBuilder; let cors = CorsBuilder::new() .allow_origin("https://app.example.com") .allow_methods(&[Method::GET, Method::POST, Method::PUT]) .allow_headers(&[http::header::CONTENT_TYPE]) .allow_credentials(true) .max_age_secs(86_400) .build(); router.plugin(cors); ``` Origin matching can also use `allow_origin_suffix(suffix)` or `allow_origin_predicate(closure)` for dynamic decisions, and `allow_private_network(true)` opts into the Private Network Access preflight. `try_build()` returns a `Result` instead of panicking on an invalid config (for example credentialed wildcards). ## Compression `plugins::compression::CompressionBuilder` negotiates response compression from the client `Accept-Encoding` header. gzip, brotli, and deflate are available by default; zstd requires the `zstd` feature. Compression is applied selectively by content type, response size, and status. ```rust use tako::plugins::compression::CompressionBuilder; let compression = CompressionBuilder::new() .enable_gzip(true) .enable_brotli(true) .enable_deflate(true) .enable_stream(true) // compress streaming bodies .min_size(1024) // skip bodies smaller than this .brotli_level(9) .build(); router.plugin(compression); ``` Per-algorithm levels are set with `gzip_level`, `brotli_level`, `deflate_level`, and `zstd_level`. `content_types(policy)` controls which MIME types are compressed, and `protect_sensitive(true)` skips compression on responses that could be vulnerable to compression side-channels. `enable_zstd(true)` only has an effect when the crate is built with the `zstd` feature; without it the zstd path is compiled out. See the [feature reference](/docs/reference/features). ## Idempotency `plugins::idempotency::IdempotencyBuilder` implements server-side idempotency for unsafe methods, keyed by a caller-supplied header (default `Idempotency-Key`). For a given key and scope it guarantees the same response within a TTL. ```rust use tako::plugins::idempotency::{IdempotencyBuilder, Scope}; let idempotency = IdempotencyBuilder::new() .ttl_secs(3_600) .scope(Scope::MethodAndPath) // or Scope::KeyOnly .coalesce_inflight(true) // concurrent dupes wait for the first to finish .verify_payload(true) // same key + different body => 409 Conflict .build(); router.plugin(idempotency); ``` Behavior: * The first request with a new key is processed and its response cached. * Concurrent requests with the same key wait for completion and receive the cached result. * Replays within the TTL return the cached result immediately. * Reusing a key with a *different* payload returns `409 Conflict` (when `verify_payload` is on). Bodies are buffered to compute a stable signature and to cache the response; `max_request_body_bytes` and `max_cached_body_bytes` bound that buffering. Storage is in-memory with periodic TTL cleanup. See the [middleware model](/docs/middleware) for how plugins fit into the chain, and [Plugins](/docs/plugins) for the plugin-vs-middleware distinction. Source: https://tako.rust-dd.com/docs/middleware/traffic --- # Metrics & Observability > Prometheus / OpenTelemetry metrics export, request-ID propagation, and upload-progress callbacks. `tako-rs-plugins` provides the observability hooks: a metrics plugin that exports to Prometheus or OpenTelemetry, request-ID propagation, and upload-progress tracking. The metrics plugin is wired through Tako's signal system; see [Observability](/docs/observability) for the broader story. ## Metrics export The metrics plugin listens to Tako's application- and route-level signals and forwards them to a backend. Two backends ship behind feature flags, each with a config type that installs the plugin and wires up the export path in one call. ### Prometheus Requires the `metrics-prometheus` feature. `PrometheusMetricsConfig::install` registers the plugin and mounts a scrape endpoint (default `/metrics`), returning the `Arc`. ```rust use tako::plugins::metrics::PrometheusMetricsConfig; use tako::router::Router; let mut router = Router::new(); let registry = PrometheusMetricsConfig::default() // endpoint_path = "/metrics" .install(&mut router); ``` `with_buckets(vec)` overrides the latency histogram bucket schedule. The scrape endpoint encodes the registry in the Prometheus text format and returns `500` (not a `200` with an error body) if encoding fails, so scraper alerting fires correctly. ### OpenTelemetry Requires the `metrics-opentelemetry` feature. `OtelMetricsConfig::install` registers the plugin with an OTLP exporter and returns the `SdkMeterProvider`, which you keep alive for the process lifetime and `shutdown()` during graceful shutdown. ```rust use tako::plugins::metrics::OtelMetricsConfig; use tako::router::Router; let mut router = Router::new(); let meter_provider = OtelMetricsConfig::default() .with_endpoint("http://localhost:4318/v1/metrics") .install(&mut router)?; // ... serve ... meter_provider.shutdown()?; ``` `with_meter_name(name)` sets the meter name (default `"tako"`). Source: [examples/metrics-opentelemetry/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/metrics-opentelemetry/src/main.rs) ```rust //! Metrics with OpenTelemetry OTLP exporter example. //! //! This example demonstrates how to use Tako's metrics plugin with //! OpenTelemetry OTLP exporter to send metrics to collectors like //! Prometheus, Jaeger, or the OpenTelemetry Collector. //! //! Run this example with: //! ```sh //! cargo run --example metrics-opentelemetry --features metrics-opentelemetry //! ``` //! //! To test with Prometheus: //! ```sh //! docker run -p 9090:9090 prom/prometheus --web.enable-otlp-receiver //! ``` //! Then update the endpoint to "http://localhost:9090/api/v1/otlp/v1/metrics" use anyhow::Result; use tako::Method; use tako::plugins::metrics::OtelMetricsConfig; use tako::responder::Responder; use tako::router::Router; use tokio::net::TcpListener; async fn hello() -> impl Responder { "Hello from metrics example".into_response() } async fn health() -> impl Responder { "OK".into_response() } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello); router.route(Method::GET, "/health", health); // Install the OpenTelemetry metrics plugin with OTLP exporter. // By default, metrics are exported to http://localhost:4318/v1/metrics let meter_provider = OtelMetricsConfig::default() .with_endpoint("http://localhost:4318/v1/metrics") .install(&mut router)?; println!("Server running on http://127.0.0.1:8080"); println!("Metrics being exported via OTLP to http://localhost:4318/v1/metrics"); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); // Shutdown the meter provider to flush remaining metrics meter_provider.shutdown()?; Ok(()) } ``` Both backends depend on the signal system, so the metrics features enable `signals` transitively. See the [feature reference](/docs/reference/features) for the full flag graph. ## Request ID `request_id::RequestId` generates or propagates a unique request identifier via the `X-Request-ID` header. If the incoming request already carries the header it is preserved (capped at 256 bytes); otherwise a UUID v4 is generated. The ID is inserted into the request extensions as `RequestIdValue` and echoed in the response header. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::request_id::RequestId; // Default: X-Request-ID with a UUID v4. let request_id = RequestId::new().into_middleware(); // Custom header name. let correlation = RequestId::new() .header_name("X-Correlation-ID") .into_middleware(); // Custom generator. let custom = RequestId::new() .generator(|| ulid_like_id()) .into_middleware(); ``` Register it globally so every request — and every log line and tracing span that reads `RequestIdValue` — is correlated: ```rust router.middleware(RequestId::new().into_middleware()); ``` ## Upload progress `upload_progress::UploadProgress` wraps the request body to track bytes received, reporting through a callback and through request extensions. The `ProgressState` passed to the callback carries `bytes_read`, `total_bytes` (from `Content-Length`, when known), and a `percent()` helper. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::upload_progress::UploadProgress; let progress = UploadProgress::new() .on_progress(|state| { let pct = state.percent().map(|p| format!("{p}%")).unwrap_or_else(|| "?%".into()); println!("{pct}: {} / {:?} bytes", state.bytes_read, state.total_bytes); }) .min_notify_interval_bytes(1024) // throttle callbacks to ~once per KiB .into_middleware(); ``` Inside a handler, the live tracker is available as a `ProgressTracker` extension, exposing `bytes_read()`, `total_bytes()`, and `percent()`: ```rust use tako::middleware::upload_progress::ProgressTracker; let pct = req .extensions() .get::() .and_then(|t| t.percent()) .unwrap_or(0); ``` Source: [examples/upload-progress/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/upload-progress/src/main.rs) ```rust use anyhow::Result; use tako::Method; use tako::middleware::IntoMiddleware; use tako::middleware::upload_progress::ProgressTracker; use tako::middleware::upload_progress::UploadProgress; use tako::responder::Responder; use tako::router::Router; use tako::types::Request; async fn upload_handler(req: Request) -> impl Responder { // Access the progress tracker from request extensions let info = req .extensions() .get::() .map(|tracker| { let bytes = tracker.bytes_read(); let total = tracker.total_bytes(); let pct = tracker.percent().unwrap_or(0); format!("Upload complete: {bytes} bytes received (total: {total:?}, {pct}%)") }) .unwrap_or_else(|| "No progress tracker found".to_string()); info } #[tokio::main] async fn main() -> Result<()> { tracing_subscriber::fmt::init(); let progress = UploadProgress::new() .on_progress(|state| { let pct = state .percent() .map(|p| format!("{p}%")) .unwrap_or_else(|| "?%".into()); println!( "Upload progress: {} / {} bytes ({pct})", state.bytes_read, state .total_bytes .map(|t| t.to_string()) .unwrap_or_else(|| "unknown".into()), ); }) .min_notify_interval_bytes(1024); // Notify at most every 1KB let mw = progress.into_middleware(); let mut router = Router::new(); router .route(Method::POST, "/upload", upload_handler) .middleware(mw); let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; println!("Upload progress server on http://127.0.0.1:8080"); println!("Test with: curl -X POST -d @somefile http://127.0.0.1:8080/upload"); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` See the [middleware model](/docs/middleware) for chain ordering and [Observability](/docs/observability) for signals, tracing, and the broader metrics pipeline. Source: https://tako.rust-dd.com/docs/middleware/metrics --- # Plugins > The tako-rs-plugins crate — bundled middleware and plugins, the plugins feature flag, and how plugins differ from middleware. `tako-rs-plugins` is the crate that ships Tako's concrete middleware and plugins. The *traits* (`IntoMiddleware`, `Next`, `TakoPlugin`) live in `tako-rs-core`; this crate hosts the ready-to-use implementations, re-exported through the umbrella crate under `tako::middleware::*` and `tako::plugins::*`. ## Middleware vs. plugin Tako has two extension mechanisms, and the catalog uses both. Knowing which is which tells you how to register it. | | Middleware | Plugin | | ----------- | ------------------------------------------------------------------- | ---------------------------------------------------- | | Trait | `IntoMiddleware` | `TakoPlugin` | | Shape | a `Fn(Request, Next) -> Future` wrapping one handler | `setup(&self, router)` that installs onto the router | | Register | `router.middleware(x.into_middleware())` or `route.middleware(...)` | `router.plugin(x)` or `route.plugin(x)` | | Typical job | inspect / transform a single request-response | install middleware, add routes, register state | A **middleware** sits directly in the request path: it gets a `Request` and a `Next`, and either calls `next.run(req).await` or returns early. A **plugin** is a setup step — its `setup` hook can register middleware *and* mount extra routes or state. That is why the metrics plugin can add a `/metrics` scrape endpoint, and why CORS/compression/rate-limiting are plugins: they wire several pieces onto the router at once. ```rust use tako::middleware::IntoMiddleware; use tako::middleware::request_id::RequestId; use tako::plugins::cors::CorsBuilder; // Middleware: build it, convert it, register it. router.middleware(RequestId::new().into_middleware()); // Plugin: build it, register it directly. router.plugin(CorsBuilder::new().build()); ``` Both can be scoped: `Router::middleware` / `Router::plugin` apply globally, while `Route::middleware` / `Route::plugin` apply to a single route. ## What's bundled **Middleware** (`tako::middleware::*`): * **Auth** — `basic_auth::BasicAuth`, `bearer_auth::BearerAuth`, `api_key_auth::ApiKeyAuth`, `jwt_auth::JwtAuth`. See [Authentication](/docs/middleware/auth). * **Security** — `csrf::Csrf`, `session::SessionMiddleware`, `security_headers::SecurityHeaders`, `body_limit::BodyLimit`. See [Security](/docs/middleware/security). * **Observability** — `request_id::RequestId`, `upload_progress::UploadProgress`, `access_log::AccessLog`, `traceparent::Traceparent`. See [Metrics & Observability](/docs/middleware/metrics). * **Cross-cutting** — `etag::Etag`, `timeout::Timeout`, `tenant::Tenant`, `circuit_breaker::CircuitBreaker`, `problem_json::ProblemJson`, `healthcheck`, plus feature-gated `ip_filter::IpFilter`, `hmac_signature::HmacSignature`, and `json_schema::JsonSchema`. **Plugins** (`tako::plugins::*`): * `cors::CorsBuilder`, `compression::CompressionBuilder`, `rate_limiter::RateLimiterBuilder`, `idempotency::IdempotencyBuilder` — see [Traffic](/docs/middleware/traffic). * `metrics::{PrometheusMetricsConfig, OtelMetricsConfig}` — see [Metrics & Observability](/docs/middleware/metrics). ## Enabling the set Most of the middleware (auth, CSRF, sessions, security headers, request ID, body limit, upload progress) compile without any extra flag. The **plugins** — CORS, compression, rate limiter, idempotency — require the `plugins` feature: ```toml [dependencies] tako-rs = { version = "2", features = ["plugins"] } ``` A few capabilities layer on top of `plugins`: | Feature | Adds | | ----------------------- | ------------------------------------------------------------------ | | `plugins` | CORS, compression (gzip/brotli/deflate), rate limiter, idempotency | | `zstd` | zstd compression in the compression plugin | | `metrics-prometheus` | Prometheus metrics export (pulls in `signals`) | | `metrics-opentelemetry` | OpenTelemetry metrics export (pulls in `signals`) | `metrics-prometheus` and `metrics-opentelemetry` both imply `plugins` and `signals`. See the [feature reference](/docs/reference/features) for the complete graph. ## Stores Stateful middleware — sessions, rate limiting, idempotency, JWKS rotation, CSRF — keep their persistence behind the `tako_rs_plugins::stores` traits. The default backend is an in-process `scc::HashMap`; companion crates can implement the same traits to move state into Redis or Postgres without changing any handler or middleware code. For how the chain is assembled and ordered, start with the [middleware model](/docs/middleware). Source: https://tako.rust-dd.com/docs/plugins --- # Streams > Server-Sent Events, WebSocket, ranged file streaming, and WebTransport over HTTP/3 in tako. Tako's streaming surface covers Server-Sent Events, WebSocket, file streaming with range / conditional GET, and WebTransport over HTTP/3. All four live in the `tako-rs-streams` crate and are re-exported under `tako::*`. ## Server-Sent Events `Sse::new(stream)` accepts any `Stream>` and frames it as `text/event-stream`. Default headers include `Cache-Control: no-cache` plus `X-Accel-Buffering: no` so reverse proxies do not buffer the response. ```rust use bytes::Bytes; use futures_util::{StreamExt, stream}; use tako::Method; use tako::responder::Responder; use tako::router::Router; use tako::sse::Sse; async fn ticker() -> impl Responder { let s = stream::unfold(0u64, |i| async move { tokio::time::sleep(std::time::Duration::from_secs(1)).await; Some((Bytes::from(format!("tick: {i}")), i + 1)) }); Sse::new(s) } #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/events", ticker); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` For structured SSE with `event:`, `id:`, and `retry:` fields, use `Sse::events(stream)` plus the `SseEvent` builder, and call `.keep_alive(Duration::from_secs(15))` to emit `:keep-alive` comments during idle periods. `tako::sse::last_event_id(headers)` reads the `Last-Event-ID` header sent by EventSource clients on reconnect. See [SSE transport](/docs/transports/sse) for the full transport reference. SSE requires the `sse` feature. ## WebSocket `TakoWs::new(req, fut)` performs the upgrade from inside an HTTP handler. The closure receives a `WebSocket` half-duplex pair (`tokio_tungstenite` on tokio, `compio_ws` on compio) that you drive to completion: ```rust use futures_util::{SinkExt, StreamExt}; use tako::Method; use tako::responder::Responder; use tako::types::Request; use tako::ws::TakoWs; use tokio_tungstenite::tungstenite::Message; async fn ws_echo(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { while let Some(Ok(msg)) = ws.next().await { if let Message::Text(t) = msg { let _ = ws.send(Message::Text(t)).await; } } }) } ``` Enable `ws` for Tokio or `compio-ws` for Compio. Configure subprotocols, frame/message limits, allowed origins, and upgrade timeout on the responder. The handler implements its own ping/pong policy; permessage-deflate is not provided. See [WebSocket transport](/docs/transports/websocket) for the full transport reference. ## File serving `tako::file_stream::FileStream` (feature `file-stream`) streams an open file and provides conditional and single-range response helpers. Metadata validators are weak ETags. `tako::r#static::ServeDirBuilder` and `ServeFile` stream assets, handle HEAD and conditional/range requests, and optionally select `.br` / `.gz` sidecars. Dotfiles are denied by default. Use an explicit fallback file path and keep served roots read-only to untrusted processes. Source: [examples/file-stream/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/file-stream/src/main.rs) ```rust use anyhow::Result; use tako::Method; use tako::file_stream::FileStream; use tako::responder::Responder; use tako::router::Router; use tako::types::Request; use tokio::fs::File; use tokio::net::TcpListener; use tokio_util::io::ReaderStream; async fn serve_file(_: Request) -> impl Responder { let file = File::open("test.txt").await.unwrap(); let stream = ReaderStream::new(file); let file_stream = FileStream::new(stream, Some("test.txt".to_string()), None); file_stream.into_response() } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/file", serve_file); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` ## WebTransport W3C WebTransport sessions run on the HTTP/3 server behind the `webtransport` feature, on either runtime: a `CONNECT` route takes the `WebTransport` extractor and gets streams and datagrams that browsers can reach. See [WebTransport](/docs/transports/webtransport). See also: * [`examples/streams`](https://github.com/rust-dd/tako/tree/main/examples/streams) and [`examples/http3-sse`](https://github.com/rust-dd/tako/tree/main/examples/http3-sse) for SSE, * [`examples/websocket`](https://github.com/rust-dd/tako/tree/main/examples/websocket), [`examples/websocket-http2`](https://github.com/rust-dd/tako/tree/main/examples/websocket-http2), and [`examples/websocket-compio`](https://github.com/rust-dd/tako/tree/main/examples/websocket-compio) for WebSocket, * [`examples/file-stream`](https://github.com/rust-dd/tako/tree/main/examples/file-stream) for ranged file serving, * [`examples/webtransport`](https://github.com/rust-dd/tako/tree/main/examples/webtransport) for browser WebTransport, and [`examples/raw-quic`](https://github.com/rust-dd/tako/tree/main/examples/raw-quic) for raw QUIC sessions. Multipart byteranges and Linux `sendfile(2)` are deferred follow-up items. Source: https://tako.rust-dd.com/docs/streams --- # Queue > In-process background job queue with retry, dead-lettering, deduplication, and delayed execution. Tako ships an in-process background job queue at `tako::queue`. It is designed for the "fire-and-forget" workloads that usually end up hand-rolled on top of `tokio::spawn` — send-email, dispatch-webhook, reindex, etc. — with retry, dead-lettering, deduplication, and delayed execution built in. ```rust use std::time::Duration; use serde::{Deserialize, Serialize}; use tako::Method; use tako::extractors::json::Json; use tako::extractors::state::State; use tako::queue::{Job, Queue, RetryPolicy}; use tako::responder::Responder; use tako::router::Router; #[derive(Serialize, Deserialize)] struct Email { to: String, subject: String } async fn enqueue( State(q): State, Json(req): Json, ) -> impl Responder { match q.0.push("send_email", &req).await { Ok(id) => (tako::StatusCode::ACCEPTED, format!("queued id={id}\n")), Err(e) => (tako::StatusCode::INTERNAL_SERVER_ERROR, e.to_string()), } } #[tokio::main] async fn main() -> anyhow::Result<()> { let queue = Queue::builder() .workers(4) .retry(RetryPolicy::exponential(3, Duration::from_millis(500))) .build(); queue.register("send_email", |job: Job| async move { let payload: Email = job.deserialize()?; println!("send_email -> {} / {}", payload.to, payload.subject); Ok(()) }); queue.start(); let mut router = Router::new(); router.with_state(queue); router.route(Method::POST, "/emails", enqueue); let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` The builder gives you a fluent way to set worker count and retry policy. Job handlers are registered by name with `queue.register(name, handler)`; each handler returns `Result<(), QueueError>`. Push jobs with `queue.push(name, payload)`, `queue.push_delayed(name, payload, duration)`, or `queue.push_dedup(name, payload, key)` to collapse duplicate pending jobs by an idempotency key. Operational helpers: * `queue.pending_count()`, `queue.inflight_count()` — gauge-style metrics for dashboards. * `queue.dead_letters()` — read the dead-letter queue. Jobs that exhaust their retry budget land here for manual inspection. * `queue.shutdown(timeout)` — drain workers gracefully on SIGTERM. * `tako::queue::cron::CronScheduler` (feature `queue-cron`) — wire a crontab spec into `queue.push`. Signals are emitted for every job lifecycle event: * `queue.job.queued`, `.started`, `.completed`, `.failed`, `.retrying`, `.dead_letter` The canonical strings live under `tako::queue::signal_ids`. The queue integrates with the in-process [signals](/docs/signals) bus, so any listener can react to job lifecycle events. The default `MemoryBackend` keeps everything in process. Companion crates `tako-stores-redis` and `tako-stores-postgres` are on the follow-up list and will implement `QueueBackend` plus the session / rate-limit / idempotency / JWKS / CSRF stores. See [`examples/job-queue`](https://github.com/rust-dd/tako/tree/main/examples/job-queue) for a full end-to-end demo including retries, delayed jobs, and the dead-letter queue. Source: https://tako.rust-dd.com/docs/queue --- # Signals > Subscribe to application, router, and route events, observe listener shutdown, and use typed RPC with the in-process signal arbiter. Signals are Tako's in-process pub/sub bus for framework-internal events and application-level RPC. The framework emits well-known signals on every request, connection, and queue job; you subscribe to the IDs that matter. ```rust use tako::Method; use tako::responder::Responder; use tako::router::Router; use tako::signals::{Signal, app_events, ids}; use tako::types::Request; async fn hello(_: Request) -> impl Responder { "hi" } fn install_listeners() { let arbiter = app_events(); // One-shot callback when the server starts up. arbiter.on(ids::SERVER_STARTED, |sig: Signal| async move { println!("server.started: {:?}", sig.metadata); }); // Long-running listener that logs every completed request. let mut rx = arbiter.subscribe(ids::REQUEST_COMPLETED); tokio::spawn(async move { while let Ok(sig) = rx.recv().await { let method = sig.metadata.get("method").cloned().unwrap_or_default(); let path = sig.metadata.get("path").cloned().unwrap_or_default(); let status = sig.metadata.get("status").cloned().unwrap_or_default(); println!("request.completed: {method} {path} -> {status}"); } }); } #[tokio::main] async fn main() -> anyhow::Result<()> { install_listeners(); let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/", hello); tako::Server::builder().build().spawn_http(listener, router).result().await?; Ok(()) } ``` `app_events()` and `app_signals()` return the same process-wide arbiter. Delivery follows the owning scope: | Event | Recipients | | -------------------------------------------------- | ------------------------------------------------------ | | `request.started`, `request.completed` | Owning router and application | | `route.request.started`, `route.request.completed` | Matching route, owning router, and application | | `server.started`, `server.stopped` | Application; per-thread workers emit once per listener | | Transport connection lifecycle events | Application | | Custom events | The arbiter on which you explicitly emit them | `request.completed` carries method, path, status, matched route, and `duration_us`. Listen with `router.on_signal(...)` for one service or `app_signals().on(...)` for all routers. `server.stopped` follows normal drain completion; cancelling the entire server future skips its final event. `ROUTER_HOT_RELOAD` is an application-defined convention, not an automatic event. ## RPC over signals The same arbiter doubles as a typed RPC bus: ```rust use std::sync::Arc; use tako::signals::SignalArbiter; #[derive(Debug)] struct AddRequest { a: i32, b: i32 } #[derive(Debug, Clone)] struct AddResponse { sum: i32 } #[tokio::main] async fn main() { let arbiter = SignalArbiter::new(); arbiter.register_rpc::( "rpc.add", |req: Arc| async move { AddResponse { sum: req.a + req.b } }, ); let res = arbiter.call_rpc::( "rpc.add", AddRequest { a: 2, b: 40 }, ).await.unwrap(); assert_eq!(res.sum, 42); } ``` `register_rpc` installs a typed handler; `call_rpc` / `call_rpc_timeout` / `call_rpc_result` invoke it. Errors are reported via the `RpcError` enum. ## Forwarding events Register a callback or exporter on an arbiter to forward selected events. The hidden experimental `signals::bus` types are not wired into routing and do not provide automatic cluster fan-out. See: * [`examples/signals-basic`](https://github.com/rust-dd/tako/tree/main/examples/signals-basic) — subscribe to `request.completed`, * [`examples/signals-route`](https://github.com/rust-dd/tako/tree/main/examples/signals-route) — per-router arbiter and custom IDs, * [`examples/signals-rpc`](https://github.com/rust-dd/tako/tree/main/examples/signals-rpc) — typed RPC handlers, * [`examples/signals-advanced`](https://github.com/rust-dd/tako/tree/main/examples/signals-advanced) and [`examples/signals-complex`](https://github.com/rust-dd/tako/tree/main/examples/signals-complex) — wildcard subscriptions and exporters. Source: https://tako.rust-dd.com/docs/signals --- # Observability > Logs, metrics, traces, and health checks with swappable defaults for production deployments. Tako covers the three usual pillars — logs, metrics, traces — plus health checks, with sensible defaults that you can swap when the deployment demands something specific. ```rust use tako::Method; use tako::middleware::IntoMiddleware; use tako::middleware::access_log::AccessLog; use tako::middleware::healthcheck::Healthcheck; use tako::middleware::request_id::RequestId; use tako::middleware::traceparent::Traceparent; use tako::plugins::metrics::PrometheusMetricsConfig; use tako::responder::Responder; use tako::router::Router; async fn hello() -> impl Responder { "ok" } #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); // Request IDs: generate X-Request-ID for inbound, propagate for outbound. let request_id = RequestId::new().into_middleware(); // Access log: one structured line per request via the `tracing` macros. let access_log = AccessLog::new().into_middleware(); // W3C Trace Context: parse `traceparent` / `tracestate`, expose them as // a TraceContext extension so handlers can propagate the span. let traceparent = Traceparent::new().into_middleware(); // Health endpoints at /live, /ready, /__drain (and a SIGTERM-friendly handle). let healthcheck = Healthcheck::new(); let _drain = healthcheck.handle(); let healthcheck_mw = healthcheck.into_middleware(); router.middleware(request_id); router.middleware(access_log); router.middleware(traceparent); router.middleware(healthcheck_mw); // Prometheus scrape endpoint at /metrics (set via PrometheusMetricsConfig). let _registry = PrometheusMetricsConfig::default().install(&mut router); router.route(Method::GET, "/", hello); tako::Server::builder().build().spawn_http(listener, router) .result().await.expect("HTTP server failed"); Ok(()) } ``` The pieces: * **Logs** — `AccessLog` middleware emits one structured line per request. The default sink is `tracing::info!`; supply a custom sink via `.sink(|record| { ... })` for JSON / OTLP / file rotation. * **Metrics** — `PrometheusMetricsConfig::install(&mut router)` (feature `metrics-prometheus`) wires a Prometheus scrape endpoint and a request-latency histogram with `with_buckets(..)` overrides. The matched-route label is bounded so cardinality stays finite. For OpenTelemetry OTLP export, use `OtelMetricsConfig::default().with_endpoint(..).install(&mut router)?` from the `metrics-opentelemetry` feature. See the [metrics middleware](/docs/middleware/metrics) reference for details. * **Tracing** — `Traceparent` parses W3C Trace Context and stores a `TraceContext` extension; outbound spans propagate automatically through the v2 client. * **Health** — `Healthcheck::new()` exposes `/live`, `/ready`, `/__drain` with configurable paths via `.live_path(..)`, `.ready_path(..)`, `.drain_path(..)`. `HealthcheckHandle::drain()` flips the gate so `/ready` starts returning `503` — a SIGTERM handler can call it before draining connections. See: * [`examples/health`](https://github.com/rust-dd/tako/tree/main/examples/health) for the readiness endpoint shape, * [`examples/metrics-opentelemetry`](https://github.com/rust-dd/tako/tree/main/examples/metrics-opentelemetry) for OTLP export. HTTP/3 qlog and `traceparent` propagation through the v2 outbound client are deferred follow-up items. Source: https://tako.rust-dd.com/docs/observability --- # Tutorials > End-to-end, copy-along guides that build a complete tako service from a single example crate in the repository. The rest of the documentation is a reference: one page per transport, per extractor, per concept. These tutorials are the opposite shape — each one builds a complete, runnable service end-to-end, wiring several pieces together the way a real application does. Every tutorial is grounded in a real crate under [`examples/`](https://github.com/rust-dd/tako/tree/main/examples) in the repository, so you can clone it, run it, and read the same code shown here. If you have not written a Tako handler yet, start with the [Quickstart](/docs/getting-started/quickstart) first — these tutorials assume you can already register a route and serve it. ## Available tutorials * **[Building a REST API](/docs/tutorials/rest-api)** — a small JSON API with path, query, and JSON-body extractors, shared application state, a custom error type that maps to HTTP status codes, and a middleware layer. Grounded in [`examples/extractors-multi`](https://github.com/rust-dd/tako/tree/main/examples/extractors-multi). * **[Realtime over WebSocket](/docs/tutorials/realtime-chat)** — an echo and a server-push feed over a full-duplex WebSocket connection, plus the SSE alternative for one-way streams. Grounded in [`examples/websocket`](https://github.com/rust-dd/tako/tree/main/examples/websocket). ## After the tutorials When you are ready to ship, the [Deployment guide](/docs/deployment) covers the single-binary, thread-per-core, and load-balancer deployment shapes, and the [feature reference](/docs/reference/features) lists every cargo flag you can opt into. Source: https://tako.rust-dd.com/docs/tutorials --- # Building a REST API > Build a small JSON REST API with path, query, and body extractors, shared state, a custom error type, and a middleware layer. This tutorial builds a small JSON REST API end-to-end. It starts from the runnable [`examples/extractors-multi`](https://github.com/rust-dd/tako/tree/main/examples/extractors-multi) crate, then layers on the pieces a real service needs: shared application state, a custom error type that maps to HTTP status codes, and a middleware layer. By the end you will have: * `GET /users/{id}/posts?page=&per_page=` — path **and** query extraction * `POST /users` — a typed JSON body, echoed back as JSON * shared state read inside a handler * a fallible handler returning `Result, ApiError>` * a `RequestId` middleware on the router If you have not served a handler before, read the [Quickstart](/docs/getting-started/quickstart) first. ## 1. The starting point The example crate is a two-route API: one route combines a path parameter with query parameters, the other accepts a JSON body and returns JSON. This is the whole file: Source: [examples/extractors-multi/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/extractors-multi/src/main.rs) ```rust use anyhow::Result; use serde::Deserialize; use serde::Serialize; use tako::Method; use tako::extractors::json::Json; use tako::extractors::params::Params; use tako::extractors::query::Query; use tako::router::Router; use tokio::net::TcpListener; #[derive(Deserialize)] struct Pagination { page: u32, per_page: u32, } #[derive(Deserialize)] struct UserPath { id: u64, } #[derive(Deserialize, Serialize, Clone)] struct CreateUser { name: String, email: String, } #[derive(Serialize)] struct Created { id: u64, name: String, email: String, } // GET /users/{id}/posts?per_page=10&page=2 // Demonstrates multiple extractors: Params + Query async fn list_user_posts(Params(user): Params, Query(p): Query) -> String { format!( "user_id={}, page={}, per_page={}", user.id, p.page, p.per_page ) } // POST /users with JSON body {"name":"...","email":"..."} // Demonstrates a body extractor async fn create(Json(user): Json) -> Json { // Normally you'd persist the user; here we just echo back with an id Json(Created { id: 1, name: user.name, email: user.email, }) } #[tokio::main] async fn main() -> Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.route(Method::GET, "/users/{id}/posts", list_user_posts); router.route(Method::POST, "/users", create); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); Ok(()) } ``` A few things to note: * Each handler argument is an **extractor**. `Params` pulls the `{id}` path segment, `Query` parses the query string, and `Json` deserializes the request body. Tako runs every extractor before calling the handler. See [Extractors](/docs/extractors) for the full catalog. * The structs deriving `Deserialize` are the typed targets for those extractors; the response struct derives `Serialize` so `Json` can encode it. * `router.route(Method::GET, "/users/{id}/posts", …)` registers the handler. `{id}` is a `matchit` path slot. (For compile-time-typed slots like `{id: u64}`, see the [`#[tako::get]` macro family](/docs/routing).) ## 2. The Cargo.toml The example depends only on the umbrella crate plus `serde` and a runtime: ```toml [dependencies] anyhow = "1" serde = { version = "1", features = ["derive"] } tako-rs = "2" tokio = { version = "1", features = ["full"] } ``` `tako-rs` is the package name; you import it as `tako`. The default feature set already covers HTTP/1.1, so no flags are needed yet. We add `plugins` later for the middleware section. ## 3. Adding shared state Most APIs need to share something with every handler — a database pool, a config struct, a cache handle. Tako exposes this through the [`State`](/docs/state) extractor backed by `Router::with_state`. Define a `Clone` config type, register it once with `with_state`, then pull it into any handler by adding a `State` argument: ```rust use std::sync::Arc; use tako::extractors::query::Query; use tako::extractors::state::State; use tako::responder::Responder; #[derive(Clone)] struct AppState { default_per_page: u32, } async fn list_user_posts( Params(user): Params, Query(p): Query, State(state): State, ) -> impl Responder { let per_page = if p.per_page == 0 { state.default_per_page } else { p.per_page }; format!("user_id={}, page={}, per_page={}", user.id, p.page, per_page) } ``` Register the state on the router before serving: ```rust let mut router = Router::new(); router.with_state(AppState { default_per_page: 20 }); router.route(Method::GET, "/users/{id}/posts", list_user_posts); ``` `State` reads from the per-router store first and surfaces the value as `Arc`. When no caller ever invoked `with_state`, the lookup short-circuits on a single `AtomicBool::Acquire`, so unused state costs nothing on the hot path. ## 4. A custom error type Real handlers fail: a user is missing, a body is invalid, a downstream call errors. Tako lets a handler return `Result` where **both** arms implement [`Responder`](/docs/routing) — the `Err` arm is turned into a response just like the `Ok` arm. Define an error enum and give it a `Responder` impl that picks the status code: ```rust use http::StatusCode; use tako::extractors::json::Json; use tako::responder::Responder; use tako::types::Response; enum ApiError { NotFound, BadRequest(String), } impl Responder for ApiError { fn into_response(self) -> Response { let (status, msg) = match self { ApiError::NotFound => (StatusCode::NOT_FOUND, "user not found".to_string()), ApiError::BadRequest(m) => (StatusCode::BAD_REQUEST, m), }; (status, msg).into_response() } } ``` Now a handler can fail with a typed error and Tako renders the right status: ```rust async fn get_user(Params(user): Params) -> Result, ApiError> { if user.id == 0 { return Err(ApiError::BadRequest("id must be non-zero".into())); } if user.id > 1_000 { return Err(ApiError::NotFound); } Ok(Json(Created { id: user.id, name: "Ada".into(), email: "ada@example.com".into() })) } ``` For machine-readable error bodies, call `router.use_problem_json()` to emit RFC 7807 `application/problem+json` for 4xx/5xx responses, or implement the `Responder` for your error type to return `tako::problem::Problem` directly. See the [migration guide](/docs/reference/migration) for the full error-handling surface (`error_handler`, `client_error_handler`, `use_problem_json`). ## 5. Adding a middleware layer Cross-cutting concerns — request IDs, auth, rate limiting — attach as middleware without touching handler code. Anything implementing `IntoMiddleware` can be attached to the whole router or a single route. The bundled middleware lives in `tako-rs-plugins`, enabled with the `plugins` feature. A `RequestId` middleware that tags every request with an `X-Request-ID` header: ```rust use tako::middleware::IntoMiddleware; use tako::middleware::request_id::RequestId; let mut router = Router::new(); router.middleware(RequestId::new().into_middleware()); router.route(Method::POST, "/users", create); ``` `router.middleware(…)` is global — every request passes through it. To scope a middleware to one route, chain `.middleware(…)` after the `route(…)` call: ```rust use tako::middleware::bearer_auth::BearerAuth; let bearer = BearerAuth::static_token("my-secret-token").into_middleware(); router .route(Method::POST, "/users", create) .middleware(bearer); ``` See [Middleware](/docs/middleware) for the full bundled set (auth, CORS, compression, rate limiting, sessions, CSRF, metrics, and more). ## 6. Running it ```bash cargo run ``` Then exercise the routes: ```bash # query + path curl 'http://127.0.0.1:8080/users/7/posts?page=2&per_page=10' # user_id=7, page=2, per_page=10 # JSON body in, JSON out curl -X POST http://127.0.0.1:8080/users \ -H 'content-type: application/json' \ -d '{"name":"Ada","email":"ada@example.com"}' # {"id":1,"name":"Ada","email":"ada@example.com"} ``` ## Next steps * [Routing](/docs/routing) — nesting, scopes, typed path slots, the macro family. * [Extractors](/docs/extractors) — the full catalog of request extractors. * [Middleware](/docs/middleware) — the bundled production middleware set. * [Realtime over WebSocket](/docs/tutorials/realtime-chat) — the next tutorial. * [Deployment](/docs/deployment) — ship the binary. Source: https://tako.rust-dd.com/docs/tutorials/rest-api --- # Realtime over WebSocket > Build a full-duplex WebSocket feature with an echo handler and a server-push ticker, then compare it to the SSE alternative. This tutorial builds a realtime feature over a full-duplex WebSocket connection. It is grounded in the runnable [`examples/websocket`](https://github.com/rust-dd/tako/tree/main/examples/websocket) crate, which serves two endpoints: * `GET /ws/echo` — reads each client message and echoes it back * `GET /ws/tick` — pushes a `tick #N` message to the client once per second These are the two halves of any chat-style feature: receiving messages from a client, and pushing messages to it on a schedule or from another source. WebSocket requires the `ws` feature on Tokio. If you have not served a handler before, read the [Quickstart](/docs/getting-started/quickstart) first. ## 1. The full example A WebSocket endpoint in Tako is an ordinary HTTP handler that returns `TakoWs::new(req, closure)`. The closure receives an upgraded socket and drives the conversation. This is the whole file: Source: [examples/websocket/src/main.rs](https://github.com/rust-dd/tako/blob/main/examples/websocket/src/main.rs) ```rust use std::time::Duration; use futures_util::SinkExt; use futures_util::StreamExt; use tako::Method; use tako::responder::Responder; use tako::types::Request; use tako::ws::TakoWs; use tokio_stream::wrappers::IntervalStream; use tokio_tungstenite::tungstenite::Message; use tokio_tungstenite::tungstenite::Utf8Bytes; pub async fn ws_echo(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { let _ = ws.send(Message::Text("Welcome to Tako WS!".into())).await; while let Some(Ok(msg)) = ws.next().await { match msg { Message::Text(txt) => { let _ = ws .send(Message::Text(Utf8Bytes::from(format!("Echo: {txt}")))) .await; } Message::Binary(bin) => { let _ = ws.send(Message::Binary(bin)).await; } Message::Ping(p) => { let _ = ws.send(Message::Pong(p)).await; } Message::Close(_) => { let _ = ws.send(Message::Close(None)).await; break; } _ => {} } } }) } pub async fn ws_tick(req: Request) -> impl Responder { TakoWs::new(req, |mut ws| async move { let mut ticker = IntervalStream::new(tokio::time::interval(Duration::from_secs(1))).enumerate(); loop { tokio::select! { msg = ws.next() => { match msg { Some(Ok(Message::Close(_))) | None => break, _ => {} } } Some((i, _)) = ticker.next() => { let _ = ws.send(Message::Text(Utf8Bytes::from(format!("tick #{i}")))).await; } } } }) } #[tokio::main] async fn main() { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080") .await .unwrap(); let mut router = tako::router::Router::new(); router.route(Method::GET, "/ws/echo", ws_echo); router.route(Method::GET, "/ws/tick", ws_tick); tako::Server::builder() .build() .spawn_http(listener, router) .result() .await .expect("HTTP server failed"); } ``` ## 2. How the upgrade works `TakoWs::new(req, |ws| async move { … })` is a [`Responder`](/docs/routing): when the handler returns it, Tako performs the WebSocket handshake on the incoming `Request` and then runs your closure with the upgraded connection. Inside the closure, `ws` is both a `Stream` of inbound `Message`s and a `Sink` you can `send` outbound `Message`s into — the `futures_util::{StreamExt, SinkExt}` traits provide `.next()` and `.send()`. The echo handler is a straight read-respond loop: ```rust while let Some(Ok(msg)) = ws.next().await { match msg { Message::Text(txt) => { let _ = ws.send(Message::Text(format!("Echo: {txt}").into())).await; } Message::Binary(bin) => { let _ = ws.send(Message::Binary(bin)).await; } Message::Ping(p) => { let _ = ws.send(Message::Pong(p)).await; } Message::Close(_) => { let _ = ws.send(Message::Close(None)).await; break; } _ => {} } } ``` The `Message` type comes from `tokio-tungstenite`. Handling `Ping`/`Close` explicitly keeps the connection healthy and shuts it down cleanly. ## 3. Server-initiated pushes A chat feed has to push messages the client did not ask for — a new message from another user, a periodic heartbeat. The `ws_tick` handler shows the pattern: `tokio::select!` races the inbound stream against a timer, so the task both reacts to the client (closing when it disconnects) and pushes on its own schedule: ```rust loop { tokio::select! { msg = ws.next() => match msg { Some(Ok(Message::Close(_))) | None => break, _ => {} }, Some((i, _)) = ticker.next() => { let _ = ws.send(Message::Text(format!("tick #{i}").into())).await; } } } ``` To turn this into a real chat fan-out, replace the `ticker` arm with a `tokio::sync::broadcast::Receiver`: every connected socket subscribes to a shared `broadcast::Sender`, and posting a message sends it to every subscriber. The `select!` skeleton stays identical — one arm reads the client, the other reads the broadcast channel. ## 4. The Cargo.toml The example pulls in the umbrella crate plus the stream/codec helpers it uses in the closures: ```toml [dependencies] futures-util = "0.3" tako-rs = { version = "2", features = ["ws"] } tokio = { version = "1", features = ["full"] } tokio-stream = "0.1" tokio-tungstenite = "0.30" ``` The `ws` feature enables `tako::ws::TakoWs`. `futures-util` provides the `Stream`/`Sink` extension traits, and `tokio-tungstenite` provides the `Message` enum used on the wire. ## 5. Running it ```bash cargo run ``` Connect with any WebSocket client. With `websocat`: ```bash # echo websocat ws://127.0.0.1:8080/ws/echo # type a line; the server replies "Echo: " # server push websocat ws://127.0.0.1:8080/ws/tick # prints "tick #0", "tick #1", … once per second ``` In the browser, the standard `WebSocket` API works unchanged: ```js const ws = new WebSocket("ws://127.0.0.1:8080/ws/echo"); ws.onmessage = (e) => console.log(e.data); ws.onopen = () => ws.send("hello"); ``` ## Tuning the connection `TakoWs` is a builder. Before returning it you can constrain the handshake and the frames — `protocols`, `max_frame_size`, `max_message_size`, `allowed_origins`, `upgrade_timeout`, and `keep_alive(WsKeepAlive)`. These guard against oversized frames and slow-loris upgrades on a public endpoint. See the [WebSocket transport reference](/docs/transports/websocket) for the full builder. WebSocket uses `ws` on Tokio and `compio-ws` on Compio. HTTP/2 WebSocket (RFC 8441) and the deflate extension are covered in the [transport reference](/docs/transports/websocket). ## When to use SSE instead If the data only flows **server → client** — a live ticker, a log tail, a progress feed — a WebSocket is more machinery than you need. [Server-Sent Events](/docs/transports/sse) give you a one-way `text/event-stream` from an ordinary HTTP handler, with browser auto-reconnect built in, behind the `sse` feature: ```rust use futures_util::stream; use tako::responder::Responder; use tako::sse::{Sse, SseEvent}; use tako::types::Request; async fn feed(_: Request) -> impl Responder { let events = stream::iter([ SseEvent::data("user joined").event("presence"), SseEvent::data("hello").event("message"), ]); Sse::events(events) } ``` Reach for a WebSocket when the client also needs to **send** — chat input, collaborative editing, game state. Reach for SSE when it only needs to listen. ## Next steps * [WebSocket transport](/docs/transports/websocket) — the full `TakoWs` builder and configuration. * [Server-Sent Events](/docs/transports/sse) — the one-way streaming alternative. * [Building a REST API](/docs/tutorials/rest-api) — the companion tutorial. * [Deployment](/docs/deployment) — ship the binary. Source: https://tako.rust-dd.com/docs/tutorials/realtime-chat --- # Deployment > The three tako deployment shapes — single binary, thread-per-core, and behind a load balancer — plus socket activation and graceful shutdown. Tako accommodates three deployment shapes, plus a few production niceties (socket activation, PROXY protocol fronting, per-thread runtimes). Pick the shape that matches your infrastructure; you can use the same router and handler code across these deployment shapes. ```rust use std::time::Duration; use tako::{Server, ServerConfig}; use tako::router::Router; use tokio::net::TcpListener; #[tokio::main] async fn main() -> anyhow::Result<()> { let listener = TcpListener::bind("0.0.0.0:8080").await?; let mut router = Router::new(); router.get("/", || async { "hi" }); let server = Server::builder() .config(ServerConfig { drain_timeout: Duration::from_secs(30), max_connections: Some(50_000), ..ServerConfig::default() }) .build(); let handle = server.spawn_http(listener, router); tako::shutdown_signal().await?; handle.shutdown(Duration::from_secs(30)).await; Ok(()) } ``` ## Single binary (the default) `Server::builder` + tokio multi-threaded runtime. `max_connections` bounds accept-spawn pressure; `drain_timeout` controls the grace period for graceful shutdown. This is the shape most services need. Consider jemalloc as the global allocator when handlers allocate. Connection tasks move between worker threads, so memory is often freed on a different thread than the one that allocated it, a pattern jemalloc usually handles better than the system allocator. A hello-world request does not allocate since 2.4, so the [benchmarks](/docs/benchmarks) show no gain there; measure your own handlers. Enable the `jemalloc` feature and declare: ```rust #[global_allocator] static GLOBAL: tako::Jemalloc = tako::Jemalloc; ``` ## Thread-per-core Two flavours, both behind cargo features: * **`per-thread`** — N current-thread tokio runtimes, one per core, sharing an `Arc` that is released when workers exit. Backed by `SO_REUSEPORT` for kernel-level fan-out across the listeners, with `balance_connections` (on by default) handing each new connection to the worker with the fewest live connections. * **`per-thread-compio`** — io\_uring (Linux), IOCP (Windows), or kqueue (macOS) via the [`compio`](https://github.com/compio-rs/compio) runtime, same per-worker shape. The entry points are re-exported from `tako` with the `per-thread` feature: ```rust use tako::{PerThreadConfig, serve_per_thread}; use tako::router::Router; fn main() -> std::io::Result<()> { let mut router = Router::new(); router.get("/", || async { "hello from per-thread" }); let config = PerThreadConfig { workers: 4, ..PerThreadConfig::default() }; serve_per_thread("0.0.0.0:8080", router, config) } ``` `serve_per_thread` waits for Ctrl+C or SIGTERM. For explicit control, use `spawn_per_thread(addr, router, config)`, await `shutdown.wait_for_bind_outcome(worker_count)`, then call `shutdown.trigger()` and join the returned OS thread handles during teardown. Bind and plugin setup errors are reported through the shutdown handle. These workers currently serve HTTP/1.1. TLS, h2c, and PROXY listeners use the regular server builder. Workers are not pinned to CPU cores by default, so the OS can move a worker off a core that another process keeps busy. On a machine the server has to itself, pinning keeps each worker's caches warm; set `pin_to_core: true` and worker *n* runs on the *n*-th CPU core. The `affinity` feature that does the pinning comes with `per-thread`. ```rust let config = PerThreadConfig { workers: 4, pin_to_core: true, ..PerThreadConfig::default() }; ``` `per-thread-compio` selects Compio consistently for workers, middleware, signals, and queues. Kernel load balancing with `SO_REUSEPORT` is intended for Linux. On macOS, where the kernel sends every connection to one listener, Tokio workers still share the load through `balance_connections`; the Compio flavour does not balance connections, so run it with one worker there. ## Behind a load balancer L4 proxies that prepend a PROXY v1/v2 header (HAProxy, AWS NLB, fly.io edges) are supported by `Server::spawn_proxy_protocol` with `proxy-protocol` enabled. The parser is TLV-aware and CRC32C-verified, and rewrites `X-Forwarded-*` from the parsed header so handlers see the real client address. See the [Transports overview](/docs/transports) and [`examples/proxy-protocol`](https://github.com/rust-dd/tako/tree/main/examples/proxy-protocol). L7 HTTP/2 proxies (Envoy, Nginx) that terminate TLS and forward cleartext h2c upstream use `Server::spawn_h2c` — see the same chapter for an h2c example. ## Socket activation Behind the `socket-activation` cargo feature. `LISTEN_FDS` / `LISTEN_PID` (and the s6 / catflap equivalents) are read by `tako::socket_activation::ListenFds::from_env()`. The returned listener types feed the existing `Server::spawn_*` methods, so the same handler code runs under systemd, runit, or supervised launch without source changes. ## Router ownership Listeners and their connection tasks own `Arc` references. Graceful shutdown stops accepting, drains active work within `drain_timeout`, then cancels remaining tasks. A runtime router-swap API is not provided; `ROUTER_HOT_RELOAD` is reserved for explicit application signals. ## Cross-references * [Transports overview](/docs/transports) — protocol-specific spawn methods. * [Observability](/docs/observability) — request IDs, access logs, metrics, traces. * [Runtime compatibility](/docs/concepts/runtimes) — the tokio vs. compio split and the per-thread variants. Source: https://tako.rust-dd.com/docs/deployment --- # Benchmarks > Hello-world throughput of tako's two serving models against Axum, Actix Web, and ntex under loopback, server-bound, and pipelined load. Hello-world throughput is not the whole story, but it shows whether the request hot path leaves performance on the table. This page compares tako's two serving models with Axum, Actix Web, and ntex under three loads in the same Linux container. Hello-world requests per second under loopback, server-bound, and pipelined load for Tako per-thread with and without core pinning, Actix Web, ntex, Tako, Tako with jemalloc, and Axum Hello-world requests per second under loopback, server-bound, and pipelined load for Tako per-thread with and without core pinning, Actix Web, ntex, Tako, Tako with jemalloc, and Axum Measured on 2026-10-01 with tako-rs 2.4.0 in a Railway container: a 24 vCPU quota on an AMD EPYC 9655P host, Linux 6.18, Rust 1.98.1, and `wrk` 4.1.0. ## Three loads * **Loopback.** 12 server workers and 12 `wrk` threads share the CPU quota, and each connection sends one request at a time. Kernel time for the loopback TCP traffic takes about 5 µs of the roughly 6 µs of server CPU per request, so the thread-per-core servers land close together. * **Server-bound.** 4 server workers face 16 `wrk` threads. The server runs out of CPU first, so its cost per request decides the result. * **Pipelined.** `wrk` writes 16 requests at a time on each connection, as in the TechEmpower plaintext test. Each syscall serves up to 16 requests, so the framework's own work dominates. ## Results Requests per second are the median of five 15-second runs. CPU per request is the server's user plus system time divided by the requests it served. Spread is half the range of the five runs, relative to the median. No run recorded a socket error or a non-2xx response. Against tako-rs 2.3.0 in the same run, each with its default settings, the default server serves 20% more requests at 1,000 connections, 21% more server-bound, and 53% more pipelined; the per-thread server serves 9% more server-bound and 30% more pipelined. ### Loopback, 100 connections | Framework | Requests/sec | p50 | p99 | CPU per request | Spread | | --------------------------- | -----------: | ------: | ------: | --------------: | -----: | | **Tako per-thread** | 1,791,799 | 0.04 ms | 0.80 ms | 6.09 µs | ±1.7% | | Actix Web | 1,730,462 | 0.04 ms | 2.78 ms | 6.42 µs | ±3.2% | | **Tako per-thread, pinned** | 1,729,861 | 0.04 ms | 2.15 ms | 6.13 µs | ±8.0% | | ntex | 1,668,935 | 0.05 ms | 1.08 ms | 6.97 µs | ±4.8% | | **Tako + jemalloc** | 567,643 | 0.14 ms | 0.94 ms | 9.23 µs | ±8.3% | | **Tako** | 553,666 | 0.15 ms | 0.90 ms | 9.08 µs | ±11.7% | | Axum | 528,781 | 0.16 ms | 1.03 ms | 11.28 µs | ±6.4% | ### Loopback, 1,000 connections | Framework | Requests/sec | p50 | p99 | CPU per request | Spread | | --------------------------- | -----------: | ------: | ------: | --------------: | -----: | | **Tako per-thread** | 1,842,063 | 0.26 ms | 1.24 ms | 5.93 µs | ±4.6% | | Actix Web | 1,786,640 | 0.43 ms | 3.02 ms | 6.32 µs | ±2.7% | | **Tako per-thread, pinned** | 1,780,810 | 0.27 ms | 3.38 ms | 6.06 µs | ±4.9% | | ntex | 1,740,750 | 0.52 ms | 1.96 ms | 6.67 µs | ±2.6% | | **Tako** | 1,401,625 | 0.61 ms | 1.96 ms | 7.58 µs | ±4.6% | | **Tako + jemalloc** | 1,386,954 | 0.62 ms | 2.16 ms | 7.79 µs | ±3.7% | | Axum | 1,132,259 | 0.76 ms | 2.57 ms | 10.10 µs | ±6.1% | ### Server-bound, 256 connections | Framework | Requests/sec | p50 | p99 | CPU per request | Spread | | --------------------------- | -----------: | ------: | ------: | --------------: | -----: | | **Tako per-thread** | 612,091 | 0.40 ms | 1.47 ms | 6.42 µs | ±1.1% | | **Tako per-thread, pinned** | 611,956 | 0.40 ms | 1.53 ms | 6.44 µs | ±1.9% | | Actix Web | 570,865 | 0.42 ms | 1.66 ms | 6.84 µs | ±1.9% | | ntex | 543,076 | 0.45 ms | 1.53 ms | 7.20 µs | ±2.9% | | **Tako** | 483,526 | 0.50 ms | 1.40 ms | 7.90 µs | ±2.8% | | **Tako + jemalloc** | 469,807 | 0.52 ms | 2.21 ms | 8.13 µs | ±3.6% | | Axum | 389,844 | 0.64 ms | 1.50 ms | 9.95 µs | ±1.8% | ### Pipelined, 256 connections | Framework | Requests/sec | CPU per request | Spread | | --------------------------- | -----------: | --------------: | -----: | | **Tako per-thread, pinned** | 15,716,822 | 0.75 µs | ±5.9% | | **Tako per-thread** | 15,033,668 | 0.77 µs | ±3.8% | | Actix Web | 12,806,477 | 0.90 µs | ±3.4% | | **Tako** | 11,512,467 | 0.91 µs | ±8.5% | | **Tako + jemalloc** | 10,853,271 | 0.97 µs | ±8.0% | | ntex | 8,749,294 | 1.31 µs | ±3.9% | With 16 requests per write, `wrk` records one latency per batch, so this table leaves latency out. Axum is left out too: `axum::serve` leaves `TCP_NODELAY` off, so pipelined responses wait on delayed ACKs, and it served 88,398 requests per second. Turning nodelay on through axum's `ListenerExt::tap_io` avoids that. ## Two serving models The default server, `Server::builder()`, runs on Tokio's multi-threaded, work-stealing runtime, the same model as Axum. It serves about 24% more requests than Axum at 1,000 connections and when the server is the bottleneck, because it spends less CPU on each request: 7.6 µs against 10.1 µs at 1,000 connections. At 100 connections the two land within 5% of each other. Since 2.4 a hello-world request does not allocate, so jemalloc no longer pulls ahead: **Tako + jemalloc** lands between 3% above and 6% below **Tako**. It can still pay off for handlers that allocate; see [Deployment](/docs/deployment#single-binary-the-default). The `per-thread` feature runs one current-thread Tokio runtime per worker, each with its own `SO_REUSEPORT` listener. It hands each new connection to the worker with the fewest live connections, and the connection then stays on that worker, so requests never cross threads. Actix Web and ntex use the same shape. Tako per-thread serves about 3% more requests than Actix Web on loopback, 7% more when the server is the bottleneck, and 17% more pipelined. The **Tako per-thread, pinned** rows turn on `pin_to_core`, which pins each worker to a CPU core and is off by default. Pinning gained about 5% pipelined, made no difference server-bound, and cost about 3% on loopback with a longer latency tail, because pinned workers compete with the co-located `wrk` threads for their cores. Every process pins from core 0, so turn it on only when the server has the machine to itself. Per-thread workers serve HTTP/1.1; TLS, h2c, and PROXY listeners use the regular server builder. The kernel balances connections across `SO_REUSEPORT` listeners only on Linux, so measure per-thread throughput there. See [Thread-per-core](/docs/deployment#thread-per-core) for the entry points. The Compio flavour (`per-thread-compio`) is not in the tables: the container blocks io\_uring, so its runtime cannot start there. ## Methodology * **Servers.** Each framework serves `GET /` → `Hello, World!` from a minimal release binary (`lto = "fat"`, `codegen-units = 1`) with its default settings; every response is the same 130 bytes. Versions: tako-rs 2.4.0, axum 0.8.9, actix-web 4.15.0, ntex 3.12.3, tokio 1.53.1. * **Load.** `wrk` runs in the same container and connects over loopback. The loopback and pipelined loads split the 24 vCPU quota into 12 server workers and 12 `wrk` threads; the server-bound load gives the server 4 workers and `wrk` 16 threads. For the pipelined load, a Lua script makes `wrk` write 16 requests at a time: ```lua init = function(args) local batch = {} for i = 1, 16 do batch[i] = wrk.format(nil, "/") end req = table.concat(batch) end request = function() return req end ``` * **Runs.** Each measurement is a 3-second warm-up followed by a 15-second `wrk --latency` run. Five rounds rotate the server order, and the tables report medians. * **CPU.** The server's user and system time come from `/proc//stat`, read before and after each run. * **jemalloc.** The `jemalloc` feature only re-exports the allocator; the benchmark binary installs it: ```rust #[global_allocator] static GLOBAL: tako::Jemalloc = tako::Jemalloc; ``` ## Caveats These are **baselines, not universal claims**. The container shares its host with other tenants, the load generator shares the CPU quota with the server, and loopback removes the real network. A hello-world route isolates the request hot path; application logic, TLS, and payload size change the picture. Reproduce on your own hardware before drawing conclusions. ## Reproducing locally The default server is [`examples/hello-world`](https://github.com/rust-dd/tako/tree/main/examples/hello-world): ```bash cargo run --release --manifest-path examples/hello-world/Cargo.toml # in another terminal: wrk -t4 -c100 -d15s --latency http://127.0.0.1:8080/ ``` Save the Lua script above as `pipeline.lua` to repeat the pipelined load: ```bash wrk -t4 -c256 -d15s -s pipeline.lua http://127.0.0.1:8080/ ``` For the per-thread server, use [`examples/bench-pt`](https://github.com/rust-dd/tako/tree/main/examples/bench-pt) on Linux. It takes the mode, the address, and the worker count: ```bash cargo run --release --manifest-path examples/bench-pt/Cargo.toml -- pt-tokio 127.0.0.1:8080 4 ``` `bench-pt` counts requests for its load-distribution report, so expect it to land somewhat below the numbers above. Source: https://tako.rust-dd.com/docs/benchmarks --- # Migrating to 2.4 > Tako 2.4 needs no code changes; review the core pinning default, the coarser HTTP/1 header deadline, and the extractor ENTRIES list. Tako 2.4.0 makes the HTTP/1 request path cheaper and needs no code changes. Per-thread workers no longer pin to cores by default, one timeout behaves slightly differently, request extensions are dropped earlier, and custom extractors can opt into the cheaper path. See the [release notes](https://github.com/rust-dd/tako/releases/tag/v2.4.0) for the complete list. ## Per-thread workers no longer pin to cores `PerThreadConfig::pin_to_core` now defaults to `false`. A pinned worker cannot move off a core that another process keeps busy, and every process pins from core 0, so two per-thread servers on one machine competed for the same cores. To keep the 2.3 behaviour on a machine the server has to itself, set it back: ```rust use tako::PerThreadConfig; let config = PerThreadConfig { pin_to_core: true, ..PerThreadConfig::default() }; ``` ## The header deadline is checked every half deadline On the Tokio server's plain HTTP/1 listener (`spawn_http` and the `serve` functions) and on per-thread workers, `header_read_timeout` no longer starts a timer for every request. Each connection checks it every half deadline instead, so a connection that stops sending request heads closes between 1 and 1.5 times `header_read_timeout` after it went idle, rather than exactly at it. A connection that is still running a handler or streaming a response is never closed by it. TLS, Unix socket, vsock, and PROXY protocol listeners keep the exact timer. ## Extractors list the entries they read `FromRequest` and `FromRequestParts` gain an associated `ENTRIES` constant of the new `tako::extractors::Entries` type. On a route without middleware or other per-request hooks, the router attaches only the entries that the handler's extractors list, such as the matched path or the connection info. The default, `Entries::ALL`, keeps existing extractors working unchanged; see [Writing an extractor](/docs/extractors#writing-an-extractor) to list fewer. ## Request extensions are dropped before the handler runs A handler hands the request's header and extension maps back for reuse as soon as its last extractor has run. A value that middleware put into the request extensions is therefore dropped before the handler body runs, not after it returns. Values that extractors read, such as a session or JWT claims, are not affected, because the extractor keeps its own copy. A middleware that stores a guard there, such as a semaphore permit, and needs it to live until the handler finishes should hold the guard itself: ```rust use std::sync::Arc; use tokio::sync::Semaphore; let limit = Arc::new(Semaphore::new(64)); router.middleware(move |req, next| { let limit = Arc::clone(&limit); async move { let _permit = limit.acquire_owned().await; next.run(req).await } }); ``` ## Performance changes These need no code changes: * Routes without middleware or other per-request hooks take a shorter dispatch path. * Request header and extension maps are reused on the thread that served them, and small handler futures are stored inline. A hello-world request on a per-thread worker's keep-alive connection makes no heap allocation, down from 11 in 2.3. * Plain HTTP/1 and per-thread connections check for shutdown with one atomic load and set no timer per request. * The per-thread accept loop runs as its own task, so it is no longer polled each time a connection wakes. Against 2.3.0 in the same [benchmark run](/docs/benchmarks), the default server serves 20% more hello-world requests at 1,000 connections and 53% more pipelined, and per-thread workers serve 9% more when the server is the bottleneck and 30% more pipelined, each with its default settings. Source: https://tako.rust-dd.com/docs/reference/migration-2-4 --- # Migrating to 2.3 > Replace the removed outbound HTTP client, review the WebTransport and gRPC changes, and use the Compio transports added in Tako 2.3. Tako 2.3.0 removes the outbound HTTP client, adds W3C WebTransport and stable streaming gRPC, and brings every transport to Compio. Most applications only need the first section. See the [release notes](https://github.com/rust-dd/tako/releases/tag/v2.3.0) for the complete list. ## The outbound HTTP client is gone The `tako::client` module (`TakoClient`, `TakoTlsClient`, `V2Client`) and the `client` and `native-certs` features were removed. A build that still enables them fails with an unknown-feature error. Use a dedicated client instead: [reqwest](https://crates.io/crates/reqwest) for async code or [ureq](https://crates.io/crates/ureq) for blocking code. ## WebTransport The `webtransport` feature now serves W3C WebTransport sessions on the HTTP/3 server, on both runtimes, and browsers can connect to it. A session starts as a `CONNECT` route that takes the `WebTransport` extractor; see [WebTransport](/docs/transports/webtransport). * `webtransport` now enables `http3`. * With `webtransport`, QUIC datagrams are always on; `ServerConfig::h3_enable_datagrams = false` no longer turns them off. * `RawQuicSession` and `serve_webtransport` are unchanged and stay Tokio-only. The raw QUIC example moved from `examples/webtransport` to `examples/raw-quic`; `examples/webtransport` is now a browser demo. ## gRPC * `GrpcServerStream` has a new public `deadline` field. A struct literal needs `deadline: None`; `GrpcServerStream::new(stream)` is unaffected. * The first `Err` item of a `GrpcServerStream` ends the call, and later items are no longer sent. Messages that are ready together share one HTTP/2 DATA frame. * A successful unary reply sends `grpc-status` in the HTTP/2 trailers instead of the response headers. grpc-go based clients, including grpcurl, rejected the old replies with `server closed the stream without sending trailers`. * `GrpcRequest` reads only as far as the first message, capped at `MAX_GRPC_MESSAGE_SIZE`, instead of buffering the whole request body. * New: the `Option` extractor and `with_deadline`, `Stream` for `GrpcClientStream` and `GrpcBidi`, `GrpcBidi::respond`, `From for GrpcStatus`, and `GrpcStatus` as a responder. See [gRPC](/docs/transports/grpc). ## Compio h2c, Unix sockets, PROXY protocol, HTTP/3, and WebTransport now run on Compio. `CompioServer` has the same `try_spawn_h2c`, `try_spawn_unix_http`, `try_spawn_proxy_protocol`, and `try_spawn_h3` methods as `Server`, and `tako::server_unix` and `tako::proxy_protocol` are available with `compio`. `UnixPeerAddr` moved to `tako::conn_info` and is still re-exported from `tako::server_unix`. `build_rustls_server_config` now installs aws-lc-rs as the process-wide rustls provider when none is installed. Builds that compile both ring and aws-lc-rs, for example with `http3`, no longer panic at startup. Source: https://tako.rust-dd.com/docs/reference/migration-2-3 --- # Migrating to 2.2 > Update PerThreadConfig literals for connection balancing and review the Tako 2.2 server performance changes. Tako 2.2.0 adds one field to `PerThreadConfig` and changes how the per-thread server spreads connections across workers. Handler, router, and transport code is unchanged. See the [release notes](https://github.com/rust-dd/tako/releases/tag/v2.2.0) for the complete release summary. ## `PerThreadConfig::balance_connections` `PerThreadConfig` gains `balance_connections: bool`, which defaults to `true`. A struct literal that lists every field without `..PerThreadConfig::default()` no longer compiles; add the field or fill the rest from the default: ```rust use tako::PerThreadConfig; let config = PerThreadConfig { workers: 4, ..PerThreadConfig::default() }; ``` With the flag set, a Tokio per-thread worker that accepts a connection while it serves more live connections than the least busy worker hands the socket to that worker. This evens out the hash-based `SO_REUSEPORT` spread on Linux, and on macOS, where the kernel sends every connection to one listener, all workers now serve traffic. Set `balance_connections: false` to keep kernel-only distribution. The Compio flavour (`per-thread-compio`) does not balance connections. ## Performance changes These need no code changes: * The Tokio HTTP/1 servers (plain, TLS, Unix, vsock, and PROXY protocol, plus per-thread workers) reuse one header-read timer per connection instead of registering a new Tokio timer for every request. On the multi-threaded runtime this removes contention on Tokio's time driver under high concurrency. `header_read_timeout` keeps its meaning. * Requests no longer update reference counts that every worker thread shares: each connection holds its own router handle, `MatchedPath` values come from a per-thread copy of the route template (the extension still holds an `Arc`), and per-thread workers watch a per-worker shutdown token. See [Benchmarks](/docs/benchmarks) for the measured effect. Source: https://tako.rust-dd.com/docs/reference/migration-2-2 --- # Migrating to 2.1 > Update handlers, middleware stores, transport features, cookies, static files, and server startup for the Tako 2.1 API. Tako 2.1.0 includes breaking API and default changes. Review this guide before upgrading from 2.0.x. See the [release notes](https://github.com/rust-dd/tako/releases/tag/v2.1.0) for the complete release summary. ## Select transport features explicitly Enable `ws`, `sse`, `proxy-protocol`, and `udp` when your application uses those transports. They are no longer part of the default build. Compio WebSockets use `compio-ws`. GraphQL subscriptions on Tokio need both `async-graphql` and `ws`. ```toml [dependencies] tako-rs = { version = "2.1", features = ["ws", "sse", "plugins", "signals"] } ``` `per-thread-compio` now enables `compio` throughout the workspace. A build has one runtime selection for routing, middleware, signals, and queues; per-thread workers follow that selection. `--all-features` selects Compio and omits the Tokio-only transport exports. See the [feature reference](/docs/reference/features). The `jemalloc` flag exposes the allocator without installing it. An application that wants it must declare: ```rust #[global_allocator] static ALLOCATOR: tako::Jemalloc = tako::Jemalloc; ``` ## Handlers, extractors, and responses Only the final handler argument may consume the body. Earlier arguments must implement `FromRequestParts`. Put `Json`, `Form`, `Bytes`, `String`, or the raw `Request` last. Implementations of `FromRequest` and `FromRequestParts` return a `Send` future; async methods can implement this contract without `async_trait` boxing. Borrowed extractor views are for manual extraction inside a handler; ordinary handler parameters must satisfy the handler's owned bounds. ```rust use tako::extractors::json::Json; use tako::extractors::state::State; #[derive(Clone)] struct Settings { prefix: String } async fn create(State(settings): State, Json(value): Json) -> String { format!("{}{}", settings.prefix, value) } ``` Buffered extractors default to a **2 MiB** limit, including chunked bodies. Use `router.body_limit(bytes)` to change it, or `router.disable_body_limit()` to explicitly allow unbounded buffering. Size failures return 413; malformed content returns the extractor's typed error. JSON accepts `application/*+json` media types and rejects unrelated content types. `Result` requires both `T` and `E` to implement `Responder`. Implement `Responder` for application errors and select an appropriate status. `anyhow` errors log their details and return a generic 500 body. String responses now carry a text content type; byte responses carry `application/octet-stream`. The tuple `(StatusCode, R)` accepts any responder body. `Next` is an opaque middleware continuation; access its behavior through `next.run(request)`. Header-only `Accept` works alone or before a raw request. `HeaderMap` is owned, `RawPath` exposes the URI path, and `Path` deserializes route parameters. ## Router state, paths, and signals Use `Router::with_state` for instance-local state. Global state functions are deprecated. Replacing an existing value of the same type now takes effect. Nested routers preserve child state, route plugins, body limits, and timeouts; child values take precedence over parent values. Child fallback and error handlers are not inherited by `nest` or `merge`; configure them on the parent. Route paths and `MatchedPath` use `Arc`. Signal IDs and metadata keys use `Cow<'static, str>`; use `.as_ref()` when comparing borrowed strings. TLS ALPN metadata uses `Bytes`, and SNI uses `Arc`. Tokio TLS fills negotiated TLS metadata; Compio's current TLS backend does not expose SNI/version getters. Request signals reach the owning router and application arbiter. Route signals also reach the matching route arbiter. `request.completed` includes the matched route and elapsed microseconds. Lifecycle signals use the application arbiter; `server.stopped` is emitted after the listener's normal drain completes. Per-thread workers emit listener lifecycle events individually. `ROUTER_HOT_RELOAD` requires explicit application emission, and the experimental `SignalBus` contract is hidden because it is not connected to dispatch. ## Shared middleware stores Each builder accepts `.store(backend)`: | Builder | Backend contract | | -------------------- | ------------------ | | `SessionMiddleware` | `SessionStore` | | `RateLimiterBuilder` | `RateLimitStore` | | `IdempotencyBuilder` | `IdempotencyStore` | | `Csrf` | `CsrfTokenStore` | | `JwtAuth` | `JwksProvider` | Import traits from `tako::stores` and reference implementations from `tako::stores::memory`. Async backend methods now return `StoreResult`. An operational failure is logged and produces 503 without exposing backend details. A shared backend must enforce TTLs and atomic operations itself. The memory implementations share state across clones within one process. Session blobs include data and an absolute-lifetime Unix timestamp. Keep replica clocks synchronized. Session rotation and destruction remove the old ID. `SessionMiddleware::handle()` administers the default memory backend; custom backends expose their own administration API. A custom rate-limit backend owns capacity and refill policy. `consume` returns `StoreResult`; accepted and rejected snapshots determine the response headers. Local `max_requests` and algorithm settings apply to the built-in limiter. `.client_ip(true)` uses the router's `IpAddrConfig` trusted proxy policy. `.ipv6_prefix(64)` groups IPv6 addresses by subnet; 128 remains the default. Idempotency stores must atomically return `IdempotencyBegin::Acquired(lease)` or `Existing(entry)`. `complete` and `remove` compare that lease token before changing data. Set the pending lease lifetime above your maximum handler duration. `IdempotencyEntry` now uses a 32-byte SHA-256 fingerprint and `Bytes` for the response body. The fingerprint covers method, path/query, content type, and body. The method/path cache key encoding changed; clear old caches during migration. A backend can override `wait` with pub/sub; its default polls every 20 ms. Unknown-length, oversized, and trailer-bearing responses pass through without caching. Replays omit `Set-Cookie` and hop-by-hop headers. Stored CSRF tokens require session middleware before CSRF, even if `bind_to_session(false)` is set. `.single_use(true)` consumes a token atomically; `.token_ttl(duration)` sets its lifetime. Stored mode uses backend-issued tokens and ignores the stateless seed hook. Stateless double-submit mode remains the default when no backend is supplied. JWT providers return `VerificationKey { algorithm, bytes }` values for a `kid`. The bundled `MultiKeyVerifier` accepts raw MAC bytes or DER public keys, bound to the declared algorithm and the verifier's allow-list. Use `.allowed_algorithms(...)` when external keys need an explicit allow-list. An empty provider result falls back to configured static keys; a provider error or a nonempty set of invalid keys fails closed. Custom verifiers must implement `verify_with_key` to use providers. Middleware constraints reach the bundled verifier before signature/time validation. Read verified claims through `extractors::jwt::JwtClaimsVerified` after `JwtAuth`. ## Cookies, streaming, and static files Session and CSRF cookies default to `Secure`. Use `.secure(false)` explicitly for local development over plain HTTP. Static files stream in bounded chunks. Metadata-generated ETags are weak; conditional requests, HEAD, single byte ranges, and precompressed sidecars are handled consistently. Dotfiles are denied by default; `.allow_dotfiles(true)` opts in. Keep served roots read-only to untrusted processes. Percent-encoded traversal paths and encoded separators are rejected. For `FileStream::try_range_response`, an inclusive end of zero means byte zero. Use `u64::MAX` for an open end. Multipart ranges remain unsupported. WebSocket upgrades validate method, HTTP version, required tokens, version 13, and the 16-byte decoded key. These are HTTP/1.1 upgrades, including over TLS; HTTP/2 Extended CONNECT is not implemented. `keep_alive(WsKeepAlive)` is deprecated and has no effect. Implement ping/pong timing in the handler that owns the socket; `max_lifetime` caps total lifetime rather than idle time. Compression skips buffered collection of open-ended streams and skips range responses. Larger buffered responses use a blocking worker; tune `CompressionBuilder::blocking_threshold` for your workload. ## Server startup and shutdown Use `Server::builder()` on Tokio or `CompioServer::builder()` on Compio. The HTTP `serve_*`, shutdown, and configuration convenience matrix is deprecated. The explicit rustls-config entry points remain available for advanced TLS configuration, as do raw-transport and per-thread functions. ```rust use tako::{Server, router::Router}; #[tokio::main] async fn main() -> Result<(), tako::types::BoxError> { let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?; let mut router = Router::new(); router.get("/", || async { "ok" }); let handle = Server::builder().build().try_spawn_http(listener, router)?; handle.result().await?; Ok(()) } ``` `try_spawn_*` returns listener, certificate, and plugin initialization errors. `ServerHandle::result()` reports errors after startup; `local_addr()` includes an ephemeral port selected by port-zero binds. `trigger()` starts graceful draining; `shutdown(timeout)` uses the smaller of that timeout and the configured drain bound. `tako::shutdown_signal()` handles Ctrl+C and Unix SIGTERM. No listener leaks its router to obtain a static reference. To advertise an existing HTTP/3 endpoint on HTTP/1 or HTTP/2 responses, install `middleware::alt_svc::AltSvc::h3(port, max_age)`. Advertising does not start a QUIC listener. Existing Alt-Svc response headers are preserved. ## Dependencies and validation The workspace stays on stable dependency releases and Rust 1.95. Major library updates include Compio, async-graphql, validation libraries, JWT dependencies, metrics exporters, OpenAPI integrations, and compression libraries. Check your application's direct imports when it shares those types with Tako. Doctests now run for every crate. Default, Tokio feature-rich, and Compio builds are checked separately because enabling all features selects the Compio path. Source: https://tako.rust-dd.com/docs/reference/migration-2-1 --- # Migrating from 1.x to 2.0 > Every breaking change between the tako 1.x line on crates.io and the 2.0.0 release, with before/after code for each. > **Status:** released. Covers every breaking change between the released > `1.x` line on crates.io and the `2.0.0` release. This guide explains *what to change in your code* to move from 1.x to 2.0. ## At a glance | Area | 1.x | 2.0 | | ----------------- | ------------------------------------------ | --------------------------------------------------------------------- | | Per-router state | `GLOBAL_STATE`, one per type | `Router::with_state` | | Handler returns | `Responder` only | `Result` where `E: Responder` | | Sub-routing | `Router::merge` | `nest()`, `scope()` | | Wrong method | 404 Not Found | 405 with an `Allow` header | | Errors | `error_handler` (5xx only) | plus `client_error_handler` and `use_problem_json()` | | Macros | `{id: u64}` only, always a `Params` struct | `{id}` and `{id: u64}`; no `Params` struct unless a typed slot exists | | Server bootstrap | `serve_*`, `serve_tls`, … | `Server::builder()` | | TLS | PEM only | PEM, DER, `Resolver`, `ReloadableResolver`, and `ClientAuth` (mTLS) | | Connection info | `SocketAddr` / `UnixPeerAddr` | `ConnInfo` (unified) | | `tako-core-local` | separate `!Send` router | removed | ## 1. Per-router typed state **1.x** ```rust let cfg = Config { db, secret }; Router::state(cfg); // one slot per TypeId, global let router = Router::new(); ``` **2.0** ```rust let router = Router::new() .with_state(Config { db, secret }); // instance-local // `State` extractor reads from the per-router store first, then falls // back to GLOBAL_STATE for backward compat. ``` Two routers in the same process can hold distinct `T` values without newtype wrappers. Hot-path overhead is one `AtomicBool::Acquire` when the feature is unused. ## 2. Handler return types **1.x** ```rust async fn handler() -> impl Responder { ... } ``` **2.0** — `Result` is supported natively: ```rust async fn handler() -> Result, ApiError> { ... } ``` `IntoResponse` is a re-export of `Responder`. New blanket impls were added for `Bytes`, `Vec`, `Cow<'static, str>`, `serde_json::Value`, `(StatusCode, HeaderMap, TakoBody)`, `(StatusCode, HeaderMap)`, `HeaderMap`, and `StatusCode`. ## 3. Sub-routing — `nest` and `scope` **1.x** ```rust let api = Router::new(); let main = Router::new(); main.merge(api); // mutates shared Arc; double-merging stacks // middleware twice ``` **2.0** ```rust let main = Router::new() .nest("/v1", v1_router) .nest("/v2", v2_router) .scope("/admin", |r| { r.layer(admin_auth).get("/", dashboard); }); ``` `nest` clones routes via `Route::cloned_with_path` so re-nesting can never double-stack middleware. `scope` carries a pending prefix consumed by every method shorthand inside the closure. ## 4. 405 with `Allow` instead of 404 **1.x** returned `404 Not Found` when the path matched but the method did not. **2.0** returns `405 Method Not Allowed` with a comma-separated `Allow` header. If you have tests asserting `404` on a path-match-method-miss, update them to expect `405` plus the appropriate `Allow` value. ## 5. RFC 7807 `application/problem+json` **1.x** had `Router::error_handler(...)` that fired only on 5xx. **2.0** adds: * `Router::client_error_handler(...)` — fires on 4xx. * `Router::use_problem_json()` — convenience that installs `default_problem_responder` for both 4xx and 5xx. * `tako::problem::Problem` struct with a `Responder` impl that emits `application/problem+json`. ## 6. Macro syntax **1.x** ```rust #[tako::route("GET", "/users/{id: u64}")] async fn get_user(id: u64) -> ... { ... } // Always emits `GetUserParams` struct, even on static paths. ``` **2.0** * `{id: u64}` and `{id}` are both accepted. The first is a typed slot, the second is `matchit` pass-through. * The `*Params` struct is **only** emitted when at least one typed slot exists. * For static paths with `name = "..."`, a unit-marker struct is still emitted so `Name::METHOD` / `Name::PATH` constants stay reachable. ## 7. Server bootstrap **1.x** ```rust serve(router, addr).await?; serve_tls(router, addr, cert, key).await?; serve_h3(router, addr, cert, key).await?; serve_unix(router, path).await?; serve_proxy_protocol(...).await?; ``` **2.0** ```rust let server = tako::Server::builder() .config(ServerConfig::default() .header_read_timeout(Duration::from_secs(30)) .keep_alive(true) .max_concurrent_streams(100) .max_connections(50_000) .drain_timeout(Duration::from_secs(60))) .tls(TlsCert::pem_paths("cert.pem", "key.pem")) .build(); let handle = server.spawn_http(listener, router); // .spawn_tls / .spawn_h2c / .spawn_h3 / .spawn_unix_http / // .spawn_proxy_protocol / .spawn_tcp_raw / .spawn_udp_raw handle.shutdown(Duration::from_secs(30)).await; ``` Changes from the original v2 shape: * The listener is handed to `spawn_*`, not the builder, so a single `Server` instance can fan out to multiple listeners. * `ServerConfig` is one flat struct instead of `HttpConfig` + `TlsConfig` + `H3Config` + `Limits`. * `ServerHandle::shutdown` returns `()` and is runtime-agnostic (Notify-based) so the same type comes back from both tokio and compio paths. ## 8. TLS **1.x** supported `TlsCert::PemPaths`. **2.0** adds: * `TlsCert::Der { certs, key, client_auth }` * `TlsCert::Resolver { resolver, client_auth }` * `ClientAuth::{Optional(roots), Required(roots)}` for mTLS, threaded through every TCP/TLS, compio-TLS, and HTTP/3 spawn path. * `ReloadableResolver` for hot-reload without a listener restart (callers wire their own file-watcher / signal trigger). * New entry points `serve_tls_with_rustls_config_and_shutdown` and `serve_h3_with_rustls_config_and_shutdown` that take a fully-built `Arc` for advanced cases. ACME (`TlsCert::Acme { ... }`) is **deferred**. ## 9. Trust store: `webpki-roots` vs OS **2.0** adds an opt-in `native-certs` feature that swaps the bundled `webpki-roots` snapshot for `rustls-native-certs` (operating-system trust store). Default behavior is unchanged — `webpki-roots` is still used unless `native-certs` is enabled. ```toml tako-rs = { version = "2", features = ["client", "native-certs"] } ``` The `client` and `native-certs` features and `tako::client` were removed in 2.3.0. Use [reqwest](https://crates.io/crates/reqwest) or [ureq](https://crates.io/crates/ureq) for outbound HTTP. ## 10. Unified `ConnInfo` **1.x** inserted a different connection-info type per transport: `SocketAddr` (TCP/TLS), `UnixPeerAddr` (Unix), something else for H3. **2.0** unifies on: ```rust struct ConnInfo { peer: PeerAddr, // Ip / Unix / Other local: PeerAddr, transport: Transport, // Http1 / Http2 / Http3 / Unix / Tcp tls: Option, // alpn, sni, version } ``` Legacy types (`SocketAddr`, `UnixPeerAddr`, `ProxyHeader`) remain in extensions for backward compatibility, alongside `ConnInfo`. ## 11. Plugin / middleware updates | Plugin / middleware | 2.0 change | | -------------------------- | -------------------------------------------------------------------------------------------------------------------- | | `session` | Idle vs absolute TTL, rolling cookie refresh, `Session::rotate()`, configurable `SameSite`/`Domain`, bulk revocation | | `rate_limiter` | Composite-key support, IETF `RateLimit-*` headers, `Algorithm::Gcra`, `UnkeyedBehavior` choice | | `idempotency` | Verified TTL = 86\_400 s, compio `inflight_wait_timeout_ms` honored | | `jwt_auth` | Iss/aud/leeway constraints, `MultiKeyVerifier`, runtime rotation/revocation, optional remote introspection | | `csrf` | Token bound to `Session`, Origin/Referer allow-list, configurable `SameSite` | | `compression` | `ContentTypePolicy` enum replaces substring filter | | `cors` | `OriginMatcher::{Exact, Suffix, Custom}`, `allow_private_network` | | `metrics` | Latency histogram with `with_buckets(..)` override | | `security_headers` | CSP + nonce, COOP/COEP/CORP, Permissions-Policy, HSTS preload toggle, `X-XSS-Protection` removed | | `request_id` | Now focused on `X-Request-ID`; `traceparent` parsing moved to a new middleware | | **NEW**: `timeout` | Per-request deadline, dynamic-per-request override | | **NEW**: `traceparent` | W3C Trace Context parser/emitter, `TraceContext` extension | | **NEW**: `access_log` | Structured one-line access log; default sink `tracing` INFO | | **NEW**: `problem+json` | Rewrites non-JSON 4xx/5xx into `application/problem+json` | | **NEW**: `circuit_breaker` | Closed/open/half-open with rolling counter | | **NEW**: `ip_filter` | CIDR allow/deny lists | | **NEW**: `healthcheck` | `/live`, `/ready`, `/__drain` with async readiness probes | | **NEW**: `etag` | SHA-1 strong validator, conditional GET | | **NEW**: `tenant` | `X-Tenant-ID` / subdomain / path-segment / custom strategies | | **NEW**: `hmac_signature` | HMAC-SHA256 signature verification | | **NEW**: `json_schema` | Request/response validator | ## 12. Backend traits `tako_rs_plugins::stores` adds five traits: * `SessionStore` * `RateLimitStore` * `IdempotencyStore` * `JwksProvider` * `CsrfTokenStore` Built-in middleware still defaults to in-memory stores. Implement these traits to back middleware with Redis / Postgres / external services. The `.store(backend)` builder connections and fallible, atomic store contracts are implemented in [2.1](/docs/reference/migration-2-1#shared-middleware-stores). > Companion crates `tako-stores-redis` and `tako-stores-postgres` are on > the follow-up list and intentionally not part of the framework > dependency surface. ## 13. Extractors * `tako::extractors::path::Path` is the new axum-style wrapper. The old zero-arg extractor was renamed `RawPath` (**breaking**). Migrate any call site that used `Path` as a no-argument extractor. * `JwtClaims` is renamed `JwtClaimsUnverified`. The old name remains as a `#[deprecated]` alias. The verifying counterpart is `tako_rs_plugins::extractors::jwt::JwtClaimsVerified`, fed by `JwtAuth` middleware. * New: `TypedHeader` (feature `typed-header`), `Extension`, `MatchedPath`, `OriginalUri`, `Host`, `Scheme`, `ConnectInfo`, `ContentLengthLimit`, `QueryMulti`, `MultipartConfig`-driven `BufferedUploadedFile`, `KeyRing`-rotated cookie extractors, `Validated` (features `validator` / `garde`). ## 14. `tako-core-local` removed The `!Send` `LocalRouter` was removed. Per-worker isolation is fully covered by `serve_per_thread` / `serve_per_thread_compio` with the thread-safe `Router`. Replace any `tako::local::*` import with the matching thread-safe equivalent. The `per-thread-local` / `per-thread-compio-local` cargo features are gone. ## 15. Streams * `Sse` gained `SseEvent` builder, `Sse::events(...)`, `Sse::keep_alive(...)`, `last_event_id(headers)` helper. * `TakoWs` builder gained `protocols`, `max_frame_size`, `max_message_size`, `allowed_origins`, `upgrade_timeout`, `keep_alive(WsKeepAlive)`. The keep-alive option never scheduled pings and is deprecated in 2.1; handlers own their ping/pong policy. Per-message deflate is not implemented. * `FileStream::with_etag(..)`, `with_last_modified(..)`, `with_content_type(..)`, `evaluate_conditional(...)`, `weak_etag_from_metadata(...)`. Multipart/byteranges and `sendfile(2)` are deferred. * `ServeDirBuilder::precompressed(...)`, `index_files([...])`, traversal rejection at parse time. SPA fallback uses the same resolver. * WebTransport is currently raw QUIC and is also exported as `RawQuicSession` so call sites can pick the honest name. W3C WebTransport sessions landed in 2.3.0: see [WebTransport](/docs/transports/webtransport). ## 16. gRPC * New: `GrpcServerStream`, `GrpcClientStream`, `GrpcBidi`. * `parse_grpc_timeout(...)` and `read_grpc_deadline(req)` insert a `GrpcDeadline(Instant)` extension. * `GrpcInterceptor` async trait + `InterceptorChain` short-circuiting on the first `Err(GrpcStatus)`. * Reflection / health scaffolding (storage layer ships now; protobuf encoders deferred until consumers don't have to ship `protoc`). * gRPC-Web byte-level decoder/encoder helpers. ## 17. GraphQL & OpenAPI * APQ via `PersistedQueryStore` trait + `MemoryPersistedQueryStore`. * Complexity / depth / cost limits builder on `async_graphql::SchemaBuilder`. * `utoipa` is now the documented primary OpenAPI integration; `vespera` remains available behind its existing cargo feature. ## 18. Queue & signals * `QueueBackend` async trait. `MemoryBackend` ships in-tree; remote brokers go in companion crates. * `Queue::push_dedup(name, payload, key)` collapses duplicate pending jobs. * Cron scheduling behind the `queue-cron` cargo feature. * New signals: `queue.job.queued / started / completed / failed / retrying / dead_letter`. Canonical strings under `tako_rs_core::queue::signal_ids`. * `signals::bus::SignalBus` and `LocalBus` are unconnected scaffolding, not automatic cluster forwarding. They are hidden from generated docs in 2.1. ## 19. v2 client `V2Client` + `V2ClientBuilder` ride on `hyper_util::client::legacy::Client`: * connection pool with idle timeout / per-host caps * per-request timeout * retry policy with exponential backoff * default `User-Agent`, `traceparent`-friendly request handling The legacy `TakoClient` / `TakoTlsClient` keep working for backward compatibility but new code should use `V2Client`. The whole `tako::client` module was removed in 2.3.0. Use [reqwest](https://crates.io/crates/reqwest) or [ureq](https://crates.io/crates/ureq) instead. ## Runtime selection `--all-features` now compiles and selects Compio. HTTP/3, h2c, WebTransport, PROXY protocol, and Unix sockets run on both runtimes. Only the raw QUIC helper (`RawQuicSession`) is Tokio-only, so use a Tokio feature subset for it. See the [2.1 migration guide](/docs/reference/migration-2-1). ## Deferred to a future 2.x release The following are tracked separately and have **not yet** landed: * `tako-stores-redis`, `tako-stores-postgres` * `TlsCert::Acme { ... }` * HTTP/3 qlog * Multipart/byteranges responder, Linux `sendfile(2)` path * Generated gRPC stubs (reflection / health protobuf) * Cluster `SignalBus` impls (Redis, NATS, Kafka) * Hot-reload `Arc` swap (current path keeps `Box::leak` for per-connection performance) When these land, this guide is updated alongside. Source: https://tako.rust-dd.com/docs/reference/migration --- # Cargo feature graph > Every cargo feature on the tako-rs umbrella crate — transports, runtimes, plugins, extractors, observability, and runtime selection. The umbrella crate `tako-rs` is the public entry point. Sub-crate features are reached through it — for example, `tako-rs/multipart` turns on `tako-rs-extractors/multipart` and `tako-rs-core/multipart` together. You should never need to depend on `tako-rs-core`, `tako-rs-extractors`, `tako-rs-server`, `tako-rs-server-pt`, `tako-rs-streams`, or `tako-rs-plugins` directly. To inspect the current shape locally: ```bash cargo metadata --format-version 1 --no-deps \ | jq '.packages[] | select(.name == "tako-rs") | .features' ``` ## Transports | Feature | Description | Gates | | ------------------- | ------------------------------------------------------------------------------------------------------------------------ | --------------------------------------------------------------------------------------------------- | | *(default)* | HTTP/1.1, raw TCP, Unix sockets, static files, core extractors and middleware. No optional flags are enabled. | always on | | `ws` | WebSocket upgrades over HTTP/1.1 on Tokio; GraphQL subscriptions also require `async-graphql`. | `tako-rs-core/ws`, `tako-rs-streams/ws` | | `sse` | Server-Sent Events on either runtime. | `tako-rs-streams/sse` | | `file-stream` | `FileStream` responder and range helpers; static file serving is available without this flag. | `tako-rs-streams/file-stream`, `tako-rs-core/file-stream` | | `proxy-protocol` | Trusted PROXY v1/v2 listener on either runtime. | `tako-rs-server/proxy-protocol` | | `udp` | Raw UDP datagrams on either runtime. | `tako-rs-server/udp` | | `socket-activation` | Inherit listener descriptors through `ListenFds`. | `tako-rs-server/socket-activation` | | `tls` | Server-side HTTPS via `rustls`. | `tako-rs-server/tls`, `tako-rs-core/tls` | | `http2` | HTTP/2 cleartext (h2c) and ALPN-negotiated h2 over TLS. | `tako-rs-server/http2`, `tako-rs-core/http2` | | `http3` | HTTP/3 over QUIC; also enables the QUIC dependency tree in `tako-rs-streams`. | `tako-rs-server/http3`, `tako-rs-streams/http3`, `tako-rs-core/http3` | | `webtransport` | W3C WebTransport sessions on the HTTP/3 server, on either runtime, plus the Tokio-only raw QUIC helper. Implies `http3`. | `http3`, `tako-rs-server/webtransport`, `tako-rs-streams/webtransport`, `tako-rs-core/webtransport` | | `vsock` | Linux vsock listener for sidecar workloads. | `tako-rs-server/vsock` | ## Runtime selection | Feature | Description | Gates | | ------------------- | ---------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------- | | `compio` | Switch the server to the `compio` runtime (io\_uring / IOCP / kqueue). Mutually exclusive with the tokio path at build time. | `tako-rs-core/compio`, `tako-rs-server/compio`, `tako-rs-streams/compio`, `tako-rs-plugins/compio`, `tako-rs-server-pt?/compio` | | `compio-tls` | TLS on compio. Implies `compio`. | `compio`, `tako-rs-server/compio-tls`, `tako-rs-core/compio-tls` | | `compio-ws` | WebSocket on compio. Implies `compio`. | `compio`, `tako-rs-streams/compio-ws`, `tako-rs-core/compio-ws` | | `per-thread` | Thread-per-core HTTP/1; Tokio by default, Compio when `compio` is selected. | `dep:tako-rs-server-pt` | | `per-thread-compio` | Thread-per-core on Compio. Implies `per-thread` and `compio` throughout the workspace. | `per-thread`, `compio`, `tako-rs-server-pt/compio` | See [Runtime compatibility](/docs/concepts/runtimes) for the tokio vs compio trade-offs. ## Plugins, middleware, signals | Feature | Description | Gates | | ---------------- | ------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------- | | `plugins` | Bundled middleware and plugins (CORS, compression, rate limiting, idempotency). | `tako-rs-core/plugins`, `tako-rs-plugins/plugins`, `tako-rs-server/plugins`, `tako-rs-server-pt?/plugins` | | `signals` | In-process pub/sub bus, queue lifecycle signals, transport signals. | `tako-rs-core/signals`, `tako-rs-server/signals`, `tako-rs-plugins/signals`, `tako-rs-server-pt?/signals` | | `ip-filter` | `IpFilter` middleware. | `tako-rs-plugins/ip-filter` | | `hmac-signature` | `HmacSignature` middleware. | `tako-rs-plugins/hmac-signature` | | `json-schema` | `JsonSchema` request-validation middleware. | `tako-rs-plugins/json-schema` | | `zstd` | Zstandard compression in `plugins::compression`. Implies `plugins`. | `tako-rs-plugins/zstd`, `tako-rs-core/zstd`, `plugins` | ## Extractors | Feature | Description | Gates | | ---------------------- | --------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------ | | `multipart` | `Multipart` / `TakoTypedMultipart` body extractors. | `tako-rs-extractors/multipart`, `tako-rs-core/multipart` | | `protobuf` | `Protobuf` extractor via `prost`. | `tako-rs-extractors/protobuf`, `tako-rs-core/protobuf` | | `simd` | Umbrella that enables both `simd-sonic` and `simd-json-impl`. | `simd-sonic`, `simd-json-impl` | | `simd-sonic` | `sonic-rs`-backed JSON path. | `tako-rs-extractors/simd-sonic`, `tako-rs-core/simd-sonic` | | `simd-json-impl` | `simd-json`-backed JSON path. | `tako-rs-extractors/simd-json-impl`, `tako-rs-core/simd-json-impl` | | `typed-header` | `TypedHeader` extractor (via `headers`). | `tako-rs-extractors/typed-header` | | `zero-copy-extractors` | Borrowed variants of `Json`, `Form`, `Query`, `HeaderMap`. | `tako-rs-extractors/zero-copy-extractors`, `tako-rs-core/zero-copy-extractors` | | `validator` | `Validated` adapter for the [`validator`](https://crates.io/crates/validator) crate. | `tako-rs-extractors/validator` | | `garde` | `Validated` adapter for the [`garde`](https://crates.io/crates/garde) crate. | `tako-rs-extractors/garde` | | `jwt-simple` | JWT decode / verify backed by [`jwt-simple`](https://crates.io/crates/jwt-simple). | `tako-rs-plugins/jwt-simple`, `tako-rs-core/jwt-simple` | | `ahash` | Swap the default request hasher for `ahash` across the workspace. | `tako-rs-core/ahash`, `tako-rs-extractors/ahash`, `tako-rs-plugins/ahash` | ## Docs, GraphQL, gRPC, OpenAPI | Feature | Description | Gates | | --------------- | ---------------------------------------------------------------------- | ------------------------------------ | | `async-graphql` | GraphQL HTTP handlers; WebSocket subscriptions additionally need `ws`. | `tako-rs-core/async-graphql` | | `graphiql` | GraphiQL IDE endpoint. | `tako-rs-core/graphiql` | | `grpc` | gRPC unary and streaming RPCs via prost. | `tako-rs-core/grpc` | | `utoipa` | OpenAPI docs via [`utoipa`](https://crates.io/crates/utoipa). | `tako-rs-core/utoipa` | | `utoipa-yaml` | YAML output for the `utoipa` integration. Implies `utoipa`. | `utoipa`, `tako-rs-core/utoipa-yaml` | | `vespera` | OpenAPI docs via [`vespera`](https://crates.io/crates/vespera). | `tako-rs-core/vespera` | ## Observability backends | Feature | Description | Gates | | ----------------------- | ------------------------------------------------------------------------ | --------------------------------------------------------------------------------------------------- | | `tako-tracing` | `tracing-subscriber` integration helpers. | `tako-rs-core/tako-tracing`, `tako-rs-server/tako-tracing` | | `metrics-prometheus` | Prometheus scrape endpoint + histogram. Implies `plugins` and `signals`. | `tako-rs-plugins/metrics-prometheus`, `tako-rs-core/metrics-prometheus`, `plugins`, `signals` | | `metrics-opentelemetry` | OpenTelemetry OTLP metrics export. Implies `plugins` and `signals`. | `tako-rs-plugins/metrics-opentelemetry`, `tako-rs-core/metrics-opentelemetry`, `plugins`, `signals` | ## Queue `queue-cron` enables cron schedules and their chrono-based time calculations in `tako-rs-core`. ## Allocator | Feature | Description | Gates | | ---------- | --------------------------------------------------------------------------------- | ----------------------- | | `jemalloc` | Re-export `tikv-jemallocator`; the application installs its own global allocator. | `dep:tikv-jemallocator` | ## Runtime and protocol combinations `--all-features` compiles and selects Compio for the HTTP server, middleware, signals, queues, and per-thread workers. The Tokio-only raw QUIC helper is absent in that build. h2c is available through `http2`, WebTransport through `webtransport`, HTTP/3 through `http3`, PROXY protocol through `proxy-protocol`, Unix sockets on every unix target, and TLS with ALPN HTTP/2 through `compio-tls` plus `http2`. Tako enables Hyper's HTTP/2 support only through `http2` (other dependencies may also enable it). The workspace selects the Tokio runtime, I/O, timer, signal, filesystem, and synchronization features it uses instead of `full`. The workspace MSRV is Rust 1.95. See the [2.1 migration guide](/docs/reference/migration-2-1) for feature changes and explicit allocator configuration. ## Removed features `client` and `native-certs` were removed in 2.3.0, together with `tako::client`. Use [reqwest](https://crates.io/crates/reqwest) or [ureq](https://crates.io/crates/ureq) for outbound HTTP. Source: https://tako.rust-dd.com/docs/reference/features --- # API reference > The in-source rustdoc on docs.rs is the canonical per-item API reference for tako — here is the per-crate link list. This site is the long-form *guide*. The canonical, per-item API reference is the in-source rustdoc published on [docs.rs](https://docs.rs). Every public type, trait, function, and feature flag is documented there, generated directly from the source, so it never drifts from the code. When this site and the rustdoc disagree about a single API, **the rustdoc wins** — it is generated from the exact released source. This site wins for intent, recommended patterns, and how the pieces fit together. ## Where to read it Application code depends only on the umbrella crate and reaches everything through the `tako::*` path, so the umbrella docs are the entry point: * **[`tako-rs`](https://docs.rs/tako-rs)** — the umbrella. Start here. Everything intended for applications is re-exported under `tako::*`. ## Per-crate rustdoc The sub-crates are implementation detail — you rarely name them directly — but their rustdoc is published for browsing the source-level layout: | Crate | docs.rs | Owns | | -------------------- | ---------------------------------------------------------------- | ---------------------------------------------------------------------------------------- | | `tako-rs` | [docs.rs/tako-rs](https://docs.rs/tako-rs) | Umbrella re-export (`tako::*`), cargo features, prelude | | `tako-rs-core` | [docs.rs/tako-rs-core](https://docs.rs/tako-rs-core) | Router, handlers, middleware/plugin traits, body/request/response, state, signals, queue | | `tako-rs-extractors` | [docs.rs/tako-rs-extractors](https://docs.rs/tako-rs-extractors) | Concrete request extractors | | `tako-rs-server` | [docs.rs/tako-rs-server](https://docs.rs/tako-rs-server) | HTTP/1.1, HTTP/2, HTTP/3, TLS, raw TCP/UDP/Unix, PROXY protocol, compio variants | | `tako-rs-streams` | [docs.rs/tako-rs-streams](https://docs.rs/tako-rs-streams) | WebSocket, SSE, file streaming, static files, WebTransport | | `tako-rs-plugins` | [docs.rs/tako-rs-plugins](https://docs.rs/tako-rs-plugins) | Bundled middleware and plugins | | `tako-rs-macros` | [docs.rs/tako-rs-macros](https://docs.rs/tako-rs-macros) | The `#[tako::route]` / `#[tako::get]` attribute family | | `tako-rs-server-pt` | [docs.rs/tako-rs-server-pt](https://docs.rs/tako-rs-server-pt) | Thread-per-core entry point | Import from the umbrella (`tako::router::Router`), not the sub-crate (`tako_rs_core::...`). Use the [migration guide](/docs/reference/migration-2-1) when upgrading to the 2.1 API. ## Building the docs locally To read the rustdoc for the exact version you depend on, build it from your own checkout: ```bash cargo doc -p tako-rs --no-deps --open ``` Add the features you use so the gated items appear — for example `--features "http2 tls plugins"`. See the [feature reference](/docs/reference/features) for the full flag list. Source: https://tako.rust-dd.com/docs/reference/api --- # Contributing > How to build and test the tako workspace, the rustfmt/MSRV/edition conventions, the crate layout, and the crates.io publish flow. Tako is a Cargo workspace of eight crates. This page covers how to build and test it, the style conventions the codebase follows, and the release flow. ## Building and testing Clone the repository and build the workspace: ```bash git clone https://github.com/rust-dd/tako cargo build --workspace ``` Run the test suite. The strictest configuration enables every feature, which activates both runtimes and every transport, extractor, and plugin: ```bash cargo test --workspace --all-features ``` `--all-features` selects Compio and excludes Tokio-only exports. Also run a Tokio feature subset to cover the Tokio side of each transport. See the [feature reference](/docs/reference/features). ## Toolchain * **MSRV 1.95**, **edition 2024**, workspace-wide. The code relies on let-chains and other edition-2024 features, so an older toolchain will not compile it. * The MSRV is recorded in the workspace `Cargo.toml` under `[workspace.package].rust-version`. CI checks this floor alongside stable and beta toolchains. ## Formatting and lints Formatting runs on **nightly**, because `rustfmt.toml` uses the nightly-only `imports_granularity` and `group_imports` options: ```bash cargo +nightly fmt --all ``` The house rustfmt config is 2-space indentation, item-granularity imports, and `StdExternalCrate` import grouping (std / external / local). Clippy must be clean at `-D warnings`. The workspace lint table sets `pedantic = warn`, so `-D warnings` catches pedantic regressions too: ```bash cargo clippy --workspace --all-features --no-deps -- -D warnings ``` ## Code style A few conventions the codebase holds to: * **No decorative comments.** No banner / divider comments, no ASCII section headers, no commented-out code, and no "added for X / used by Y / changed in PR Z" notes. A comment earns its place only when it explains a non-obvious *why* — a hidden constraint, a workaround, an invariant, a surprising performance choice. If a section feels big enough to need a header, split it into a submodule instead. * **Don't restate the code.** Skip comments that paraphrase the next line; good names carry the narration. * **Doc comments stay short.** A one-line summary plus an example covers most public items. Walls of doc text get skimmed and go stale. * **No `#[non_exhaustive]` marker.** The codebase does not annotate public types with `#[non_exhaustive]`. Document source-breaking changes in the migration guide, including added enum variants and required struct fields. ## Commit messages Use bare Conventional-Commit prefixes with **no scope**: `fix:`, `feat:`, `perf:`, `docs:`, `refactor:`, `chore:` — never the `fix(scope):` form. ## Workspace layout Eight crates, published as `tako-rs` plus seven `tako-rs-*` sub-crates. The umbrella re-exports everything under `tako::*`; the sub-crates are implementation detail. | Crate | Owns | | -------------------- | -------------------------------------------------------------------------------------------------------------------------------- | | `tako-rs` | Umbrella re-export, cargo feature surface, prelude, allocator hook | | `tako-rs-core` | Router, handlers, middleware/plugin traits, body/request/response types, state, signals, queue, GraphQL / gRPC / OpenAPI helpers | | `tako-rs-extractors` | Concrete request extractors | | `tako-rs-server` | HTTP/1.1, HTTP/2, HTTP/3, WebTransport, TLS, raw TCP / UDP / Unix, PROXY protocol, compio variants | | `tako-rs-streams` | WebSocket, SSE, file streaming, static files, raw QUIC sessions | | `tako-rs-plugins` | Bundled middleware and plugins | | `tako-rs-macros` | The `#[tako::route]` / `#[tako::get]` attribute family | | `tako-rs-server-pt` | Thread-per-core entry point | For more on what each crate owns, see [Project layout](/docs/getting-started/project-layout). The unprefixed `tako-*` names on crates.io are owned by an unrelated party, so the workspace publishes under the `tako-rs-*` prefix. Application code depends on the umbrella `tako-rs` and imports it as `tako`. ## Publishing Releases go out through `publish.sh` at the repository root. It runs a pre-publish gate (fmt + clippy + workspace tests), then publishes each crate to crates.io in topological order — each sub-crate before any crate that depends on it, with the umbrella `tako-rs` last. ```bash ./publish.sh --dry-run # validate without uploading (still runs the gate) ./publish.sh # publish for real ``` Other flags: `--allow-dirty` to publish with uncommitted changes, and `--skip-gate` to bypass the fmt/clippy/test gate (not recommended). The script is re-runnable — crates already published at the current version are skipped, so a partial failure can be retried safely. Source: https://tako.rust-dd.com/docs/contributing