Skip to main content

ntex_io/
lib.rs

1//! Asynchronous I/O abstractions for the ntex ecosystem.
2//!
3//! [`Io`] wraps an underlying [`IoStream`] and coordinates buffered reads,
4//! writes, backpressure, timeouts, and shutdown. Protocol transforms can be
5//! composed through [`Filter`] layers, while [`Framed`] combines an I/O stream
6//! with an `ntex-codec` encoder and decoder.
7//!
8//! Use [`IoConfig`] to configure buffer thresholds and connection timeouts.
9#![deny(clippy::pedantic)]
10#![allow(
11    clippy::missing_fields_in_debug,
12    clippy::missing_errors_doc,
13    clippy::missing_panics_doc,
14    clippy::must_use_candidate
15)]
16use std::io::{Error as IoError, Result as IoResult};
17use std::{any::Any, any::TypeId, fmt, task::Poll};
18
19pub mod cfg;
20pub mod testing;
21pub mod types;
22
23mod buf;
24mod ctx;
25mod filter;
26mod filterptr;
27mod flags;
28mod framed;
29mod io;
30mod ioref;
31mod macros;
32mod ops;
33mod seal;
34mod utils;
35mod waiters;
36
37use ntex_codec::Decoder;
38
39pub use self::buf::{FilterBuf, FilterCtx};
40pub use self::cfg::IoConfig;
41pub use self::ctx::IoContext;
42pub use self::filter::{Base, Filter, Layer};
43pub use self::framed::Framed;
44pub use self::io::{Io, IoRef};
45pub use self::ops::{Id, TimerHandle};
46pub use self::seal::{IoBoxed, Sealed};
47pub use self::utils::Decoded;
48pub use self::waiters::Waiter;
49
50#[doc(hidden)]
51pub use self::flags::Flags;
52
53/// Filter readiness state.
54#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
55pub enum Readiness {
56    /// The I/O task may proceed with I/O operations.
57    Ready,
58    /// The transport must be closed gracefully.
59    ///
60    /// The I/O task must close both directions of the connection and then
61    /// release it. For a socket this is `shutdown(SHUT_RDWR)` followed by
62    /// `close()`. Any operation still in flight should be canceled.
63    ///
64    /// During a graceful shutdown buffered output is drained before this is
65    /// reported. If the connection ended because of a failure or an expired
66    /// shutdown deadline, output may instead be undeliverable and discarded.
67    /// In every case the I/O task must not attempt another flush.
68    Close,
69    /// The transport must be released immediately.
70    ///
71    /// This is reported only for an explicit force close through
72    /// [`IoRef::terminate`](crate::IoRef::terminate), so whatever is still
73    /// buffered is discarded on purpose. The I/O task must not perform a
74    /// graceful close: no receive queue drain and no `shutdown(SHUT_RDWR)`,
75    /// just release the connection. That keeps an aborted stream
76    /// distinguishable from one that ended normally, instead of terminating a
77    /// truncated response with a clean `FIN`.
78    ///
79    /// A connection that ends because of an I/O failure, a filter failure or an
80    /// expired shutdown deadline reports [`Close`](Self::Close) instead: the
81    /// transport is gone or unusable, so there is nothing to gain from
82    /// aborting it.
83    Terminate,
84}
85
86impl Readiness {
87    /// Merges two readiness states without regard to argument order.
88    ///
89    /// `Terminate` overrides every other state, `Close` overrides `Pending` and
90    /// `Ready`, and `Pending` overrides `Ready`.
91    pub fn merge(val1: Poll<Readiness>, val2: Poll<Readiness>) -> Poll<Readiness> {
92        match (val1, val2) {
93            (Poll::Ready(Readiness::Terminate), _) | (_, Poll::Ready(Readiness::Terminate)) => {
94                Poll::Ready(Readiness::Terminate)
95            }
96            (Poll::Ready(Readiness::Close), _) | (_, Poll::Ready(Readiness::Close)) => {
97                Poll::Ready(Readiness::Close)
98            }
99            (Poll::Pending, _) | (_, Poll::Pending) => Poll::Pending,
100            (Poll::Ready(Readiness::Ready), Poll::Ready(Readiness::Ready)) => {
101                Poll::Ready(Readiness::Ready)
102            }
103        }
104    }
105}
106
107/// A processing layer that transforms an I/O stream's read and write buffers.
108///
109/// Read processing runs from the transport toward the application. Write and
110/// shutdown processing run from the application toward the transport. Each
111/// callback receives the buffers immediately before and after this layer.
112///
113/// Implementations must move or transform all bytes they consume. Bytes left
114/// in a source buffer remain available to the layer on a later callback.
115#[allow(unused_variables)]
116pub trait FilterLayer: fmt::Debug + 'static {
117    /// Returns type-indexed information exposed by this layer.
118    ///
119    /// Returning `None` allows the query to continue through the remaining
120    /// filter chain.
121    fn query(&self, id: TypeId) -> Option<Box<dyn Any>> {
122        None
123    }
124
125    /// Processes incoming data from the transport-facing source buffer into
126    /// the application-facing destination buffer.
127    ///
128    /// This is also called once after clean transport read EOF, with
129    /// [`IoRef::is_read_eof`] returning `true`.
130    fn process_read_buf(&self, buf: &FilterBuf<'_>) -> IoResult<()>;
131
132    /// Processes outgoing data from the application-facing source buffer into
133    /// the transport-facing destination buffer.
134    fn process_write_buf(&self, buf: &FilterBuf<'_>) -> IoResult<()>;
135
136    /// Performs graceful filter shutdown.
137    ///
138    /// Returning `Poll::Pending` keeps the filter active and causes shutdown to
139    /// be polled again after the I/O task is notified. A ready result allows
140    /// shutdown to continue toward the transport.
141    ///
142    /// A filter that waits for input from the peer must check
143    /// [`IoRef::is_read_eof`] and return a ready result once it is set: after a
144    /// clean read EOF no further input can arrive, so pending forever would
145    /// only stall the close until the shutdown timeout expires. The runtime
146    /// also ends the shutdown phase itself in that case, but it cannot know
147    /// whether the filter considers the shutdown complete.
148    fn shutdown(&self, buf: &FilterBuf<'_>) -> IoResult<Poll<()>> {
149        Ok(Poll::Ready(()))
150    }
151}
152
153/// An underlying transport that can be managed by [`Io`].
154///
155/// [`start`](IoStream::start) is called exactly once when the transport is
156/// wrapped in [`Io`]. The implementation must start its read and write tasks,
157/// use the supplied [`IoContext`] to exchange buffers and readiness state, and
158/// return a handle that remains valid for the connection's lifetime.
159pub trait IoStream {
160    /// Starts transport-specific I/O tasks and returns their control handle.
161    fn start(self, _: IoContext) -> Box<dyn Handle>;
162}
163
164#[doc(hidden)]
165/// Callbacks invoked around filter-chain processing.
166pub trait IoCallbacks {
167    /// Called before processing the read or write filter chain.
168    fn before_processing(&self, io: &IoRef);
169
170    /// Called after processing the read or write filter chain.
171    fn after_processing(&self, io: &IoRef);
172}
173
174/// Control handle for transport-specific I/O tasks.
175///
176/// The handle is called synchronously by the connection state and must not
177/// block.
178pub trait Handle {
179    /// Returns type-indexed transport information.
180    fn query(&self, _: TypeId) -> Option<Box<dyn Any>> {
181        None
182    }
183
184    #[inline]
185    /// Requests that the transport start or resume a write operation.
186    fn write(&self, _: &IoContext) {}
187}
188
189/// Current status of the I/O state.
190#[derive(Copy, Clone, Debug, PartialEq, Eq, Hash)]
191pub enum IoTaskStatus {
192    /// Work remains, the task should perform another I/O operation.
193    ///
194    /// This reports that the connection still has work to do, not that the
195    /// transport can accept or supply bytes right now. A task must re-arm
196    /// transport readiness before the next operation. On the write side it is
197    /// returned whenever output is still buffered, including after an attempt
198    /// that made no progress.
199    Io,
200    /// Pause the task until the context wakes it.
201    Pause,
202    /// Stop the task and release its transport resources.
203    Stop,
204}
205
206/// I/O status update events.
207#[derive(Debug)]
208pub enum IoStatusUpdate {
209    /// The dispatcher timer has expired.
210    Timeout,
211    /// Write backpressure is currently active.
212    WriteBackpressure,
213    /// The connection is no longer usable.
214    ///
215    /// Reported once the connection can no longer dispatch frames, whether
216    /// because transport shutdown has begun, the peer disconnected, an error
217    /// occurred, or shutdown was started locally with
218    /// [`IoRef::close`](crate::IoRef::close) or
219    /// [`IoRef::terminate`](crate::IoRef::terminate). The transport backend
220    /// may still be finishing teardown.
221    ///
222    /// Carries the connection error, which may originate from the transport,
223    /// a filter, or shutdown. `None` means no error was recorded.
224    PeerGone(Option<IoError>),
225}
226
227/// Errors that can occur while receiving data.
228pub enum RecvError<U: Decoder> {
229    /// The dispatcher timer has expired.
230    Timeout,
231    /// Write backpressure is currently active.
232    WriteBackpressure,
233    /// Failed to decode an incoming frame.
234    Decoder(U::Error),
235    /// The connection is no longer usable.
236    ///
237    /// Reported once the connection can no longer dispatch frames, whether
238    /// because transport shutdown has begun, the peer disconnected, an error
239    /// occurred, or shutdown was started locally with
240    /// [`IoRef::close`](crate::IoRef::close) or
241    /// [`IoRef::terminate`](crate::IoRef::terminate). The transport backend
242    /// may still be finishing teardown.
243    ///
244    /// Carries the connection error, which may originate from the transport,
245    /// a filter, or shutdown. `None` means no error was recorded.
246    PeerGone(Option<IoError>),
247}
248
249impl<U> fmt::Debug for RecvError<U>
250where
251    U: Decoder,
252    <U as Decoder>::Error: fmt::Debug,
253{
254    fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
255        match *self {
256            RecvError::Timeout => {
257                write!(fmt, "RecvError::Timeout")
258            }
259            RecvError::WriteBackpressure => {
260                write!(fmt, "RecvError::WriteBackpressure")
261            }
262            RecvError::Decoder(ref e) => {
263                write!(fmt, "RecvError::Decoder({e:?})")
264            }
265            RecvError::PeerGone(ref e) => {
266                write!(fmt, "RecvError::PeerGone({e:?})")
267            }
268        }
269    }
270}
271
272#[cfg(test)]
273mod tests {
274    use super::*;
275    use ntex_codec::BytesCodec;
276    use std::io;
277
278    #[test]
279    fn test_fmt() {
280        assert!(format!("{:?}", IoStatusUpdate::Timeout).contains("Timeout"));
281        assert!(format!("{:?}", RecvError::<BytesCodec>::Timeout).contains("Timeout"));
282        assert!(
283            format!("{:?}", RecvError::<BytesCodec>::WriteBackpressure)
284                .contains("WriteBackpressure")
285        );
286        assert!(
287            format!(
288                "{:?}",
289                RecvError::<BytesCodec>::Decoder(io::Error::other("err"))
290            )
291            .contains("RecvError::Decoder")
292        );
293        assert!(
294            format!(
295                "{:?}",
296                RecvError::<BytesCodec>::PeerGone(Some(io::Error::other("err")))
297            )
298            .contains("RecvError::PeerGone")
299        );
300    }
301
302    #[test]
303    fn readiness_merge() {
304        let states = [
305            Poll::Pending,
306            Poll::Ready(Readiness::Ready),
307            Poll::Ready(Readiness::Close),
308            Poll::Ready(Readiness::Terminate),
309        ];
310
311        for val1 in states {
312            for val2 in states {
313                assert_eq!(Readiness::merge(val1, val2), Readiness::merge(val2, val1));
314            }
315        }
316
317        assert_eq!(
318            Readiness::merge(Poll::Pending, Poll::Ready(Readiness::Ready)),
319            Poll::Pending
320        );
321        assert_eq!(
322            Readiness::merge(Poll::Pending, Poll::Ready(Readiness::Close)),
323            Poll::Ready(Readiness::Close)
324        );
325        assert_eq!(
326            Readiness::merge(Poll::Ready(Readiness::Ready), Poll::Ready(Readiness::Close)),
327            Poll::Ready(Readiness::Close)
328        );
329        assert_eq!(
330            Readiness::merge(Poll::Pending, Poll::Ready(Readiness::Terminate)),
331            Poll::Ready(Readiness::Terminate)
332        );
333        assert_eq!(
334            Readiness::merge(
335                Poll::Ready(Readiness::Close),
336                Poll::Ready(Readiness::Terminate)
337            ),
338            Poll::Ready(Readiness::Terminate)
339        );
340    }
341}