Skip to main content

ntex_io/
filter.rs

1use std::{any, cell::Cell, io, task::Context, task::Poll};
2
3use crate::{FilterCtx, FilterLayer, IoRef, Readiness, io::IoState};
4
5#[derive(Debug)]
6/// Base filter that connects a filter chain to the underlying transport.
7pub struct Base(IoRef);
8
9impl Base {
10    pub(crate) fn new(inner: IoRef) -> Self {
11        Base(inner)
12    }
13}
14
15#[derive(Debug)]
16/// One processing layer wrapped around an existing filter chain.
17///
18/// Values of this type are created by [`Io::add_filter`](crate::Io::add_filter).
19/// `F` is the outer, newly added layer and `L` is the previously installed
20/// inner chain.
21pub struct Layer<F, L = Base>(pub(crate) F, L, Cell<bool>);
22
23impl<F: FilterLayer, L: Filter> Layer<F, L> {
24    pub(crate) fn new(f: F, l: L) -> Self {
25        Self(f, l, Cell::new(false))
26    }
27}
28
29pub(crate) struct NullFilter;
30
31const NULL: NullFilter = NullFilter;
32
33impl NullFilter {
34    pub(super) const fn get() -> &'static dyn Filter {
35        &NULL
36    }
37}
38
39/// Complete filter-chain interface used by [`Io`](crate::Io).
40///
41/// Most filters should implement [`FilterLayer`] and be installed with
42/// [`Io::add_filter`](crate::Io::add_filter). Implement this trait directly
43/// only when wrapping or replacing a complete chain.
44pub trait Filter: 'static {
45    /// Returns type-indexed information exposed by this chain.
46    fn query(&self, id: any::TypeId) -> Option<Box<dyn any::Any>>;
47
48    /// Processes incoming data from the transport toward the application.
49    fn process_read_buf(&self, ctx: &mut FilterCtx<'_>) -> io::Result<()>;
50
51    /// Processes outgoing data from the application toward the transport.
52    fn process_write_buf(&self, ctx: &mut FilterCtx<'_>) -> io::Result<()>;
53
54    /// Performs graceful shutdown from the outermost layer toward the
55    /// transport.
56    fn shutdown(&self, ctx: &mut FilterCtx<'_>) -> io::Result<Poll<()>>;
57
58    /// Checks whether transport read operations may proceed.
59    ///
60    /// Reads continue through the filter shutdown phase so that filters can
61    /// complete theirs, and are paused for the transport shutdown phase, so
62    /// [`Readiness::Close`] is resolved only once the connection is
63    /// terminated. A force close reports [`Readiness::Terminate`] instead,
64    /// which releases the connection without a graceful close. That decision is
65    /// made by [`IoContext`](crate::IoContext) rather than by the chain, so
66    /// that it survives the chain being dropped along with
67    /// [`Io`](crate::Io).
68    fn poll_read_ready(&self, cx: &mut Context<'_>) -> Poll<Readiness>;
69
70    /// Checks whether transport write operations may proceed.
71    ///
72    /// Resolves to [`Readiness::Close`] once the connection enters a graceful
73    /// shutdown and all buffered output has reached the transport, or as soon
74    /// as it ends because of a failure. See [`poll_read_ready`] for the
75    /// force-close case.
76    ///
77    /// [`poll_read_ready`]: Self::poll_read_ready
78    fn poll_write_ready(&self, cx: &mut Context<'_>) -> Poll<Readiness>;
79}
80
81/// Transport read readiness as decided by the io state, without the filter
82/// chain.
83pub(crate) fn read_readiness(st: &IoState) -> Poll<Readiness> {
84    if st.flags.is_force_closing() {
85        // Only an explicit `IoRef::terminate()` aborts the connection. A
86        // transport failure, a filter failure or an expired shutdown
87        // deadline end the connection too, but the transport still closes
88        // it gracefully.
89        Poll::Ready(Readiness::Terminate)
90    } else if st.flags.is_aborted() {
91        // The connection ended because of a failure, so no further input
92        // can be used; the transport closes it gracefully.
93        Poll::Ready(Readiness::Close)
94    } else if st.flags.is_read_eof() {
95        // The transport read side is closed, no further input can
96        // arrive. This outranks filter shutdown below: a filter that
97        // waits for input would otherwise keep the transport polling a
98        // closed read side.
99        Poll::Pending
100    } else if st.flags.is_stopping() {
101        // Transport shutdown phase. The filters are done, so no further
102        // input can be used and the read task pauses. The receive queue
103        // is drained by the transport itself, just before it closes the
104        // connection.
105        Poll::Pending
106    } else if st.flags.is_stopping_filters() {
107        // A filter may still need input to complete its shutdown, so
108        // keep reading even though the application paused reads.
109        Poll::Ready(Readiness::Ready)
110    } else if st.flags.is_read_paused_or_backpressure() || st.flags.is_read_wr_backpressure() {
111        // read buffer is full or is not processed by dispatcher yet,
112        // or output produced by reading has not drained
113        Poll::Pending
114    } else {
115        Poll::Ready(Readiness::Ready)
116    }
117}
118
119/// Transport write readiness as decided by the io state, without the filter
120/// chain.
121pub(crate) fn write_readiness(st: &IoState) -> Poll<Readiness> {
122    if st.flags.is_force_closing() {
123        // see `read_readiness`
124        Poll::Ready(Readiness::Terminate)
125    } else if st.flags.is_aborted() {
126        // The connection ended because of a failure, so there is nothing
127        // left to drain; the transport closes it gracefully.
128        Poll::Ready(Readiness::Close)
129    } else if st.flags.is_stopping() {
130        // Transport shutdown phase. Buffered output is drained into the
131        // transport first; `Readiness::Close` is reported only once
132        // nothing is left to write.
133        if st.buffer.write_buf_size() != 0 {
134            Poll::Ready(Readiness::Ready)
135        } else if st.wr_inflight.get() != 0 {
136            // the transport still holds output that has not reached
137            // the peer, its completion wakes the write task
138            Poll::Pending
139        } else {
140            Poll::Ready(Readiness::Close)
141        }
142    } else if st.flags.is_write_paused() {
143        Poll::Pending
144    } else {
145        Poll::Ready(Readiness::Ready)
146    }
147}
148
149impl Filter for Base {
150    fn query(&self, id: any::TypeId) -> Option<Box<dyn any::Any>> {
151        if let Some(hnd) = self.0.0.handle.take() {
152            let res = hnd.query(id);
153            self.0.0.handle.set(Some(hnd));
154            res
155        } else {
156            None
157        }
158    }
159
160    fn poll_read_ready(&self, cx: &mut Context<'_>) -> Poll<Readiness> {
161        let st = &self.0.0;
162        let res = read_readiness(st);
163        if !matches!(res, Poll::Ready(Readiness::Close | Readiness::Terminate)) {
164            st.read_task.register(cx.waker());
165        }
166        res
167    }
168
169    fn poll_write_ready(&self, cx: &mut Context<'_>) -> Poll<Readiness> {
170        let st = &self.0.0;
171        let res = write_readiness(st);
172        if !matches!(res, Poll::Ready(Readiness::Close | Readiness::Terminate)) {
173            st.write_task.register(cx.waker());
174        }
175        res
176    }
177
178    #[inline]
179    fn process_read_buf(&self, _: &mut FilterCtx<'_>) -> io::Result<()> {
180        Ok(())
181    }
182
183    #[inline]
184    fn process_write_buf(&self, _: &mut FilterCtx<'_>) -> io::Result<()> {
185        Ok(())
186    }
187
188    #[inline]
189    fn shutdown(&self, _: &mut FilterCtx<'_>) -> io::Result<Poll<()>> {
190        Ok(Poll::Ready(()))
191    }
192}
193
194impl<F, L> Filter for Layer<F, L>
195where
196    F: FilterLayer,
197    L: Filter,
198{
199    #[inline]
200    fn query(&self, id: any::TypeId) -> Option<Box<dyn any::Any>> {
201        self.0.query(id).or_else(|| self.1.query(id))
202    }
203
204    #[inline]
205    fn shutdown(&self, ctx: &mut FilterCtx<'_>) -> io::Result<Poll<()>> {
206        if !self.2.get() {
207            if ctx.with_buffer(|buf| self.0.shutdown(buf))?.is_ready() {
208                self.process_write_buf(ctx)?;
209                self.2.set(true);
210
211                // Discard the write buffer; it won't be processed
212                ctx.clear_write_buf();
213            } else {
214                // Output produced by this filter's shutdown sits in the buffer
215                // of the inner filter, which only runs its write processing
216                // here, so move it towards the transport while waiting.
217                ctx.with_next(|ctx| self.1.process_write_buf(ctx))?;
218                return Ok(Poll::Pending);
219            }
220        }
221        ctx.with_next(|ctx| self.1.shutdown(ctx))
222    }
223
224    #[inline]
225    fn process_read_buf(&self, ctx: &mut FilterCtx<'_>) -> io::Result<()> {
226        ctx.with_next(|ctx| self.1.process_read_buf(ctx))?;
227        if self.2.get() {
228            Ok(())
229        } else {
230            ctx.with_buffer(|buf| self.0.process_read_buf(buf))
231        }
232    }
233
234    #[inline]
235    fn process_write_buf(&self, ctx: &mut FilterCtx<'_>) -> io::Result<()> {
236        if !self.2.get() {
237            ctx.with_buffer(|buf| self.0.process_write_buf(buf))?;
238        }
239        ctx.with_next(|ctx| self.1.process_write_buf(ctx))
240    }
241
242    #[inline]
243    fn poll_read_ready(&self, cx: &mut Context<'_>) -> Poll<Readiness> {
244        self.1.poll_read_ready(cx)
245    }
246
247    #[inline]
248    fn poll_write_ready(&self, cx: &mut Context<'_>) -> Poll<Readiness> {
249        self.1.poll_write_ready(cx)
250    }
251}
252
253impl Filter for NullFilter {
254    #[inline]
255    fn query(&self, _: any::TypeId) -> Option<Box<dyn any::Any>> {
256        None
257    }
258
259    // The filter chain is gone once `Io` has been dropped, so nothing is left
260    // that could process buffered data. The connection is closed gracefully;
261    // `IoContext` reports `Readiness::Terminate` on top of this when the
262    // application asked for a force close, or when the drop had to discard
263    // output that never reached the peer.
264    #[inline]
265    fn poll_read_ready(&self, _: &mut Context<'_>) -> Poll<Readiness> {
266        Poll::Ready(Readiness::Close)
267    }
268
269    #[inline]
270    fn poll_write_ready(&self, _: &mut Context<'_>) -> Poll<Readiness> {
271        Poll::Ready(Readiness::Close)
272    }
273
274    #[inline]
275    fn process_read_buf(&self, _: &mut FilterCtx<'_>) -> io::Result<()> {
276        Ok(())
277    }
278
279    #[inline]
280    fn process_write_buf(&self, _: &mut FilterCtx<'_>) -> io::Result<()> {
281        Ok(())
282    }
283
284    #[inline]
285    fn shutdown(&self, _: &mut FilterCtx<'_>) -> io::Result<Poll<()>> {
286        Ok(Poll::Ready(()))
287    }
288}