ntex framework guide
Learn how ntex services fit together, run servers, manage state and I/O, and build web applications.
- Service Model
- Service Pipelines
- Service State
- Runtime
- Server
- Worker and Request State
- I/O Abstraction Layer
- Web Application Framework
- Byte Buffers
API reference
Migration
Component and Service Model
The term component model can make a simple idea sound more complicated than it really is. It often brings to mind IoC containers, layers of dependency injection, and plenty of indirection. That is not what we are trying to build.
Our goal is much simpler: reusable components that are easy to understand and, most importantly, easy to combine.
Imagine a service that receives an operation, performs some work, and returns a result. Its implementation might be quite complex, but callers should not need to know about any of that. From the outside, the service should have a small, predictable interface.
In Rust, the most natural way to express this is with a function:
async fn execute(op: Operation) -> Result<OperationResult, Error> {
// Execute the operation.
// ...
}
That is all the caller needs to see. Give the function an Operation, and it
either returns an OperationResult or reports an Error. How it produces that
result is an implementation detail.
This small abstraction already gives us almost everything we need: one input, one output, and a clear contract. There is no hidden framework behavior or unnecessary ceremony.
It is also easy to understand, straightforward to test, and naturally composable.
Building an HTTP Endpoint
Suppose we want to expose our execution service through an HTTP endpoint. What else do we need?
At a minimum, we need a thin layer to handle the HTTP-specific details. It must deserialize the incoming request, call the execution service, and serialize the result into an HTTP response.
Once again, the shape can stay simple. The endpoint can be just another function:
async fn endpoint(req: HttpRequest) -> Result<HttpResponse, Error> {
// Extract an operation from the request.
let operation = load_operation(req).await?;
// Execute the operation.
let result = execute(operation).await?;
// Convert the result into an HTTP response.
into_response(result).await
}
There is nothing particularly fancy happening here. The endpoint is simply glue code:
- Turn the HTTP request into a domain value—an
Operation. - Pass the operation to the execution service.
- Turn the resulting
OperationResultinto an HTTP response.
Each part has a single responsibility, and the boundaries are explicit. The execution engine knows nothing about HTTP, while the HTTP layer does not need to know how the operation is executed.
More importantly, these pieces do not depend on one another directly. Each one simply transforms one type into another.
That is the key idea.
Composing Services
To compose two services, the output of one simply needs to match the input of the next. When the types line up, the services naturally form a transformation chain.
In this example, the overall transformation is from an HTTP request to an HTTP response:
HttpRequest -> HttpResponse
Everything that happens in between is an implementation detail.
We can add more steps, such as authentication and authorization, without changing the overall shape of the system:
async fn endpoint(req: HttpRequest) -> Result<HttpResponse, Error> {
// Authenticate the request.
let req = authenticate(req).await?;
// Extract the operation from the request.
let operation = load_operation(req).await?;
// Check whether the operation is allowed.
let operation = authorize(operation).await?;
// Execute the operation.
let result = execute(operation).await?;
// Convert the result into an HTTP response.
into_response(result).await
}
The processing flow now looks like this:
HttpRequest
-> authenticate
-> load operation
-> authorize
-> execute
-> build response
-> HttpResponse
Each step is small and focused. It receives a value, performs its work, and returns the value expected by the next step. There is no tight coupling between the components.
Once we start thinking in these terms, the component model becomes almost trivial: components are functions, and composition means connecting compatible outputs and inputs.
Despite its simplicity, this model scales surprisingly far.
Defining a Service
Now that we have the basic idea, we can formalize what a service actually is.
At its core, a service accepts an input and produces an output. The operation may be asynchronous, and it may fail. That should already sound familiar—it is the same shape as the functions in the previous examples.
We can describe this idea with a small Service trait that resembles Rust’s
Fn traits:
trait Service<Req> {
/// Response produced by the service.
type Res;
/// Error produced by the service.
type Error;
async fn call(&self, req: Req) -> Result<Self::Res, Self::Error>;
}
This trait gives us a common way to describe any transformation from a request to a response.
The real value lies in what this common interface allows us to build. We are no longer limited to ordinary functions. A service can be:
- an asynchronous function;
- a struct with its own configuration or internal state;
- a wrapper around another service; or
- a chain composed of several services.
As long as every component follows the same Service contract, we can write
generic tools that work with all of them. This is what makes middleware,
combinators, and reusable pipelines possible.
Instead of calling a collection of unrelated functions, we now have a system in which every component is composable by design.
Services in Networking
This model is not limited to HTTP endpoints. If you look closely at networking code—especially the generic, reusable parts—you will find the same pattern almost everywhere.
Consider a TCP connection handler:
impl Service<TcpStream, Res = (), Error = io::Error>
It receives a TcpStream and either handles the connection successfully or
returns an error. Once connection handling follows this interface, we can build
a generic server that works with any compatible TCP service. This is the basic
idea behind ntex-server.
A TCP connector fits the same model:
impl Service<net::SocketAddr, Res = TcpStream, Error = io::Error>
It receives a socket address and returns an established TCP connection. From the caller’s point of view, it is simply another transformation:
SocketAddr -> TcpStream
A TLS handshake is another example:
impl<T: Stream> Service<T, Res = TlsStream<T>, Error = io::Error>
It receives a plain stream and returns a TLS-enabled stream:
T -> TlsStream<T>
The TLS service does not need to know whether the stream came from a server, a client, or an in-memory transport. It only needs a compatible stream as input.
All these components follow the same basic pattern:
input -> service -> output
The types differ, but the model remains the same. A connector creates a stream, a TLS service transforms it, and a connection handler consumes it. Because these components share a common interface, they can be wrapped, reused, and composed without needing to understand one another’s implementation.
This is the same idea we explored with HTTP requests, applied to the rest of the networking stack.
From a Simple Trait to a Practical Framework
The ntex-service crate turns this model into a practical abstraction for real applications. The entire ntex framework builds on it, from low-level connection handling to high-level HTTP and web services.
Of course, a production-ready service model needs more than the minimal
Service trait shown here. It must also support concerns such as:
- service initialization and configuration;
- readiness and backpressure;
- graceful shutdown;
- middleware;
- service composition; and
- pipeline state.
These features add some complexity, but they do not change the core idea. A service still accepts an input and produces an output. Everything else exists to make that simple model reliable, reusable, and practical at scale.
Service Pipelines
In the previous section, we described what a Service is and how to define one.
Consider the HTTP endpoint from our earlier example:
async fn endpoint(req: HttpRequest) -> Result<HttpResponse, Error> {
// Extract an operation from the request.
let operation = load_operation(req).await?;
// Execute the operation.
let result = execute(operation).await?;
// Convert the result into an HTTP response.
into_response(result).await
}
This function is a sequence of fallible, asynchronous transformations. Each step takes the output of the previous step and either produces the next value or returns an error.
That structure maps naturally to a service pipeline. We can compose the same
steps with an and_then combinator, similar to Result::and_then():
let pipeline = service(authenticate)
.and_then(load_operation)
.and_then(authorize)
.and_then(execute)
.and_then(into_response);
Each service runs only if the previous one succeeds. If any stage returns an error, the remaining stages are skipped and the pipeline returns that error.
The completed pipeline is itself a service:
Service<HttpRequest, Res = HttpResponse>
From the outside, it behaves just like the original endpoint function. It
accepts an HttpRequest and produces an HttpResponse. Everything that happens
between those two types remains an implementation detail.
The benefit of this approach is that each stage remains independent. We can replace, remove, or insert a service without rewriting the rest of the pipeline.
For example, adding request throttling is a local change:
let pipeline = service(authenticate)
.and_then(throttle) // <- throttling service
.and_then(load_operation)
.and_then(authorize)
.and_then(execute)
.and_then(into_response);
The external interface has not changed. The pipeline still represents the same overall transformation:
HttpRequest -> HttpResponse
The only requirement is that the types line up: the output of one service must match the input expected by the next.
A throttling service can be designed specifically for HTTP requests:
Service<HttpRequest, Res = HttpRequest>
Such a service inspects a request and returns it unchanged when the request is allowed to continue.
The service can also be generic over its input:
Service<R, Res = R>
A generic throttling service can act as a reusable concurrency or rate-limiting boundary anywhere in a pipeline. It does not need to know anything about HTTP or the application’s domain types.
This is what makes service pipelines useful: they let us build larger behavior from small, focused transformations while preserving a simple interface for the result.
Implementing and_then
Given the Service trait, we can describe and_then as a generic operation. It
accepts two services and connects them so that the output of the first becomes
the input of the second:
fn and_then<Req, A, B>(first: A, second: B) -> impl Service<Req, Res = B::Res, Error = A::Error>
where
A: Service<Req>,
B: Service<A::Res, Error = A::Error>,
{
AndThen { first, second }
}
The type constraints describe the rules of composition:
firstaccepts the original request.secondaccepts the response produced byfirst.- Both services use the same error type.
- The combined service returns the response produced by
second.
The implementation is equally straightforward:
struct AndThen<A, B> {
first: A,
second: B,
}
impl<Req, A, B> Service<Req> for AndThen<A, B>
where
A: Service<Req>,
B: Service<A::Res, Error = A::Error>,
{
type Res = B::Res;
type Error = A::Error;
async fn call(&self, req: Req) -> Result<Self::Res, Self::Error> {
let intermediate = self.first.call(req).await?;
self.second.call(intermediate).await
}
}
The first service processes the request and produces an intermediate value. If
it succeeds, that value is passed to the second service. If it fails, the ?
operator returns the error immediately and the second service is never called.
This is ordinary sequential composition with short-circuiting on error. Other
familiar combinators, such as then, map, and map_err, follow the same
general pattern.
In practice, this machinery is already available in the
ntex-service crate.
Its
ServiceChain
type provides a fluent, strongly typed API for composing services:
let chain = service(authenticate)
.and_then(load_operation)
.and_then(authorize)
.and_then(execute)
.and_then(into_response);
Each call wraps the existing chain in another service combinator. The resulting
ServiceChain is itself a service, so it can be extended further, wrapped in
middleware, or placed in a Pipeline.
The entire chain remains strongly typed. If one service produces a value that the next service cannot accept, the code fails to compile. Invalid pipelines are therefore caught while the application is being built rather than after it starts running.
Readiness and Backpressure
So far, Service::call() describes what a service does, but it says nothing
about when the service is able to accept more work.
Real services are not always ready. A service may have reached its limit for in-flight requests, filled an internal queue, or be waiting for an external resource such as a database connection.
If callers can invoke call() at any time, the service must either buffer an
unlimited amount of work or implement backpressure through a separate,
non-composable mechanism.
We can make backpressure part of the service contract by adding a readiness check:
trait Service<Req> {
type Res;
type Error;
async fn ready(&self) -> Result<(), Self::Error>;
async fn call(&self, req: Req) -> Result<Self::Res, Self::Error>;
}
Conceptually, ready() acts as a gate. It waits until the service can accept
more work without exceeding its internal limits. If the service cannot become
ready—for example, because an underlying resource has failed—it returns an
error.
A caller would follow this pattern:
service.ready().await?;
let response = service.call(request).await?;
For a composed service such as AndThen<A, B>, readiness is no longer local to
one component. A request may pass through both services, so the combined service
is ready only when both components are ready.
A simplified implementation looks like this:
impl<Req, A, B> Service<Req> for AndThen<A, B>
where
A: Service<Req>,
B: Service<A::Res, Error = A::Error>,
{
type Res = B::Res;
type Error = A::Error;
async fn ready(&self) -> Result<(), Self::Error> {
self.first.ready().await?;
self.second.ready().await?;
Ok(())
}
async fn call(&self, req: Req) -> Result<Self::Res, Self::Error> {
let intermediate = self.first.call(req).await?;
self.second.call(intermediate).await
}
}
This captures the basic rule: the composed service should not accept a request unless every stage needed to process it is ready.
A production implementation can check the services concurrently rather than waiting for them one at a time. More importantly, it must coordinate readiness across all concurrent calls to the same service graph.
That coordination is the responsibility of the Pipeline.
Pipeline Readiness
A Pipeline turns a service chain into a runnable service graph. It owns the
graph’s state and coordinates readiness, calls, and shutdown.
This matters because several requests may be in flight at the same time. Whether the pipeline can accept more work may depend on active calls, queue capacity, connection limits, or external resources used by any of its inner services.
Readiness is therefore a cross-cutting concern. It cannot be managed reliably by looking at each call or service in isolation.
Before dispatching a request, the pipeline waits until the required services are ready. If one of them reaches capacity, processing pauses until that service can accept more work.
This gives us a useful distinction:
ServiceChaindescribes how services are connected.Pipelineowns the runnable service graph and coordinates its lifecycle.
Shared Readiness
Concurrent calls use the same underlying pipeline and therefore share its readiness state. Each active call receives its own pipeline binding, which identifies the call while readiness is being coordinated.
A normal call creates this binding internally:
let response = pipeline.call(request).await?;
Sometimes a call must be stored, moved into another task, or allowed to outlive
the borrow of pipeline. In that case, Pipeline::call_static() returns an
owned future that keeps its binding alive:
let call = pipeline.call_static(request);
let response = call.await?;
Both methods use the same service graph and shared readiness state. The difference is how the lifetime of the call is managed.
This coordination prevents multiple callers from independently driving the same readiness check. It also avoids hiding excess work in unbounded internal queues: when the pipeline is not ready, callers wait.
A successful pipeline readiness check, such as Pipeline::poll_ready(), is
remembered. The next call skips its own readiness check, later calls check
readiness again. A dispatcher that waits for readiness before each call
therefore checks the service only once per request.
Readiness Between Services
Readiness must also be respected while a request moves through the service chain. A service should not invoke the next service directly, because doing so would bypass that service’s readiness check.
Instead, services call one another through an execution context:
let result = ctx.call(&next_service, request).await?;
The context waits for next_service to become ready before invoking it. This
allows backpressure to propagate through the chain, even if readiness changes
while a request is being processed.
In ntex-service, this execution context is represented by Ctx. It carries the
pipeline-level information needed to coordinate otherwise independent service
instances.
A useful mental model is:
Servicedefines what happens at each step.ServiceChaindescribes how the steps are connected.Pipelineowns the runnable graph and coordinates shared readiness.Ctxcontrols how execution moves safely between services.
The Complete Service Trait
Once we include state, readiness, and shutdown, a simplified version of the ntex
Service trait looks like this:
trait Service<St, Req> {
type Res;
type Error;
async fn ready(&self, ctx: Ctx<'_, Self, St>) -> Result<(), Self::Error>;
async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<Self::Res, Self::Error>;
async fn shutdown(&self, ctx: Ctx<'_, Self, St>);
}
The context serves several related purposes:
- It provides access to the pipeline state.
- It connects services to the pipeline that owns them.
- It coordinates readiness between services.
- It allows one service to call another without bypassing backpressure.
- It participates in orderly shutdown.
The important point is that readiness is not merely a method that callers are expected to use correctly. It is built into the way requests move through a pipeline. This makes backpressure explicit, composable, and enforceable without relying on hidden queues or coordination outside the service model.
The St parameter represents the service state made available through the
context. We will explore service and pipeline state in the next section.
Shutdown
The shutdown() method represents the final stage of the service lifecycle:
async fn shutdown(&self, ctx: Ctx<'_, Self, St>);
It is called by the service’s owner—typically a Pipeline—when the service is
being taken out of operation.
A service can use shutdown() to:
- stop background tasks gracefully;
- release external resources;
- flush buffered data; and
- allow in-flight work to finish cleanly.
Simple services may not need to do anything during shutdown. Services that own resources or run background tasks, however, should use this method to clean up properly rather than stopping abruptly.
Composite services should also propagate shutdown to their inner services so that the entire pipeline stops cleanly.
Service State
Services often need access to information that is not part of an individual request. Examples include connection metadata, a database pool, application configuration, or a channel to another subsystem.
The Service trait represents this information with its St type parameter:
trait Service<St, Req> {
type Res;
type Error;
async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<Self::Res, Self::Error>;
async fn ready(&self, ctx: Ctx<'_, Self, St>) -> Result<(), Self::Error>;
async fn shutdown(&self, ctx: Ctx<'_, Self, St>);
}
A service borrows the pipeline state through Ctx instead of owning it. The
service can still own its own fields, such as configuration or inner services.
The pipeline state is available during calls, readiness checks, and shutdown:
async fn call(&self, req: Req, ctx: Ctx<'_, Self, AppState>) -> Result<Self::Res, Self::Error> {
let state: &AppState = ctx.st();
// ...
}
Because the state type is part of Service<St, Req>, Rust verifies that a
service is used with compatible state. A service requiring AppState cannot be
placed directly in a pipeline that supplies an unrelated type.
State-aware functions
An asynchronous function whose first argument is &St can be converted into a
state-aware service. This is the simplest way to define one:
use std::convert::Infallible;
use ntex::Pipeline;
struct AppState {
prefix: &'static str,
}
async fn format_value(state: &AppState, value: usize) -> Result<String, Infallible> {
Ok(format!("{}-{value}", state.prefix))
}
#[ntex::main]
async fn main() {
let service = Pipeline::new(AppState { prefix: "item" }, format_value);
assert_eq!(service.call(10).await.unwrap(), "item-10");
}
The function receives &AppState for each call. The IntoService
implementation performs this conversion automatically. You can call
fn_service_st explicitly when the service value must be named or passed
through an API before creating the pipeline.
Testing state-aware function services
A state-aware function service can be tested without starting a server. Create
the service with fn_service_st, place it and the test state in a
Pipeline::new, then call it with the request value. This exercises the same
state access and readiness handling used in production.
When the state contains external dependencies, a mocking library can verify how the service uses them:
test-mockalldefines the state behavior as a trait. Mockall creates a mock state and checks calls to its accessor and operation methods.test-shimforgekeeps the concrete state type. Shimforge replaces its methods during the test, including a method on an object nested in the state, and checks the receiver’s attributes and call arguments.
Both examples configure call counts and ordering, invoke the service through a pipeline, and assert its response.
Pipeline-owned state
Pipeline::new stores one service and one state value together. Every call,
readiness check, and shutdown operation on that pipeline uses the same state
instance.
The state belongs to that pipeline rather than to the process. Creating another
pipeline creates another state instance unless the application explicitly
shares a resource. For example, several worker-local states can each hold a
clone of the same Arc<Pool>.
Pipeline bindings share the original pipeline and its state; cloning a
PipelineBinding does not clone the state value. Bindings also participate
in the pipeline’s coordinated readiness and shutdown handling.
Calling nested services
Ctx does more than provide ctx.st(). It connects a service call to the
pipeline’s readiness and shutdown machinery. A service that wraps another
service should call it through the context:
async fn call(&self, req: Req, ctx: Ctx<'_, Self, St>) -> Result<Self::Res, Self::Error> {
ctx.call(&self.inner, req).await
}
Ctx::call waits for the inner service to become ready before calling it.
Ctx::call_nowait skips that check and should only be used when readiness has
already been established. Ctx also provides corresponding ready() and
shutdown() operations for implementing service combinators and middleware.
A wrapping service should not call an inner service’s trait methods directly.
Use Ctx so the pipeline can supply the required context and coordinate the
service chain.
Substituting state
Sometimes one service in a larger chain needs a different state type.
map_state wraps that service with a fixed state value while allowing the
outer pipeline to use another state type:
use std::convert::Infallible;
use ntex::{Pipeline, service::map_state};
struct Metrics {
prefix: &'static str,
}
async fn record(metrics: &Metrics, value: usize) -> Result<String, Infallible> {
Ok(format!("{}-{value}", metrics.prefix))
}
#[ntex::main]
async fn main() {
let service = map_state(Metrics { prefix: "metric" }, record);
// The wrapped service uses Metrics even though the outer pipeline uses ().
let pipeline = Pipeline::new((), service);
assert_eq!(pipeline.call(10).await.unwrap(), "metric-10");
}
State supplied per operation
PipelineState stores a service without owning its state. The caller supplies
a state reference for each readiness check, call, or shutdown operation:
use std::convert::Infallible;
use ntex::service::pipeline::PipelineState;
struct RequestContext {
prefix: &'static str,
}
async fn format_value(state: &RequestContext, value: usize) -> Result<String, Infallible> {
Ok(format!("{}-{value}", state.prefix))
}
#[ntex::main]
async fn main() {
let service = PipelineState::new(format_value);
let first = RequestContext { prefix: "first" };
let second = RequestContext { prefix: "second" };
assert_eq!(service.call(1, &first).await.unwrap(), "first-1");
assert_eq!(service.call(2, &second).await.unwrap(), "second-2");
}
Use Pipeline when one state value should remain attached to the service. Use
PipelineState when the same service and readiness machinery must operate with
a state selected by each caller. PipelineState::bind_state() can attach an
owned state value and produce a normal PipelineBinding, the state type must
implement Clone.
Carrying state with a request
Some protocol boundaries receive both a request and the state that a nested
pipeline should use. The RequestState trait represents such an input.
State<St, Req> and (St, Req) implement it by splitting the value into a
state and a request:
use ntex::service::{RequestState, State};
let input = State {
state: "connection-1",
req: "request",
};
let (connection, request) = input.unpack();
assert_eq!(connection, "connection-1");
assert_eq!(request, "request");
For example, ntex HTTP services use RequestState to extract state associated
with an accepted connection and then provide that state to the HTTP request and
control-service pipelines. This keeps connection setup state separate from the
HTTP request value while preserving its concrete type.
A plain Io also implements RequestState with () as its state, so an
ordinary server can pass accepted connections to HttpService without
wrapping them. See
Passing Connection State to the HTTP Service
for an example that wraps an accepted Io in State.
Runtime
ntex uses a single-threaded execution model on each runtime thread. Tasks stay
on the thread where they were started, so services and futures can use Rc,
Cell, RefCell, and other types that are not Send or Sync.
That does not limit an application to one CPU core. A server normally starts
several worker threads, each with its own single-threaded runtime. State shared
between workers must still use thread-safe types such as Arc, atomics, or
locks.
Most applications only need #[ntex::main], #[ntex::test], and
ntex::rt::spawn. The lower-level System and Arbiter APIs are useful for
explicit startup, shutdown, and cross-thread coordination.
Choosing a runtime backend
DefaultRuntime selects the async runtime and matching ntex network reactor
together. Cargo features choose the backend:
tokiouses Tokio. When a Tokio runtime handle is entered on the current thread, the adapter uses it; otherwise it creates a current-thread Tokio runtime. ntex tasks run in a TokioLocalSet.compiocreates a Compio runtime on each runtime thread. Compio selects the I/O driver for the host platform and configuration.- without either feature, ntex uses its native runtime. On Linux it tries io_uring and falls back to polling; other Unix platforms use polling, and Windows uses IOCP.
If both tokio and compio are enabled, Tokio takes precedence. With the
native backend, neon-polling forces polling on Unix and neon-uring requires
io_uring on Linux. Do not enable both explicit native reactor features.
For example:
[dependencies]
ntex = { version = "4", features = ["tokio"] }
The Tokio and Compio backends are convenient when an application also uses libraries from those ecosystems. The native backend avoids an additional general-purpose runtime and lets ntex choose its platform reactor directly.
SystemRunner::block_on blocks the current thread. Do not call it from inside
an async Tokio task to try to nest one runtime inside another. Usually the
cleanest boundary is to let #[ntex::main] own the process entry point.
Spawning tasks
ntex::rt::spawn starts a future on the current runtime thread:
#[ntex::main]
async fn main() {
let task = ntex::rt::spawn(async {
// `Rc`, `RefCell`, and other `!Send` values may be used here.
10usize
});
assert_eq!(task.await.unwrap(), 10);
}
The future and its result do not need to implement Send, but the future must
be 'static. Use async move when it captures owned values. spawn() panics
when called outside an active ntex runtime.
The returned JoinHandle can be awaited, canceled with cancel(), detached
with detach(), or queried with is_finished(). A task panic or cancellation
is reported as JoinError, not as the task’s normal output.
Dropping a local join handle detaches the task; it does not cancel it. Keep and await the handle when the result matters. A detached task is not guaranteed to finish before its arbiter or the whole system shuts down.
To submit work to another runtime thread, obtain an arbiter’s handle and call
spawn() on it. The future and its output must both be Send + 'static
because they cross thread boundaries. Once the future is running on that
arbiter, it can start more non-Send work with ntex::rt::spawn.
Do not use a remote join handle as a portable abort mechanism. With the Tokio and Compio backends, canceling work sent through another arbiter abandons the result but does not stop the remote task. Use an explicit cancellation signal when cross-thread work must be stoppable.
Timers and timeouts
Use ntex::time instead of a backend-specific timer so the code works with
all three runtime backends:
use ntex::time::{Millis, sleep, timeout};
#[ntex::main]
async fn main() {
sleep(Millis(10)).await;
let result = timeout(Millis(100), async { "done" }).await;
assert_eq!(result.unwrap(), "done");
}
ntex timers are intended for scheduling and timeouts, not high-resolution
measurement. Their granularity is roughly 16 milliseconds. A zero-duration
sleep or timeout still waits for at least one timer tick.
timeout_checked is the variant that treats a zero timeout as disabled.
System lifecycle
System is the top-level runtime context. It owns the shared runtime
configuration, tracks arbiters, manages signal delivery, and provides the
blocking thread pool. #[ntex::main] builds a system automatically.
Building one manually returns a SystemRunner:
fn main() {
let result = ntex::rt::System::build()
.name("worker")
.build(ntex::rt::DefaultRuntime)
.block_on(async { 10usize });
assert_eq!(result, 10);
}
SystemRunner has two common ownership models:
block_on(future)starts a system, drives one root future, and returns its output when that future completes;run(callback)invokes a synchronous startup callback inside the running system, then keeps the event loop alive untilSystem::stop()orSystem::stop_with_code()is called.
For a process driven by an explicit stop request:
use std::io;
use ntex::rt::{self, DefaultRuntime, System};
use ntex::time::{Millis, sleep};
fn main() -> io::Result<()> {
System::build()
.name("my-app")
.build(DefaultRuntime)
.run(|| {
rt::spawn(async {
sleep(Millis(10)).await;
System::current().stop();
});
Ok(())
})
}
run_until_stop() is run() without a startup callback. All these methods
consume the runner, so choose one model rather than calling block_on() and
then run_until_stop() on the same value. A non-zero stop code is returned as
an I/O error.
Signal and panic handling are disabled by default. Enabling signal handling
makes process signals available through ntex::rt::signals; it does not by
itself define the application’s shutdown policy. ntex servers listen for
process signals while they are running and coordinate their own shutdown.
Arbiters and runtime threads
An Arbiter represents one runtime event-loop thread. Every system has a
primary arbiter. Arbiter::new() or Arbiter::with_name() starts another OS
thread using the same system runtime configuration. ntex server workers are
also arbiter threads.
Arbiter::current() returns the arbiter for the current thread and panics when
called outside one. For an arbiter created with new() or with_name(),
stop() requests shutdown and join() waits for its thread to exit.
The primary arbiter runs on the thread that started the system. It cannot be
stopped independently, has no owned thread handle, and join() returns
immediately. Stop the System to end the primary arbiter.
The two levels have different shutdown scopes:
Arbiter::stop()stops one separately created runtime thread;System::stop()stops every registered arbiter and ends the system.
The runtime also offers two typed storage scopes. System::get_value() stores
a thread-safe value shared by the complete system, while
Arbiter::get_value() stores a cloneable value local to one runtime thread.
The latter is useful for per-worker clients or caches that should not be
shared.
Blocking work
Blocking an arbiter thread pauses connection handling, timers, and every other
task on that thread. Move CPU-heavy work and blocking system calls to
spawn_blocking():
#[ntex::main]
async fn main() {
let total = ntex::rt::spawn_blocking(|| (0..1_000u64).sum::<u64>())
.await
.unwrap();
assert_eq!(total, 499_500);
}
The closure and its result must implement Send. Inside a system, the work
runs on a dynamically sized blocking pool. The pool allows up to 256 workers
by default, and an idle worker exits after 60 seconds. Configure those values
with thread_pool_limit() and thread_pool_recv_timeout():
use std::time::Duration;
fn main() {
ntex::rt::System::build()
.thread_pool_limit(32)
.thread_pool_recv_timeout(Duration::from_secs(30))
.build(ntex::rt::DefaultRuntime)
.block_on(async {});
}
If spawn_blocking() is called outside a running system, it does not
create a background worker. The closure runs immediately on the calling thread
and the returned future is already ready. Treat that as a fallback, not as a
way to start asynchronous work before the runtime.
Dropping the returned future prevents queued work from starting, but cannot
interrupt a closure that is already running. Use detach() when queued work
should continue even if nobody needs its result. Cancellation and closure
panics are reported as BlockingError.
Runtime diagnostics
The system pings spawned arbiters every two seconds and keeps their ten most
recent round-trip records. System::list_arbiter_pings() exposes those
records. Set ping_interval(0) to disable the checks, or tune
ping_interval() and ping_threshold() on System::build().
On Linux, System::set_latency_callback() can receive a backtrace when an
arbiter misses the configured threshold, which defaults to one second.
Backtrace capture requires ntex signal handling and uses SIGUSR2; an
embedding application should not reserve that signal for another purpose.
Custom runners
System::build().build(...) accepts any implementation of ntex_rt::Runner.
This is an advanced integration point, not merely an executor switch: ntex
networking also expects a compatible reactor on every runtime thread. A custom
runner must establish everything expected by the selected backend.
Use DefaultRuntime unless the application deliberately provides both sides
of that integration.
Runtime attributes
#[ntex::main] turns an async function into a synchronous entry point, builds
a System, and drives the function’s future:
#[ntex::main]
async fn main() {
// Start the application.
}
Supported options are:
name = "..."for the system name;signals = trueorfalse;panic_handling = trueorfalse;ping_interval = Nin milliseconds, where zero disables arbiter pings;rt = Typeto select a custom runtime runner.
For example:
#[ntex::main(
name = "my-service",
signals = true,
panic_handling = true,
ping_interval = 250,
)]
async fn main() {
// Start the application.
}
The async function may return a value such as Result<(), E>; the generated
synchronous function returns the root future’s output.
#[ntex::test] creates a fresh system named after the test, marks it as a
testing system, disables signal and panic handling, and initializes ntex test
logging:
#[ntex::test]
async fn runtime_test() {
let value = ntex::rt::spawn(async { 10usize }).await.unwrap();
assert_eq!(value, 10);
}
Test logging initialization is skipped when the no-test-logging feature is
enabled.
The useful boundaries to remember are:
- local tasks may be non-
Send, but cross-arbiter work must beSend; - dropping a join handle does not stop its task;
- blocking closures belong in
spawn_blocking; - ntex timers keep code independent of the selected backend;
- stop a separately created
Arbiterfor one runtime thread and theSystemfor the whole runtime; - use
DefaultRuntimeunless the application also provides the matching ntex network reactor.
Server
ntex servers use several single-threaded workers to make use of multiple CPU cores. A dedicated accept thread owns the listening sockets and sends each connection to an available worker:
listeners -> accept thread -> worker-local service
The server does not know whether a connection speaks HTTP, MQTT, or a custom
protocol. It manages sockets and workers; the service created for each worker
decides how to handle the resulting Io.
Starting a server
Create a server inside an active ntex System, register at least one listener,
and call run():
use ntex::http::{HttpService, Response};
#[ntex::main]
async fn main() -> std::io::Result<()> {
ntex::server::build()
.bind(
"http",
"127.0.0.1:8080",
ntex::SharedCfg::default(),
async |_| {
HttpService::new(async |_| {
Ok::<_, std::io::Error>(
Response::Ok().body("Hello world!"),
)
})
},
)?
.run()
.await
}
server::build() needs a current System, which #[ntex::main] creates in
this example. run() panics if no listener has been registered.
bind() receives:
- a service name used in logs and to associate listeners with services;
- an address that implements
ToSocketAddrs; - a
SharedCfgattached to accepted connections; - an asynchronous factory that builds the connection service for each worker.
run() returns a cloneable Server controller. Awaiting it waits for the
server to stop. The future always resolves to Ok(()); it is a completion
notification, not a health report for every connection or worker.
Binding and existing listeners
bind() creates and owns the listening socket. If an address resolves to
several socket addresses, ntex tries all of them and registers every listener
that binds successfully. The call succeeds when at least one address binds.
Use a concrete SocketAddr when the application needs exactly one listener.
Use listen() when the application already owns a TcpListener, for example
with socket activation or custom socket options. The builder’s backlog()
setting applies only to listeners created by later bind() and configure()
calls; it cannot change the backlog of an existing listener.
On Unix, bind_uds() and listen_uds() provide the same choices for Unix
domain sockets. bind_uds() removes an existing file at the socket path before
binding and removes the socket file again when the server stops.
A server may register several listeners and several protocols. Keep service
names clear and unique when using modular configure() callbacks, because
those callbacks attach worker services to listeners by name.
Workers and service factories
By default, the server starts one worker for each available logical CPU. Set a
different number with workers():
let builder = ntex::server::build()
.workers(4);
Use at least one worker. With no workers, the server has nowhere to dispatch connections and remains paused.
Each worker runs on its own arbiter thread and owns its own service instance.
The factory passed to bind() is called for every worker. It may be called
again if a worker restarts or a service is recreated after a readiness
failure, so factory setup should be safe to repeat.
The factory must be Send, Clone, and 'static, but the service it creates
stays on one worker and does not need to implement Send or Sync. This is
why a worker-local service may use Rc, Cell, or RefCell.
Data shared across workers must be shared explicitly with thread-safe types
such as Arc, atomics, or locks. For richer process-wide and worker-local
state, use build_with_config() and ServerAppConfig; the next chapter
explains that model in detail.
Connections are distributed across workers that currently report ready. If a worker reaches its connection limit, or one of its registered services is not ready, that worker temporarily stops receiving new connections. If no worker is available, the accept loop pauses until one becomes ready again.
The default limit is 25,600 concurrent connections per worker:
let builder = ntex::server::build()
.workers(4)
.max_connections(10_000);
This limit is process-wide and shared by every ntex server. Configure it before starting servers; each worker takes the current value when its services are created.
Server configuration
The commonly useful ServerBuilder settings are:
name()sets the accept and worker thread-name prefix. It defaults to the current system name.workers()sets the worker count.backlog()sets the listen backlog for listeners created afterward.max_connections()sets the process-wide per-worker connection limit.enable_affinity()pins workers to CPU cores when the platform exposes suitable core IDs.graceful_shutdown_timeout()bounds how long a graceful stop waits for workers. The default is 30 seconds.stop_on_panic()stops the server instead of restarting a worker that fails.graceful_shutdown()makesstop_on_panicworker failures,SIGQUIT, fatal signals, and application panics use graceful rather than immediate server shutdown.disable_signals()disables the server’s built-in signal handling.stop_runtime()stops the complete ntexSystemafter this server stops.
For example:
use ntex::time::Seconds;
let builder = ntex::server::build()
.name("api")
.workers(4)
.backlog(1024)
.max_connections(20_000)
.graceful_shutdown_timeout(Seconds(15))
.stop_on_panic();
stop_runtime() is useful when one server owns the complete process
lifecycle. Do not enable it when other independent work or servers must keep
using the same System.
status_handler() can connect listener readiness to a supervisor or metrics
system. It runs on the accept thread and reports when listeners pause or
resume. The status describes acceptance state, not whether every worker is
healthy, and the same status may be reported more than once. Keep the handler
quick and non-blocking.
Controlling a running server
Clone the controller when another task needs to pause, resume, or stop the server:
use ntex::server::Server;
async fn manage(server: Server) -> std::io::Result<()> {
let control = server.clone();
ntex::rt::spawn(async move {
control.pause().await;
// Perform maintenance or wait for an external readiness condition.
control.resume().await;
control.stop(true).await;
});
server.await
}
The server starts in a paused state and resumes once its first worker becomes
ready. pause() stops accepting new connections without closing existing
ones; new connection attempts normally wait in the operating system’s listen
backlog. resume() enables the listeners again.
Use stop(true) for a graceful stop. The accept loop closes its listeners,
then the server asks workers currently available for dispatch to finish active
work and waits up to the configured graceful-shutdown timeout. A worker still
initializing or recreating its service is not part of that wait, so long-lived
factory setup should have its own cancellation or ownership boundary.
Use stop(false) when the controller should not wait for graceful worker
draining. Available worker services are still asked to shut down and get a
default three-second local timeout, but that cleanup may continue after the
server controller reports completion. An immediate stop is therefore not a
guarantee that every in-flight operation has finished.
Repeated stop requests join the stop already in progress. Await either
stop(...) or the original Server when later work must begin only after the
server has reached its stopped state.
Signals and shutdown
Signal handling is enabled by default:
| Event | Default behavior |
|---|---|
SIGTERM | graceful stop |
SIGINT / Ctrl-C | immediate stop |
SIGQUIT | immediate stop, or graceful with graceful_shutdown() |
SIGHUP | ignored by the server |
| fatal signal or application panic | immediate stop, or graceful with graceful_shutdown() |
Fatal-signal notifications are installed with Unix signal handling. Application-panic notifications additionally require runtime panic handling. A fatal-signal shutdown is best effort because the operating system may terminate the process before graceful cleanup finishes. On Windows, Ctrl-C is delivered as the interrupt event.
Signal delivery belongs to the ntex System, not to one isolated server. If
several servers share a system, or the application installs its own signal
policy, call disable_signals() consistently and stop the appropriate server
controllers yourself.
The graceful server timeout and connection shutdown timeout solve different
problems. The server timeout bounds a worker as a whole. Each connection uses
its own IoConfig::set_shutdown_timeout() while flushing and closing. Leave
enough server-level time for those connection shutdowns to complete.
Connection and protocol configuration
The SharedCfg passed to bind() or listen() is cloned onto each accepted
Io. The I/O layer and protocol services retrieve the configuration types
they understand from the same container:
use ntex::http::{HttpService, HttpServiceConfig, Response};
use ntex::io::IoConfig;
use ntex::time::Seconds;
use ntex::SharedCfg;
#[ntex::main]
async fn main() -> std::io::Result<()> {
let cfg = SharedCfg::new("HTTP")
.add(
IoConfig::new()
.set_shutdown_timeout(Seconds(2)),
)
.add(
HttpServiceConfig::new()
.set_keepalive_timeout(Seconds(30)),
);
ntex::server::build()
.bind("http", "127.0.0.1:8080", cfg, async |_| {
HttpService::new(async |_| {
Ok::<_, std::io::Error>(
Response::Ok().body("Hello world!"),
)
})
})?
.run()
.await
}
IoConfig controls connection-level behavior such as shutdown, memory, and
read-rate limits. HttpServiceConfig controls HTTP behavior such as
keep-alive, request-head limits, and protocol timeouts. TLS and other protocols
add their own configuration types to the same SharedCfg.
The I/O chapter explains buffer limits, read rates, write backpressure, keep-alive, and shutdown timing in detail.
Test servers
test_server() starts one worker on a separate thread, binds an available
local port, and disables signal handling:
use ntex::http::{HttpService, Response};
use ntex::{client::Client, server};
#[ntex::test]
async fn test_server_response() {
let server = server::test_server(async || {
HttpService::new(async |_| {
Ok::<_, std::io::Error>(
Response::Ok().body("Hello world!"),
)
})
});
let url = format!("http://{}/", server.addr());
let response = Client::new()
.get(url.as_str())
.send()
.await
.unwrap();
assert!(response.status().is_success());
}
Use addr() for HTTP clients or connect() to open a client Io with the
test server’s client configuration. server() returns the underlying server
controller.
TestServerBuilder can set separate server and client SharedCfg values.
build_test_server() exposes the full ServerBuilder for custom listeners;
when using it, set the address on the returned TestServer before calling
connect().
Dropping the last TestServer clone stops its server and runtime. That drop
briefly blocks the current thread while the test thread shuts down, so drop it
outside timing-sensitive assertions.
Worker and request state
Each ntex worker runs its own single-threaded runtime. To make use of multiple CPU cores, the server starts several workers and distributes incoming connections among them. Each worker then creates and runs its own application instance.
Starting an application usually happens in two stages.
First, the server loads configuration that is shared across the entire process. This normally happens on the main thread and may include reading environment variables, parsing command-line arguments, loading configuration files, or fetching settings from an external service.
Next, ntex uses that configuration to initialize each worker. The
ServerAppConfig trait controls this part of the process. You provide the
server with an object that implements ServerAppConfig, and the server calls
its create() method inside each worker every time that worker starts. If a
worker fails and is restarted, create() runs again and the restarted worker
gets fresh state.
The configuration object must implement Send + Sync because the server shares
it across worker threads. The state returned by create() is different: it
belongs to one worker and stays on that worker’s thread. As a result, worker
state can contain single-threaded types such as Rc and RefCell.
Worker state must also implement Clone. The server clones it for every
listener, and servers such as web::server_with_config() clone it again for
every connection. Put data that should be shared by all connections of a worker
behind a shared handle, such as Rc<RefCell<T>> or Rc<Cache>, so that
the clones refer to the same value. Plain fields are copied, and changes to
them are not visible to other clones.
The create() method is asynchronous, so it can do more than simply construct
a struct. It can open database connections, create client instances, initialize
caches, or prepare any other resources the worker needs. If something goes
wrong, it can return an error rather than starting the worker with invalid
state.
struct AppBuilder {
// Configuration shared across the process
}
#[derive(Clone)]
struct AppState {
// Resources owned by one worker, usually behind `Rc`
}
impl ServerAppConfig for AppBuilder {
type State = AppState;
// Called each time a worker starts.
async fn create(&self) -> io::Result<Self::State> {
Ok(AppState {
// Initialize this worker's resources.
})
}
}
A separate type is not always necessary. ServerAppConfig is also implemented
for async closures that return io::Result<State>, and NoConfig provides
unit state for servers that do not need any:
let server = ntex::server::build_with_config(async || {
Ok(AppState {
// ...
})
});
Creating the Worker Application
Once the worker state is ready, the server calls an application factory to
build the worker’s application instance. The factory runs once per worker and
receives a reference to the state created by ServerAppConfig::create().
#[ntex::main]
async fn main() -> std::io::Result<()> {
// Load the process-wide configuration.
let builder = AppBuilder {
// ...
};
web::server_with_config(builder, async |_state: &AppState| {
web::App::new().service(
web::resource("/").to(async || {
web::HttpResponse::Ok()
}),
)
})
.bind("127.0.0.1:8080", SharedCfg::default())?
.run()
.await
}
Here, AppBuilder holds the process-wide configuration and implements
ServerAppConfig. For every worker, ntex calls AppBuilder::create() to
produce a new AppState. It then passes that state to the application factory,
which builds the worker’s web::App factory. The App service itself is
created for each connection, with its own clone of the worker state.
The factory can use the state while setting up the application and its services. The same state is also available to the worker’s connection-handler pipelines throughout their lifecycle.
Worker-Local and Process-Wide State
A worker creates its state once and reuses it for every connection it handles.
It does not call create() for every connection. Connections receive clones of
the worker state, so values behind Rc are shared by every connection of that
worker.
This means that connections handled by the same worker see the same shared data. Connections handled by different workers, however, use different state instances. Updating local state in one worker does not automatically update the state in another worker.
This separation is often useful. Because worker-local state never crosses thread boundaries, it can use simple single-threaded data structures without paying the cost of synchronization.
Sometimes state really does need to be shared across all workers. In that case,
it must be shared explicitly. One common approach is to wrap the shared value in
an Arc and clone it into each worker state. If the value is mutable, it will
also need a suitable synchronization mechanism, such as an atomic type,
Mutex, or RwLock.
The best choice depends on how the data is used. Prefer worker-local state when workers do not need to coordinate. Use process-wide shared state only when changes must be visible across workers.
Complete process/worker configuration example
Service and Pipeline State
ntex-service 5 adds a state parameter to the Service trait. This gives every
service access to state while it is checking readiness, handling requests, or
shutting down.
A simplified version of the trait looks like this:
trait Service<St, Req> {
type Res;
type Error;
async fn call(
&self,
req: Req,
ctx: Ctx<'_, Self, St>,
) -> Result<Self::Res, Self::Error>;
async fn ready(
&self,
ctx: Ctx<'_, Self, St>,
) -> Result<(), Self::Error>;
async fn shutdown(&self, ctx: Ctx<'_, Self, St>);
}
Here, St is the type of state the service expects. Since it is part of the
Service definition, Rust can check that a service is always used with the
right kind of state.
The same state is available in all three lifecycle methods:
ready()can use it while checking whether the service is ready for more work.call()can use it while handling a request.shutdown()can use it while the service is stopping.
You can think of a stateful service as an asynchronous function that receives both the state and the request:
async fn my_service(
state: &St,
req: Req,
) -> Result<Res, Error> {
// ...
}
The Service trait does not pass state as a separate argument. Instead, it
provides state through Ctx. The context also connects the service to its
pipeline, which manages readiness, request processing, and shutdown.
Constructing a Pipeline
To call a service, you first place it in a Pipeline. A pipeline keeps the
service and its state together. It can contain a single service or a chain of
services, all of which can use the same state.
Here is a small example:
async fn my_service(
state: &AppState,
req: usize,
) -> Result<String, io::Error> {
// Use the pipeline state to handle the request.
todo!()
}
#[ntex::main]
async fn main() -> io::Result<()> {
// Create a pipeline containing one state instance and one service.
let svc = ntex::Pipeline::new(AppState::new(), my_service);
// Both calls use the same AppState instance.
let first = svc.call(1).await?;
let second = svc.call(2).await?;
println!("{first}:{second}");
Ok(())
}
The AppState is created when the pipeline is built. It stays with the pipeline
and is reused for every call. Calling svc.call() does not create a new state
instance.
If the pipeline contains a chain of services, each compatible service can access the same state. This makes it easy to share resources across a service chain without adding a copy of those resources to every service.
Worker state
Earlier, we saw that ntex creates a separate state instance for each server worker. That state also becomes the state of the worker’s connection-handler pipeline.
When a worker starts, the server calls the application factory and passes it a reference to the worker state. The factory builds a connection-handler service, and the server places that service in a pipeline together with the worker state.
/// Handles an I/O connection using worker-local state.
async fn handle_io(
state: &AppState,
io: ntex::io::Io,
) -> Result<(), io::Error> {
// Process the connection using the worker state.
Ok(())
}
#[ntex::main]
async fn main() -> std::io::Result<()> {
// Load the process-wide configuration.
let builder = AppBuilder {
// ...
};
ntex::server::build_with_config(builder)
.bind(
"handler",
"127.0.0.1:8080",
SharedCfg::new("S"),
async |_state: &AppState| {
// Build the connection-handler service.
ntex::service(handle_io)
},
)?
.run()
.await
}
The factory is called once for each worker and receives that worker’s
AppState. It can use the state to configure the service during construction.
The server places the returned service and a clone of the worker state in a
pipeline. When
a connection is dispatched to handle_io service, the pipeline passes access to the
same AppState instance alongside the connection.
Each worker has its own state. Connections handled by the same worker share the data behind the state’s shared handles, but connections handled by different workers do not automatically share state.
State Accumulation
The state we have discussed so far is long-lived. Worker state lives for as long as the worker, and pipeline state lives for as long as the pipeline. In both cases, the state is created once and reused across many calls.
Sometimes we need state with a shorter lifetime. In particular, we may want to collect information while a request or connection moves through a service chain. Each service can inspect what has already been collected, add more information, and pass the updated state to the next service.
Imagine a pipeline that accepts an incoming connection. The first service might
record basic details such as the connection ID and peer address. The next
service could read the TLS ClientHello and extract the Server Name Indication
(SNI). Other services might load configuration for that server name, validate
the client, calculate resource usage, apply throttling, and finally perform the
TLS handshake before handing the connection to an HTTP server.
The pipeline might look like this:
accept connection
→ collect connection information
→ read SNI
→ load client information
→ validate and throttle
→ negotiate TLS
→ handle HTTP
ntex-service provides the RequestState trait and the State type for this
kind of state. Together, they let us pass a message and its accumulated state
through a service chain.
State and types that implement RequestState are different from pipeline state.
Pipeline state is shared across calls, while accumulated state belongs to
a single request. Every connection moving through the pipeline carries
its own state.
State<St, Req> is the most common implementation, but it is not the only
one. A (St, Req) tuple also implements RequestState, which is convenient
for small services. Io and IoBoxed implement it with () state, so a plain
connection can be passed to any service that expects a RequestState.
The ntex protocol servers support RequestState, including ntex::http,
ntex-h2, ntex-mqtt, and ntex-amqp.
Collecting Connection Information
Let’s start with a service that receives a raw I/O connection and collects some basic information about it:
use std::{io, net::SocketAddr, time::Instant};
use ntex::{io::Io, service::State};
use uuid::Uuid;
struct Connection {
id: Uuid,
created: Instant,
peer_addr: SocketAddr,
}
async fn connect(_: &AppState, io: Io) -> io::Result<State<Connection, Io>> {
let peer_addr = load_peer_addr(&io)?;
Ok(State {
req: io,
state: Connection {
id: Uuid::new_v4(),
created: Instant::now(),
peer_addr,
},
})
}
The next service can unpack these values, perform the TLS handshake, and add more information:
use ntex::service::{RequestState, State};
struct ConnectionWithTls {
id: Uuid,
created: Instant,
peer_addr: SocketAddr,
server_name: String,
}
async fn tls(_: &AppState, msg: State<Connection, Io>) -> io::Result<State<ConnectionWithTls, TlsIo>> {
let (connection, io) = msg.unpack();
let server_name = load_server_name(&io).await?;
let io = accept_tls(io, &server_name).await?;
Ok(State {
req: io,
state: ConnectionWithTls {
id: connection.id,
created: connection.created,
peer_addr: connection.peer_addr,
server_name,
},
})
}
This service changes both parts of the input. It turns the raw Io connection
into a TLS-enabled TlsIo, and it extends the connection state with the server
name extracted during TLS negotiation.
The connection ID, creation time, and peer address are carried forward. By the time the connection reaches the HTTP server, its state contains everything collected by the earlier services.
Passing Connection State to the HTTP Service
We can now connect these services into a single chain:
use ntex::{server, service, SharedCfg};
use ntex::http::{self, Request, Response};
/// Handles an HTTP request using information about its connection.
async fn handle_request(st: &ConnectionWithTls, _req: Request) -> io::Result<Response> {
Ok(Response::Ok()
.body(format!("server name: {}", st.server_name)),
)
}
#[ntex::main]
async fn main() -> io::Result<()> {
let builder = AppBuilder {
// ...
};
server::build_with_config(builder)
.bind(
"HTTP",
"127.0.0.1:8080",
SharedCfg::new("S"),
async |_: &AppState| {
// Collect the initial connection information.
service(connect)
// Perform the TLS handshake.
.and_then(tls)
// Pass the established connection to the HTTP service.
.and_then(http::HttpService::new(handle_request))
},
)?
.run()
.await
}
The worker’s AppState is still available to connect and tls. It contains
worker-level resources that can be reused across many connections. At the same
time, each connection carries its own state.
Complete state accumulation example
Implementation Details
ntex supports several network protocols, and each one handles state a little differently. The important distinction is between state that belongs to a connection and state that belongs to a single request.
ntex::http
HttpService expects a handler with roughly the following shape:
async fn handler(state: &St, req: http::Request) -> Result<http::Response, Err> {
// ...
}
Here, state is the state of the HTTP connection. It is extracted from the
value passed to HttpService through the RequestState trait.
For example, an earlier service can return a State<ConnectionWithTls, Io>.
When HttpService receives that value, it separates the connection state
from the I/O stream. The I/O stream is used to run the HTTP protocol,
while ConnectionWithTls becomes the state provided to the HTTP request handler.
This is the mechanism used in the previous section to make connection-specific information, such as the peer address and TLS server name, available while handling HTTP requests.
ntex::web
A web::App can act as the handler for HttpService. When it does, the state
extracted by HttpService becomes the application’s state.
A web service looks roughly like this:
async fn handler(state: &St, req: web::WebRequest) -> Result<web::WebResponse, Err> {
// ...
}
The state argument contains the long-lived application or connection state.
The request can also carry its own state through the generic
parameter of WebRequest:
web::WebRequest<RequestState>
This gives a web application two separate kinds of state:
Stis provided by the service pipeline and remains available across requests.RequestStatebelongs to one request and moves through the request-processing chain.
A filter or middleware can inspect a WebRequest, perform some work, and
return a new WebRequest with a different state type. This makes it possible
to build up request-specific information as the request moves through
the application.
For example, a request-processing chain might look like this:
validate request
→ load and authenticate the user
→ apply throttling
→ call the handler
The request starts without any additional state:
web::WebRequest<()>
After the authentication service loads and verifies the user, it can return:
web::WebRequest<AuthenticatedUser>
The next service can then require WebRequest<AuthenticatedUser>. This means it
cannot be called until authentication has completed and the request contains an
authenticated user.
Another service could validate the request and return a different state type:
web::WebRequest<ValidatedRequest>
In this way, the request type shows how far the request has moved through the processing chain. Each service declares the state it expects and the state it produces, and Rust checks that the services are connected in the right order.
Keeping pipeline state and request state separate lets the application reuse connection-level resources while building up request-specific information as each request moves through the service chain.
I/O Abstraction Layer
ntex provides an I/O abstraction layer that keeps protocol and service code independent of a specific runtime or reactor implementation, such as Tokio, Compio, or Neon.
In addition to presenting a consistent interface across these backends, the I/O layer provides the building blocks needed to manage buffered reads and writes, backpressure, timeouts, and graceful shutdown.
Socket abstractions
ntex achieves independence from the underlying socket implementation through
dependency inversion. The I/O subsystem does not call Tokio, Compio, or Neon
socket APIs directly. Instead, a runtime adapter moves bytes between its socket
and the buffers owned by Io. One way to picture it: the read half carries
bytes from the transport to the application and the write half carries them
back, with Io sitting between the two. Each half is an in-memory buffer that
Io owns, which is what lets it apply backpressure and run filter
transformations on bytes as they pass through.
An active connection has three cooperating execution responsibilities:
- A protocol or dispatcher task consumes decoded input and queues encoded
output through
IoorIoRef. This task typically runs the protocol service and, indirectly, application code. - A transport read task waits for read readiness, reads bytes from the socket, and submits them to the I/O subsystem.
- A transport write task takes queued bytes from the I/O subsystem and writes them to the socket.
The runtime adapter decides how to schedule these responsibilities. The
built-in adapters typically use cooperating read and write tasks, but this is
an implementation detail rather than part of the IoStream contract.
The read task cooperates with ntex backpressure. It reads only while
IoContext::poll_read_ready permits more input. It takes a buffer through
IoContext::take_read_buf and releases it with the result through
IoContext::release_read_buf. This allows the I/O
subsystem to process filters, wake the dispatcher, and pause further reads when
the configured high-water mark is reached.
A backend that reads synchronously, without holding the buffer across an await,
should use IoContext::with_read_buf instead. It combines the two steps and
borrows the read buffer in place, so no buffer is detached and appended back.
The write task follows the same pattern. It waits for
IoContext::poll_write_ready, obtains queued data with
IoContext::with_write_dst, and reports through
IoContext::update_write_status how many bytes reached the peer. ntex can
then apply write backpressure, resume waiting services, and coordinate
graceful shutdown.
A transport does not have to write the bytes before it reports. Completion based backends take ownership of whole pages and keep them until the operation completes. Those bytes leave the write buffer but have not reached the peer, so ntex counts them as outstanding output until the transport either reports them as written or returns them to the buffer. Flush completion, write backpressure, and the shutdown drain all account for them.
Both methods return an IoTaskStatus. Io means work remains, not that the
transport is ready, so a task re-arms transport readiness before its next
operation. On the write side it is returned whenever output is still buffered,
including after an attempt that made no progress. Output the transport already
owns does not keep the task running, because its completion wakes the task
again.
Both readiness methods resolve to a Readiness value rather than a plain
ready signal. Readiness::Ready allows the task to perform its next operation.
Readiness::Close means the connection must be taken down: the task closes
both directions and releases the socket, cancelling anything still in flight.
For a socket this is shutdown(SHUT_RDWR) followed by close(). Whatever the
peer still has in the receive queue should be discarded first, because closing
a socket with unread input aborts the connection with an RST and loses the
output that was just drained.
Readiness::Close covers a graceful shutdown as well as a connection that ends
because of an I/O failure, a filter failure, or an expired shutdown deadline.
On the graceful path the I/O subsystem reports it only once buffered output
has been drained, so nothing is left to flush.
Readiness::Terminate is reported only for an explicit force close through
IoRef::terminate. Whatever is still buffered is discarded on purpose, and
the task must not close gracefully: no receive queue drain and no
shutdown(SHUT_RDWR). It aborts the connection instead, for a socket by
setting SO_LINGER to zero before close(), so that the peer sees an RST and
cannot mistake a truncated stream for a complete one.
Teardown is reported back through two methods. A transport that fails or
decides to abort calls IoContext::stop with the error, which terminates the
connection without draining pending work. Once transport teardown has actually
finished, and regardless of which side initiated it, the task calls
IoContext::stopped. That is what completes Io::shutdown, which resolves
only after the backend reports the connection stopped, not merely when filter
shutdown and flushing are done.
Socket types integrate with the I/O subsystem by implementing IoStream.
Its start() method receives an IoContext, starts the transport-specific
tasks, and returns a Handle used to control or query the transport.
A simplified Tokio-style adapter looks like this:
impl IoStream for TcpStream {
fn start(self, ctx: IoContext) -> Box<dyn Handle> {
let socket = Rc::new(self);
tokio::task::spawn_local(read_task(socket.clone(), ctx.clone()));
tokio::task::spawn_local(write_task(socket.clone(), ctx));
Box::new(SocketHandle(socket))
}
}
async fn read_task(socket: Rc<TcpStream>, ctx: IoContext) {
loop {
// resolves `IoContext::poll_read_ready`
match wait_for_read_readiness(&ctx).await {
Readiness::Ready => {}
Readiness::Close | Readiness::Terminate => break,
}
let mut buf = ctx.take_read_buf();
let result = read_from_socket(&socket, &mut buf);
match ctx.release_read_buf(buf, result) {
IoTaskStatus::Io => {}
IoTaskStatus::Pause => wait_for_read_resume(&ctx).await,
IoTaskStatus::Stop => break,
}
}
}
async fn write_task(socket: Rc<TcpStream>, ctx: IoContext) {
let terminate = loop {
// resolves `IoContext::poll_write_ready`
match wait_for_write_readiness(&ctx).await {
Readiness::Ready => {}
Readiness::Close => break false,
Readiness::Terminate => break true,
}
let result = ctx.with_write_dst(|buf| {
write_to_socket(&socket, buf)
});
// `result` carries the number of bytes that reached the peer
match ctx.update_write_status(result) {
IoTaskStatus::Io => {}
IoTaskStatus::Pause => wait_for_queued_output(&ctx).await,
IoTaskStatus::Stop => break false,
}
};
if terminate {
// force close, abort so that the peer sees an RST
abort_socket(&socket);
} else {
// close both directions
shutdown_socket(&socket).await;
}
// report teardown as complete
ctx.stopped(None);
}
The example is illustrative; runtime adapters use their native readiness and
buffer APIs. The important boundary is that only the adapter accesses the
socket, while the application and protocol layers interact with Io and
IoRef. Buffer limits and read/write backpressure are therefore applied
consistently across all supported runtimes.
An Io object is normally passed to a protocol service such as
ntex::http::HttpService or ntex_mqtt::Server. The protocol service
decodes incoming bytes into protocol messages and encodes its responses back
into the I/O write buffer without depending on the concrete socket type.
The runtime-specific implementations are provided by the ntex-net crate,
which supports Tokio, Compio, and Neon backends.
Transport handles and metadata
Because the underlying socket is hidden behind the I/O abstraction, code using
Io cannot access transport-specific methods directly. The Handle returned
by IoStream::start() provides the bridge to the underlying transport. It can
resume the transport write operation on request and expose transport-specific
metadata through typed queries.
Each backend decides which query types it supports. For example, the built-in
network backends expose the remote socket address as PeerAddr.
IoRef::query is also available on Io through dereferencing and returns a
QueryItem, from which the value can be retrieved with get():
use ntex::io::{Io, types::PeerAddr};
fn log_peer_addr(io: &Io) {
if let Some(addr) = io.query::<PeerAddr>().get() {
println!("Peer address {:?}", addr.into_inner());
}
}
Queries travel through the filter stack before reaching the transport handle, so filters may also expose their own typed metadata.
Configuration and timeouts
I/O settings are stored in IoConfig and are normally added to a
SharedCfg. Io::new() retrieves the IoConfig from that shared
configuration, using the default settings when it is not present.
use ntex::{
SharedCfg,
io::IoConfig,
time::{Millis, Seconds},
util::BytePageSize,
};
let cfg = SharedCfg::new("my-protocol")
.add(
IoConfig::new()
.set_connect_timeout(Millis(5_000))
.set_keepalive_timeout(Seconds(30))
.set_shutdown_timeout(Seconds(2))
.set_frame_read_rate(Seconds(2), Seconds(10), 1_024)
.set_write_timeout(Seconds(10))
.set_read_size(BytePageSize::Size4, BytePageSize::Size32)
.set_read_backpressure(32 * 1024)
.set_write_backpressure(32 * 1024)
.set_write_buf_threshold(8 * 1024),
)
.build();
These settings are used by different parts of the stack:
- The connection timeout is applied by
ntex-netwhile resolving and opening an outgoing connection. - The keep-alive timeout and frame read-rate limits are interpreted by
protocol dispatchers. A frame read-rate limit protects a decoder from peers
that send one incomplete frame too slowly. It also applies to the first
frame of a new connection, so a peer that connects and stays silent is
closed once the limit expires. A write timeout protects against peers
that stop reading: it bounds how long write backpressure may stay enabled,
from the moment it is enabled until it is disabled, before the dispatcher
stops with a write timeout. Keep-alive and read-rate timers do not run
during that time. Output left after backpressure is disabled is bounded
only by keep-alive. The HTTP/1 dispatcher does not use these settings; it is
configured by
HttpServiceConfig, including its ownset_write_timeout(). - The graceful-shutdown timeout bounds both phases of shutdown together: the filter shutdown and the transport drain of pending output. It cannot be disabled; a zero timeout is rejected.
- The read and write high-water marks enable backpressure. Write backpressure is released after outstanding output falls to half its high-water mark, counting both buffered output and output a transport has taken ownership of but not yet written to the peer.
- The min and max read page sizes bound the adaptive page size of new read
buffers. Before another socket read, a buffer grows once less than
BytePageSize::lowof its page remains free. Output is held inBytePagesof the write page size, so the write backpressure setting takes only a high-water mark. - The write buffer page size controls newly allocated
BytePages, while the write threshold controls when supported transports attempt an early direct write.
Connection and keep-alive timeouts are disabled by default. Frame read-rate limits and the write timeout are also disabled. The default graceful-shutdown timeout is one second, and the default read and write high-water marks are approximately 32 KiB and 16 KiB.
Start with those defaults. Change one setting because of an observed workload, not because larger values sound faster:
- Read page sizes control allocation and read-call frequency, not the amount of unread data a connection may queue. Increase the maximum for sustained large frames or bodies; keep the minimum small when many mostly idle connections are expected.
- Backpressure watermarks control when producers or transports should pause. They are not hard memory limits: one socket read or one encoded item can cross a watermark.
- The write threshold is a latency hint for starting transport work earlier. It does not flush data and does not replace write backpressure.
- A write timeout is especially important for untrusted peers. Without one, a peer that stops reading can keep a backpressured connection and its output alive indefinitely.
An established connection can switch to another shared configuration with
Io::set_config. This is useful when a protocol upgrade changes timeout or
buffer requirements. The method is unsafe: replacing the configuration may
release the allocation that IoRef::cfg hands out, so no reference obtained
from it may be live across the call or used afterwards.
Filter subsystem
Applications often need to transform a byte stream before a protocol service processes it. TLS must decrypt incoming records and encrypt outgoing data. A protocol may also be tunneled through another framing layer, such as MQTT over WebSocket.
ntex implements these transformations with FilterLayer. A filter operates
on in-memory byte buffers: FilterLayer::process_read_buf transforms data
received from the next inner layer, while FilterLayer::process_write_buf
transforms data queued by the application before passing it toward the
transport. Filters do not perform socket I/O themselves, so the same filter
can be used with any supported runtime backend.
Filters are composable. Io::add_filter adds a layer and allocates the
intermediate read and write buffers that separate it from adjacent layers. For
example, an MQTT service can receive its byte stream through either of these
stacks:
socket <-> TLS <-> MQTT
socket <-> TLS <-> WebSocket <-> MQTT
On reads, bytes move from the socket through the inner filters toward the application. On writes, they move in the opposite direction. In the second stack, the WebSocket filter removes and creates WebSocket framing, while the TLS filter decrypts and encrypts the resulting byte stream. The MQTT service still reads and writes MQTT bytes and does not need to know which transport filters are installed below it.
Filters may maintain protocol state, expose typed metadata through query(),
emit output while processing input, and participate in graceful shutdown. For
example, a TLS filter can expose the negotiated protocol or peer certificate,
and a WebSocket filter can generate a close frame during shutdown. Output a
filter writes while processing a read is pushed toward the transport as soon as
that processing finishes, without waiting for the application to write.
Most byte transformations only need FilterLayer. Lower-level concerns that
must observe or control the entire filter chain can instead wrap the current
chain with Filter by using Io::map_filter.
In addition to processing buffers, queries, and shutdown, Filter participates
in read and write readiness decisions. A wrapper can therefore delay readiness
to implement custom throttling and wake the I/O tasks when work may resume. It
can also observe buffer processing for metrics or accounting without changing
the byte stream.
A custom Filter normally stores the filter it wraps and delegates every
operation it does not intentionally override. ntex provides forwarding macros
for readiness, queries, and shutdown to make this pattern less error-prone.
Typed versus erased filter stacks
The filter stack is encoded in the type parameter of Io. A new connection
starts as Io<Base>, using the Base filter. Calling add_filter(layer)
consumes the current value and returns Io<Layer<U, F>>, where the Layer
marker pairs the new outer layer U with the previous stack F. Keeping this
concrete type provides static dispatch and allows Io::filter to return the
concrete outer filter.
let io: Io<Base> = create_io();
let io: Io<Layer<MyFilter, Base>> = io.add_filter(MyFilter::new());
At service boundaries, different connections may have different concrete
filter stacks. Io::seal erases the stack type and returns Io<Sealed>,
using the Sealed marker, while Io::boxed returns the IoBoxed
convenience wrapper. Both operations consume the original Io value and retain
the same connection state and filter behavior behind a dynamically dispatched
Filter.
let io: IoBoxed = io.boxed();
start_protocol(io);
Additional typed layers can still be added to a sealed stream and the result can be erased again when necessary. Type erasure is therefore normally performed at the boundary where a protocol or service needs one uniform I/O type, rather than while constructing the filter stack.
Read/write streams
Incoming and outgoing bytes use separate buffer paths.
Reading
The transport adapter places bytes read from the socket into a BytesMut.
The bytes pass through the filter chain and arrive in the application-facing
read buffer, where a codec or protocol service can inspect and consume them.
BytesMut is a contiguous, growable buffer. A decoder can split immutable
Bytes values from it without copying larger payloads; short values are
copied into the Bytes handle itself. This is useful when a decoded message
must retain part of the input after the decoder continues processing later
data.
Read buffers are ntex-bytes pages. Each connection picks the page size of
new read buffers between the min and max set with
IoConfig::set_read_size, Size4 and Size64 by default. It starts at the
min, so idle and light connections pin small pages. Consecutive reads that
fill their buffer form one batch with the read that ends it; a batch larger
than the page grows the page size to fit it, up to the max. After four
batches in a row that fit in half of the next smaller page, the page size
shrinks by one step, down to the min. Io::set_config restarts it at the
min. Equal min and max sizes fix the page size.
The page size does not affect read backpressure: the read high-water mark is
the Size32 capacity, 32,752 bytes, by default;
IoConfig::set_read_backpressure sets another one. An empty buffer goes
back to the per-thread page cache of its size, shared with write pages and any
other pooled BytesMut; set_page_cache_size tunes it. A frame split from
the read buffer keeps the page alive, the page returns to the cache once the
last frame is dropped. Before another socket read, the adapter obtains a
buffer from IoContext. ntex calls BytesMut::reserve_more once less than
BytePageSize::low of the page remains free, compacting the data within its
page when it still fits there. If larger input must be buffered, the buffer moves to
bigger page sizes and, past the largest one, to a plain allocation that is
freed instead of cached.
Read backpressure is based on the size of the application-facing read buffer.
When it reaches the configured high-water mark, ntex pauses the transport read
task. Consuming input through IoRef::decode, IoRef::with_buf,
IoRef::with_read_src or IoRef::with_read_dst wakes the read task once
the buffered input falls to half that mark.
Io::recv and Io::read_exact wait for more input, so they release
backpressure regardless of how much is still buffered.
Choose the highest-level read API that matches the protocol:
- Use
recv(codec)for framed messages. It decodes buffered input and waits when the codec says the frame is incomplete. - Use
read_exact()for a fixed-size header or field. - For a custom parser, inspect or consume the read buffer through
IoRef, and callIo::read_moreonly after deciding that more bytes are required.read_more()is an active request: it resumes transport reads and releases read backpressure even if the buffer is still large.IoRef::read_dst_sizeis the passive size check.
Writing
Application output is queued in BytePages, a growable collection of
byte pages. Internally allocated pages use the size configured by IoConfig.
Owned buffers passed to IoRef::encode_bytes can also become pages directly,
avoiding a copy when their storage can be retained. Write filters consume the
pages in order, transform their contents, and place the result into the next
buffer toward the transport.
IoRef::encode, IoRef::encode_slice, and IoRef::encode_bytes queue
output but do not wait for every byte to reach the socket. Queueing makes the
data available to the transport write task and schedules that task when it
needs to be resumed. On transports that support direct writes, the configured
write threshold can trigger an earlier write while the application is still
producing output, reducing latency for large responses.
Use Io::flush to apply the configured write wait policy. flush(false) is
a backpressure wait, not a request to drain everything: it returns immediately
while the outstanding output is below the high-water mark. If the high-water
mark has been reached, it waits until the outstanding output falls to half that
mark. flush(true) waits until all queued data has reached the peer, including
data a transport has taken ownership of but not yet written. Io::send
combines codec encoding with a full flush.
Write backpressure is advisory. The encode methods still accept more output
after the high-water mark is reached, because refusing half of a protocol
message would be difficult to recover from. A producer that is not managed by
a dispatcher should await IoRef::write_ready before queueing each next
chunk, or use flush(false) between batches:
use ntex::{
io::Io,
util::Bytes,
};
async fn queue_chunk(io: &Io, chunk: Bytes) -> std::io::Result<()> {
io.write_ready().await?;
io.encode_bytes(chunk)
}
This bounds normal streaming output around the configured watermark. It does not split a single large chunk or turn the watermark into a strict per-connection memory cap.
This separation allows codecs and application services to work with bytes without depending on socket readiness, while the I/O subsystem coordinates buffering and backpressure.
Connection lifecycle and shutdown
A connection stays usable until the application closes it, the peer disconnects, or the transport reports an error.
A service that is waiting on something other than input, such as an unready
dependency or a slow response, still needs to notice that the connection
requires attention. Io::poll_status_update reports the next status as an
IoStatusUpdate value:
Timeoutwhen the dispatcher timer has expired, for example the configured keep-alive timeout.WriteBackpressurewhen outstanding output has reached the write high-water mark, so the producer should stop and flush.PeerGoneonce the connection has closed, whether the peer disconnected, the transport failed, or the shutdown was started locally. It carries the transport error if one occurred, andNoneafter a clean close.
Use Io::poll_read_pause instead when the service should also stop accepting
more transport input while it waits. The pause is cooperative: decoding,
touching the application read buffer, or explicitly asking for more input
resumes reads. This makes it suitable for a dispatcher waiting on readiness,
not for permanently disabling the read half.
The same conditions reach a codec-driven service as RecvError from
Io::poll_recv, which additionally reports decoder failures. Code that only
needs to be woken when the connection goes away can await the Waiter
future returned by IoRef::on_disconnect. IoRef::is_active becomes
false as soon as closing starts; IoRef::is_closed becomes true only after
the backend has released the transport and teardown has finished.
A clean EOF from the peer ends the read direction but leaves the write half open, so a service can still finish encoding and flushing its response before closing.
IoRef::close requests a graceful shutdown and returns immediately.
Io::shutdown drives that shutdown to completion, in two phases.
In the first phase both directions stay open and every filter gets a chance to
emit its own closing data through FilterLayer::shutdown, such as a TLS
close_notify or a WebSocket close frame. Reads keep running and are processed
normally, so closing data sent by the peer still reaches the filters. A read
failure does not abort the shutdown; it only means the filter handshake cannot
finish.
The second phase belongs to the transport. It drains the remaining output into the connection and then closes both directions. The read side is paused: the filters are done, so no further input can be used, and the transport discards whatever is left in the receive queue just before it closes. Input is no longer delivered to the application in this phase.
A single deadline bounds both phases. If the filters finish early, the
transport drain gets the remaining time. If the deadline expires while filters
are still pending, ntex ends that phase and still runs transport shutdown with
the already-expired deadline; when the transport phase observes it, undrained
output is discarded. Io::shutdown reports the timed-out error after
transport teardown.
IoRef::terminate skips the process entirely and drops the connection
without flushing pending output.
A graceful close, including one forced forward by a timeout or failure, reaches
the transport as Readiness::Close. An explicit IoRef::terminate reaches
it as Readiness::Terminate, so a socket backend can abort rather than send a
clean end-of-stream. The transport reports either teardown through
IoContext::stopped, which is what allows Io::shutdown to resolve.
Dropping the owning Io value is not a substitute for awaiting shutdown. If
no output would be lost, dropping it starts a normal close. If buffered output
can no longer be delivered because the filter stack is being dropped, ntex
aborts the transport so the peer does not mistake a truncated stream for a
complete one. Call shutdown().await when delivery and a clean close matter.
Testing
IoTest provides a pair of interconnected in-memory transports for testing
codecs, filters, and protocol services without opening sockets. Each endpoint
implements IoStream and can be wrapped in Io. Writing to one IoTest
endpoint supplies input to the other endpoint, while read() collects bytes
written back by the peer.
use ntex::codec::BytesCodec;
use ntex::io::{Io, testing::IoTest};
use ntex::util::Bytes;
#[ntex::test]
async fn protocol_io() {
let (client, server) = IoTest::create();
// Allow the server transport to write to the client.
client.remote_buffer_cap(1024);
let io = Io::from(server);
client.write(b"request");
let request = io.recv(&BytesCodec).await.unwrap().unwrap();
assert_eq!(request, Bytes::from_static(b"request"));
io.send(Bytes::from_static(b"response"), &BytesCodec)
.await
.unwrap();
assert_eq!(client.read().await.unwrap(), b"response"[..]);
}
Tests can also force pending reads, inject read or write errors, close either
side, constrain write capacity to exercise backpressure, and attach a
PeerAddr value for transport-query tests.
Web Applications and Routing
The ntex::web module provides everything needed to build an HTTP
application: routing, request extractors, responses, middleware, and testing
tools.
web::App is where an application comes together. You add routes, resources,
scopes, middleware, filters, and fallback services, then ntex builds them into
the service that handles incoming requests.
The server evaluates the application factory independently in each worker:
use ntex::web::{self, HttpResponse};
#[ntex::main]
async fn main() -> std::io::Result<()> {
web::server(async |_| {
web::App::new()
.route(
"/",
web::get().to(async || HttpResponse::Ok().body("Hello")),
)
})
.bind("127.0.0.1:8080", ntex::SharedCfg::default())?
.run()
.await
}
The factory returns an application factory rather than a running service.
HttpService uses it to create the application service for each connection,
so application services are never shared between workers. A worker can invoke
the outer factory again if its service has to be recreated, so keep factory
setup repeatable rather than relying on it to run exactly once.
To learn more about workers, server configuration, and application state, see Server and Worker and request state.
Handlers
The quickest way to register a handler is App::route():
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new()
.route("/users", web::get().to(async || HttpResponse::Ok()))
.route(
"/users",
web::post().to(async || HttpResponse::Created()),
);
A handler takes request extractors that implement FromRequest and returns
a value that implements Responder.
Each call to App::route() creates a separate Resource containing one
Route. The route’s method and custom guards become resource guards. This
means you can register the same path several times with different guards. It
also means that when those guards fail, ntex skips that resource and uses the
application fallback, which returns 404 Not Found by default.
When several routes belong to the same path, group them with
App::service() and web::resource(). A resource can also have its own name,
middleware, filters, and fallback:
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new().service(
web::resource("/users/{id}")
.name("user")
.route(web::get().to(async || HttpResponse::Ok()))
.route(web::delete().to(async || HttpResponse::NoContent())),
);
Resource::to() adds a route without a method guard, so it accepts every HTTP
method:
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new().service(
web::resource("/health").to(async || HttpResponse::Ok()),
);
ntex also provides attribute macros for common routes. The generated service
can be passed directly to App::service():
use ntex::web::{self, HttpResponse};
#[web::get("/health")]
async fn health() -> HttpResponse {
HttpResponse::Ok().into()
}
let app = web::App::<()>::new().service(health);
Method-specific route helpers include get, post, put, delete, patch,
head, and query. Routes for other methods are created with web::method(),
for example web::method(Method::OPTIONS).
Method matching is exact. A GET route does not also register HEAD, and
ntex does not add an OPTIONS route for you. Register those methods explicitly
when clients need them.
Choosing a Registration Style
For a small endpoint, App::route() is usually the clearest choice. It keeps
the path, method, and handler together:
use ntex::web::{self, App};
let app = App::default()
.route("/health", web::get().to(async || "ready"));
Use an explicit Resource when one path has several methods, needs a name, or
has its own middleware, filter, or fallback. Use a Scope when several
resources share a path prefix or policy. This distinction is more than style:
an App::route() method mismatch continues through the application router,
while a route mismatch inside a selected Resource reaches that resource’s
fallback.
web::service() is the lower-level option for code that already implements an
ntex service over WebRequest. Most handler-based applications do not need it.
Builder Order
App is configured in two phases. Application-wide settings come first:
middleware(), filter(), case_insensitive_routing(), with_config(), and
external_resource(). The first call to route(), service(),
configure(), or default_service() switches the builder to service
registration. After that, only
route(), service(), and default_service() are available.
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new()
// Application-wide settings
.case_insensitive_routing()
.external_resource("docs", "https://docs.example.com/{page}")
// Service registration
.route("/", web::get().to(async || HttpResponse::Ok()))
.default_service(web::to(async || HttpResponse::NotFound()));
Calling a setting such as external_resource() after route() does not
compile. Scope follows the same rule for its settings: guard(),
middleware(), filter(), and case_insensitive_routing().
How Routing Works
A request moves through the application in several stages:
- Application middleware wraps the application service.
- Application filters process the incoming
WebRequest. - The application router finds a
Resourcewhose path and resource guards match. - Resource middleware wraps the selected resource’s filter and routes.
- The resource filter processes the request.
- The resource checks its routes in registration order. A route matches only if all its method and custom guards pass.
- The first matching route calls its handler.
Middleware can return a response without calling the service it wraps. When that happens, ntex skips every later stage.
Where routing stops determines which fallback runs:
- If no resource path-and-guard combination matches, the application default
service runs. The built-in application default returns
404 Not Found. - If a resource path matches but none of its routes match, the resource
default service runs. The built-in resource default returns
405 Method Not Allowed. - Resource fallbacks are independent of application and scope fallbacks.
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new()
.service(
web::resource("/reports")
.route(web::get().to(async || HttpResponse::Ok()))
.default_service(
web::to(async || HttpResponse::MethodNotAllowed()),
),
)
.default_service(
web::to(async || HttpResponse::NotFound().body("Unknown path")),
);
Resource guards help the application router choose a resource. Route guards are checked later, after a resource has been selected. Guards can check methods, headers, or custom conditions:
use ntex::http::Method;
use ntex::web::{self, guard, HttpResponse};
let app = web::App::<()>::new().service(
web::resource("/events").route(
web::route()
.method(Method::POST)
.guard(guard::Header("content-type", "application/json"))
.to(async || HttpResponse::Accepted()),
),
);
Guards only receive the request head. They are a good fit for methods, headers, and other request metadata, but they cannot inspect the request body or values produced by extractors. Use a filter, extractor, or handler for those checks.
This matters when a request uses the wrong method.
App::route("/reports", web::get().to(handler)) uses the application fallback
for a non-GET request because the method guard belongs to the generated
resource. In contrast,
web::resource("/reports").route(web::get().to(handler)) first selects the
resource by path, then uses its default 405 Method Not Allowed response when
the method guard fails. Scope::route() behaves like App::route().
Order Matters
Resources and scopes are considered in registration order, and the first path whose guards pass wins. Routes inside a resource follow the same rule. Put specific patterns before broad dynamic or remainder patterns:
use ntex::web::{self, App};
let app = App::default()
.service(web::resource("/files/status").to(async || "ready"))
.service(web::resource("/files/{path}*").to(async || "file"));
Reversing these registrations would let the remainder resource handle
/files/status. For the same reason, add an unguarded route with
Resource::to() after guarded routes; once it matches, later routes cannot be
reached.
Use App::middleware() to wrap the whole application. Use App::filter() to
transform every incoming WebRequest before routing. Resources and scopes
offer the same methods when the behavior should apply to only part of the
application. See Web Application Filters and Middleware for more
information.
Route Path Format
Route patterns contain static and dynamic segments separated by /. ntex adds
a leading / to resource and scope paths when needed, but including it makes
the full route easier to read.
Static Paths
A static pattern matches the same path:
/users
/users/profile
Trailing slashes matter: /users and /users/ are different paths. Query
strings are not included in path matching.
Routing is case-sensitive by default.
App::case_insensitive_routing() makes static path segments ASCII
case-insensitive, but does not change how dynamic segments are matched. Scopes
can enable the same behavior for their own nested routes.
Dynamic Segments
Write a dynamic segment as {name}. It matches one non-empty path segment:
/users/{id}
/teams/{team}/users/{user}
You can combine variables with static text or use several variables in one segment:
/releases/v{major}.{minor}
/files/{name}.{extension}
Add a regular expression as {name:regex} to restrict what a variable
accepts:
/users/{id:[0-9]+}
/releases/{version:v[0-9]+\.[0-9]+}
The regular expression must match the whole segment. Invalid route patterns and regular expressions panic while the application is built, so define routes as trusted application configuration rather than from user input.
Remainder Matches
Add * after a dynamic variable to capture the rest of the path, including
/ separators:
/files/{path}*
This matches both /files/readme.txt and /files/images/logo.svg. The path
value is readme.txt in the first case and images/logo.svg in the second.
The default expression is .*, so the captured value may be empty. Remainder
matches cannot use a custom regular expression.
Unlike normal dynamic segments, a remainder value is not percent-decoded. A
request for /files/my%20file.txt captures my%20file.txt. Decode the value
yourself if the handler needs the decoded path, and validate it before using
it as a file system path.
Static tail patterns such as /files/* also work, but a named remainder is
usually clearer and works naturally with typed path extraction.
Multiple Patterns
A resource or scope can use an array or vector to accept several patterns:
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new()
.service(
web::resource(["/health", "/status"])
.to(async || HttpResponse::Ok()),
);
Every pattern points to the same service.
When a multi-pattern resource is named, URL generation uses its first pattern. Put the canonical path first and treat the remaining patterns as aliases.
Accessing Path Variables
ntex percent-decodes variables matched by normal dynamic segments and stores
them in the request’s match information. Remainder values are stored as they
appear in the request path. The easiest way to read them is with the typed Path
extractor:
use ntex::web::{self, HttpResponse};
async fn user(path: web::types::Path<(u32,)>) -> HttpResponse {
let user_id = path.0;
HttpResponse::Ok().body(format!("user {user_id}"))
}
let app = web::App::<()>::new().route(
"/users/{id:[0-9]+}",
web::get().to(user),
);
Tuples deserialize variables in path order. To access variables by name, use a
struct that derives serde::Deserialize:
use ntex::web::{self, HttpResponse};
#[derive(serde::Deserialize)]
struct UserPath {
organization: String,
user_id: u32,
}
async fn user(path: web::types::Path<UserPath>) -> HttpResponse {
HttpResponse::Ok().body(format!(
"organization: {}, user: {}",
path.organization,
path.user_id,
))
}
let app = web::App::<()>::new().route(
"/organizations/{organization}/users/{user_id:[0-9]+}",
web::get().to(user),
);
The struct fields must have the same names as the route variables. This
example requires the derive feature from serde.
You can also read variables directly from HttpRequest::match_info():
use ntex::web::{self, HttpRequest};
async fn file(req: HttpRequest) -> String {
req.match_info()
.get("path")
.unwrap_or_default()
.to_owned()
}
let app = web::App::<()>::new().route(
"/files/{path}*",
web::get().to(file),
);
ntex splits path segments before percent-decoding them. As a result, an
encoded slash such as %2F can appear inside a normal dynamic variable after
decoding. Escapes that are not valid UTF-8 remain percent-encoded.
Path values are still untrusted input. In particular, a variable that looked
like one URL segment can decode to ../ or contain /. Validate it before
using it as a file name, database key, or other identifier with stricter
rules.
Scopes
A Scope groups related services under a shared path prefix:
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new().service(
web::scope("/api")
.service(
web::resource("/users")
.route(web::get().to(async || HttpResponse::Ok())),
)
.service(
web::resource("/users/{id}")
.route(web::get().to(async || HttpResponse::Ok())),
),
);
These resources match /api/users and /api/users/{id}. You can nest scopes
and use variables in their prefixes. Nested handlers can access those
variables:
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new().service(
web::scope("/organizations/{org}")
.route(
"/users/{user}",
web::get().to(
async |path: web::types::Path<(String, String)>| {
HttpResponse::Ok()
.body(format!("{}:{}", path.0, path.1))
},
),
),
);
A scope prefix matches complete path segments rather than arbitrary text.
scope("/api") can contain routes below /api/, but it does not by itself
provide a handler for the bare /api path and it never matches /apix.
Inside the scope, an empty resource pattern ("") matches /api, while "/"
matches /api/:
use ntex::web::{self, HttpResponse};
let app = web::App::<()>::new().service(
web::scope("/api")
.route("", web::get().to(async || HttpResponse::Ok().body("/api")))
.route("/", web::get().to(async || HttpResponse::Ok().body("/api/"))),
);
A scope can have its own guards, middleware, filters, case-sensitivity setting,
and fallback service. Once a scope matches, unmatched paths inside it use the
scope fallback. Without a custom scope fallback, ntex returns its built-in
404 Not Found; it does not use the application fallback.
Named and External Resources
Give a resource a name when handlers need to link to it without hard-coding its path:
use ntex::web::{self, HttpRequest, HttpResponse, error::UrlGenerationError};
async fn index(req: HttpRequest) -> Result<HttpResponse, UrlGenerationError> {
let url = req.url_for("user", ["42"])?;
Ok(HttpResponse::Ok().body(url.to_string()))
}
let app = web::App::<()>::new()
.service(
web::resource("/users/{id}")
.name("user")
.route(web::get().to(async || HttpResponse::Ok())),
)
.route("/", web::get().to(index));
Here, the index handler generates an absolute URL such as
http://some-host-name/users/42.
HttpRequest::url_for() returns a ntex::url::Url. Values fill dynamic
segments in pattern order. For a resource inside a scope, values for dynamic
scope segments come before values for the resource itself. Use
HttpRequest::url_for_static() when the pattern has no dynamic segments.
Each value is percent-encoded as one literal segment, so URL delimiters in a
value cannot change the generated URL’s structure.
Keep resource names unique across the application. Names are application-wide lookup keys, not names local to a scope, and duplicate names make generated links depend on which definition is found.
Generation can fail if the name is unknown, too few values are supplied, or
the completed URL is invalid. Returning UrlGenerationError from a handler,
as above, avoids panicking on those errors.
For an application resource, the scheme and host come from the request’s
ConnectionInfo, while the path comes from the resource and its containing
scopes. ConnectionInfo can use Forwarded, X-Forwarded-Proto,
X-Forwarded-Host, or Host, so a front-end proxy should replace untrusted
forwarding headers before generated URLs are exposed to clients.
Use App::external_resource() to register a named URL that is not handled by
the application. Register it before the first route or service, or register it
inside App::configure() with ServiceConfig::external_resource():
use ntex::web::{self, HttpRequest, HttpResponse, error::UrlGenerationError};
async fn docs(req: HttpRequest) -> Result<HttpResponse, UrlGenerationError> {
let url = req.url_for("documentation", ["page1.html"])?;
Ok(HttpResponse::Ok().body(url.to_string()))
}
let app = web::App::<()>::new()
.external_resource(
"documentation",
"https://docs.example.com/{page}",
)
.route("/docs", web::get().to(docs));
The docs handler generates https://docs.example.com/page1.html. Because
the external pattern already contains a scheme and host, it does not use the
request’s connection information. External resources are URL templates only;
they do not participate in request routing.
Modular Configuration
As an application grows, keeping every route in one builder chain becomes hard
to read. App::configure() provides a ServiceConfig that lets another
function register part of the application:
use ntex::web::{self, HttpResponse};
fn configure_api(cfg: &mut web::ServiceConfig<()>) {
cfg.route(
"/health",
web::get().to(async || HttpResponse::Ok()),
);
cfg.service(
web::resource("/version")
.to(async || HttpResponse::Ok().body("4")),
);
}
let app = web::App::<()>::new()
.configure(configure_api)
.route("/", web::get().to(async || HttpResponse::Ok()));
ServiceConfig can register routes, services, and external resources. It does
not create another routing boundary. The registered items join the application
at the point where configure() is called.
Both applications and scopes support modular configuration.
Testing Routes
The web::test module lets you initialize an App, send requests to it, and
inspect responses without opening a network socket:
use ntex::http::StatusCode;
use ntex::web::{self, HttpResponse};
use ntex::web::test::{TestRequest, call_service, init_service};
#[ntex::test]
async fn user_route() {
let service = init_service(
web::App::new().route(
"/users/{id}",
web::get().to(async || HttpResponse::Ok()),
),
)
.await;
let request = TestRequest::get()
.uri("/users/42")
.to_request();
let response = call_service(&service, request).await;
assert_eq!(response.status(), StatusCode::OK);
}
It is worth testing method mismatches, trailing slashes, dynamic-variable constraints, scope boundaries, and fallback services. Each one exercises a different part of route selection.
Lower-Level HTTP Service
For most applications, web::server(factory) is the simplest way to start a
server. It is equivalent to web::HttpServer::new(factory).
After you register routes and services, the application builder can become a
service factory that converts http::Request values into http::Response
values. HttpService uses this factory to handle HTTP/1.1 and HTTP/2
connections.
Use the lower-level server builder when you need to add connection-level
services before HttpService. These services can inspect or transform the
Io object and map the state passed to the HTTP and application services.
The same server can be built at the lower level like this:
use ntex::{http, server, web, SharedCfg};
use ntex::web::HttpResponse;
#[ntex::main]
async fn main() -> std::io::Result<()> {
server::build()
.bind("http", "127.0.0.1:8080", SharedCfg::new("S"), async |_| {
http::HttpService::new(
web::App::new()
.route("/", web::get().to(async || {
HttpResponse::Ok().body("Hello")
}))
.build(),
)
})?
.run()
.await
}
To give the web application its own state instead of the surrounding service pipeline’s state, see Web Application State.
Web Application State
A web::App is parameterized by an application state type. It usually
inherits this state from the surrounding server or service pipeline, making
the same value available to handlers, extractors, filters, middleware, and
error responses.
There are two kinds of state to keep in mind:
- Application state is the state supplied by the service pipeline. A
worker’s value is cloned for each connection, and handlers borrow that
connection’s clone. Data that must be shared between connections belongs
behind a shared handle such as
Rc. - Request-local state belongs to one request. It is owned by
WebRequest, can change type as the request moves through filters and middleware, and can be moved into a state-aware handler.
This guide explains how both kinds of state work with web::App. See
Worker and request state for worker initialization
and process-wide versus worker-local data, and Service state
for the underlying service-state model.
The web::State Trait
To use a type as application state, implement web::State:
pub trait State: 'static {
type Error;
}
The trait has no methods. Its associated Error type defines the application’s
error domain. Handlers, filters, middleware, and fallback services use it
through WebError and WebResponseError.
Most applications can use DefaultError:
use ntex::web;
#[derive(Clone)]
struct ApplicationState {
greeting: String,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
web::State itself does not require Clone, but in practice application
state used with the built-in HTTP server must implement it.
ServerAppConfig::State and state passed to build_with() must be Clone,
and HttpService clones the connection state while creating each
connection’s application service. A non-Clone state type can still be used
with an application factory directly, but it cannot be carried through these
server paths.
Inheriting Worker State
web::server_with_config() creates one state value per worker and passes it
to that worker’s application factory. The returned App uses the same state
type:
use std::io;
use ntex::{SharedCfg, web};
/// State created independently for each worker.
#[derive(Clone)]
struct ApplicationState {
greeting: String,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
/// Process-wide factory used to create worker state.
struct ApplicationConfig;
impl ntex::server::ServerAppConfig for ApplicationConfig {
type State = ApplicationState;
async fn create(&self) -> io::Result<Self::State> {
Ok(ApplicationState {
greeting: "Hello".to_owned(),
})
}
}
async fn index(state: &ApplicationState, _request_state: ()) -> String {
state.greeting.clone()
}
#[ntex::main]
async fn main() -> io::Result<()> {
web::server_with_config(ApplicationConfig, async |_| {
web::App::new().route("/", web::get().to_with_state(index))
})
.bind("127.0.0.1:8080", SharedCfg::default())?
.run()
.await
}
Every request handled by a worker can access that worker’s state through the
service context. Each connection receives its own clone of the worker state,
so values that all connections of a worker should share, such as caches or
counters, belong behind Rc. See
Worker-Local and Process-Wide State
for details.
Design State for Cheap Cloning
Application state is usually a small collection of handles rather than a large data structure copied for every connection:
use std::{cell::RefCell, collections::HashMap, rc::Rc, sync::Arc};
use ntex::web;
#[derive(Clone)]
struct ApplicationState {
worker_cache: Rc<RefCell<HashMap<String, String>>>,
shared_settings: Arc<Settings>,
}
struct Settings {
service_name: String,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
Here, connections on the same worker share worker_cache. All workers can
also share the same thread-safe shared_settings value when the configuration
factory gives each worker a clone of the same Arc.
Handlers receive &ApplicationState, so mutation normally happens through a
client handle or an interior-mutable value such as Cell, RefCell, a lock,
or an atomic. Keep RefCell borrows short and do not hold them across an
.await, where unrelated work could try to borrow the same value and panic.
Accessing Application State in Handlers
Use Route::to_with_state() when a handler needs direct access to state.
Resource::to_with_state() does the same for a resource route without a
method guard. The handler receives the arguments in this order:
- A shared reference to the application state.
- The request-local state.
- Any request extractors.
use ntex::web;
#[derive(Clone)]
struct ApplicationState {
service_name: &'static str,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
async fn user(
state: &ApplicationState,
_request_state: (),
user_id: web::types::Path<u32>,
) -> String {
format!("{} user {}", state.service_name, user_id.into_inner())
}
fn main() {
let app = web::App::<ApplicationState>::new().route(
"/users/{user_id}",
web::get().to_with_state(user),
);
}
Filters, middleware, and lower-level services access application state through
their service Ctx.
Handlers registered with Route::to() receive only request extractors. Use
to() when the handler does not need application or request-local state.
Extractors also receive a shared reference to application state through
FromRequest, so a custom extractor can use application services or
configuration. Extractors do not receive request-local state; that value is a
separate handler argument supplied only by to_with_state().
Request-Local State
Application state is long-lived and shared across requests. Request-local
state is different: it belongs to one request and is carried by WebRequest.
Its type is the parameter in WebRequest<RequestState>.
An ordinary web request starts with () as its request-local state. Each
filter or middleware can keep that value, mutate it through
WebRequest::st_mut(), or replace it with another type through
map_state().
This is useful for data produced while processing a request, such as
authentication details. Filters and middleware can transform the value with
WebRequest::map_state(), and a handler registered with to_with_state()
receives the result as its second argument:
use std::convert::Infallible;
use ntex::web::{self, WebRequest};
#[derive(Clone)]
struct ApplicationState {
service_name: &'static str,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
async fn index(state: &ApplicationState, user_id: usize) -> String {
format!("{} user {user_id}", state.service_name)
}
fn main() {
let app = web::App::<ApplicationState>::new()
.filter(async |req: WebRequest<()>| {
// Authentication could produce this request-local value.
Ok::<_, Infallible>(req.map_state(|()| 42usize))
})
.route("/", web::get().to_with_state(index));
}
Request-local state moves through the filter and middleware chain. Changing
the state also changes the WebRequest type expected by the next service.
This lets Rust check the pipeline at compile time: a handler expecting
AuthInfo, for example, can only follow a filter or middleware that produces
WebRequest<AuthInfo>.
Inside filters and middleware, WebRequest::st() and WebRequest::st_mut()
return the request-local state, not the application state. Application state
comes from the service context instead.
map_state() consumes the previous value. If later stages need some of the
old information, include it in the new state:
struct Authenticated {
user_id: usize,
}
struct Validated {
user: Authenticated,
request_id: String,
}
Application filters affect every route. A scope or resource filter can create state required only by that part of the application, which keeps unrelated handlers from depending on it.
Request-local state is not stored in HttpRequest extensions, and request
extractors cannot read it. Use to_with_state() when the final handler needs
the value. A handler registered with to() ignores and drops it.
Typed Web Configuration Is Separate
WebAppConfig also has a typed value store populated with set_state().
Despite the similar name, these values are not the application’s service state:
to_with_state()does not pass them as its first argument.- They are read from
HttpRequest::app_state()orWebRequest::app_state(). - They must be
Send + Sync, unlike worker-local state, which may containRcandRefCell. - There is one value for each concrete type.
This store is useful for HTTP-facing configuration, especially extractor limits and settings that are naturally read from a request:
use ntex::web::{self, App, HttpRequest, WebAppConfig};
struct Limits {
max_items: usize,
}
async fn limits(req: HttpRequest) -> String {
let max_items = req
.app_state::<Limits>()
.map_or(100, |limits| limits.max_items);
format!("maximum items: {max_items}")
}
let config = WebAppConfig::new().set_state(Limits { max_items: 50 });
let app = App::default()
.with_config(config)
.route("/limits", web::get().to(limits));
Use application state for services and resources that participate in the
service lifecycle. Use WebAppConfig values for typed web configuration.
Supplying Fixed State with .build_with()
After the first call to route(), service(), configure(), or
default_service(), the application builder becomes AppServices (see
Builder Order). You can finish building it
with either build() or AppServices::build_with().
build() uses the state from the surrounding service pipeline.
build_with(state) gives the application its own fixed state instead, allowing
the outer pipeline to use a different state type.
Use build_with() when the HTTP or server pipeline has different state, or no
state at all, but the web application needs its own:
use ntex::web::{self, HttpResponse};
#[derive(Clone)]
struct ApplicationState {
greeting: &'static str,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
async fn index(state: &ApplicationState, _request_state: ()) -> HttpResponse {
HttpResponse::Ok().body(state.greeting)
}
fn main() {
let app = web::App::<ApplicationState>::new()
.route("/", web::get().to_with_state(index))
.build_with::<()>(ApplicationState {
greeting: "Hello",
});
}
In this standalone example, () is the outer pipeline’s state type. It is
normally inferred when the application factory is passed to HttpService.
The application state must implement Clone because build_with() clones it
when creating application service instances.
build_with() replaces the state for the application; it does not add another
state layer. Code outside the application continues to use the outer state,
while web handlers and application middleware receive the fixed state.
Use server_with_config() when state needs asynchronous per-worker
initialization or should be recreated with a worker. Use build_with() when a
ready value should be embedded in an application factory, especially when the
surrounding HTTP pipeline uses another state type.
The AppState<T> Wrapper
If the application uses DefaultError, AppState<T> saves you from writing
a web::State implementation. It works with any 'static T and dereferences
to the wrapped value. It also implements Clone and Default whenever T
does:
use ntex::web;
#[derive(Clone)]
struct Settings {
service_name: &'static str,
}
fn main() {
let state = web::AppState::new(Settings {
service_name: "users",
});
assert_eq!(state.service_name, "users");
assert_eq!(state.st().service_name, "users");
}
Use a dedicated state type and implement web::State directly when you need a
custom error type or other state-specific behavior.
AppState<T> does not change how values are cloned or shared. If T::clone()
copies a plain field, connections receive separate copies; if it clones an
Rc or Arc, they share the value behind that handle.
Testing Stateful Routes
test::init_service_st() builds an application with an explicit state value.
Keep a clone of a shared handle when the test needs to inspect what the handler
changed:
use std::{cell::Cell, rc::Rc};
use ntex::web::{self, test};
#[derive(Clone)]
struct ApplicationState {
hits: Rc<Cell<usize>>,
}
impl web::State for ApplicationState {
type Error = web::DefaultError;
}
async fn index(state: &ApplicationState, (): ()) -> &'static str {
state.hits.set(state.hits.get() + 1);
"ok"
}
#[ntex::test]
async fn state_is_used() {
let state = ApplicationState {
hits: Rc::new(Cell::new(0)),
};
let service = test::init_service_st(
state.clone(),
web::App::<ApplicationState>::new()
.route("/", web::get().to_with_state(index)),
)
.await;
let request = test::TestRequest::get().uri("/").to_request();
let _response = test::call_service(&service, request).await;
assert_eq!(state.hits.get(), 1);
}
Web Application Filters and Middleware
Like the rest of ntex, the web framework is built around Service and
Middleware. A web application is simply a service that accepts an HTTP
request, runs it through the configured pipeline, and returns a
WebResponse.
Filters and middleware let you customize this pipeline in different ways:
- A filter transforms an incoming
WebRequestbefore processing continues. - Middleware wraps a service, so it can run code both before and after that service.
Both can use application state. Middleware services receive it through
Ctx, and a filter function can take it as its first argument. For the
distinction between application state and request-local state, see
Web Application State.
Filters
App::filter(), Scope::filter(), and Resource::filter() append a
service to the inbound request pipeline. A filter accepts
WebRequest<In> and must return WebRequest<Out>:
WebRequest<In> -> filter -> WebRequest<Out>
Filters run in registration order, with each filter receiving the output of
the previous one. If a filter returns an error, ntex skips the remaining
filters, router, and handler, then renders the error through
WebResponseError.
Because request state is part of the WebRequest type, filters can enforce
requirements at compile time. For example, an authentication filter can turn
WebRequest<()> into WebRequest<AuthenticatedUser>. A handler that expects
AuthenticatedUser can only be registered after a filter that provides it:
use std::convert::Infallible;
use ntex::web::{self, WebRequest};
struct AuthenticatedUser {
id: u64,
}
async fn authenticate(
req: WebRequest<()>,
) -> Result<WebRequest<AuthenticatedUser>, Infallible> {
Ok(req.map_state(|()| AuthenticatedUser { id: 42 }))
}
async fn profile(_app: &(), user: AuthenticatedUser) -> String {
format!("User {}", user.id)
}
web::App::default()
.filter(authenticate)
.route("/profile", web::get().to_with_state(profile));
Here the application state is (), so profile receives &() as its first
argument and the AuthenticatedUser produced by the filter as its second.
WebRequest::map_state() consumes the request, transforms its request-local
state, and preserves the HTTP request and payload.
A filter can also take the application state as its first argument and return
any error that implements WebResponseError. The error stops processing and
becomes the response:
use ntex::web::{self, WebRequest, error};
fn main() {
let app = web::App::default()
.filter(async |_state: &(), req: WebRequest<()>| {
if req.headers().contains_key("x-api-key") {
Ok(req)
} else {
Err(error::ErrorForbidden("missing API key"))
}
})
.route("/", web::get().to(async || "Hello"));
}
Filters only run on the way in. They cannot inspect or modify the handler’s response. Use middleware when you need to wrap the complete request-response operation.
Filters also cannot return a successful response directly: their successful
output must be another WebRequest. Return an error to reject a request, or
use middleware when a cache hit, maintenance page, or authorization decision
should produce a response without running the inner service.
Middleware
A middleware component receives a service and wraps it with another service.
The wrapper can inspect or modify the request, call the inner service through
Ctx, and then inspect or modify the response. It can also participate in
readiness and shutdown.
Middleware can change the request-state type before calling the inner service,
but a filter is usually simpler when you only need to process the request.
Middleware can also short-circuit the pipeline by returning a WebResponse
without calling the inner service. WebRequest::into_response() is the
convenient way to keep the original request attached to that response.
Middleware can be installed with App::middleware(),
Scope::middleware(), or Resource::middleware(). Built-in middleware
can be installed directly:
use ntex::web::{self, App, middleware};
App::default()
.middleware(
middleware::DefaultHeaders::new()
.header("x-application", "example"),
)
.middleware(middleware::Logger::default())
.route("/", web::get().to(async || "Hello"));
Middleware runs in registration order on the way in. In this example,
DefaultHeaders receives the request first, followed by Logger. The response
takes the opposite path: Logger processes it first, followed by
DefaultHeaders.
DefaultHeaders only inserts a header when the response does not already have
one, so a handler can override the default. Logger writes through the log
facade; an application must install and configure a logger to see its output.
A custom middleware follows the normal ntex service model:
use ntex::http::header::{HeaderName, HeaderValue};
use ntex::web::{self, WebRequest, WebResponse};
use ntex::{Ctx, Middleware, Service};
struct ResponseHeader;
struct ResponseHeaderService<S> {
service: S,
}
impl<S, St> Middleware<S, St> for ResponseHeader {
type Service = ResponseHeaderService<S>;
fn create(&self, _state: &St, service: S) -> Self::Service {
ResponseHeaderService { service }
}
}
impl<S, St, ReqSt> Service<St, WebRequest<ReqSt>>
for ResponseHeaderService<S>
where
S: Service<St, WebRequest<ReqSt>, Res = WebResponse>,
{
type Res = WebResponse;
type Error = S::Error;
ntex::forward_ready!(St, service);
ntex::forward_shutdown!(St, service);
async fn call(
&self,
req: WebRequest<ReqSt>,
ctx: Ctx<'_, Self, St>,
) -> Result<Self::Res, Self::Error> {
let mut res = ctx.call(&self.service, req).await?;
res.headers_mut().insert(
HeaderName::from_static("x-application"),
HeaderValue::from_static("example"),
);
Ok(res)
}
}
web::App::default()
.middleware(ResponseHeader)
.route("/", web::get().to(async || "Hello"));
Middleware::create() runs when the application service is constructed. It
receives the application state and the service to wrap, and returns the
per-service wrapper. Its call() method then receives a Ctx for each
request. If the wrapper needs to retain something from application state,
clone an owned handle during create() rather than trying to retain the
borrowed &St.
Always call the inner service through Ctx::call() so the pipeline handles
readiness and lifecycle events correctly. The forward_ready! and
forward_shutdown! macros forward readiness checks and shutdown to the wrapped
service. Without them, the default implementations report the middleware as
always ready and never shut down the inner service. The macros take the state
type and a named field, so store the inner service in a named field rather
than a tuple field. See Service pipelines for more on
readiness and shutdown.
Configure Layers Before Routes
At each level, add filters and middleware before the first route, service,
configuration callback, or default service. Those registration methods move
the builder to its service-registration stage, where filter() and
middleware() are no longer available:
use ntex::web::{self, middleware, App, WebRequest};
App::default()
.middleware(middleware::Logger::default())
.filter(async |req: WebRequest<()>| {
Ok::<_, std::convert::Infallible>(req)
})
.route("/", web::get().to(async || "Hello"));
The same rule applies to Scope and Resource. See
Builder Order.
Middleware and filter calls are two separate stacks, even if their builder calls are interleaved. All middleware at a level wraps that level’s complete filter-and-router service, so it runs before every filter on the inbound path.
Where Filters and Middleware Run
You can install filters and middleware on an application, scope, or resource:
| Level | Runs for |
|---|---|
App | Every request entering the application |
Scope | Requests whose scope prefix and scope guards match |
Resource | Requests whose resource path and resource guards match |
For example, a request handled by a resource inside a scope follows this path:
application middleware
-> application filters
-> application routing and scope guards
-> scope middleware
-> scope filters
-> scope routing and resource guards
-> resource middleware
-> resource filters
-> route guards
-> route handler
The response travels back through the middleware layers in reverse. At each level, middleware runs in registration order for the request and reverse registration order for the response. Filters are not part of the response path.
The router checks a scope or resource’s own guards before entering its middleware and filters. Application filters run before the application router chooses a scope or resource, and scope filters run before the scope router chooses a nested resource. Route guards are later: the resource filter runs first, then routes and their guards are checked in registration order.
Fallback services run inside the layers of their level. When no resource
matches, the application default service runs after application middleware
and filters. When a scope matches but none of its services do, the scope
default service runs after the scope’s middleware and filters. When a resource
matches but no route does, its default service runs after resource middleware
and filters; without a custom default, it returns 405 Method Not Allowed.
See Web Applications and Routing for guard and fallback
behavior.
Errors and the Response Path
Extractor failures and errors returned by handler Result values are
converted into WebResponse values inside the route handler. Response
middleware can therefore inspect or modify those error responses normally.
Errors returned directly by filters, middleware services, or other web
services take a different path. They remain service errors while unwinding
through middleware and are rendered through WebResponseError outside the
application middleware stack. A middleware that propagates an inner error
with ? does not see the final rendered response. If it must process that
case, it needs to catch the error and convert it into a WebResponse itself.
See Web Application Errors for the error model.
Choosing Between Them
A useful rule of thumb is to use a filter for request-only work and middleware when you also need the response.
Choose a filter when:
- Only the incoming request needs processing.
- The request-local state type should change.
- Failure should stop processing before routing or handler execution.
Choose middleware when:
- Both the request and response need to be observed or modified.
- Processing may finish early with a successful response.
- Logic must wrap an entire application, scope, or resource.
- The component must control how the wrapped service’s readiness or shutdown is reported.
Filters are services too. A filter implemented as a full Service takes part
in readiness and shutdown like any other stage of the pipeline, but it cannot
see the response or the wrapped service.
Web Application Extractors and Responders
Most handlers do not need to parse requests or build responses by hand. ntex handles both sides of the process:
- Extractors provide the values passed to a handler.
- Responders turn the handler’s result into an HTTP response.
This leaves the handler free to focus on application logic.
Extractors
Handler arguments are extractors. Before calling a handler registered with
Route::to(), ntex creates each argument from the request in the order it
appears in the function signature. Each argument type must implement
FromRequest. A handler can take up to 16 extractor arguments.
The built-in extractors cover common request data:
| Extractor | Reads |
|---|---|
HttpRequest | A cheap clone of the request handle |
Path<T> | Dynamic path segments |
Query<T> | URL query parameters |
Json<T> | A JSON request body |
Form<T> | A URL-encoded request body |
Payload | The streaming request body |
Bytes | The complete request body as bytes |
String | The complete request body as decoded text |
Option<T> | Some(value), or None if T fails |
Result<T, T::Error> | The value, or the error produced by T |
(A, B, ...) | Several extractors combined into one argument |
For example, this handler reads a user ID from the path and an optional flag from the query string:
use ntex::web::{self, App};
#[derive(serde::Deserialize)]
struct Options {
details: Option<bool>,
}
async fn user(
user_id: web::types::Path<u32>,
options: web::types::Query<Options>,
) -> String {
format!(
"User {}, details: {}",
user_id.into_inner(),
options.details.unwrap_or(false),
)
}
App::default().route(
"/users/{user_id}",
web::get().to(user),
);
If any extractor fails, the handler is not called. ntex turns the error into a
response through WebResponseError.
Path<T> and Query<T> both deserialize with serde, but they represent
different shapes:
- a path tuple reads captured segments by position;
- a path struct reads them by route variable name;
- a query string is a set of
key=valuepairs, so use a struct or map rather than a tuple.
Path values are percent-decoded after the URL has been split into segments.
That means an encoded slash such as %2F can appear inside one extracted
value. Treat path values as input, not as safe file-system paths.
Request Bodies
A request has only one body stream, and all extractors share it. Json<T>,
Form<T>, Bytes, and String consume that stream, so a handler should
normally have only one body extractor. Their default in-memory limits are:
| Extractor | Default limit |
|---|---|
Json<T> | 32 KiB |
Form<T> | 16 KiB |
Bytes and String | 256 KiB |
Json<T>, Form<T>, and Query<T> use serde to deserialize values.
Json<T> accepts JSON media types, including types with a +json suffix.
Form<T> requires application/x-www-form-urlencoded. String decodes the
body according to the request charset; Bytes leaves it unchanged.
Use JsonConfig, FormConfig, and PayloadConfig to change those
limits. JsonConfig can accept additional content types. PayloadConfig can
require a content type for Bytes and String; without that setting, those
extractors accept any content type. FormConfig changes only the size limit.
Store extractor configuration in WebAppConfig with set_state() and
install it with App::with_config():
use ntex::web::{self, App, WebAppConfig};
fn main() {
let config = WebAppConfig::new()
.set_state(web::types::JsonConfig::default().limit(4096))
.set_state(web::types::PayloadConfig::new(64 * 1024));
let app = App::default()
.with_config(config)
.route("/", web::post().to(async |body: String| body));
}
Values stored with set_state() must be Send + Sync. If an application does
not call with_config(), it uses the WebAppConfig from the connection’s
shared configuration, and extractors fall back to their defaults when no
configuration of the right type is stored.
Use Payload when the handler should process the body incrementally instead
of buffering it in memory. It takes ownership of the remaining body stream and
does not apply PayloadConfig; the handler decides how much data to read and
how to handle it.
Optional Extractors
Sometimes invalid input should not stop the request. Wrapping an extractor in
Option<T> turns a successful extraction into Some(value) and a failure
into None. The failure is also logged:
use ntex::web::{self, App};
#[derive(serde::Deserialize)]
struct Search {
query: String,
}
async fn search(params: Option<web::types::Query<Search>>) -> String {
match params {
Some(params) => format!("Searching for {}", params.query),
None => "No search query".to_owned(),
}
}
App::default().route("/search", web::get().to(search));
If the handler needs the actual error, use Result<T, T::Error> instead. Both
forms allow the handler to run when extraction fails.
These wrappers do not rewind the request body. If a wrapped body extractor reads part or all of the stream before failing, a later extractor sees only what remains. In practice, keep it as the handler’s only body extractor.
Custom Extractors
You can turn application-specific request data into a handler argument by
implementing FromRequest. A custom extractor receives the application
state, the HTTP request, and mutable access to the request body.
The body argument is the raw http::Payload stream. The Payload
extractor in web::types is a different type: a handler argument that wraps
this stream. Custom extractors use http::Payload directly:
use ntex::http::Payload;
use ntex::web::{self, App, FromRequest, HttpRequest, InternalError};
struct ClientName(String);
impl<St: web::State> FromRequest<St> for ClientName {
type Error = InternalError<&'static str>;
async fn from_request(_: &St, req: &HttpRequest, _: &mut Payload) -> Result<Self, Self::Error> {
req.headers()
.get("x-client-name")
.and_then(|value| value.to_str().ok())
.map(|value| ClientName(value.to_owned()))
.ok_or_else(|| web::error::ErrorBadRequest("Missing client name"))
}
}
async fn hello(client: ClientName) -> String {
format!("Hello, {}!", client.0)
}
App::default().route("/", web::get().to(hello));
Handlers registered with Route::to_with_state() receive the application
state and request-local state before their extractor arguments. Extractors
otherwise behave in the same way. They can use application state, but they do
not receive the request-local state carried by WebRequest. See
Web Application State for details.
Responders
Handler return values are responders. Any type that implements Responder
can be returned from a handler, and ntex turns it into an HTTP response.
Responder conversion happens after the handler completes and receives both
the application state and the original request.
Common responder types include:
| Responder | Result |
|---|---|
HttpResponse or HttpResponseBuilder | The response as configured |
String, &String, or &'static str | 200 OK, text/plain; charset=utf-8 |
Bytes, BytesMut, or &'static [u8] | 200 OK, application/octet-stream |
Json<T> | 200 OK, application/json |
Form<T> | 200 OK, application/x-www-form-urlencoded |
Option<T> | The inner response, or 404 Not Found for None |
Result<T, E> | The inner response, or the error response produced by E |
(T, StatusCode) | The inner response with a different status code |
Either<A, B> | The response of whichever variant is returned |
InternalError<T> | The error response with its configured status |
() | 200 OK with an empty body, for applications with () state |
The error type in Result<T, E> must implement WebResponseError. See
Web Application Errors for how errors become responses.
Either is available as ntex::util::Either.
Use Responder::with_status() and Responder::with_header() when you only
need to change the status or set a header:
use ntex::http::StatusCode;
use ntex::web::{self, App, Responder};
#[derive(serde::Serialize)]
struct User {
id: u32,
}
async fn create_user() -> impl Responder {
web::types::Json(User { id: 42 })
.with_status(StatusCode::CREATED)
.with_header("x-api-version", "1")
}
App::default().route(
"/users",
web::post().to(create_user),
);
These methods wrap the original responder. They keep its body and other headers, then apply the requested status and headers. A configured header replaces values with the same name that the inner responder produced.
For simple handlers, returning a concrete responder such as Json<T> keeps
the signature clear. Use impl Responder when the concrete wrapper type is
unimportant, and Either<A, B> when different branches naturally produce
different responder types.
Custom Responders
Implement Responder when one of your own types should be returned directly
from a handler. The conversion is asynchronous and can inspect application
state or request metadata:
use ntex::http::Response;
use ntex::web::{self, App, HttpRequest, Responder};
struct Greeting(&'static str);
impl<St: web::State> Responder<St> for Greeting {
async fn respond_to(self, _: &St, _: &HttpRequest) -> Response {
Response::Ok()
.content_type("text/plain; charset=utf-8")
.body(self.0)
}
}
async fn hello() -> Greeting {
Greeting("Hello!")
}
App::default().route("/", web::get().to(hello));
Web Application Errors
Web applications often need a consistent error format, even when failures come
from reusable components or external libraries. ntex handles this through an
error domain. Each application state selects its domain with the
State::Error associated type, and every error that can reach the application
must implement WebResponseError for that domain.
The domain type is a marker: it names a set of rendering rules, and it does not have to be an error that handlers return. It can be the application’s own error enum, as in the example below, or a separate empty type.
This means handlers can return ordinary Result values instead of building
error responses by hand. The exact conversion point depends on where the error
comes from: extractor and handler errors become responses inside the route,
while errors returned by filters, middleware, and lower-level services remain
service errors until they leave the application stack.
DefaultError provides plain-text responses for common ntex errors. If your
API needs JSON bodies, extra headers, error codes, or other metadata, define a
custom error domain and choose how each error is rendered. The compiler checks
that every error used by the application has a renderer for the selected
domain, so a missing implementation prevents the application from compiling.
Default Error Handling
Every application state has an error type. You choose it through the State
trait:
use ntex::web;
#[derive(Clone)]
struct AppState;
impl web::State for AppState {
type Error = web::DefaultError;
}
Here, the application uses DefaultError. It already knows how to handle
common ntex errors. The response body is the error’s Display text, sent as
text/plain:
- malformed request data, such as invalid query strings, JSON or form bodies,
and a wrong
Content-Type, returns400 Bad Request; - JSON and form bodies that exceed their configured limit return
413 Payload Too Large; - a malformed form
Content-Lengthreturns411 Length Required; - path segments that fail to deserialize return
404 Not Found; - response serialization failures, such as a
Json<T>responder failing to encode its value, return500 Internal Server Error.
The Bytes and String extractors use their own PayloadError; its default
mapping is 400 Bad Request, including when the configured body limit is
exceeded.
If your application uses () or AppState<T> as its state, DefaultError
is selected automatically. In that case, there is usually nothing else to
configure. DefaultError is only a marker selecting these implementations; it
is not constructed at runtime, and it does not make every Rust error
renderable. A custom error returned by your code still needs a
WebResponseError<St, DefaultError> implementation or an InternalError
wrapper.
Returning Errors from Handlers
A handler can return Result just like any other async Rust function. ntex
uses the successful value as the response, or turns the error into an error
response.
For simple cases, the helpers in web::error let you choose an HTTP status
without defining a new error type:
use ntex::web::{self, App, InternalError};
async fn create_user(
) -> Result<&'static str, InternalError<&'static str>> {
Err(web::error::ErrorBadRequest("Invalid user"))
}
fn main() {
App::default().route(
"/users",
web::post().to(create_user),
);
}
web::error has 39 Error* helpers, one for each common 4xx and 5xx status,
such as ErrorBadRequest, ErrorUnauthorized, ErrorForbidden,
ErrorNotFound, ErrorConflict, and ErrorInternalServerError. Each returns
an InternalError<T>, which renders the given status with T’s Display text
as the body.
These helpers are convenient at an application boundary, but their body is public. Avoid wrapping database errors, tokens, internal paths, or other sensitive details directly. Log the original error and return a message that is safe for the client.
Under the hood, the successful value must implement Responder, while the
error must implement WebResponseError for the application’s domain. The
error is rendered as soon as the handler returns, inside the Result
responder.
Some error types implement WebResponseError for every domain, so they work
with DefaultError and custom domains alike:
InternalError<T>, including everyError*helper;WebError, an error that has already been converted for the domain;ntex::util::Either<A, B>, when bothAandBimplement it;std::convert::Infallible.
InternalError::from_response() is useful when a call site needs a completely
prepared response rather than a status and plain-text body:
use ntex::web::{self, HttpResponse, InternalError};
async fn create_user(
) -> Result<&'static str, InternalError<&'static str>> {
Err(InternalError::from_response(
"duplicate user",
HttpResponse::Conflict()
.header("x-error-code", "user-exists")
.body("A user with that name already exists"),
))
}
Custom Error Domains
For a larger application, you may want all errors to follow rules of your own.
Set State::Error to an application-specific type, then implement
WebResponseError for each error that the application can produce.
This gives you full control over how errors are rendered, including errors
from external components.
WebResponseError requires std::error::Error + 'static, so every error type
must implement the standard Error trait. The example below derives it with
the thiserror crate, which must be added to your Cargo.toml.
The following example returns 409 Conflict when a user already exists and
400 Bad Request when the request contains invalid JSON:
use ntex::web::{self, App, HttpResponse, WebResponseError, error::JsonPayloadError};
#[derive(Clone)]
struct AppState;
#[derive(Debug, thiserror::Error)]
enum ApiError {
#[error("User already exists")]
UserExists,
}
impl web::State for AppState {
type Error = ApiError;
}
impl WebResponseError<AppState, ApiError> for ApiError {
fn error_response(&self, _state: &AppState) -> HttpResponse {
HttpResponse::Conflict().body(self.to_string())
}
}
impl WebResponseError<AppState, ApiError> for JsonPayloadError {
fn error_response(&self, _state: &AppState) -> HttpResponse {
match self {
JsonPayloadError::Overflow => {
HttpResponse::PayloadTooLarge().body("Request body is too large")
}
JsonPayloadError::ContentType => {
HttpResponse::UnsupportedMediaType().body("Expected a JSON request")
}
_ => HttpResponse::BadRequest().body("Invalid JSON"),
}
}
}
#[derive(serde::Deserialize)]
struct User {
name: String,
}
async fn create_user(
user: web::types::Json<User>,
) -> Result<String, ApiError> {
if user.name == "admin" {
Err(ApiError::UserExists)
} else {
Ok(format!("Created {}", user.name))
}
}
fn main() {
App::<AppState>::new().route(
"/users",
web::post().to(create_user),
);
}
There are two possible errors to account for:
- the handler can return
ApiError; - the
Json<User>extractor can returnJsonPayloadErrorbefore the handler is called.
JsonPayloadError is defined by ntex, yet the application can still implement
WebResponseError<AppState, ApiError> for it. Rust’s orphan rules allow this
because AppState and ApiError are local types.
Selecting a custom domain replaces the DefaultError implementations for that
application. Add mappings for every extractor, responder, filter, middleware,
or service error the application uses. This is deliberate: it prevents an API
from silently falling back to a response format or status it did not choose.
Errors can also happen while ntex is creating the response. For example, if a
handler returns Json<T>, your error domain must know how to handle
serde_json::Error. A Form<T> responder similarly requires a mapping for
serde_urlencoded::ser::Error.
The same rule applies outside handlers. Errors from filters and middleware,
at the application, scope, or resource level, must also implement
WebResponseError for the domain. Unlike handler errors, they are wrapped in a
WebError and rendered when they reach the application service.
If a required conversion is missing, the route does not compile. This catches unhandled error cases while you are building the application rather than when a request arrives.
The error_response() method receives the application state, so it can use
shared configuration while building the response. Its default response is
500 Internal Server Error with the error’s Display text as a plain-text
body; override it to choose the status, headers, and body that make sense for
the error. For a JSON API, a response builder’s json() method can serialize a
small error-body struct here.
Where Error Responses Travel
Extractor errors, handler Result errors, and responder serialization errors
are converted while the route is running. They return through resource, scope,
and application middleware as ordinary responses, so response middleware can
inspect or modify them.
Filters, middleware services, and custom web services return errors through
the service layer. ntex type-erases those errors into WebError while they
unwind, then the outer application service calls error_response(). As a
result, middleware using ctx.call(...).await? sees an error, not the final
rendered response. If it must add headers or otherwise inspect that response,
it needs to catch the error and convert it itself.
For normal handlers, return Result<T, E> and let the responder handle the
conversion. Construct WebError directly only when writing lower-level web
services that use the service error path.
Test the Rendered Response
Error handling is part of the HTTP contract, so test the status, headers, and body rather than only testing the Rust error value:
use ntex::http::StatusCode;
use ntex::web::{self, App};
#[derive(Clone)]
struct AppState;
impl web::State for AppState {
type Error = web::DefaultError;
}
async fn create_user(
) -> Result<&'static str, web::InternalError<&'static str>> {
Err(web::error::ErrorConflict("User already exists"))
}
#[ntex::test]
async fn duplicate_user_response() {
let app = web::test::init_service_st(
AppState,
App::<AppState>::new().route(
"/users",
web::post().to(create_user),
),
)
.await;
let request = web::test::TestRequest::post()
.uri("/users")
.to_request();
let response = web::test::call_service(&app, request).await;
assert_eq!(response.status(), StatusCode::CONFLICT);
}
Byte Buffers
When a socket read contains two messages, a decoder needs to hand the first
message to the application while keeping the second for later. Copying each
message into a new Vec<u8> works, but adds allocations and memory copies to
every request. ntex-bytes lets the decoder hand out an owned view instead.
The view keeps the data alive without borrowing the decoder’s read buffer.
The same idea helps on the write side: a small header and a large body can be queued separately, rather than copied into one growing buffer. For tiny values, sharing would cost more than copying, so the crate stores short byte strings directly in their handles. For frequently reused I/O buffers, per-thread page caches reduce trips to the allocator.
The ntex crate re-exports the main types in ntex::util: Bytes,
BytesMut, ByteString, BytePages, BytePage, BytePageSize,
Buf and BufMut. Page cache tuning, set_page_cache_size, is
available from ntex_bytes directly.
The mental model
Start with three types. Use BytesMut while filling or changing a contiguous
buffer, Bytes once data is ready to share, and BytePages when output can
remain in separate chunks. ByteString is the text version of Bytes: it
adds the guarantee that the bytes are valid UTF-8.
The important distinction is between a handle and its allocation.
Several handles can keep one allocation alive, but they do not all have
permission to change it. A BytesMut owns the writable tail; Bytes views
can own earlier, read-only portions. Splitting off a message gives it an
independent lifetime, not necessarily independent storage.
Here, “zero-copy” means avoiding a copy between application buffers. It does not promise that the operating system or a TLS layer will avoid copying. It is also an optimization, not a rule: copying a handful of bytes is often cheaper than maintaining another shared reference.
If you remember only one rule from this chapter, make it this:
Share data while it is moving through the pipeline; copy or trim the small part that needs to live much longer than the buffer it came from.
If you are unsure where to start, use a BytesMut for input and message
construction, hand completed data out as Bytes, and reach for BytePages
only when the consumer can accept several chunks. The rest of this chapter
explains when those defaults are worth changing.
Types at a glance
| Type | Mutability | Clone cost | Typical use |
|---|---|---|---|
Bytes | Immutable | Usually shares storage or copies inline data; external storage decides | Decoded frames, payloads, header values |
BytesMut | Unique | Copies the data | Read buffers, building output |
ByteString | Immutable | Same as Bytes | UTF-8 text: header names, paths, topics |
BytePages | Unique | Shares or copies according to page storage and append rules | Write queues |
BytePage | Immutable | Shares data, Vec pages copy | One chunk of a BytePages queue |
BytePageSize | - | - | Allocation size classes |
Bytes
Bytes is an immutable view into contiguous memory. Heap-backed values
normally share a reference-counted allocation when cloned, sliced or split.
Small values and static data use simpler representations, described below.
Bytes is Send + Sync, so it can be passed between threads.
use ntex_bytes::Bytes;
let mut msg = Bytes::copy_from_slice(&[1u8; 1024]);
// These two views share the original allocation.
let first = msg.split_to(256);
assert_eq!(first.len(), 256);
assert_eq!(msg.len(), 768);
let part = msg.slice(100..200); // shares the allocation too
let cloned = msg.clone(); // another handle, not another 768-byte buffer
drop(msg);
assert_eq!(part.len(), 100); // the other handles still own their data
assert_eq!(cloned.len(), 768);
A Bytes value uses one of four storage kinds:
| Kind | Created by | Clone |
|---|---|---|
| Inline | Data of at most 23 bytes (11 on 32-bit targets) | Copies the bytes, no allocation |
| Static | Bytes::from_static, From<&'static [u8]>, From<&'static str> | Copies the pointer |
| Shared | A heap buffer, usually from BytesMut | Increments the reference count |
| External | Bytes::from_ext, Arc<str> via ByteString | Delegated to the owner’s vtable |
Inline storage keeps small values inside the Bytes struct itself.
Operations such as Bytes::slice, Bytes::split_to,
BytesMut::split_to and BytesMut::freeze copy small results inline
instead of taking another reference to the heap buffer. A short header name
can therefore outlive a request without keeping its read buffer alive.
Bytes::is_inline reports the storage kind. Static constructors keep their
static representation even for short strings.
The examples below assume a 64-bit target, where the inline limit is 23 bytes. On 32-bit targets it is 11 bytes. These limits and the layouts later in the chapter describe the current implementation, not a stable memory-layout API.
When sharing retains too much memory
A view that is too large to inline keeps the whole allocation alive, even if it uses only a small part. Keeping a 100-byte token from a 16 KiB read buffer in a long-lived cache can therefore retain far more than 100 bytes. That is not a leak: the allocation is still owned. It is a lifetime tradeoff.
Bytes::trimdown helps at this boundary. It moves data inline when it
fits, or copies it into an exactly sized buffer when at least 64 bytes of
the backing capacity are unused. Inline and static values are left alone.
use ntex_bytes::Bytes;
let packet = Bytes::copy_from_slice(&[b'x'; 16 * 1024]);
let mut token = packet.slice(100..200);
token.trimdown(); // copy 100 bytes so `token` no longer retains the packet
drop(packet);
assert_eq!(token.len(), 100);
Trimming one view does not release the original allocation if other handles still own it. Use this deliberately for long-lived values, rather than trimming every frame and losing the benefit of sharing.
Externally owned data
External storage lets Bytes share memory owned by something else. A type
implements the unsafe StorageExt trait and provides a static vtable with
as_ptr, len, clone and drop functions. A clone that returns None
makes clones copy the data into native storage instead. Arc<str>
implements it, so ByteString::from(Arc<str>) shares the string without
copying it. This is an integration API, not something a normal decoder needs
to implement. The unsafe contract requires the exposed bytes to remain valid
and immutable for every live handle, including across threads; the vtable
must clone and release ownership correctly.
BytesMut
BytesMut is a unique, growable view into a heap buffer, similar to
Vec<u8>. Only the BytesMut handle can write into its part of the buffer,
but the same allocation can be shared with Bytes values split off its
front. Think of the allocation as a line moving through a decoder:
| completed frame | completed frame | unread data | spare capacity |
|<------ immutable Bytes views ---->|<--------- BytesMut -------->|
Only one BytesMut can own the writable part of a given allocation. Once a
prefix is handed out as Bytes, the mutable view advances past it and cannot
overwrite it. The decoder can continue filling the tail while the application
reads earlier frames, without locks around the data.
use ntex_bytes::{BufMut, BytesMut};
let mut buf = BytesMut::with_capacity(1024);
buf.put_slice(b"hello ");
buf.put_u16(0x1234);
buf.extend_from_slice(b"world");
// 6 bytes fit inline, so `head` is a copy and `buf` stays unique
let head = buf.split_to(6);
assert!(head.is_inline());
assert!(buf.is_unique());
buf.extend_from_slice(&[0u8; 100]);
let body = buf.take(); // shares the allocation with `buf`
assert!(!buf.is_unique());
drop(body);
assert!(buf.is_unique());
// a unique buffer gets its whole capacity back
buf.clear();
assert_eq!(buf.capacity(), 1024);
The main operations:
| Method | Effect |
|---|---|
BytesMut::split_to | Returns the first at bytes as Bytes; shares large results and copies small ones inline |
BytesMut::take | Returns all data as Bytes, keeping the spare capacity for self |
BytesMut::freeze | Consumes the mutable handle; reuses heap storage or copies a small result inline |
BytesMut::advance_to | Drops the first cnt bytes, O(1) |
BytesMut::clear | Empties the buffer, reclaims the full capacity if unique |
BytesMut::reserve | Ensures spare capacity, see Growth |
BytesMut::is_unique | true if no Bytes shares the allocation |
BytesMut::from(Bytes) | Reuses a unique shared heap allocation; copies inline, static, external or still-shared data |
Three similar-looking operations cover most handoff points:
- Use
split_to(n)after decoding a complete prefix and keep reading from the remainder. - Use
take()when all current bytes form one message but the mutable handle should stay with the connection and reuse its remaining capacity. - Use
freeze()when the mutable buffer itself is finished. It consumes the handle and turns its current data intoBytes.
None is universally faster. Small results are copied inline; larger results usually share the allocation. Choose the operation that matches ownership first, then measure if that path matters.
BytesMut::new currently starts with 112 bytes of data capacity (a
128-byte allocation including its header),
BytesMut::with_capacity allocates exactly the requested capacity, and
BytesMut::with_page_size takes a pooled page, see
Pooled buffers. Cloning a BytesMut always copies the
data.
Length, capacity and spare space
len() counts initialized bytes. capacity() counts what the current
mutable view can hold, including those bytes, and the difference is spare
space at the end. Capacity is not necessarily the size of the whole
allocation: splitting or advancing the front shrinks the mutable view.
Dropping a frame makes its space eligible for reuse, but does not immediately
move the read cursor backward. A later clear() reclaims the whole allocation
if it is unique; reserve() can move unread data to the front when more space
is needed. If another handle is still alive, that memory cannot be overwritten.
clear() means “discard my data”, not “free my allocation”. This is useful
for a buffer that will be filled again. Drop the buffer when it is no longer
needed; whether its memory is freed or cached depends on how it was allocated.
A buffer’s typical journey
In a protocol decoder, one allocation often goes through the same cycle many times:
- The transport writes socket data into the spare tail of a
BytesMut. - The decoder examines the initialized prefix.
- A complete frame is removed with
split_to. - The returned
Bytestravels through services and application code. - The
BytesMutstays with the connection and receives more input. - When all shared frames are dropped,
reserveorclearcan reclaim the space at the front.
This is why the mutable buffer can remain useful even while immutable frames are alive: they own disjoint views. It is also why holding one frame for a long time can change allocation behavior for that connection. The decoder remains correct, but it may need another allocation when the tail fills.
BytesMut implements both the ntex BufMut trait and the BufMut trait
of the bytes crate. They differ in one detail: ntex’s
BufMut::remaining_mut() reports the current spare capacity, while the
bytes version reports usize::MAX - len because the buffer grows on demand.
All put_* methods reserve capacity as needed in both cases.
Zero-copy decoding
Splitting frames off the read buffer is the core pattern of ntex codecs. The decoder below parses a two-byte, big-endian length followed by a payload. Large frames share the read buffer; tiny ones are copied inline:
use ntex_bytes::{Bytes, BytesMut};
fn decode(src: &mut BytesMut) -> Option<Bytes> {
if src.len() < 2 {
return None;
}
let len = u16::from_be_bytes([src[0], src[1]]) as usize;
if src.len() < 2 + len {
// make room for the rest of the frame
src.reserve(2 + len - src.len());
return None;
}
src.advance_to(2);
Some(src.split_to(len))
}
let mut src = BytesMut::copy_from_slice(b"\0\x03one");
assert_eq!(decode(&mut src).unwrap(), b"one"[..]);
assert!(src.is_empty());
The length field is treated as untrusted input. The decoder checks that the
complete frame is present before calling advance_to and split_to, because
those methods panic when asked to move past the end. If an index is not already
validated, use BytesMut::split_to_checked instead:
use ntex_bytes::{Bytes, BytesMut};
fn take_prefix(src: &mut BytesMut, len: usize) -> Option<Bytes> {
src.split_to_checked(len)
}
let mut src = BytesMut::copy_from_slice(b"hello");
assert_eq!(take_prefix(&mut src, 2).unwrap(), b"he"[..]);
assert_eq!(src, b"llo"[..]);
Bounds checks are not a size limit. A real decoder should reject a declared
length above its protocol or application limit before calling reserve;
otherwise a peer can make the connection retain a very large buffer without
ever sending the promised frame.
While a frame is alive, the read buffer is not unique, and the frame keeps
the whole allocation alive, unless the frame was small enough to inline.
Once all shared frames are dropped, the read buffer is unique again, and a
later clear() or reserve() can reuse the space in front of its data.
If it must grow before then, it copies only its remaining data into another
buffer; existing frames remain valid in the old one.
ByteString
ByteString is an immutable UTF-8 string backed by Bytes. It has the same
storage kinds and cost model, and dereferences to &str.
use ntex_bytes::{ByteString, Bytes};
let header = ByteString::from("content-type: text/plain");
let name = header.slice(0..12); // copies this short result inline on 64-bit targets
assert_eq!(name, "content-type");
// validates UTF-8, takes over the bytes without copying
let s = ByteString::try_from(Bytes::from_static(b"utf-8 text")).unwrap();
assert_eq!(s.as_str(), "utf-8 text");
Indices are byte offsets, not character positions. slice, split_to and
split_off panic if an index is not on a UTF-8 character boundary, so take
care when working with non-ASCII text. ByteString::from_bytes_unchecked
is unsafe: it skips validation, and the caller must guarantee valid UTF-8.
Use the checked conversion unless that guarantee is already established.
With the simd feature, UTF-8 validation uses SIMD instructions.
Converting an owned Bytes or BytesMut into ByteString validates UTF-8
and then reuses the bytes, apart from the usual small-value inlining.
Converting back with ByteString::into_bytes is free. Converting a
String or Vec<u8> copies into native byte storage;
ByteString::from(Arc<str>) is the exception and shares the Arc.
use ntex_bytes::{ByteString, BytesMut};
let bytes = BytesMut::copy_from_slice("hello".as_bytes()).freeze();
let text = ByteString::try_from(bytes).unwrap();
assert_eq!(text, "hello");
let bytes = text.into_bytes();
assert_eq!(bytes, b"hello"[..]);
ByteString can retain a larger allocation for the same reason as Bytes.
Call ByteString::trimdown when a short string is being promoted from a
request buffer into long-lived state.
Buf and BufMut
Buf reads from a buffer with a cursor, BufMut writes into one. They
provide the get_* and put_* helpers for integers and floats in big- and
little-endian order. The unsuffixed integer methods use big-endian order;
methods ending in _le use little-endian order.
use ntex_bytes::{Buf, BufMut, BytesMut};
let mut buf = BytesMut::with_capacity(16);
buf.put_u32(0xDEAD_BEEF);
buf.put_u8(1);
let mut b = buf.freeze();
assert_eq!(b.get_u32(), 0xDEAD_BEEF);
assert_eq!(b.get_u8(), 1);
assert!(!b.has_remaining());
Reading advances the cursor: after get_u32(), those four bytes are no
longer part of the remaining input. Reads panic if there are not enough bytes,
so a decoder should check remaining() or len() before reading a field.
The traits work with more than these owned buffer types: a &[u8] can be a
read cursor too, and a generic Buf need not store all its data contiguously.
Bytes and BytesMut also implement the Buf trait of the bytes
crate, and BytesMut its BufMut trait, so they work with libraries built
on bytes. BytesMut and BytePages implement std::io::Write, and
BytesMut implements std::fmt::Write, so write! can format into them.
Bytes and ByteString implement serde’s Serialize and Deserialize.
Panicking and checked operations
Methods such as slice, split_to, split_off, advance_to and the
get_* family assume their indices or lengths have already been validated.
They panic on an invalid range or insufficient input. This keeps a codec’s
hot path simple after it has performed its bounds checks.
For values derived directly from input, prefer the checked variants where
available: Bytes::slice_checked, Bytes::split_to_checked,
Bytes::split_off_checked and BytesMut::split_to_checked. A failed
check then becomes None instead of a process-level panic.
Internal organization
You do not need the layout details to use the types, but they explain two design choices: why tiny values are copied and why a mutable buffer can share an allocation without sharing its writable bytes.
Handles
Bytes is three machine words: a data pointer, a length, and an offset
word. The two low bits of offset select the storage kind, the remaining
bits hold the distance from the start of the heap buffer to the data
pointer. Inline storage reuses all three words for data and keeps the kind
and length in one byte, which is where the 23-byte inline capacity comes
from.
BytesMut is a single pointer to its heap buffer. Its view, the start and
length of its data, is stored in the buffer header, so BytesMut is one
word in size.
Heap buffers
Each native shared heap buffer is one allocation: a 16-byte header followed
by the data. External storage and Vec-backed pages have their own layouts.
header (16 bytes) data
+--------+-----+----------+-----------+---------+---------+------------+-------+
| offset | len | capacity | ref_count | Bytes 1 | Bytes 2 | BytesMut | spare |
+--------+-----+----------+-----------+---------+---------+------------+-------+
^ ^ ^ ^
allocation data start offset offset + len
| Field | Meaning |
|---|---|
offset | Start of the BytesMut view, from the beginning of the allocation |
len | Length of the BytesMut view |
capacity | Data capacity of the whole allocation, never changes while it is shared |
ref_count | Number of handles: the BytesMut plus every Bytes view |
Everything else is derived from these fields. The spare capacity after the
BytesMut view is 16 + capacity - offset - len, and the
BytePageSize class follows from the allocation size 16 + capacity,
see Pooled buffers.
Lengths, capacities and offsets use u32; the reference count is an
AtomicU32. A native heap buffer holds at most u32::MAX - 16 bytes of data,
just under 4 GiB. Requests for larger capacities panic.
Each Bytes view stores its own pointer and length, and computes the header
address from its offset. Views split off a BytesMut cover memory in front
of the BytesMut view, so the BytesMut never writes to memory that a
Bytes can read.
Reference counting and threads
The reference count is atomic and follows the same memory ordering as
std::sync::Arc: a Release decrement, and an Acquire fence before the
last handle frees the buffer. BytesMut::is_unique uses an Acquire
load, so after it returns true, all accesses through dropped handles,
including handles dropped on other threads, happen before the buffer is
reused. Apart from the reference count, Bytes handles read only the
capacity field of the header, which does not change while the buffer is
shared, so the BytesMut handle can update its view without
synchronization.
Allocation strategy
Regular buffers
BytesMut::with_capacity makes one allocation of 16 + capacity bytes,
and BytesMut::capacity is exactly the requested capacity. Buffers
created this way, by copy_from_slice, and by conversions such as
Bytes::from(Vec<u8>) have no page
size, unless the requested capacity is exactly a page capacity. They are
freed when their last handle is dropped.
In particular, Bytes::from(Vec<u8>) copies the data into the crate’s storage
(or inline for small values); it does not adopt the vector allocation.
The shared-buffer header needs space before the data. If you are building
data specifically to share as Bytes, starting with a BytesMut avoids
that conversion copy.
Growth
Appending data reserves capacity as needed. reserve(additional) asks for
room for that many bytes after the current data, not for a total capacity.
BytesMut::reserve tries these options in order:
- If the buffer has enough spare capacity, nothing happens.
- If the buffer is unique and the whole allocation is large enough, the data
is moved to the front of the allocation, reusing the space of dropped
Bytesviews. - If the buffer is unique and has no page size, the allocation is grown with
realloc, which can often extend it in place. - If the buffer has a page size, its data moves to a pooled page, see Pooled buffers.
- Otherwise, the data is copied into a new allocation without a page size.
Bytesviews split off the buffer keep the old allocation alive.
When growth needs more storage, the target capacity is at least twice the
current length, so repeatedly appending small amounts does not allocate on
every append. realloc may extend an allocation in place, but it may also
move it; the buffer makes no promise that its address stays the same.
The three reservation methods serve different purposes:
| Method | Use it when |
|---|---|
BytesMut::reserve | You want room for more bytes and expect the buffer to keep growing. |
BytesMut::reserve_exact | You know how many more bytes you need and want to avoid the doubling policy. The allocation is exactly the new capacity plus the 16-byte header; pooled buffers are not rounded up to a page class and keep one only if the new capacity is a page capacity. |
BytesMut::reserve_more | You don’t know how much more is coming, for example the next read. If less than BytePageSize::low remains (1 KiB for Size16, 4 KiB without a page size), the allocation is reused when it holds the data plus half_capacity(): a unique buffer is compacted in place, a shared page moves to a page of the same class. Otherwise a pooled buffer moves to the next page class and any other buffer grows by its capacity, by at least 112 bytes and at most 64 KiB. |
Neither reserve nor reserve_exact is a general-purpose shrinking
operation: if there is already enough spare capacity, it leaves the buffer
alone.
Page sizes
BytePageSize defines the size classes for pooled buffers. A page
allocation requests exactly the class size, including the header, rather
than a class-sized payload plus extra metadata. This avoids overshooting
those useful size boundaries, though an allocator’s actual size classes
are its own implementation detail. Data capacity is the class size minus
the 16-byte header. These are allocation categories, not operating-system
virtual-memory pages.
| Class | Allocation | Capacity | half_capacity() | low() | Cached pages by default |
|---|---|---|---|---|---|
Size4 | 4 KiB | 4,080 bytes | 2 KiB | 256 bytes | 128 |
Size8 | 8 KiB | 8,176 bytes | 4 KiB | 512 bytes | 64 |
Size16 | 16 KiB | 16,368 bytes | 8 KiB | 1 KiB | 64 |
Size24 | 24 KiB | 24,560 bytes | 12 KiB | 1.5 KiB | 32 |
Size32 | 32 KiB | 32,752 bytes | 16 KiB | 2 KiB | 16 |
Size48 | 48 KiB | 49,136 bytes | 16 KiB | 3 KiB | 8 |
Size64 | 64 KiB | 65,520 bytes | 16 KiB | 4 KiB | 16 |
Size128 | 128 KiB | 131,056 bytes | 16 KiB | 8 KiB | 2 |
Size256 | 256 KiB | 262,128 bytes | 16 KiB | 16 KiB | 1 |
Unset | - | 65,520 bytes | 16 KiB | 4 KiB | never cached |
Size16 is the default class. BytePageSize::for_capacity returns the
smallest class that holds a given capacity, or Unset above the largest
data capacity (262,128 bytes).
BytePageSize::next and BytePageSize::prev step between classes.
half_capacity() is the recommended write-buffer threshold for a page size,
low(), 2/32 of the class size, is the free-capacity threshold of
BytesMut::reserve_more.
The enum is #[non_exhaustive], so more classes may be added.
use ntex_bytes::BytePageSize;
assert_eq!(BytePageSize::for_capacity(100), BytePageSize::Size4);
assert_eq!(BytePageSize::for_capacity(20_000), BytePageSize::Size24);
assert_eq!(BytePageSize::Size16.next(), BytePageSize::Size24);
assert_eq!(BytePageSize::Size256.next(), BytePageSize::Unset);
Pooled buffers
BytesMut::with_page_size takes a page of the given class from the
current thread’s page cache, or allocates a new one if the cache is empty.
with_page_size(BytePageSize::Unset) takes a Size64 page, the two have
the same allocation size. BytesMut::page_size reports the class of a
buffer.
The header does not store the class. A buffer belongs to a class when its
allocation, header included, is exactly the class size, so the class is
looked up from the capacity field. Any buffer of a page size is pooled,
including one created by with_capacity or copy_from_slice with exactly
a page capacity, or grown into one by realloc.
When the last handle to a page is dropped, whether it is the BytesMut or
a Bytes view split off it, the page returns to the cache of its class on
the thread that drops it. If that cache is full, the page is freed. Only
buffers with a page size are cached. Buffers with Unset are always freed.
A pooled buffer that grows moves to a page of the smallest class that fits
the new capacity, but never to a smaller class than its current one. The
data is copied, and the old page returns to the cache once its Bytes
views are dropped. If the growth target exceeds the largest page’s data
capacity, the buffer becomes a regular buffer without a page size. From then
on it grows like any regular buffer and is freed when its last handle drops.
use ntex_bytes::{BytePageSize, BytesMut};
let mut buf = BytesMut::with_page_size(BytePageSize::Size4);
assert_eq!(buf.capacity(), 4096 - 16);
buf.extend_from_slice(&[0; 5000]);
assert_eq!(buf.page_size(), BytePageSize::Size8);
buf.extend_from_slice(&[0; 10_000]);
assert_eq!(buf.page_size(), BytePageSize::Size16);
// beyond the largest class, a regular buffer
buf.reserve(300 * 1024);
assert_eq!(buf.page_size(), BytePageSize::Unset);
Small results are inlined, which affects pooled buffers too: freezing a
pooled buffer with at most 23 bytes of data copies them inline. If no earlier
Bytes views still share the allocation, the page can return to the cache
right away.
use ntex_bytes::{BytePageSize, BytesMut};
let mut buf = BytesMut::with_page_size(BytePageSize::Size32);
buf.extend_from_slice(b"small");
let small = buf.freeze();
assert!(small.is_inline());
assert_eq!(small, b"small"[..]);
// `small` no longer needs the page; a later allocation can reuse it.
Pages handed to another thread return to that thread’s cache. A thread that only receives data, for example a thread that writes data produced on other threads, fills its cache up to the limits and frees the rest. While a thread’s thread-local storage is being destroyed, the cache is unavailable, and pages released then are freed.
Tuning the page cache
Caching trades retained memory for fewer allocations. An idle worker can still hold cached pages, but those pages are ready for reuse; they are not live messages.
set_page_cache_size sets the number of cached pages of one class for the
current thread. The defaults, listed in the table above, cache more pages of
small classes and fewer of large ones. The largest shares go to the classes
that I/O uses most: 4 to 16 KiB pages, where connection reads start and the
default write page; and 64 KiB pages, where reads of busy connections grow
to. With all caches full, a thread retains about 5 MiB.
use ntex_bytes::{BytePageSize, set_page_cache_size};
// retain more large pages for a thread that handles large messages
set_page_cache_size(BytePageSize::Size256, 4);
// stop returning 4 KiB pages to this thread's cache
set_page_cache_size(BytePageSize::Size4, 0);
The setting affects only the calling thread, so call it on every thread that
uses pooled buffers, for example at the start of each worker thread. A
smaller limit does not free pages already in the cache. They are reused, and
the limit applies when pages are released. The older set_pages_cache(),
which sets one limit for all classes, is deprecated.
These limits count cached pages, not pages in use. They do not bound the memory held by live requests or slow consumers. Start with the defaults, then adjust them using the sizes and concurrency of your actual workload. More workers also means more independent caches.
For example, the default cache can retain roughly 5 MiB per worker when every size class is full. Eight otherwise idle workers could therefore retain about 41 MiB in page caches. That may be a good trade when traffic returns quickly; for sparse workloads or many workers, smaller limits may be better.
BytePages
Suppose a response contains a short header, a large existing body and a
trailer. A contiguous buffer may have to move earlier output as it grows,
and copying the body into it would be wasteful. BytePages keeps a queue
of chunks instead: new bytes fill a current writable page, while owned
payloads can be queued separately. When the current page fills, another
page is started. Existing output stays where it is.
use ntex_bytes::{BufMut, BytePageSize, BytePages, Bytes};
let mut pages = BytePages::new(BytePageSize::Size4);
pages.extend_from_slice(b"HTTP/1.1 200 OK\r\n\r\n");
let body = Bytes::copy_from_slice(&[b'x'; 8192]);
pages.append(body); // moves the existing body into the queue, no extra copy
pages.put_slice(b"trailer");
assert_eq!(pages.num_pages(), 3);
// the transport consumes pages in order
while let Some(page) = pages.take() {
// write `page` to the socket
assert!(!page.is_empty());
}
How data enters the queue:
put_slice(),extend_from_slice()and the otherBufMutmethods copy data into the current page. New pages are taken from the page cache with the queue’s page size.BytePages::appendadds an owned buffer. If the current page holds data, buffers of up to 4 KiB are copied into it. Larger buffers become a separate page without copying, and the current page keeps its spare capacity for later writes. If the current page is empty, the buffer is added without copying, and a uniqueBytesMutwith spare capacity becomes the new current page.BytePages::prependinserts a page at the front.
How data leaves it:
BytePages::takereturns the nextBytePage, with the current page last.BytePages::split_toandBytePages::split_intomove a byte prefix to another queue, splitting a page if needed.BytePages::freezereturns all data as oneBytes. A single page is converted according to its storage kind: native heap storage can be reused, tiny data may be inlined, andVecstorage is copied. Several pages are copied into one buffer.BytePages::copy_toleaves the source unchanged and appends cloned pages to another queue.BytePages::move_toleaves the source empty and appends its pages instead. Both followappend’s small-buffer copy rules;copy_toalso copiesVec-backed pages when cloning them.
A BytePage holds a Bytes, the storage of a BytesMut, or a Vec<u8>.
Splitting or advancing a Vec page copies its data, so prefer the other
kinds when output will be split or shared. Queuing an owned Vec directly
can still avoid an initial copy; that is different from converting the
vector to Bytes first.
Keep data paged for as long as the consumer accepts chunks. Calling
freeze() just before passing it to a chunk-aware transport would undo the
main benefit by combining all pages into one allocation.
Paging avoids copies; it does not bound queued output. A producer can still
outpace the transport and retain many pages, including large buffers appended
without copying. Protocol code should stop producing when its I/O layer
reports write backpressure instead of treating BytePages as an unlimited
queue.
The page size of a BytePages cannot be Unset, BytePages::new and
BytePages::set_page_size panic on it. BytePages::default() uses
Size16. The page list itself is also reused: each thread keeps up to 128
empty page lists for new BytePages values.
Changing the page size affects future page allocations, not data already queued or the current page. Page buffers are allocated lazily, so creating an empty queue does not immediately allocate a payload page.
Buffers in the I/O layer
The I/O Abstraction Layer builds on these types. Reads and writes share the same per-thread page cache:
- Reads. The transport reads into a pooled
BytesMutread buffer. Codecs split frames off it withsplit_to, so large decoded messages share the read buffer while tiny ones are copied inline. New read buffers use a per-connection page size that adapts to the read load between the min and max read sizes ofIoConfig,Size4andSize64by default. A buffer grows once less thanBytePageSize::low()of its page remains free, larger input moves it to bigger page sizes, beyondSize256it becomes an unpooled buffer that is freed when empty. - Writes. Encoders write into
BytePageswith the page size ofIoConfig::write_size,Size16by default.IoRef::encode_bytesappends owned buffers, so large payloads are not copied. The write threshold,half_capacity()of the page size by default, controls when a transport may start writing while output is still being produced.
An empty read buffer is released immediately, but its page returns to the
cache only when the last frame split from it is dropped. Holding a whole
request or frame therefore keeps the full page alive; copy the field instead
when only a small part needs to survive. set_page_cache_size tunes the
cache for read and write pages alike.
What memory usage means in practice
When investigating memory growth, separate three categories:
- Live data: bytes still owned by requests, responses, transports or application state.
- Retained allocations: a small live view keeps a larger shared allocation
alive.
trimdowncan help when the view is intentionally long-lived. - Cached allocations: empty pages or read buffers kept for reuse. Cache limits control this category, not live data.
Only the first category is directly proportional to pending work. The other two are performance tradeoffs and can make resident memory stay high after a traffic spike without indicating an ownership leak.
Performance guidelines
The fastest choice depends on how long the data lives, not just how large it is. A useful starting point is:
| Situation | Start with |
|---|---|
| Hand a decoded frame to its immediate consumer | split_to or take, so large payloads can share the input allocation. |
| Keep one field in a long-lived cache | A slice followed by trimdown, or an explicit copy, to avoid retaining an entire input buffer. |
| Build a contiguous message of a known size | BytesMut::with_capacity, or one reserve for the bytes still to be written. |
| Repeatedly build short-lived I/O buffers | with_page_size with a class that fits the usual size; cache hits avoid a fresh payload allocation. |
| Send an existing large body | BytePages::append, rather than copying it with put_slice. |
Queue an owned Vec<u8> for output | Append the Vec directly to BytePages; converting it to Bytes first copies it. |
| Use constant protocol bytes | from_static, which neither allocates nor copies the bytes. |
Drop shared frames when they are no longer needed, and do not equate “cheap to clone” with “cheap to retain”. A clone may cost only an atomic increment yet extend the lifetime of a large allocation. Conversely, an inline clone copies a few bytes but needs neither allocation nor reference counting. Measure allocation volume and retained memory alongside throughput when deciding whether to share, copy or cache.
Migrating from ntex 3 to ntex 4
ntex 4 requires Rust 1.97 and upgrades the service stack to ntex-service 5.
Most migration work involves the new service state and lifecycle APIs.
Server
Server service factories now create a Service directly. They no longer
create a ServiceFactory whose initialization configuration is SharedCfg.
The factory is called once per worker and receives a reference to that
worker’s application state.
For a server without custom application state, continue to use
ntex::server::build(). The factory receives &():
use std::io;
use ntex::http::{HttpService, Response};
use ntex::SharedCfg;
#[ntex::main]
async fn main() -> io::Result<()> {
ntex::server::build()
.bind(
"http",
"127.0.0.1:8080",
SharedCfg::default(),
async |_| {
HttpService::new(async |_| {
Ok::<_, io::Error>(Response::Ok().body("Hello"))
})
},
)?
.run()
.await
}
Use ntex::server::build_with_config() when each worker needs application
state. Its argument implements ServerAppConfig and creates the state for each
worker. An asynchronous closure can be used directly:
use std::io;
use ntex::SharedCfg;
#[derive(Clone)]
struct WorkerState;
#[ntex::main]
async fn main() -> io::Result<()> {
ntex::server::build_with_config(
async || Ok::<_, io::Error>(WorkerState),
)
.bind(
"service",
"127.0.0.1:8080",
SharedCfg::default(),
async |state: &WorkerState| {
let state = state.clone();
ntex::service::fn_service(move |_| {
let _state = state.clone();
async { Ok::<_, io::Error>(()) }
})
},
)?
.run()
.await
}
ntex-service 5 also removes Service::poll(). The Service trait now uses
the asynchronous ready() and shutdown() lifecycle methods, while service
chains provide readiness() and shutdown() callbacks. Pipeline and
middleware APIs have been updated to bind and propagate service state.
Pipeline::call_nowait() and PipelineBinding::call_nowait() have been
removed; use call(), which skips the readiness check when the last pipeline
readiness check has succeeded and no other call has consumed it.
Server builder
Several ServerBuilder and web HttpServer methods have been renamed or
removed:
| ntex 3 | ntex 4 |
|---|---|
maxconn() | max_connections() |
shutdown_timeout() | graceful_shutdown_timeout() |
HttpServer::maxconnrate() | HttpServer::max_tls_handshakes() |
ServerBuilder::config(), HttpServer::config() | pass SharedCfg to bind() or listen() |
on_worker_start(), on_accept() | build_with_config() or HttpServer::with_config() |
WorkerPool::shutdown_timeout() is also renamed to
graceful_shutdown_timeout().
Runtime features
The deprecated neon feature has been removed from ntex, ntex-rt and
ntex-net, the native runtime is used without it. Remove it from ntex and
direct ntex-rt or ntex-net dependencies. The neon-iocp feature has been
removed as well, Windows always uses IOCP. Select the tokio, compio,
neon-polling or neon-uring feature when a specific runtime backend is
required.
HTTP services
HttpService::openssl() and HttpService::rustls() have been replaced by the
http::openssl() and http::rustls() service wrappers.
OpenSSL:
let service = ntex::http::openssl(
acceptor,
ntex::http::HttpService::new(handler),
);
rustls now accepts the ALPN protocol names separately:
let service = ntex::http::rustls(
config,
&["h2", "http/1.1"],
ntex::http::HttpService::new(handler),
);
HTTP configuration
HttpServiceConfig::set_enable_headers_vec()is replaced byset_headers_vec(bool).h1::Codec::upgrade()has been removed; useRequest::upgrade()to detect upgrade andCONNECTrequests.- New settings:
set_half_close(),set_host_validation(),set_max_start_line_size(), andset_write_timeout().
URLs
ntex 4 uses urly::Url, re-exported as ntex::url, instead of http::Uri:
ntex::http::Uriand thentex::http::urimodule have been removed.RequestHead::uri,Request::uri(),HttpRequest::uri(), andWebRequest::uri()areUrlvalues.Url::path()andUrl::query()returnntex::url::Pathandntex::url::Query; useas_str()to get string slices.HttpRequest::path()andHttpRequest::query_string()still return&str.HttpRequest::url_for()returnsntex::url::Urlinstead ofurl::Url.- The
urlfeature has been removed; URL support is always available. - The HTTP and WebSocket clients accept any type that converts to
Url, including&str,String, andhttp::Uri. Client connectors useConnect<Url>instead ofConnect<Uri>. - Server request targets are normalized: dot segments are removed, a fragment
is dropped, and an origin-form path starting with
//is a path, not an authority.
HTTP client
The HTTP client builder, connector, and connection pool have been redesigned:
ClientBuildernow owns the connector configuration; the separate client connector builder has been removed.- Custom connector factories are replaced by connector
Servicevalues. ClientBuilder::build()is synchronous and takes aSharedCfg.- Client settings are stored in
ClientConfiginsideSharedCfg.
The default connector is configured automatically:
use ntex::SharedCfg;
use ntex::client::{Client, ClientConfig};
fn create_client() -> Client {
let cfg = SharedCfg::new("client").add(
ClientConfig::new()
.set_h1_connection_limit(8)
.set_response_payload_limit(256 * 1024),
);
Client::builder().build(cfg)
}
Use ClientBuilder::connector() or ClientBuilder::secure_connector() to
install a custom connector service.
Request defaults have moved from ClientBuilder to ClientConfig:
ntex 3 ClientBuilder | ntex 4 ClientConfig |
|---|---|
header() | set_header() |
basic_auth(), bearer_auth() | set_basic_auth(), set_bearer_auth() |
response_timeout(), disable_timeout() | set_response_timeout(), disable_timeout() |
response_payload_limit() | set_response_payload_limit() |
response_payload_timeout() | set_response_payload_timeout() |
The ClientConfig getters timeout(), payload_limit(), and
payload_timeout() are now response_timeout(), response_payload_limit(),
and response_payload_timeout().
ClientBuilder::disable_redirects(), max_redirects(), and
no_default_headers() have been removed.
The error variants ClientError::TunnelNotSupported, ConnectError::Timeout,
ConnectError::SslError, ConnectError::SslHandshakeError, and
EncodeError::Fmt have been removed. InvalidUrl::Http is replaced by
InvalidUrl::Parse.
WebSocket client
WebSocket client settings have moved to WsClientConfig. Construct
WsClient directly with the URI and configuration; a separate builder is no
longer required:
use ntex::ws::{WsClient, WsClientConfig};
let client = WsClient::new(
"ws://127.0.0.1:8080/ws",
WsClientConfig::new().set_max_frame_size(128 * 1024),
);
let connection = client.connect().await?;
URL validation errors are now reported by connect() rather than by
WsClient::new(), as WsConfigError::Parse instead of WsConfigError::Http.
Custom connectors and TLS are still selected with connector(), openssl(),
or rustls() on WsClient.
WsSink::on_disconnect() now returns ntex::io::Waiter<'static>; the
OnDisconnect future type has been removed.
Web applications
Web application state and error handling have been redesigned.
Application state in handlers
The application state type implements web::State. For common state that uses
the default web error type, web::AppState<T> provides a ready-made wrapper.
State is created once per worker with web::server_with_config() and is
passed to state-aware handlers by reference.
The web::types::State<St> extractor has been removed. Replace handlers that
use it with Route::to_with_state(). A state-aware handler receives:
- A shared reference to the application state.
- The request-local state.
- Any request extractors.
use std::io;
use ntex::{web, SharedCfg};
#[derive(Clone)]
struct TestAppState {
value: &'static str,
}
impl web::State for TestAppState {
type Error = web::DefaultError;
}
async fn index(
state: &TestAppState,
_request_state: (),
) -> web::HttpResponse {
web::HttpResponse::Ok().body(state.value)
}
#[ntex::main]
async fn main() -> io::Result<()> {
web::server_with_config(
async || {
Ok::<_, io::Error>(TestAppState {
value: "Hello",
})
},
async |_| {
web::App::<TestAppState>::new()
.route("/", web::get().to_with_state(index))
},
)
.bind("127.0.0.1:8080", SharedCfg::default())?
.run()
.await
}
The state-aware variants are available as web::to_with_state(),
Route::to_with_state(), Resource::to_with_state(), and
ResourceServices::to_with_state(). Continue to use to() when a handler only
needs request extractors.
The old App::state() model is replaced by service state for the primary
application state. Additional configuration values can be stored with
WebAppConfig::set_state() and retrieved with HttpRequest::app_state().
Request-local state is available through WebRequest::st() and
WebRequest::st_mut(). Filters and middleware can change its type with
WebRequest::map_state().
Custom implementations of the Handler trait must change call() from a
method returning impl Future to an async fn. Ordinary async fn handlers
do not need this change.
Error handling
The ErrorRenderer API has been removed. The state’s associated Error type
defines the application’s error type, and errors are rendered through
WebResponseError<St, Err>. Error rendering receives the application state
instead of an HttpRequest. Service initialization errors now use
ntex::error::Failure and IntoFailure.
Middleware order
Middleware registered with App::middleware(), Scope::middleware(), and
Resource::middleware() now runs in registration order. The first registered
middleware is the outermost one: it receives the request first and the
response last. In ntex 3 the order was reversed, and the last registered
middleware was the outermost one.
use ntex::web::{self, App, middleware};
App::default()
.middleware(middleware::DefaultHeaders::new().header("x-app", "example"))
.middleware(middleware::Logger::default())
.route("/", web::get().to(async || "Hello"));
In ntex 4, DefaultHeaders receives the request before Logger, and Logger
processes the response before DefaultHeaders. In ntex 3 the same code ran
Logger first. To keep the ntex 3 behavior, reverse the order of
middleware() calls.
Custom codecs and I/O
This section applies to code that implements codecs, filters, or dispatchers
directly on top of ntex::io and ntex::codec.
Codecs
ntex-codec 2 writes encoded data into BytePages. The deprecated
encode(BytesMut) method has been removed and encodev() has been renamed to
encode(), which is now required:
impl Encoder for MyCodec {
type Item = Bytes;
type Error = io::Error;
fn encode(&self, item: Bytes, dst: &mut BytePages) -> Result<(), io::Error> {
dst.append(item);
Ok(())
}
}
Decoder::decode_eof() is called when the peer closes the stream. Its default
implementation calls decode(); override it to decode a final frame that has
no terminator. If undecodable bytes remain after a clean EOF, Io::recv() and
the dispatcher report an io::ErrorKind::UnexpectedEof error.
IoConfig
set_disconnect_timeout()/disconnect_timeout()are renamed toset_shutdown_timeout()/shutdown_timeout(). A zero timeout panics.- Read and write buffers are
ntex-bytespages and share one page cache, configured per thread withntex_bytes::set_page_cache_size().BufConfigandIoConfig::read_buf()/write_buf()have been removed. - Each connection adapts its read page size to its reads, between the sizes set
with
set_read_size(min, max)(4 KiB to 64 KiB by default). - Backpressure is configured with a single high watermark and is released at half of it:
| ntex 3 | ntex 4 |
|---|---|
set_read_buf(high, low, cache_size) | set_read_backpressure(high) and set_read_size(min, max) |
set_write_buf(high, low, cache_size) | set_write_backpressure(high) |
set_write_page_size(), write_page_size() | set_write_size(), write_size() |
set_write_timeout()closes connections whose peer stops reading.
Io and IoRef
| ntex 3 | ntex 4 |
|---|---|
IoRef::force_close() | IoRef::terminate() |
IoRef::wants_shutdown() | IoRef::close() |
IoRef::with_read_buf() | IoRef::with_read_dst() |
IoRef::with_read_src_buf() | IoRef::with_read_src() |
IoRef::with_write_buf() | IoRef::with_write_src() |
IoRef::with_write_dst_buf() | IoRef::with_write_dst() |
IoRef::on_disconnect() returning OnDisconnect | returns Waiter<'static> |
Io::read_ready(), Io::poll_read_ready() | Io::read_more(), Io::poll_read_more() |
Io::poll_dispatch() | Io::register_dispatch() |
Io::pause() | removed; reads resume via poll_read_more() |
Io::set_config() | pass the configuration to Io::new() |
Io::take() | unsafe Io::take() |
IoRef::resize_read_buf(), IoContext::resize_read_buf() | BytesMut::reserve_more() |
IoStatusUpdate::KeepAlive | IoStatusUpdate::Timeout |
IoRef::is_closed() now reports whether closing has finished. Use the new
IoRef::is_active() to check whether the connection is still usable.
Dispatcher
Reason::KeepAliveTimeout is renamed to Reason::KeepAlive, and the new
Reason::WriteTimeout is reported when IoConfig::set_write_timeout()
expires.
Other API changes
ntex::rt:System::stop_on_panic()has been removed andBuilder::stop_on_panic()no longer has an effect. UseBuilder::panic_handling()or#[ntex::main(panic_handling = true)].System::set_latency_callback()requires aSend + Synccallback.ntex::time:query_system_time()has been removed; usesystem_time().ntex_util::channel::bstream::Receiver::max_buffer_size()is deprecated in favor ofset_watermarks().ntex::http::HeaderMapno longer implementsFromIterator; build maps withinsert()orappend().HeaderValue::to_str()accepts any valid UTF-8 value, not only visible ASCII.
Router
ntex::router (ntex-router 2) renames several APIs:
| ntex 3 | ntex 4 |
|---|---|
Router::recognize_mut_checked() | Router::recognize_checked_mut() |
ResourceDef::resource_path(), resource_path_named() | ResourceDef::build_path(), build_path_named() |
ResourceDef::name_mut() | ResourceDef::set_name() |
RouterBuilder::rdef() | RouterBuilder::resource() |
Path::unprocessed() | Path::path() |
Path::query() | Path::get() |
Router::build(), RouterBuilder::finish() | removed |
RouterBuilder registration methods return &mut RouterEntry instead of a
tuple; use set_id(), set_name(), set_check_value(), and
resource_mut(). Path::skip() takes a u32. Path segments that do not
decode to valid UTF-8 are kept percent-encoded.
Connection and protocol configuration
SharedCfg remains the container for connection and protocol configuration,
but it is separate from worker application state:
- Pass
SharedCfgto serverbind()orlisten()methods. - Pass
SharedCfgtoClientBuilder::build(). - Use
build_with_config()orweb::server_with_config()for per-worker application state.
See Connection and Protocol Configuration for a complete example.