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

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.