# 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::