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}