I/O Abstraction Layer
ntex provides an I/O abstraction layer that keeps protocol and service code independent of a specific runtime or reactor implementation, such as Tokio, Compio, or Neon.
In addition to presenting a consistent interface across these backends, the I/O layer provides the building blocks needed to manage buffered reads and writes, backpressure, timeouts, and graceful shutdown.
Socket abstractions
ntex achieves independence from the underlying socket implementation through
dependency inversion. The I/O subsystem does not call Tokio, Compio, or Neon
socket APIs directly. Instead, a runtime adapter moves bytes between its socket
and the buffers owned by Io. One way to picture it: the read half carries
bytes from the transport to the application and the write half carries them
back, with Io sitting between the two. Each half is an in-memory buffer that
Io owns, which is what lets it apply backpressure and run filter
transformations on bytes as they pass through.
An active connection has three cooperating execution responsibilities:
- A protocol or dispatcher task consumes decoded input and queues encoded
output through
IoorIoRef. This task typically runs the protocol service and, indirectly, application code. - A transport read task waits for read readiness, reads bytes from the socket, and submits them to the I/O subsystem.
- A transport write task takes queued bytes from the I/O subsystem and writes them to the socket.
The runtime adapter decides how to schedule these responsibilities. The
built-in adapters typically use cooperating read and write tasks, but this is
an implementation detail rather than part of the IoStream contract.
The read task cooperates with ntex backpressure. It reads only while
IoContext::poll_read_ready permits more input. It takes a buffer through
IoContext::take_read_buf and releases it with the result through
IoContext::release_read_buf. This allows the I/O
subsystem to process filters, wake the dispatcher, and pause further reads when
the configured high-water mark is reached.
A backend that reads synchronously, without holding the buffer across an await,
should use IoContext::with_read_buf instead. It combines the two steps and
borrows the read buffer in place, so no buffer is detached and appended back.
The write task follows the same pattern. It waits for
IoContext::poll_write_ready, obtains queued data with
IoContext::with_write_dst, and reports through
IoContext::update_write_status how many bytes reached the peer. ntex can
then apply write backpressure, resume waiting services, and coordinate
graceful shutdown.
A transport does not have to write the bytes before it reports. Completion based backends take ownership of whole pages and keep them until the operation completes. Those bytes leave the write buffer but have not reached the peer, so ntex counts them as outstanding output until the transport either reports them as written or returns them to the buffer. Flush completion, write backpressure, and the shutdown drain all account for them.
Both methods return an IoTaskStatus. Io means work remains, not that the
transport is ready, so a task re-arms transport readiness before its next
operation. On the write side it is returned whenever output is still buffered,
including after an attempt that made no progress. Output the transport already
owns does not keep the task running, because its completion wakes the task
again.
Both readiness methods resolve to a Readiness value rather than a plain
ready signal. Readiness::Ready allows the task to perform its next operation.
Readiness::Close means the connection must be taken down: the task closes
both directions and releases the socket, cancelling anything still in flight.
For a socket this is shutdown(SHUT_RDWR) followed by close(). Whatever the
peer still has in the receive queue should be discarded first, because closing
a socket with unread input aborts the connection with an RST and loses the
output that was just drained.
Readiness::Close covers a graceful shutdown as well as a connection that ends
because of an I/O failure, a filter failure, or an expired shutdown deadline.
On the graceful path the I/O subsystem reports it only once buffered output
has been drained, so nothing is left to flush.
Readiness::Terminate is reported only for an explicit force close through
IoRef::terminate. Whatever is still buffered is discarded on purpose, and
the task must not close gracefully: no receive queue drain and no
shutdown(SHUT_RDWR). It aborts the connection instead, for a socket by
setting SO_LINGER to zero before close(), so that the peer sees an RST and
cannot mistake a truncated stream for a complete one.
Teardown is reported back through two methods. A transport that fails or
decides to abort calls IoContext::stop with the error, which terminates the
connection without draining pending work. Once transport teardown has actually
finished, and regardless of which side initiated it, the task calls
IoContext::stopped. That is what completes Io::shutdown, which resolves
only after the backend reports the connection stopped, not merely when filter
shutdown and flushing are done.
Socket types integrate with the I/O subsystem by implementing IoStream.
Its start() method receives an IoContext, starts the transport-specific
tasks, and returns a Handle used to control or query the transport.
A simplified Tokio-style adapter looks like this:
impl IoStream for TcpStream {
fn start(self, ctx: IoContext) -> Box<dyn Handle> {
let socket = Rc::new(self);
tokio::task::spawn_local(read_task(socket.clone(), ctx.clone()));
tokio::task::spawn_local(write_task(socket.clone(), ctx));
Box::new(SocketHandle(socket))
}
}
async fn read_task(socket: Rc<TcpStream>, ctx: IoContext) {
loop {
// resolves `IoContext::poll_read_ready`
match wait_for_read_readiness(&ctx).await {
Readiness::Ready => {}
Readiness::Close | Readiness::Terminate => break,
}
let mut buf = ctx.take_read_buf();
let result = read_from_socket(&socket, &mut buf);
match ctx.release_read_buf(buf, result) {
IoTaskStatus::Io => {}
IoTaskStatus::Pause => wait_for_read_resume(&ctx).await,
IoTaskStatus::Stop => break,
}
}
}
async fn write_task(socket: Rc<TcpStream>, ctx: IoContext) {
let terminate = loop {
// resolves `IoContext::poll_write_ready`
match wait_for_write_readiness(&ctx).await {
Readiness::Ready => {}
Readiness::Close => break false,
Readiness::Terminate => break true,
}
let result = ctx.with_write_dst(|buf| {
write_to_socket(&socket, buf)
});
// `result` carries the number of bytes that reached the peer
match ctx.update_write_status(result) {
IoTaskStatus::Io => {}
IoTaskStatus::Pause => wait_for_queued_output(&ctx).await,
IoTaskStatus::Stop => break false,
}
};
if terminate {
// force close, abort so that the peer sees an RST
abort_socket(&socket);
} else {
// close both directions
shutdown_socket(&socket).await;
}
// report teardown as complete
ctx.stopped(None);
}
The example is illustrative; runtime adapters use their native readiness and
buffer APIs. The important boundary is that only the adapter accesses the
socket, while the application and protocol layers interact with Io and
IoRef. Buffer limits and read/write backpressure are therefore applied
consistently across all supported runtimes.
An Io object is normally passed to a protocol service such as
ntex::http::HttpService or ntex_mqtt::Server. The protocol service
decodes incoming bytes into protocol messages and encodes its responses back
into the I/O write buffer without depending on the concrete socket type.
The runtime-specific implementations are provided by the ntex-net crate,
which supports Tokio, Compio, and Neon backends.
Transport handles and metadata
Because the underlying socket is hidden behind the I/O abstraction, code using
Io cannot access transport-specific methods directly. The Handle returned
by IoStream::start() provides the bridge to the underlying transport. It can
resume the transport write operation on request and expose transport-specific
metadata through typed queries.
Each backend decides which query types it supports. For example, the built-in
network backends expose the remote socket address as PeerAddr.
IoRef::query is also available on Io through dereferencing and returns a
QueryItem, from which the value can be retrieved with get():
use ntex::io::{Io, types::PeerAddr};
fn log_peer_addr(io: &Io) {
if let Some(addr) = io.query::<PeerAddr>().get() {
println!("Peer address {:?}", addr.into_inner());
}
}
Queries travel through the filter stack before reaching the transport handle, so filters may also expose their own typed metadata.
Configuration and timeouts
I/O settings are stored in IoConfig and are normally added to a
SharedCfg. Io::new() retrieves the IoConfig from that shared
configuration, using the default settings when it is not present.
use ntex::{
SharedCfg,
io::IoConfig,
time::{Millis, Seconds},
util::BytePageSize,
};
let cfg = SharedCfg::new("my-protocol")
.add(
IoConfig::new()
.set_connect_timeout(Millis(5_000))
.set_keepalive_timeout(Seconds(30))
.set_shutdown_timeout(Seconds(2))
.set_frame_read_rate(Seconds(2), Seconds(10), 1_024)
.set_write_timeout(Seconds(10))
.set_read_size(BytePageSize::Size4, BytePageSize::Size32)
.set_read_backpressure(32 * 1024)
.set_write_backpressure(32 * 1024)
.set_write_buf_threshold(8 * 1024),
)
.build();
These settings are used by different parts of the stack:
- The connection timeout is applied by
ntex-netwhile resolving and opening an outgoing connection. - The keep-alive timeout and frame read-rate limits are interpreted by
protocol dispatchers. A frame read-rate limit protects a decoder from peers
that send one incomplete frame too slowly. It also applies to the first
frame of a new connection, so a peer that connects and stays silent is
closed once the limit expires. A write timeout protects against peers
that stop reading: it bounds how long write backpressure may stay enabled,
from the moment it is enabled until it is disabled, before the dispatcher
stops with a write timeout. Keep-alive and read-rate timers do not run
during that time. Output left after backpressure is disabled is bounded
only by keep-alive. The HTTP/1 dispatcher does not use these settings; it is
configured by
HttpServiceConfig, including its ownset_write_timeout(). - The graceful-shutdown timeout bounds both phases of shutdown together: the filter shutdown and the transport drain of pending output. It cannot be disabled; a zero timeout is rejected.
- The read and write high-water marks enable backpressure. Write backpressure is released after outstanding output falls to half its high-water mark, counting both buffered output and output a transport has taken ownership of but not yet written to the peer.
- The min and max read page sizes bound the adaptive page size of new read
buffers. Before another socket read, a buffer grows once less than
BytePageSize::lowof its page remains free. Output is held inBytePagesof the write page size, so the write backpressure setting takes only a high-water mark. - The write buffer page size controls newly allocated
BytePages, while the write threshold controls when supported transports attempt an early direct write.
Connection and keep-alive timeouts are disabled by default. Frame read-rate limits and the write timeout are also disabled. The default graceful-shutdown timeout is one second, and the default read and write high-water marks are approximately 32 KiB and 16 KiB.
Start with those defaults. Change one setting because of an observed workload, not because larger values sound faster:
- Read page sizes control allocation and read-call frequency, not the amount of unread data a connection may queue. Increase the maximum for sustained large frames or bodies; keep the minimum small when many mostly idle connections are expected.
- Backpressure watermarks control when producers or transports should pause. They are not hard memory limits: one socket read or one encoded item can cross a watermark.
- The write threshold is a latency hint for starting transport work earlier. It does not flush data and does not replace write backpressure.
- A write timeout is especially important for untrusted peers. Without one, a peer that stops reading can keep a backpressured connection and its output alive indefinitely.
An established connection can switch to another shared configuration with
Io::set_config. This is useful when a protocol upgrade changes timeout or
buffer requirements. The method is unsafe: replacing the configuration may
release the allocation that IoRef::cfg hands out, so no reference obtained
from it may be live across the call or used afterwards.
Filter subsystem
Applications often need to transform a byte stream before a protocol service processes it. TLS must decrypt incoming records and encrypt outgoing data. A protocol may also be tunneled through another framing layer, such as MQTT over WebSocket.
ntex implements these transformations with FilterLayer. A filter operates
on in-memory byte buffers: FilterLayer::process_read_buf transforms data
received from the next inner layer, while FilterLayer::process_write_buf
transforms data queued by the application before passing it toward the
transport. Filters do not perform socket I/O themselves, so the same filter
can be used with any supported runtime backend.
Filters are composable. Io::add_filter adds a layer and allocates the
intermediate read and write buffers that separate it from adjacent layers. For
example, an MQTT service can receive its byte stream through either of these
stacks:
socket <-> TLS <-> MQTT
socket <-> TLS <-> WebSocket <-> MQTT
On reads, bytes move from the socket through the inner filters toward the application. On writes, they move in the opposite direction. In the second stack, the WebSocket filter removes and creates WebSocket framing, while the TLS filter decrypts and encrypts the resulting byte stream. The MQTT service still reads and writes MQTT bytes and does not need to know which transport filters are installed below it.
Filters may maintain protocol state, expose typed metadata through query(),
emit output while processing input, and participate in graceful shutdown. For
example, a TLS filter can expose the negotiated protocol or peer certificate,
and a WebSocket filter can generate a close frame during shutdown. Output a
filter writes while processing a read is pushed toward the transport as soon as
that processing finishes, without waiting for the application to write.
Most byte transformations only need FilterLayer. Lower-level concerns that
must observe or control the entire filter chain can instead wrap the current
chain with Filter by using Io::map_filter.
In addition to processing buffers, queries, and shutdown, Filter participates
in read and write readiness decisions. A wrapper can therefore delay readiness
to implement custom throttling and wake the I/O tasks when work may resume. It
can also observe buffer processing for metrics or accounting without changing
the byte stream.
A custom Filter normally stores the filter it wraps and delegates every
operation it does not intentionally override. ntex provides forwarding macros
for readiness, queries, and shutdown to make this pattern less error-prone.
Typed versus erased filter stacks
The filter stack is encoded in the type parameter of Io. A new connection
starts as Io<Base>, using the Base filter. Calling add_filter(layer)
consumes the current value and returns Io<Layer<U, F>>, where the Layer
marker pairs the new outer layer U with the previous stack F. Keeping this
concrete type provides static dispatch and allows Io::filter to return the
concrete outer filter.
let io: Io<Base> = create_io();
let io: Io<Layer<MyFilter, Base>> = io.add_filter(MyFilter::new());
At service boundaries, different connections may have different concrete
filter stacks. Io::seal erases the stack type and returns Io<Sealed>,
using the Sealed marker, while Io::boxed returns the IoBoxed
convenience wrapper. Both operations consume the original Io value and retain
the same connection state and filter behavior behind a dynamically dispatched
Filter.
let io: IoBoxed = io.boxed();
start_protocol(io);
Additional typed layers can still be added to a sealed stream and the result can be erased again when necessary. Type erasure is therefore normally performed at the boundary where a protocol or service needs one uniform I/O type, rather than while constructing the filter stack.
Read/write streams
Incoming and outgoing bytes use separate buffer paths.
Reading
The transport adapter places bytes read from the socket into a BytesMut.
The bytes pass through the filter chain and arrive in the application-facing
read buffer, where a codec or protocol service can inspect and consume them.
BytesMut is a contiguous, growable buffer. A decoder can split immutable
Bytes values from it without copying larger payloads; short values are
copied into the Bytes handle itself. This is useful when a decoded message
must retain part of the input after the decoder continues processing later
data.
Read buffers are ntex-bytes pages. Each connection picks the page size of
new read buffers between the min and max set with
IoConfig::set_read_size, Size4 and Size64 by default. It starts at the
min, so idle and light connections pin small pages. Consecutive reads that
fill their buffer form one batch with the read that ends it; a batch larger
than the page grows the page size to fit it, up to the max. After four
batches in a row that fit in half of the next smaller page, the page size
shrinks by one step, down to the min. Io::set_config restarts it at the
min. Equal min and max sizes fix the page size.
The page size does not affect read backpressure: the read high-water mark is
the Size32 capacity, 32,752 bytes, by default;
IoConfig::set_read_backpressure sets another one. An empty buffer goes
back to the per-thread page cache of its size, shared with write pages and any
other pooled BytesMut; set_page_cache_size tunes it. A frame split from
the read buffer keeps the page alive, the page returns to the cache once the
last frame is dropped. Before another socket read, the adapter obtains a
buffer from IoContext. ntex calls BytesMut::reserve_more once less than
BytePageSize::low of the page remains free, compacting the data within its
page when it still fits there. If larger input must be buffered, the buffer moves to
bigger page sizes and, past the largest one, to a plain allocation that is
freed instead of cached.
Read backpressure is based on the size of the application-facing read buffer.
When it reaches the configured high-water mark, ntex pauses the transport read
task. Consuming input through IoRef::decode, IoRef::with_buf,
IoRef::with_read_src or IoRef::with_read_dst wakes the read task once
the buffered input falls to half that mark.
Io::recv and Io::read_exact wait for more input, so they release
backpressure regardless of how much is still buffered.
Choose the highest-level read API that matches the protocol:
- Use
recv(codec)for framed messages. It decodes buffered input and waits when the codec says the frame is incomplete. - Use
read_exact()for a fixed-size header or field. - For a custom parser, inspect or consume the read buffer through
IoRef, and callIo::read_moreonly after deciding that more bytes are required.read_more()is an active request: it resumes transport reads and releases read backpressure even if the buffer is still large.IoRef::read_dst_sizeis the passive size check.
Writing
Application output is queued in BytePages, a growable collection of
byte pages. Internally allocated pages use the size configured by IoConfig.
Owned buffers passed to IoRef::encode_bytes can also become pages directly,
avoiding a copy when their storage can be retained. Write filters consume the
pages in order, transform their contents, and place the result into the next
buffer toward the transport.
IoRef::encode, IoRef::encode_slice, and IoRef::encode_bytes queue
output but do not wait for every byte to reach the socket. Queueing makes the
data available to the transport write task and schedules that task when it
needs to be resumed. On transports that support direct writes, the configured
write threshold can trigger an earlier write while the application is still
producing output, reducing latency for large responses.
Use Io::flush to apply the configured write wait policy. flush(false) is
a backpressure wait, not a request to drain everything: it returns immediately
while the outstanding output is below the high-water mark. If the high-water
mark has been reached, it waits until the outstanding output falls to half that
mark. flush(true) waits until all queued data has reached the peer, including
data a transport has taken ownership of but not yet written. Io::send
combines codec encoding with a full flush.
Write backpressure is advisory. The encode methods still accept more output
after the high-water mark is reached, because refusing half of a protocol
message would be difficult to recover from. A producer that is not managed by
a dispatcher should await IoRef::write_ready before queueing each next
chunk, or use flush(false) between batches:
use ntex::{
io::Io,
util::Bytes,
};
async fn queue_chunk(io: &Io, chunk: Bytes) -> std::io::Result<()> {
io.write_ready().await?;
io.encode_bytes(chunk)
}
This bounds normal streaming output around the configured watermark. It does not split a single large chunk or turn the watermark into a strict per-connection memory cap.
This separation allows codecs and application services to work with bytes without depending on socket readiness, while the I/O subsystem coordinates buffering and backpressure.
Connection lifecycle and shutdown
A connection stays usable until the application closes it, the peer disconnects, or the transport reports an error.
A service that is waiting on something other than input, such as an unready
dependency or a slow response, still needs to notice that the connection
requires attention. Io::poll_status_update reports the next status as an
IoStatusUpdate value:
Timeoutwhen the dispatcher timer has expired, for example the configured keep-alive timeout.WriteBackpressurewhen outstanding output has reached the write high-water mark, so the producer should stop and flush.PeerGoneonce the connection has closed, whether the peer disconnected, the transport failed, or the shutdown was started locally. It carries the transport error if one occurred, andNoneafter a clean close.
Use Io::poll_read_pause instead when the service should also stop accepting
more transport input while it waits. The pause is cooperative: decoding,
touching the application read buffer, or explicitly asking for more input
resumes reads. This makes it suitable for a dispatcher waiting on readiness,
not for permanently disabling the read half.
The same conditions reach a codec-driven service as RecvError from
Io::poll_recv, which additionally reports decoder failures. Code that only
needs to be woken when the connection goes away can await the Waiter
future returned by IoRef::on_disconnect. IoRef::is_active becomes
false as soon as closing starts; IoRef::is_closed becomes true only after
the backend has released the transport and teardown has finished.
A clean EOF from the peer ends the read direction but leaves the write half open, so a service can still finish encoding and flushing its response before closing.
IoRef::close requests a graceful shutdown and returns immediately.
Io::shutdown drives that shutdown to completion, in two phases.
In the first phase both directions stay open and every filter gets a chance to
emit its own closing data through FilterLayer::shutdown, such as a TLS
close_notify or a WebSocket close frame. Reads keep running and are processed
normally, so closing data sent by the peer still reaches the filters. A read
failure does not abort the shutdown; it only means the filter handshake cannot
finish.
The second phase belongs to the transport. It drains the remaining output into the connection and then closes both directions. The read side is paused: the filters are done, so no further input can be used, and the transport discards whatever is left in the receive queue just before it closes. Input is no longer delivered to the application in this phase.
A single deadline bounds both phases. If the filters finish early, the
transport drain gets the remaining time. If the deadline expires while filters
are still pending, ntex ends that phase and still runs transport shutdown with
the already-expired deadline; when the transport phase observes it, undrained
output is discarded. Io::shutdown reports the timed-out error after
transport teardown.
IoRef::terminate skips the process entirely and drops the connection
without flushing pending output.
A graceful close, including one forced forward by a timeout or failure, reaches
the transport as Readiness::Close. An explicit IoRef::terminate reaches
it as Readiness::Terminate, so a socket backend can abort rather than send a
clean end-of-stream. The transport reports either teardown through
IoContext::stopped, which is what allows Io::shutdown to resolve.
Dropping the owning Io value is not a substitute for awaiting shutdown. If
no output would be lost, dropping it starts a normal close. If buffered output
can no longer be delivered because the filter stack is being dropped, ntex
aborts the transport so the peer does not mistake a truncated stream for a
complete one. Call shutdown().await when delivery and a clean close matter.
Testing
IoTest provides a pair of interconnected in-memory transports for testing
codecs, filters, and protocol services without opening sockets. Each endpoint
implements IoStream and can be wrapped in Io. Writing to one IoTest
endpoint supplies input to the other endpoint, while read() collects bytes
written back by the peer.
use ntex::codec::BytesCodec;
use ntex::io::{Io, testing::IoTest};
use ntex::util::Bytes;
#[ntex::test]
async fn protocol_io() {
let (client, server) = IoTest::create();
// Allow the server transport to write to the client.
client.remote_buffer_cap(1024);
let io = Io::from(server);
client.write(b"request");
let request = io.recv(&BytesCodec).await.unwrap().unwrap();
assert_eq!(request, Bytes::from_static(b"request"));
io.send(Bytes::from_static(b"response"), &BytesCodec)
.await
.unwrap();
assert_eq!(client.read().await.unwrap(), b"response"[..]);
}
Tests can also force pending reads, inject read or write errors, close either
side, constrain write capacity to exercise backpressure, and attach a
PeerAddr value for transport-query tests.