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}