Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

ntex framework guide

Learn how ntex services fit together, run servers, manage state and I/O, and build web applications.

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:

  1. Turn the HTTP request into a domain value—an Operation.
  2. Pass the operation to the execution service.
  3. Turn the resulting OperationResult into 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:

  • first accepts the original request.
  • second accepts the response produced by first.
  • 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:

  • ServiceChain describes how services are connected.
  • Pipeline owns 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:

  • Service defines what happens at each step.
  • ServiceChain describes how the steps are connected.
  • Pipeline owns the runnable graph and coordinates shared readiness.
  • Ctx controls 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-mockall defines the state behavior as a trait. Mockall creates a mock state and checks calls to its accessor and operation methods.
  • test-shimforge keeps 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:

  • tokio uses 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 Tokio LocalSet.
  • compio creates 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 until System::stop() or System::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 = true or false;
  • panic_handling = true or false;
  • ping_interval = N in milliseconds, where zero disables arbiter pings;
  • rt = Type to 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 be Send;
  • 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 Arbiter for one runtime thread and the System for the whole runtime;
  • use DefaultRuntime unless 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:

  1. a service name used in logs and to associate listeners with services;
  2. an address that implements ToSocketAddrs;
  3. a SharedCfg attached to accepted connections;
  4. 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() makes stop_on_panic worker 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 ntex System after 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:

EventDefault behavior
SIGTERMgraceful stop
SIGINT / Ctrl-Cimmediate stop
SIGQUITimmediate stop, or graceful with graceful_shutdown()
SIGHUPignored by the server
fatal signal or application panicimmediate 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:

  • St is provided by the service pipeline and remains available across requests.
  • RequestState belongs 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:

  1. A protocol or dispatcher task consumes decoded input and queues encoded output through Io or IoRef. This task typically runs the protocol service and, indirectly, application code.
  2. A transport read task waits for read readiness, reads bytes from the socket, and submits them to the I/O subsystem.
  3. 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-net while 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 own set_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::low of its page remains free. Output is held in BytePages of 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 call Io::read_more only 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_size is 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:

  • Timeout when the dispatcher timer has expired, for example the configured keep-alive timeout.
  • WriteBackpressure when outstanding output has reached the write high-water mark, so the producer should stop and flush.
  • PeerGone once 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, and None after 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:

  1. Application middleware wraps the application service.
  2. Application filters process the incoming WebRequest.
  3. The application router finds a Resource whose path and resource guards match.
  4. Resource middleware wraps the selected resource’s filter and routes.
  5. The resource filter processes the request.
  6. The resource checks its routes in registration order. A route matches only if all its method and custom guards pass.
  7. 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:

  1. A shared reference to the application state.
  2. The request-local state.
  3. 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() or WebRequest::app_state().
  • They must be Send + Sync, unlike worker-local state, which may contain Rc and RefCell.
  • 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 WebRequest before 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:

LevelRuns for
AppEvery request entering the application
ScopeRequests whose scope prefix and scope guards match
ResourceRequests 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:

ExtractorReads
HttpRequestA 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
PayloadThe streaming request body
BytesThe complete request body as bytes
StringThe 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=value pairs, 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:

ExtractorDefault limit
Json<T>32 KiB
Form<T>16 KiB
Bytes and String256 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:

ResponderResult
HttpResponse or HttpResponseBuilderThe response as configured
String, &String, or &'static str200 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, returns 400 Bad Request;
  • JSON and form bodies that exceed their configured limit return 413 Payload Too Large;
  • a malformed form Content-Length returns 411 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, return 500 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 every Error* helper;
  • WebError, an error that has already been converted for the domain;
  • ntex::util::Either<A, B>, when both A and B implement 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 return JsonPayloadError before 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

TypeMutabilityClone costTypical use
BytesImmutableUsually shares storage or copies inline data; external storage decidesDecoded frames, payloads, header values
BytesMutUniqueCopies the dataRead buffers, building output
ByteStringImmutableSame as BytesUTF-8 text: header names, paths, topics
BytePagesUniqueShares or copies according to page storage and append rulesWrite queues
BytePageImmutableShares data, Vec pages copyOne 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:

KindCreated byClone
InlineData of at most 23 bytes (11 on 32-bit targets)Copies the bytes, no allocation
StaticBytes::from_static, From<&'static [u8]>, From<&'static str>Copies the pointer
SharedA heap buffer, usually from BytesMutIncrements the reference count
ExternalBytes::from_ext, Arc<str> via ByteStringDelegated 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:

MethodEffect
BytesMut::split_toReturns the first at bytes as Bytes; shares large results and copies small ones inline
BytesMut::takeReturns all data as Bytes, keeping the spare capacity for self
BytesMut::freezeConsumes the mutable handle; reuses heap storage or copies a small result inline
BytesMut::advance_toDrops the first cnt bytes, O(1)
BytesMut::clearEmpties the buffer, reclaims the full capacity if unique
BytesMut::reserveEnsures spare capacity, see Growth
BytesMut::is_uniquetrue 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 into Bytes.

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:

  1. The transport writes socket data into the spare tail of a BytesMut.
  2. The decoder examines the initialized prefix.
  3. A complete frame is removed with split_to.
  4. The returned Bytes travels through services and application code.
  5. The BytesMut stays with the connection and receives more input.
  6. When all shared frames are dropped, reserve or clear can 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
FieldMeaning
offsetStart of the BytesMut view, from the beginning of the allocation
lenLength of the BytesMut view
capacityData capacity of the whole allocation, never changes while it is shared
ref_countNumber 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:

  1. If the buffer has enough spare capacity, nothing happens.
  2. 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 Bytes views.
  3. If the buffer is unique and has no page size, the allocation is grown with realloc, which can often extend it in place.
  4. If the buffer has a page size, its data moves to a pooled page, see Pooled buffers.
  5. Otherwise, the data is copied into a new allocation without a page size. Bytes views 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:

MethodUse it when
BytesMut::reserveYou want room for more bytes and expect the buffer to keep growing.
BytesMut::reserve_exactYou 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_moreYou 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.

ClassAllocationCapacityhalf_capacity()low()Cached pages by default
Size44 KiB4,080 bytes2 KiB256 bytes128
Size88 KiB8,176 bytes4 KiB512 bytes64
Size1616 KiB16,368 bytes8 KiB1 KiB64
Size2424 KiB24,560 bytes12 KiB1.5 KiB32
Size3232 KiB32,752 bytes16 KiB2 KiB16
Size4848 KiB49,136 bytes16 KiB3 KiB8
Size6464 KiB65,520 bytes16 KiB4 KiB16
Size128128 KiB131,056 bytes16 KiB8 KiB2
Size256256 KiB262,128 bytes16 KiB16 KiB1
Unset-65,520 bytes16 KiB4 KiBnever 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 other BufMut methods copy data into the current page. New pages are taken from the page cache with the queue’s page size.
  • BytePages::append adds 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 unique BytesMut with spare capacity becomes the new current page.
  • BytePages::prepend inserts a page at the front.

How data leaves it:

  • BytePages::take returns the next BytePage, with the current page last.
  • BytePages::split_to and BytePages::split_into move a byte prefix to another queue, splitting a page if needed.
  • BytePages::freeze returns all data as one Bytes. A single page is converted according to its storage kind: native heap storage can be reused, tiny data may be inlined, and Vec storage is copied. Several pages are copied into one buffer.
  • BytePages::copy_to leaves the source unchanged and appends cloned pages to another queue. BytePages::move_to leaves the source empty and appends its pages instead. Both follow append’s small-buffer copy rules; copy_to also copies Vec-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 BytesMut read buffer. Codecs split frames off it with split_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 of IoConfig, Size4 and Size64 by default. A buffer grows once less than BytePageSize::low() of its page remains free, larger input moves it to bigger page sizes, beyond Size256 it becomes an unpooled buffer that is freed when empty.
  • Writes. Encoders write into BytePages with the page size of IoConfig::write_size, Size16 by default. IoRef::encode_bytes appends 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. trimdown can 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:

SituationStart with
Hand a decoded frame to its immediate consumersplit_to or take, so large payloads can share the input allocation.
Keep one field in a long-lived cacheA slice followed by trimdown, or an explicit copy, to avoid retaining an entire input buffer.
Build a contiguous message of a known sizeBytesMut::with_capacity, or one reserve for the bytes still to be written.
Repeatedly build short-lived I/O bufferswith_page_size with a class that fits the usual size; cache hits avoid a fresh payload allocation.
Send an existing large bodyBytePages::append, rather than copying it with put_slice.
Queue an owned Vec<u8> for outputAppend the Vec directly to BytePages; converting it to Bytes first copies it.
Use constant protocol bytesfrom_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 3ntex 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 by set_headers_vec(bool).
  • h1::Codec::upgrade() has been removed; use Request::upgrade() to detect upgrade and CONNECT requests.
  • New settings: set_half_close(), set_host_validation(), set_max_start_line_size(), and set_write_timeout().

URLs

ntex 4 uses urly::Url, re-exported as ntex::url, instead of http::Uri:

  • ntex::http::Uri and the ntex::http::uri module have been removed.
  • RequestHead::uri, Request::uri(), HttpRequest::uri(), and WebRequest::uri() are Url values. Url::path() and Url::query() return ntex::url::Path and ntex::url::Query; use as_str() to get string slices. HttpRequest::path() and HttpRequest::query_string() still return &str.
  • HttpRequest::url_for() returns ntex::url::Url instead of url::Url.
  • The url feature has been removed; URL support is always available.
  • The HTTP and WebSocket clients accept any type that converts to Url, including &str, String, and http::Uri. Client connectors use Connect<Url> instead of Connect<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:

  • ClientBuilder now owns the connector configuration; the separate client connector builder has been removed.
  • Custom connector factories are replaced by connector Service values.
  • ClientBuilder::build() is synchronous and takes a SharedCfg.
  • Client settings are stored in ClientConfig inside SharedCfg.

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

  1. A shared reference to the application state.
  2. The request-local state.
  3. 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 to set_shutdown_timeout() / shutdown_timeout(). A zero timeout panics.
  • Read and write buffers are ntex-bytes pages and share one page cache, configured per thread with ntex_bytes::set_page_cache_size(). BufConfig and IoConfig::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 3ntex 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 3ntex 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 OnDisconnectreturns 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::KeepAliveIoStatusUpdate::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 and Builder::stop_on_panic() no longer has an effect. Use Builder::panic_handling() or #[ntex::main(panic_handling = true)]. System::set_latency_callback() requires a Send + Sync callback.
  • ntex::time: query_system_time() has been removed; use system_time().
  • ntex_util::channel::bstream::Receiver::max_buffer_size() is deprecated in favor of set_watermarks().
  • ntex::http::HeaderMap no longer implements FromIterator; build maps with insert() or append().
  • HeaderValue::to_str() accepts any valid UTF-8 value, not only visible ASCII.

Router

ntex::router (ntex-router 2) renames several APIs:

ntex 3ntex 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 SharedCfg to server bind() or listen() methods.
  • Pass SharedCfg to ClientBuilder::build().
  • Use build_with_config() or web::server_with_config() for per-worker application state.

See Connection and Protocol Configuration for a complete example.