Skip to main content

ntex_io/
ioref.rs

1use std::{any, fmt, hash, io, ptr};
2
3use ntex_bytes::{BytePage, BytePages, BytesMut};
4use ntex_codec::{Decoder, Encoder};
5use ntex_service::cfg::SharedCfg;
6use ntex_util::time::Seconds;
7
8use crate::ops::{Id, Iops, TimerHandle};
9use crate::waiters::{TAG_DISCONNECT, TAG_WRITE, Waiter};
10use crate::{Decoded, Filter, FilterBuf, Flags, Handle, IoConfig, IoContext, IoRef, types};
11
12impl IoRef {
13    #[inline]
14    /// Gets the ID.
15    pub fn id(&self) -> Id {
16        self.0.id()
17    }
18
19    #[inline]
20    /// Gets the I/O tag.
21    pub fn tag(&self) -> &'static str {
22        self.0.tag()
23    }
24
25    #[doc(hidden)]
26    /// Gets the state flags. (for debug purpose only)
27    pub fn flags(&self) -> Flags {
28        self.0.flags.clone()
29    }
30
31    #[inline]
32    /// Gets the current filter.
33    pub(crate) fn filter(&self) -> &dyn Filter {
34        self.0.filter()
35    }
36
37    #[inline]
38    /// Gets the configuration.
39    pub fn cfg(&self) -> &IoConfig {
40        &self.0.cfg
41    }
42
43    #[inline]
44    /// Gets the shared configuration.
45    pub fn shared(&self) -> SharedCfg {
46        self.0.cfg.shared()
47    }
48
49    #[inline]
50    /// Checks whether the I/O stream is active.
51    ///
52    /// This becomes `false` as soon as the connection leaves its active
53    /// state, whether it was closed locally, force-terminated, or the
54    /// transport reported the peer as gone. Closing is not instantaneous, so
55    /// buffered output may still be flushing and buffered input stays
56    /// readable after this goes `false`; use
57    /// [`is_closed`](Self::is_closed) to ask whether closing has finished.
58    pub fn is_active(&self) -> bool {
59        self.0.flags.is_active()
60    }
61
62    #[inline]
63    /// Checks whether the transport read half reached clean EOF.
64    ///
65    /// Buffered input remains available and the write half may still be used.
66    pub fn is_read_eof(&self) -> bool {
67        self.0.flags.is_read_eof()
68    }
69
70    #[inline]
71    /// Checks whether the I/O stream is closed.
72    ///
73    /// This becomes `true` once the backend released the underlying socket
74    /// and transport teardown has finished, so nothing further can be read
75    /// from or delivered to the peer. Every way a connection can end reaches
76    /// this state, whether it closed gracefully, was force-terminated or the
77    /// peer disappeared. Use [`is_active`](Self::is_active) to also cover a
78    /// close that is still in progress.
79    ///
80    /// Buffered input that was already received stays readable.
81    pub fn is_closed(&self) -> bool {
82        self.0.flags.is_closed()
83    }
84
85    #[inline]
86    /// Checks whether read back-pressure is enabled.
87    ///
88    /// This becomes `true` once unread data in the application-facing read
89    /// buffer reaches the configured high watermark, which parks the transport
90    /// read task.
91    ///
92    /// Two different paths release it. Consuming through
93    /// [`decode`](Self::decode), [`with_buf`](Self::with_buf),
94    /// [`with_read_src`](Self::with_read_src) or
95    /// [`with_read_dst`](Self::with_read_dst) releases it once the buffer has
96    /// fallen to at most half the high watermark. Asking for more input through
97    /// [`Io::poll_read_more`](crate::Io::poll_read_more), and the methods built
98    /// on it, releases it immediately however much data is still buffered.
99    pub fn is_rd_backpressure(&self) -> bool {
100        self.0.flags.is_rd_backpressure()
101    }
102
103    #[inline]
104    /// Checks whether write back-pressure is enabled.
105    ///
106    /// This becomes `true` once outstanding output reaches the configured high
107    /// watermark. Outstanding output includes buffered data and data that the
108    /// transport owns but has not yet written to the peer. It is released once
109    /// the outstanding size falls to half the high watermark.
110    ///
111    /// Nothing enforces the signal: encoding continues to succeed while it is
112    /// set. Producers that are not driven by a dispatcher should check this
113    /// before encoding more, or the write buffer grows without bound. See
114    /// [`encode`](Self::encode) and [`write_ready`](Self::write_ready).
115    pub fn is_wr_backpressure(&self) -> bool {
116        self.0.flags.is_wr_backpressure()
117    }
118
119    #[inline]
120    /// Checks whether transport reads are paused by the filter chain.
121    ///
122    /// This is `true` while a filter is not ready for transport reads although
123    /// the io state allows them, for example a filter that waits for its own
124    /// resources. No input arrives during the pause, it is not caused by the
125    /// peer, so read timeouts should not run. The dispatcher is notified when
126    /// the pause starts and when it ends.
127    pub fn is_read_filter_paused(&self) -> bool {
128        self.0.flags.is_read_filter_paused()
129    }
130
131    #[inline]
132    /// Checks whether transport writes are paused by the filter chain.
133    ///
134    /// This is `true` while buffered output is waiting for a filter that is
135    /// not ready for transport writes. Output does not drain during the pause,
136    /// it is not caused by the peer, so write timeouts should not run. The
137    /// dispatcher is notified when the pause starts and when it ends.
138    pub fn is_write_filter_paused(&self) -> bool {
139        self.0.flags.is_write_filter_paused()
140    }
141
142    /// Waits until the write buffer can accept more output.
143    ///
144    /// Completes immediately unless write back-pressure is enabled. While it
145    /// is, waits until the outstanding output falls to the release threshold
146    /// (half of the high watermark), the level at which the dispatcher
147    /// releases back-pressure. Any number of tasks can wait at once, the
148    /// write task wakes them without involving the dispatcher.
149    ///
150    /// Producers that are not driven by a dispatcher can await this before
151    /// encoding more, so the write buffer does not grow without bound.
152    ///
153    /// Fails once the connection is closing or closed, and with
154    /// [`io::ErrorKind::TimedOut`] if the
155    /// [write timeout](crate::IoConfig::set_write_timeout) is set and expires
156    /// first.
157    pub async fn write_ready(&self) -> io::Result<()> {
158        self.0.write_ready().await
159    }
160
161    /// Gracefully closes the connection.
162    ///
163    /// Initiates the I/O stream shutdown process.
164    pub fn close(&self) {
165        self.0.start_shutdown();
166    }
167
168    /// Force-closes the connection.
169    ///
170    /// The dispatcher does not wait for incomplete responses. The I/O stream is
171    /// terminated without any graceful period, and whatever is still buffered
172    /// is discarded.
173    ///
174    /// The transport aborts the connection instead of closing it gracefully, so
175    /// the peer most likely observes an `RST` rather than a clean end of
176    /// stream, and output that has not been acknowledged yet is lost. That is
177    /// what keeps a truncated response distinguishable from a complete one, but
178    /// it also means this must not be used to end a connection normally. Use
179    /// [`close`](Self::close) for that.
180    pub fn terminate(&self) {
181        log::trace!("{}: Terminate io stream object", self.tag());
182        self.0.force_close_connection();
183    }
184
185    /// Queries filter-specific data.
186    pub fn query<T: 'static>(&self) -> types::QueryItem<T> {
187        let _borrow = self.0.buffer.borrow();
188        types::QueryItem::new(self.filter().query(any::TypeId::of::<T>()))
189    }
190
191    #[inline]
192    /// Encodes an item into the write buffer.
193    ///
194    /// This method reports codec errors only. Any `io::Error` produced while
195    /// buffering is discarded: if the connection is already closing or closed
196    /// the item is not encoded, and a transport or filter error raised by an
197    /// eager backend write is dropped. Such errors remain observable later
198    /// through [`crate::Io::poll_flush`] or [`crate::Io::poll_recv`]. Use
199    /// [`encode_slice`](Self::encode_slice) or
200    /// [`encode_bytes`](Self::encode_bytes) when they must be observed at the
201    /// call site.
202    ///
203    /// # Back-pressure is advisory
204    ///
205    /// Encoding never blocks and never refuses. Once buffered output reaches
206    /// the configured high watermark this arms write back-pressure and wakes
207    /// the dispatch task, but the item is still buffered and `Ok` is still
208    /// returned. A caller that keeps encoding without consulting
209    /// [`is_wr_backpressure`](Self::is_wr_backpressure), or awaiting
210    /// [`Io::poll_status_update`](crate::Io::poll_status_update) or
211    /// [`Io::poll_flush`](crate::Io::poll_flush), will grow the write buffer
212    /// without bound, because a slow peer cannot slow the producer down on its
213    /// own. Honouring the signal is the caller's responsibility.
214    pub fn encode<U>(&self, item: U::Item, codec: &U) -> Result<(), <U as Encoder>::Error>
215    where
216        U: Encoder,
217    {
218        self.with_write_src(|buf| codec.encode(item, buf))
219            .unwrap_or_else(|_| Ok(()))
220    }
221
222    #[inline]
223    /// Encodes the slice into the write buffer.
224    ///
225    /// If this triggers an eager backend write, any transport or filter error
226    /// from that write is returned immediately.
227    ///
228    /// Write back-pressure is advisory here too; see [`encode`](Self::encode).
229    pub fn encode_slice(&self, src: &[u8]) -> io::Result<()> {
230        self.with_write_src(|buf| buf.extend_from_slice(src))
231    }
232
233    #[inline]
234    /// Writes bytes to the write buffer.
235    ///
236    /// If this triggers an eager backend write, any transport or filter error
237    /// from that write is returned immediately.
238    ///
239    /// Write back-pressure is advisory here too; see [`encode`](Self::encode).
240    pub fn encode_bytes<B>(&self, src: B) -> io::Result<()>
241    where
242        BytePage: From<B>,
243    {
244        self.with_write_src(|buf| buf.append(src))
245    }
246
247    /// Attempts to decode a frame from the read buffer.
248    ///
249    /// Once the transport reached eof this uses [`Decoder::decode_eof`]
250    /// instead of [`Decoder::decode`].
251    ///
252    /// This mutates the read state: it clears read readiness, and consuming
253    /// enough bytes may release read backpressure. It also cancels a pause
254    /// installed by [`Io::poll_read_pause`](crate::Io::poll_read_pause) and wakes the transport
255    /// read task.
256    ///
257    /// Decoded frames that share the read buffer's allocation keep the whole
258    /// buffer alive, see [`IoConfig::set_read_size`](crate::IoConfig::set_read_size).
259    pub fn decode<U>(
260        &self,
261        codec: &U,
262    ) -> Result<Option<<U as Decoder>::Item>, <U as Decoder>::Error>
263    where
264        U: Decoder,
265    {
266        self.0.buffer.with_read_dst(self, |buf| {
267            let res = self.decode_buf(codec, buf);
268            self.0.flags.unset_read_ready();
269            self.update_read_destination(buf);
270            res
271        })
272    }
273
274    /// Attempts to decode a frame from the read buffer.
275    ///
276    /// `Decoded::consumed` reports the bytes taken by this attempt and
277    /// `Decoded::remains` the bytes left in the application-facing read
278    /// buffer. Once the transport reached eof this uses
279    /// [`Decoder::decode_eof`] instead of [`Decoder::decode`].
280    ///
281    /// Like [`decode`](Self::decode), this mutates the read state: it clears
282    /// read readiness, may release read backpressure, and cancels a pause
283    /// installed by [`Io::poll_read_pause`](crate::Io::poll_read_pause).
284    pub fn decode_item<U>(
285        &self,
286        codec: &U,
287    ) -> Result<Decoded<<U as Decoder>::Item>, <U as Decoder>::Error>
288    where
289        U: Decoder,
290    {
291        self.0.buffer.with_read_dst(self, |buf| {
292            let len = buf.len();
293            let res = self.decode_buf(codec, buf).map(|item| Decoded {
294                item,
295                remains: buf.len(),
296                consumed: len - buf.len(),
297            });
298            self.0.flags.unset_read_ready();
299            self.update_read_destination(buf);
300            res
301        })
302    }
303
304    fn decode_buf<U: Decoder>(
305        &self,
306        codec: &U,
307        buf: &mut BytesMut,
308    ) -> Result<Option<U::Item>, U::Error> {
309        if self.0.flags.is_read_eof() {
310            codec.decode_eof(buf)
311        } else {
312            codec.decode(buf)
313        }
314    }
315
316    /// Sends the write buffer to the I/O layer.
317    ///
318    /// Requires the underlying runtime to implement `.write()`;
319    /// otherwise, no action is taken.
320    pub fn send_buf(&self) -> io::Result<()> {
321        self.consolidate_write_state(true)
322    }
323
324    pub(crate) fn ops_send_buf(&self) {
325        let st = &self.0;
326        #[cfg(feature = "trace")]
327        log::trace!(
328            "{}: ops-send == buf:{} flags:{:?}",
329            st.tag(),
330            st.buffer.write_buf_size(),
331            st.flags
332        );
333
334        if st.flags.is_wr_send_scheduled() {
335            st.flags.unset_wr_send_scheduled();
336
337            if st.flags.is_write_paused() {
338                // call `Handle::write()`.
339                // if write task is not paused, io write is pending
340                // need to wake write task for io completeion
341                if self.call_write() == WakeWriteTask::Yes {
342                    st.wake_write_task();
343                    st.flags.unset_write_paused();
344                }
345            } else {
346                st.wake_write_task();
347            }
348        }
349    }
350
351    /// Provides temporary access to the outermost filter buffers.
352    ///
353    /// Filter callbacks run before and after `f`, and any produced write data
354    /// is scheduled for delivery after the closure returns. Errors from an
355    /// eager backend write are returned to the caller.
356    ///
357    /// The destination exposed by
358    /// [`FilterBuf::with_read_buffers`](crate::FilterBuf::with_read_buffers) is
359    /// the application-facing read destination, so consuming enough of it
360    /// releases read backpressure and cancels an installed read pause.
361    ///
362    /// Buffer access is not reentrant. The closure must not use this
363    /// connection to access an overlapping read or write buffer, or change the
364    /// filter chain.
365    pub fn with_buf<F, R>(&self, f: F) -> io::Result<R>
366    where
367        F: FnOnce(&mut FilterBuf<'_>) -> R,
368    {
369        self.with_callbacks(|cb| cb.before_processing(self));
370        let result = self.0.buffer.with_filter(self, |ctx| ctx.with_buffer(f));
371        self.with_callbacks(|cb| cb.after_processing(self));
372        self.release_read_destination();
373
374        self.consolidate_write_state(false)?;
375        Ok(result)
376    }
377
378    /// Provides mutable access to the application-facing read destination.
379    ///
380    /// This holds the decoded bytes the application consumes; see
381    /// [`with_read_src`](Self::with_read_src) for the transport-facing source.
382    ///
383    /// This mutates the read state whether or not `f` consumes anything. While
384    /// read back-pressure is active nothing is released until the buffer has
385    /// fallen to at most half the high watermark, so until then read readiness
386    /// and any installed read pause are left in place. Once it has, or when
387    /// back-pressure was not active, read readiness is cleared and a pause
388    /// installed by [`Io::poll_read_pause`](crate::Io::poll_read_pause) is
389    /// cancelled, waking the transport read task.
390    ///
391    /// Use [`crate::Io::poll_read_more`] rather than this method to check
392    /// whether data is available.
393    ///
394    /// The closure must not access this connection's application-facing read
395    /// destination again. Nested access is unsupported and may terminate the
396    /// connection or lose nested buffer changes.
397    pub fn with_read_dst<F, R>(&self, f: F) -> R
398    where
399        F: FnOnce(&mut BytesMut) -> R,
400    {
401        self.0.buffer.with_read_dst(self, |buf| {
402            let res = f(buf);
403            self.update_read_destination(buf);
404            res
405        })
406    }
407
408    #[inline]
409    /// Returns the size of the application-facing read destination.
410    ///
411    /// Unlike [`with_read_dst`](Self::with_read_dst) this does not change the
412    /// read state, read readiness, read back-pressure and an installed read
413    /// pause are left in place, and no buffer is allocated. Returns `0` when
414    /// called from inside [`with_read_dst`](Self::with_read_dst).
415    pub fn read_dst_size(&self) -> usize {
416        self.0.buffer.read_dst_size()
417    }
418
419    /// Provides mutable access to the application-facing write source.
420    ///
421    /// This holds the bytes the application produces; see
422    /// [`with_write_dst`](Self::with_write_dst) for the transport-facing
423    /// destination.
424    ///
425    /// Returns an error without invoking `f` if the connection is closing or
426    /// closed. Data appended by `f` is scheduled for delivery. If that starts
427    /// an eager backend write, its transport or filter error is returned.
428    ///
429    /// # Panics
430    ///
431    /// Panics if the closure accesses the same application-facing write
432    /// buffer again.
433    pub fn with_write_src<F, R>(&self, f: F) -> io::Result<R>
434    where
435        F: FnOnce(&mut BytePages) -> R,
436    {
437        let st = &self.0;
438
439        if st.flags.is_active() {
440            let result = st.buffer.with_write_src(f);
441            self.consolidate_write_state(false)?;
442            Ok(result)
443        } else if st.flags.is_peer_gone() {
444            Err(st.error_or_disconnected())
445        } else {
446            Err(io::Error::other("I/O stream is closing"))
447        }
448    }
449
450    #[inline]
451    /// Provides mutable access to the transport-facing read source.
452    ///
453    /// This is the buffer the transport fills; it is the counterpart of the
454    /// application-facing destination exposed by
455    /// [`with_read_dst`](Self::with_read_dst). Primarily intended for transport
456    /// and filter implementations.
457    ///
458    /// Without a filter installed this is the same buffer as the
459    /// application-facing destination, so consuming enough of it releases read
460    /// backpressure and cancels an installed read pause. Unlike
461    /// [`with_read_dst`](Self::with_read_dst) it never clears read readiness.
462    ///
463    /// The closure must not access this connection's transport-facing read
464    /// source again. Nested access is unsupported and may terminate the
465    /// connection or lose nested buffer changes.
466    pub fn with_read_src<F, R>(&self, f: F) -> R
467    where
468        F: FnOnce(&mut BytesMut) -> R,
469    {
470        let result = self.0.buffer.with_read_src(self, f);
471        self.release_read_destination();
472        result
473    }
474
475    #[inline]
476    /// Provides mutable access to the transport-facing write destination.
477    ///
478    /// This is the buffer the transport drains; it is the counterpart of the
479    /// application-facing source exposed by
480    /// [`with_write_src`](Self::with_write_src). Primarily intended for
481    /// transport and filter implementations.
482    ///
483    /// # Panics
484    ///
485    /// Panics if the closure accesses the same transport-facing write buffer
486    /// again.
487    pub fn with_write_dst<F, R>(&self, f: F) -> R
488    where
489        F: FnOnce(&mut BytePages) -> R,
490    {
491        self.0.buffer.with_write_dst(f)
492    }
493
494    /// Schedules buffered output for delivery and updates write state.
495    ///
496    /// When output is buffered and the write task is paused, this either
497    /// performs an eager in-place write through the transport handle or
498    /// schedules a write operation. An eager write requires direct writes to be
499    /// enabled and, unless `force` is set, at least
500    /// [`IoConfig::write_buf_threshold`](crate::IoConfig::write_buf_threshold)
501    /// bytes to be buffered; `force` makes any non-empty buffer eligible. A
502    /// write operation is scheduled when an eager write leaves data behind, or
503    /// when no eager write was attempted.
504    ///
505    /// Returns the connection error once the connection is stopping or
506    /// terminating with an error set. This is how a failed eager write or
507    /// filter reaches the caller, and it also prevents further eager writes on
508    /// a connection that is already gone.
509    ///
510    /// Finally enables write back-pressure and wakes the dispatcher if buffered
511    /// output has reached the configured high watermark. The size is re-read
512    /// first because an eager write may have drained it.
513    pub(crate) fn consolidate_write_state(&self, force: bool) -> io::Result<()> {
514        let st = &self.0;
515
516        // wake write task if needsed
517        let size = st.buffer.write_buf_size();
518
519        #[cfg(feature = "trace")]
520        log::trace!("{}: write-upd == buf:{size} flags:{:?}", st.tag(), st.flags);
521
522        if size > 0 && st.flags.is_write_paused() {
523            // The app encodes data in response to incoming data,
524            // continuing to fill the write buffer until all data
525            // has been processed. Only then can the runtime wake
526            // the write task to send the buffered data.
527            //
528            // By that time, the buffer may have accumulated a large
529            // amount of data, causing it to be sent in large bursts,
530            // which introduces latency. To prevent this behavior and
531            // flatten data delivery to the peer, IoRef can initiate
532            // out-of-order writes based on a configured threshold.
533            if st.flags.is_direct_wr_enabled() && (force || size >= st.cfg.write_buf_threshold()) {
534                // Send data in-place
535                if self.call_write() == WakeWriteTask::Yes {
536                    #[cfg(feature = "trace")]
537                    log::trace!(
538                        "{}: write-upd == schedule(more):{} flags:{:?}",
539                        st.tag(),
540                        st.buffer.write_buf_size(),
541                        st.flags
542                    );
543                    if !st.flags.is_wr_send_scheduled() {
544                        // More data needs to be sent
545                        st.flags.set_wr_send_scheduled();
546                        Iops::schedule_write(st.id());
547                    }
548                } else {
549                    st.flags.unset_wr_send_scheduled();
550                }
551            } else if !st.flags.is_wr_send_scheduled() {
552                #[cfg(feature = "trace")]
553                log::trace!("{}: write-upd == schedule(too small)", st.tag());
554                st.flags.set_wr_send_scheduled();
555                Iops::schedule_write(st.id());
556            }
557        }
558
559        if !st.flags.is_active()
560            && let Some(err) = st.error()
561        {
562            return Err(err);
563        }
564
565        // A direct write may have changed the amount of buffered data.
566        // In-flight output counts too: it has not reached the peer yet.
567        let size = st.write_outstanding();
568
569        // Enable backpressure
570        if !st.flags.is_wr_backpressure() && st.is_wr_backpressure_needed(size) {
571            st.flags.set_wr_backpressure();
572            st.wake_dispatch_task();
573        }
574        Ok(())
575    }
576
577    /// Updates read state after the application-facing destination was accessed.
578    ///
579    /// While read back-pressure is active nothing is released until `buf` has
580    /// fallen to at most half the high watermark. Until then read readiness and
581    /// any installed read pause are deliberately left in place, keeping the
582    /// transport read task parked. Once it has, read readiness and
583    /// back-pressure are cleared together; without back-pressure only read
584    /// readiness is cleared.
585    ///
586    /// Whenever something is released, a pause installed by
587    /// [`Io::poll_read_pause`](crate::Io::poll_read_pause) is cancelled and the
588    /// transport read task is woken.
589    ///
590    /// See [`release_read_destination`](Self::release_read_destination) for the
591    /// variant used when the caller may have drained a different buffer of the
592    /// filter chain.
593    fn update_read_destination(&self, buf: &mut BytesMut) {
594        let st = &self.0;
595
596        #[cfg(feature = "trace")]
597        log::trace!(
598            "{}: read-upd == buf:{} flags:{:?}",
599            st.tag(),
600            buf.len(),
601            st.flags
602        );
603
604        if st.flags.is_rd_backpressure() {
605            // Keep reads paused until enough buffered data has been consumed.
606            if !st.should_disable_rd_backpressure(buf.len()) {
607                return;
608            }
609            st.flags.unset_read_ready_and_backpressure();
610        } else {
611            st.flags.unset_read_ready();
612        }
613
614        if st.flags.is_read_paused() {
615            st.wake_read_task();
616            st.flags.unset_read_paused();
617        }
618    }
619
620    /// Releases read backpressure and any installed read pause.
621    ///
622    /// Used by the accessors that can drain the application-facing read
623    /// destination without going through
624    /// [`with_read_dst`](Self::with_read_dst). Unlike
625    /// `update_read_destination()` this never clears read readiness, because
626    /// the caller may have touched a different buffer of the chain. It only
627    /// removes a stale pause, so it can never suppress a wakeup.
628    fn release_read_destination(&self) {
629        let st = &self.0;
630
631        if st.flags.is_rd_backpressure() {
632            if !st.should_disable_rd_backpressure(st.buffer.read_dst_size()) {
633                return;
634            }
635            st.flags.unset_read_ready_and_backpressure();
636        }
637
638        if st.flags.is_read_paused() {
639            st.wake_read_task();
640            st.flags.unset_read_paused();
641        }
642    }
643
644    /// Wakeup dispatcher
645    pub fn notify_dispatcher(&self) {
646        log::trace!("{}: Timer, notify dispatcher", self.tag());
647        self.0.wake_dispatch_task();
648    }
649
650    /// Wakeup dispatcher and send Timeout error
651    pub fn notify_timeout(&self) {
652        self.0.notify_timeout();
653    }
654
655    /// Returns the currently registered dispatcher timer handle.
656    ///
657    /// [`TimerHandle::ZERO`] is returned when no timer is registered.
658    pub fn timer_handle(&self) -> TimerHandle {
659        self.0.timeout.get()
660    }
661
662    /// Starts or updates the dispatcher timer.
663    ///
664    /// The timer uses second-granularity deadlines. When it expires,
665    /// [`poll_status_update`](crate::Io::poll_status_update) reports
666    /// [`IoStatusUpdate::Timeout`](crate::IoStatusUpdate::Timeout).
667    ///
668    /// A zero timeout cancels the current timer but does not consume a timeout
669    /// notification that has already been delivered. Use
670    /// [`stop_timer`](Self::stop_timer) when leaving a protocol phase to also
671    /// clear such a notification.
672    pub fn start_timer(&self, timeout: Seconds) -> TimerHandle {
673        let cur_hnd = self.0.timeout.get();
674
675        if timeout.is_zero() {
676            if cur_hnd.is_set() {
677                self.0.timeout.set(TimerHandle::ZERO);
678                cur_hnd.unregister(self);
679            }
680            TimerHandle::ZERO
681        } else if cur_hnd.is_set() {
682            let hnd = cur_hnd.update(timeout, self);
683            if hnd != cur_hnd {
684                log::trace!("{}: Update timer {:?}", self.tag(), timeout);
685                self.0.timeout.set(hnd);
686            }
687            hnd
688        } else {
689            log::trace!("{}: Start timer {:?}", self.tag(), timeout);
690            let hnd = TimerHandle::register(timeout, self);
691            self.0.timeout.set(hnd);
692            hnd
693        }
694    }
695
696    /// Stops the timer and clears any pending timeout notification.
697    pub fn stop_timer(&self) {
698        self.0.flags.check_dispatcher_timeout();
699
700        let hnd = self.0.timeout.get();
701        if hnd.is_set() {
702            log::trace!("{}: Stop timer", self.tag());
703            self.0.timeout.set(TimerHandle::ZERO);
704            hnd.unregister(self);
705        }
706    }
707
708    /// Returns a future that resolves when the complete I/O stream disconnects.
709    ///
710    /// A clean peer read EOF does not resolve this future because the write
711    /// half remains usable. It resolves once the transport backend reports
712    /// that teardown has finished, which happens after local shutdown or
713    /// force termination. [`terminate`](Self::terminate) requests that
714    /// teardown but does not itself resolve the future.
715    pub fn on_disconnect(&self) -> Waiter<'static> {
716        Waiter::new_static(self.clone(), TAG_DISCONNECT)
717    }
718
719    /// Wakes all [`Waiter`](crate::Waiter) waiters of the tag.
720    ///
721    /// Tags reserved for internal use are ignored.
722    pub fn wake(&self, tag: usize) {
723        if tag < TAG_WRITE {
724            self.0.extensions.wake(tag);
725        }
726    }
727
728    /// Creates a [`Waiter`](crate::Waiter) for the tag.
729    ///
730    /// The waiter registers on its first poll.
731    ///
732    /// # Panics
733    ///
734    /// Panics in debug builds if the tag is reserved for internal use,
735    /// `usize::MAX - 1` and `usize::MAX - 2` are reserved.
736    pub fn waiter(&self, tag: usize) -> Waiter<'_> {
737        Waiter::new(self, tag)
738    }
739
740    #[doc(hidden)]
741    /// Register filter callbacks
742    ///
743    /// The callbacks are discarded once the connection is closed or its `Io`
744    /// has been dropped. They are released together with the `Io`, so storing
745    /// them afterwards could keep the connection state alive through an
746    /// `IoRef` they hold.
747    pub fn register_filter_callbacks<F: crate::IoCallbacks + 'static>(&self, f: F) {
748        if !self.0.flags.is_closed() && !self.0.is_io_dropped() {
749            self.0.extensions.register_filter_callbacks(f);
750        }
751    }
752
753    /// Call handle write method, returns true if `write-paused` is still set
754    fn call_write(&self) -> WakeWriteTask {
755        if let Some(hnd) = self.0.handle.take() {
756            self.0.flags.unset_write_paused();
757            #[cfg(feature = "trace")]
758            log::trace!(
759                "{}: call-write ({}), flags:{:?}",
760                self.tag(),
761                self.0.buffer.write_buf_size(),
762                self.0.flags
763            );
764            let ctx = unsafe { &*(ptr::from_ref(self).cast::<IoContext>()) };
765            hnd.write(ctx);
766            self.restore_handle(hnd);
767        }
768        if self.0.flags.is_write_paused() {
769            WakeWriteTask::No
770        } else {
771            WakeWriteTask::Yes
772        }
773    }
774
775    /// Reinstalls the transport handle after a reentrant transport callback.
776    ///
777    /// The handle is taken for the duration of the call so that a nested
778    /// `call_write()` cannot reenter the transport. If the
779    /// callback terminated the connection, `terminate_connection()` and
780    /// `stop_connection()` found the slot empty and could not release the
781    /// transport, so the handle is dropped here instead of being reinstalled.
782    /// A graceful shutdown keeps it: the write task still needs the transport
783    /// to shut it down.
784    fn restore_handle(&self, hnd: Box<dyn Handle>) {
785        if self.0.flags.is_terminating() || self.0.flags.is_closed() {
786            drop(hnd);
787        } else {
788            self.0.handle.set(Some(hnd));
789        }
790    }
791
792    pub(crate) fn with_callbacks<F>(&self, f: F)
793    where
794        F: FnOnce(&dyn crate::IoCallbacks),
795    {
796        self.0.extensions.with_callbacks(f);
797    }
798}
799
800#[derive(Copy, Clone, PartialEq, Eq, Debug)]
801enum WakeWriteTask {
802    Yes,
803    No,
804}
805
806impl Eq for IoRef {}
807
808impl PartialEq for IoRef {
809    #[inline]
810    fn eq(&self, other: &Self) -> bool {
811        self.0.eq(&other.0)
812    }
813}
814
815impl hash::Hash for IoRef {
816    #[inline]
817    fn hash<H: hash::Hasher>(&self, state: &mut H) {
818        self.0.hash(state);
819    }
820}
821
822impl fmt::Debug for IoRef {
823    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
824        f.debug_struct("IoRef")
825            .field("state", self.0.as_ref())
826            .finish()
827    }
828}
829
830#[cfg(test)]
831mod tests {
832    use std::cell::{Cell, RefCell};
833    use std::{future::Future, future::poll_fn, pin::Pin, rc::Rc, task::Poll};
834
835    use ntex_bytes::Bytes;
836    use ntex_codec::BytesCodec;
837    use ntex_util::future::{Either, lazy};
838    use ntex_util::time::{Millis, sleep, timeout};
839
840    use super::*;
841
842    /// The three public lifecycle predicates answer different questions and
843    /// must not be read as stages of one another.
844    #[ntex::test]
845    async fn lifecycle_predicates_are_independent() {
846        // a force termination starts: it ends abnormally and is closing, but
847        // the backend still holds the socket
848        let (_client, server) = IoTest::create();
849        let io = Io::from(server);
850        assert!(io.is_active() && !io.flags().is_terminating() && !io.is_closed());
851
852        io.terminate();
853        assert!(!io.is_active());
854        assert!(io.flags().is_terminating());
855        assert!(!io.is_closed(), "socket released before teardown ran");
856
857        // a graceful shutdown also releases the socket, without ever being an
858        // abnormal ending
859        let (client, server) = IoTest::create();
860        client.remote_buffer_cap(1024);
861        let io = Io::from(server);
862        io.shutdown().await.unwrap();
863        assert!(io.is_closed());
864        assert!(!io.is_active());
865        assert!(
866            !io.flags().is_terminating(),
867            "a graceful shutdown reported as an abort"
868        );
869    }
870    use crate::{FilterCtx, Io, testing::IoTest};
871
872    const BIN: &[u8] = b"GET /test HTTP/1\r\n\r\n";
873    const TEXT: &str = "GET /test HTTP/1\r\n\r\n";
874
875    #[ntex::test]
876    async fn utils() {
877        let (client, server) = IoTest::create();
878        client.remote_buffer_cap(1024);
879        client.write(TEXT);
880
881        let state = Io::from(server);
882        assert_eq!(state.get_ref(), state.get_ref());
883
884        let msg = state.recv(&BytesCodec).await.unwrap().unwrap();
885        assert_eq!(msg, Bytes::from_static(BIN));
886        assert_eq!(state.get_ref(), state.as_ref().clone());
887        assert!(format!("{state:?}").find("Io {").is_some());
888        assert!(format!("{:?}", state.get_ref()).find("IoRef {").is_some());
889
890        let res = poll_fn(|cx| Poll::Ready(state.poll_recv(&BytesCodec, cx))).await;
891        assert!(res.is_pending());
892        client.write(TEXT);
893        sleep(Millis(50)).await;
894        let res = poll_fn(|cx| Poll::Ready(state.poll_recv(&BytesCodec, cx))).await;
895        if let Poll::Ready(msg) = res {
896            assert_eq!(msg.unwrap(), Bytes::from_static(BIN));
897        }
898
899        client.read_error(io::Error::other("err"));
900        let msg = state.recv(&BytesCodec).await;
901        assert!(msg.is_err());
902        assert!(state.flags().is_closed());
903
904        let (client, server) = IoTest::create();
905        client.remote_buffer_cap(1024);
906        let state = Io::from(server);
907
908        client.read_error(io::Error::other("err"));
909        let res = poll_fn(|cx| Poll::Ready(state.poll_recv(&BytesCodec, cx))).await;
910        if let Poll::Ready(msg) = res {
911            assert!(msg.is_err());
912            assert!(state.flags().is_closed());
913        }
914
915        let (client, server) = IoTest::create();
916        client.remote_buffer_cap(1024);
917        let state = Io::from(server);
918        assert_eq!(0, state.with_write_dst(|b| b.len()));
919        state.encode_slice(b"test").unwrap();
920        assert_eq!(4, state.with_write_dst(|b| b.len()));
921        let buf = client.read().await.unwrap();
922        assert_eq!(buf, Bytes::from_static(b"test"));
923
924        client.write(b"test");
925        state.read_more().await.unwrap();
926        let buf = state.decode(&BytesCodec).unwrap().unwrap();
927        assert_eq!(buf, Bytes::from_static(b"test"));
928
929        client.write_error(io::Error::other("err"));
930        let err = state
931            .send(Bytes::from_static(b"test"), &BytesCodec)
932            .await
933            .unwrap_err();
934        assert!(matches!(
935            err,
936            Either::Right(ref err)
937                if err.kind() == io::ErrorKind::Other && err.to_string() == "err"
938        ));
939        assert!(state.flags().is_closed());
940
941        let res = state.send(Bytes::from_static(b"test"), &BytesCodec).await;
942        assert!(res.is_err());
943
944        let (client, server) = IoTest::create();
945        client.remote_buffer_cap(1024);
946        let state = Io::from(server);
947        state.terminate();
948        assert!(state.flags().is_terminating());
949        assert!(!state.flags().is_stopping());
950        assert!(!state.flags().is_closed());
951        state.shutdown().await.unwrap();
952        assert!(state.flags().is_closed());
953    }
954
955    #[ntex::test]
956    async fn zero_byte_write_reports_write_zero() {
957        let (client, server) = IoTest::create();
958        client.remote_buffer_cap(1024);
959        let state = Io::from(server);
960
961        client.write_zero();
962        let err = state
963            .send(Bytes::from_static(b"test"), &BytesCodec)
964            .await
965            .unwrap_err();
966
967        assert!(matches!(
968            err,
969            Either::Right(ref err) if err.kind() == io::ErrorKind::WriteZero
970        ));
971        assert!(state.flags().is_closed());
972    }
973
974    #[ntex::test]
975    #[allow(clippy::unit_cmp)]
976    async fn on_disconnect() {
977        let (client, server) = IoTest::create();
978        let state = Io::from(server);
979        let mut waiter = state.on_disconnect();
980        assert_eq!(
981            lazy(|cx| Pin::new(&mut waiter).poll(cx)).await,
982            Poll::Pending
983        );
984        let mut waiter2 = waiter.clone();
985        assert_eq!(
986            lazy(|cx| Pin::new(&mut waiter2).poll(cx)).await,
987            Poll::Pending
988        );
989        client.close().await;
990        assert!(state.is_read_eof());
991        assert!(state.is_active());
992        assert_eq!(
993            lazy(|cx| Pin::new(&mut waiter).poll(cx)).await,
994            Poll::Pending
995        );
996        assert_eq!(
997            lazy(|cx| Pin::new(&mut waiter2).poll(cx)).await,
998            Poll::Pending
999        );
1000
1001        timeout(Millis(1000), state.shutdown())
1002            .await
1003            .expect("stream shutdown did not complete")
1004            .unwrap();
1005        timeout(Millis(1000), waiter)
1006            .await
1007            .expect("disconnect waiter was not notified");
1008        timeout(Millis(1000), waiter2)
1009            .await
1010            .expect("cloned disconnect waiter was not notified");
1011
1012        let mut waiter = state.on_disconnect();
1013        assert_eq!(
1014            lazy(|cx| Pin::new(&mut waiter).poll(cx)).await,
1015            Poll::Ready(())
1016        );
1017
1018        let (client, server) = IoTest::create();
1019        let state = Io::from(server);
1020        let mut waiter = state.on_disconnect();
1021        assert_eq!(
1022            lazy(|cx| Pin::new(&mut waiter).poll(cx)).await,
1023            Poll::Pending
1024        );
1025        client.read_error(io::Error::other("err"));
1026        assert_eq!(waiter.await, ());
1027
1028        let mut waiter = state.on_disconnect();
1029        assert_eq!(
1030            lazy(|cx| Pin::new(&mut waiter).poll(cx)).await,
1031            Poll::Ready(())
1032        );
1033    }
1034
1035    #[ntex::test]
1036    async fn write_to_closed_io() {
1037        let (_client, server) = IoTest::create();
1038        let state = Io::from(server);
1039        state.terminate();
1040
1041        assert!(!state.is_active());
1042        assert!(state.encode_slice(TEXT.as_bytes()).is_err());
1043        assert!(state.encode_bytes(Bytes::from_static(BIN)).is_err());
1044        assert!(
1045            state
1046                .with_write_src(|buf| buf.extend_from_slice(BIN))
1047                .is_err()
1048        );
1049    }
1050
1051    #[derive(Debug)]
1052    struct Counter<F> {
1053        layer: F,
1054        idx: usize,
1055        out_bytes: Rc<Cell<usize>>,
1056        read_order: Rc<RefCell<Vec<usize>>>,
1057        write_order: Rc<RefCell<Vec<usize>>>,
1058    }
1059
1060    impl<F: Filter> Filter for Counter<F> {
1061        fn process_read_buf(&self, ctx: &mut FilterCtx<'_>) -> io::Result<()> {
1062            self.read_order.borrow_mut().push(self.idx);
1063            self.layer.process_read_buf(ctx)
1064        }
1065
1066        fn process_write_buf(&self, ctx: &mut FilterCtx<'_>) -> io::Result<()> {
1067            self.write_order.borrow_mut().push(self.idx);
1068            ctx.with_buffer(|buf| {
1069                buf.with_write_buffers(|src, _| {
1070                    self.out_bytes.set(self.out_bytes.get() + src.len());
1071                });
1072            });
1073            self.layer.process_write_buf(ctx)
1074        }
1075
1076        crate::forward_ready!(layer);
1077        crate::forward_query!(layer);
1078        crate::forward_shutdown!(layer);
1079    }
1080
1081    #[ntex::test]
1082    async fn filter() {
1083        let out_bytes = Rc::new(Cell::new(0));
1084        let read_order = Rc::new(RefCell::new(Vec::new()));
1085        let write_order = Rc::new(RefCell::new(Vec::new()));
1086
1087        let (client, server) = IoTest::create();
1088        let io = Io::from(server)
1089            .map_filter(|layer| Counter {
1090                layer,
1091                idx: 1,
1092                out_bytes: out_bytes.clone(),
1093                read_order: read_order.clone(),
1094                write_order: write_order.clone(),
1095            })
1096            .seal();
1097
1098        client.remote_buffer_cap(1024);
1099        client.write(TEXT);
1100        let msg = io.recv(&BytesCodec).await.unwrap().unwrap();
1101        assert_eq!(msg, Bytes::from_static(BIN));
1102
1103        io.send(Bytes::from_static(b"test"), &BytesCodec)
1104            .await
1105            .unwrap();
1106        let buf = client.read().await.unwrap();
1107        assert_eq!(buf, Bytes::from_static(b"test"));
1108
1109        client.write(TEXT);
1110        let msg = io.recv(&BytesCodec).await.unwrap().unwrap();
1111        assert_eq!(msg, Bytes::from_static(BIN));
1112
1113        // the number of write passes depends on how often the runtime polls
1114        // `send()`, each pass sees the 4 queued bytes until the transport
1115        // writes them, at least once from `send()` and once from the transport
1116        assert!(out_bytes.get() >= 8, "{}", out_bytes.get());
1117        assert_eq!(out_bytes.get() % 4, 0, "{}", out_bytes.get());
1118    }
1119
1120    #[ntex::test]
1121    async fn boxed_filter() {
1122        let out_bytes = Rc::new(Cell::new(0));
1123        let read_order = Rc::new(RefCell::new(Vec::new()));
1124        let write_order = Rc::new(RefCell::new(Vec::new()));
1125
1126        let (client, server) = IoTest::create();
1127        let state = Io::from(server)
1128            .map_filter(|layer| Counter {
1129                layer,
1130                idx: 2,
1131                out_bytes: out_bytes.clone(),
1132                read_order: read_order.clone(),
1133                write_order: write_order.clone(),
1134            })
1135            .map_filter(|layer| Counter {
1136                layer,
1137                idx: 1,
1138                out_bytes: out_bytes.clone(),
1139                read_order: read_order.clone(),
1140                write_order: write_order.clone(),
1141            });
1142        let state = state.seal();
1143
1144        client.remote_buffer_cap(1024);
1145        client.write(TEXT);
1146        let msg = state.recv(&BytesCodec).await.unwrap().unwrap();
1147        assert_eq!(msg, Bytes::from_static(BIN));
1148
1149        state
1150            .send(Bytes::from_static(b"test"), &BytesCodec)
1151            .await
1152            .unwrap();
1153        let buf = client.read().await.unwrap();
1154        assert_eq!(buf, Bytes::from_static(b"test"));
1155
1156        // the number of write passes depends on how often the runtime polls
1157        // `send()`, each pass sees the 4 queued bytes in both layers until the
1158        // transport writes them, at least once from `send()` and once from the
1159        // transport
1160        assert!(out_bytes.get() >= 16, "{}", out_bytes.get());
1161        assert_eq!(out_bytes.get() % 8, 0, "{}", out_bytes.get());
1162        assert_eq!(state.0.buffer.with_write_dst(|b| b.len()), 0);
1163
1164        // refs
1165        assert_eq!(Rc::strong_count(&out_bytes), 3);
1166        drop(state);
1167        assert_eq!(Rc::strong_count(&out_bytes), 1);
1168        assert_eq!(*read_order.borrow(), &[1, 2][..]);
1169        let write_order = write_order.borrow();
1170        assert!(write_order.len() >= 6, "{write_order:?}");
1171        assert!(
1172            write_order.chunks(2).all(|c| c == [1, 2]),
1173            "{write_order:?}"
1174        );
1175    }
1176
1177    #[ntex::test]
1178    async fn timer_start_update_and_stop() {
1179        let (_client, server) = IoTest::create();
1180        let io = Io::new(server, SharedCfg::new("TIMER"));
1181        assert_eq!(io.shared().tag(), "TIMER");
1182        assert_eq!(io.timer_handle(), TimerHandle::ZERO);
1183
1184        let hnd = io.start_timer(Seconds(5));
1185        assert!(hnd.is_set());
1186        assert_eq!(io.timer_handle(), hnd);
1187        assert!(hnd.remains() >= Seconds(4) && hnd.remains() <= Seconds(5));
1188        assert!(hnd.instant() > TimerHandle::ZERO.instant());
1189
1190        // the same timeout keeps the registration
1191        assert_eq!(io.start_timer(Seconds(5)), hnd);
1192        assert_eq!(io.timer_handle(), hnd);
1193
1194        // a different timeout moves it
1195        let hnd2 = io.start_timer(Seconds(30));
1196        assert_ne!(hnd2, hnd);
1197        assert_eq!(io.timer_handle(), hnd2);
1198        assert!(hnd2.remains() >= Seconds(29));
1199
1200        // a second io can share the deadline slot
1201        let (_client2, server2) = IoTest::create();
1202        let io2 = Io::from(server2);
1203        assert!(io2.start_timer(Seconds(30)).is_set());
1204
1205        // zero timeout cancels the timer
1206        assert_eq!(io.start_timer(Seconds::ZERO), TimerHandle::ZERO);
1207        assert_eq!(io.timer_handle(), TimerHandle::ZERO);
1208        assert_eq!(io.start_timer(Seconds::ZERO), TimerHandle::ZERO);
1209        io2.stop_timer();
1210
1211        let hnd = TimerHandle::ZERO + Seconds(3);
1212        assert!(hnd.is_set());
1213        assert_eq!(
1214            hnd.instant() - TimerHandle::ZERO.instant(),
1215            std::time::Duration::from_secs(3)
1216        );
1217    }
1218
1219    #[ntex::test]
1220    async fn notify_dispatcher_wakes_status_poll() {
1221        use std::sync::{Arc, atomic::AtomicBool, atomic::Ordering};
1222
1223        struct Flag(AtomicBool);
1224
1225        impl std::task::Wake for Flag {
1226            fn wake(self: Arc<Self>) {
1227                self.0.store(true, Ordering::Relaxed);
1228            }
1229        }
1230
1231        let (_client, server) = IoTest::create();
1232        let io = Io::from(server);
1233
1234        let flag = Arc::new(Flag(AtomicBool::new(false)));
1235        let waker = std::task::Waker::from(flag.clone());
1236        let mut cx = std::task::Context::from_waker(&waker);
1237        assert!(io.poll_status_update(&mut cx).is_pending());
1238
1239        io.notify_dispatcher();
1240        assert!(flag.0.load(Ordering::Relaxed));
1241        // a plain wakeup reports no status update
1242        assert!(lazy(|cx| io.poll_status_update(cx)).await.is_pending());
1243    }
1244
1245    #[ntex::test]
1246    #[allow(clippy::mutable_key_type)]
1247    async fn io_ref_hash() {
1248        let (_client, server) = IoTest::create();
1249        let io = Io::from(server);
1250        let (_client2, server2) = IoTest::create();
1251        let io2 = Io::from(server2);
1252
1253        let mut set = std::collections::HashSet::new();
1254        assert!(set.insert(io.get_ref()));
1255        assert!(!set.insert(io.get_ref()));
1256        assert!(set.insert(io2.get_ref()));
1257
1258        let mut set = std::collections::HashSet::new();
1259        assert!(set.insert(&io));
1260        assert!(!set.insert(&io));
1261        assert!(set.insert(&io2));
1262    }
1263}