Skip to main content

IoBoxed

Struct IoBoxed 

Source
pub struct IoBoxed(/* private fields */);
Expand description

An Io object whose filter-chain type has been erased.

Implementations§

Source§

impl IoBoxed

Source

pub fn take(&mut self) -> Self

Transfers the live I/O state into a new object.

This does not clone the connection. The current object is replaced with a stopped placeholder and should no longer be used for I/O.

Methods from Deref<Target = Io<Sealed>>§

Source

pub fn get_ref(&self) -> IoRef

Returns a cloneable reference to this connection’s shared state.

Source

pub unsafe fn take(&self) -> Self

Transfers the live connection state into a new Io object.

This does not clone the connection. self is replaced with a stopped placeholder and should no longer be used for I/O.

§Safety

No reference derived from self may be alive across this call, the filter returned by Io::filter, the IoRef it dereferences to and the config returned by IoRef::cfg all borrow the transferred state, which is dropped together with the returned Io.

Source

pub unsafe fn set_config<T: Into<SharedCfg>>(&self, cfg: T)

Replaces this connection’s shared I/O configuration.

The write-buffer page size and eager-write enablement are updated immediately, the read page size restarts at the new min read size. Existing allocated buffers and an already registered timer are not recreated.

§Safety

No reference obtained from IoRef::cfg for this connection may be live when this method is called or used afterward. Replacing the configuration may release the allocation backing those references.

Source

pub async fn recv<U>( &self, codec: &U, ) -> Result<Option<U::Item>, Either<U::Error, Error>>
where U: Decoder,

Reads and decodes the next item from the incoming stream.

Returns Ok(None) when the connection closed before another item could be decoded and nothing was left undecoded, whether the peer disconnected or the shutdown was started locally.

If the peer closed its write half while the codec still held a partial item, the stream was truncated and this returns io::ErrorKind::UnexpectedEof in Either::Right rather than Ok(None), so that a cut-off frame is not mistaken for a clean end of stream. Undecodable bytes left after a locally started shutdown are not treated as truncation.

Codec errors are returned in Either::Left. Dispatcher timeouts and connection errors are returned in Either::Right; a connection error may originate from the transport, a filter, or shutdown.

If write backpressure prevents further reads, this method first waits for the write buffer to fall below its configured threshold. A dispatcher timeout that fires during this wait is returned as well.

Source

pub async fn read_exact(&self, dst: &mut [u8]) -> Result<()>

Reads exactly enough bytes from this I/O stream to fill dst.

If there is not enough data available, waits for incoming data. If clean EOF or an error-free shutdown occurs before dst is filled, this returns io::ErrorKind::UnexpectedEof. Transport errors are passed through unchanged.

Each wait goes through read_more, so this releases read backpressure unconditionally rather than waiting for the read buffer to drain to half the high watermark.

Source

pub async fn read_more(&self) -> Result<Option<()>>

Waits until application-facing data is available and allows the transport to read more data.

If reads are paused or under backpressure, calling this method resumes the read task. This is not a passive check of the current buffer, and read backpressure is released however much data is still buffered.

Returns Ok(Some(())) when input that has not been reported yet is available, Ok(None) when no further input will be reported, and Err if the transport failed. See poll_read_more for what None means after a clean EOF.

Source

pub async fn read_notify(&self) -> Result<Option<()>>

Waits for the next read from the transport.

Use this when a filter needs more source bytes. Unlike read_more, data already waiting in the application buffer does not complete this wait. If the read task is paused, this method wakes it.

Returns Some(()) when the transport provides more input. If clean EOF leaves final data in the application buffer, it returns Some(()) once and None afterward. A transport error is returned unchanged.

Source

pub async fn send<U>( &self, item: U::Item, codec: &U, ) -> Result<(), Either<U::Error, Error>>
where U: Encoder,

Encodes an item and sends it to the peer, fully flushing the write buffer.

The flush is bounded by the write timeout, if one is set; when it expires this returns io::ErrorKind::TimedOut in Either::Right.

Source

pub async fn flush(&self, full: bool) -> Result<()>

Wakes the write task and requests a flush of queued output.

This is the asynchronous counterpart to poll_flush. A full flush completes once all output has reached the peer, including output a transport has taken ownership of but not yet written.

The wait is bounded by the write timeout, if one is set; when it expires this returns io::ErrorKind::TimedOut.

Source

pub async fn shutdown(&self) -> Result<()>

Gracefully shuts down the I/O stream.

Shutdown runs in two phases, bounded together by a single IoConfig::set_shutdown_timeout. First the filters shut down while both directions stay open, so a filter can emit its closing data and read the peer’s. Then the transport drains the remaining output, pauses the read side, and closes the connection.

This completes once the transport backend has finished its shutdown operation, not merely once the output has been drained.

If the shutdown deadline expires in either phase, this returns an io::ErrorKind::TimedOut error once the transport has stopped.

Source

pub fn poll_read_more(&self, cx: &mut Context<'_>) -> Poll<Result<Option<()>>>

Polls for application-facing data and allows the transport to read more.

If reads are paused or under backpressure, this resumes the read task. It therefore changes the read state and should not be used as a passive buffer check.

The release is unconditional: unlike consumption through IoRef::decode, IoRef::with_buf, IoRef::with_read_src or IoRef::with_read_dst, which waits for the read buffer to fall to at most half the high watermark, this releases read backpressure however much data is still buffered. Asking for more input is taken as the dispatcher declaring itself able to accept it. This also resumes reads paused because output produced by reading, for example replies to peer pings, has not drained yet.

§Returns
  • Poll::Pending while waiting for more data.

  • Poll::Ready(Ok(Some(()))) when input that has not been reported yet is available.

  • Poll::Ready(Ok(None)) when no further input will be reported: after a clean EOF once the available input has been reported, or when the stream closes without an error. Clean EOF leaves the write half open.

    This reports arrivals, not buffer contents. Input that has already been reported stays in the read buffer and remains decodable through IoRef::decode and IoRef::with_read_dst, so None does not imply that the read buffer is empty.

  • Poll::Ready(Err(e)) if the transport failed.

Source

pub fn poll_read_notify(&self, cx: &mut Context<'_>) -> Poll<Result<Option<()>>>

Polls for the next read from the transport.

This is the polling version of read_notify. Existing application data does not make it ready. When another transport read is needed, this wakes the read task and registers the current waker.

Some(()) means that more input arrived. Clean EOF may produce one last Some(()) when filters leave final application data; later polls return None. Transport errors are returned unchanged.

Paused or back-pressured reads are resumed through poll_read_more, which releases read backpressure unconditionally.

Source

pub fn poll_recv<U>( &self, codec: &U, cx: &mut Context<'_>, ) -> Poll<Result<U::Item, RecvError<U>>>
where U: Decoder,

Decodes the next item from the incoming byte stream.

Returns Poll::Pending when the codec needs more input, after going through poll_read_more, which wakes the read task and releases read backpressure unconditionally.

An error return does not register the waker.

Source

pub fn poll_recv_decode<U>( &self, codec: &U, cx: &mut Context<'_>, ) -> Result<Decoded<U::Item>, RecvError<U>>
where U: Decoder,

Attempts to decode an item and reports buffer progress.

Decoded::consumed is the number of bytes consumed by this decode attempt and Decoded::remains is the number left in the application-facing read buffer. If the codec needs more input, this returns Ok with item set to None after arranging for cx to be woken when progress is possible.

An error return does not register the waker. While the connection is open, an expired dispatcher timer and then active write backpressure are reported before the read buffer is decoded, so an error never follows a decode attempt and Decoded is never lost. The caller decides whether input received together with a timeout is decoded first, see IoRef::decode_item. Once the connection is closing neither is reported, the buffered input is decoded and RecvError::PeerGone is returned when no item is left.

When the codec needs more input this goes through poll_read_more, which releases read backpressure unconditionally.

Source

pub fn poll_flush(&self, cx: &mut Context<'_>, full: bool) -> Poll<Result<()>>

Wakes the write task and instructs it to flush data.

A full flush waits until all output has reached the peer.

Otherwise this returns immediately while the outstanding size is below the configured high watermark. Reaching that watermark enables write backpressure, and the call then waits until the outstanding size falls to half of it.

Output that a completion based transport has taken ownership of counts as outstanding until it reaches the peer, so a full flush does not complete while a write is still in flight.

Source

pub fn poll_shutdown(&self, cx: &mut Context<'_>) -> Poll<Result<()>>

Polls graceful shutdown through transport completion.

Poll::Ready is returned only after the transport backend marks the connection stopped, not merely when filter shutdown and flushing finish.

Source

pub fn poll_read_pause(&self, cx: &mut Context<'_>) -> Poll<IoStatusUpdate>

Pauses the read task and polls for a status update.

The transport stops reading until the pause is cancelled. There is no explicit resume: the pause is cancelled implicitly by any operation that touches the read buffer or asks for more input, namely read_more, poll_read_more, IoRef::decode, IoRef::decode_item, and IoRef::with_read_dst. Releasing read backpressure cancels it as well. Because those methods are available through every IoRef clone, the pause holds only while no other holder touches the read buffer.

See poll_status_update for the reported updates.

Source

pub fn poll_status_update(&self, cx: &mut Context<'_>) -> Poll<IoStatusUpdate>

Polls for available status updates.

Timeout consumes the pending dispatcher-timeout notification. WriteBackpressure is reported while backpressure is active. The poll that observes the write buffer falling below its release threshold releases backpressure and reports no status update, matching poll_flush. PeerGone is returned once the connection has closed, whether the peer disconnected, the transport failed, or the shutdown was started locally.

Source

pub fn register_dispatch(&self, cx: &mut Context<'_>)

Registers a dispatch task.

Methods from Deref<Target = IoRef>§

Source

pub fn id(&self) -> Id

Gets the ID.

Source

pub fn tag(&self) -> &'static str

Gets the I/O tag.

Source

pub fn cfg(&self) -> &IoConfig

Gets the configuration.

Source

pub fn shared(&self) -> SharedCfg

Gets the shared configuration.

Source

pub fn is_active(&self) -> bool

Checks whether the I/O stream is active.

This becomes false as soon as the connection leaves its active state, whether it was closed locally, force-terminated, or the transport reported the peer as gone. Closing is not instantaneous, so buffered output may still be flushing and buffered input stays readable after this goes false; use is_closed to ask whether closing has finished.

Source

pub fn is_read_eof(&self) -> bool

Checks whether the transport read half reached clean EOF.

Buffered input remains available and the write half may still be used.

Source

pub fn is_closed(&self) -> bool

Checks whether the I/O stream is closed.

This becomes true once the backend released the underlying socket and transport teardown has finished, so nothing further can be read from or delivered to the peer. Every way a connection can end reaches this state, whether it closed gracefully, was force-terminated or the peer disappeared. Use is_active to also cover a close that is still in progress.

Buffered input that was already received stays readable.

Source

pub fn is_rd_backpressure(&self) -> bool

Checks whether read back-pressure is enabled.

This becomes true once unread data in the application-facing read buffer reaches the configured high watermark, which parks the transport read task.

Two different paths release it. Consuming through decode, with_buf, with_read_src or with_read_dst releases it once the buffer has fallen to at most half the high watermark. Asking for more input through Io::poll_read_more, and the methods built on it, releases it immediately however much data is still buffered.

Source

pub fn is_wr_backpressure(&self) -> bool

Checks whether write back-pressure is enabled.

This becomes true once outstanding output reaches the configured high watermark. Outstanding output includes buffered data and data that the transport owns but has not yet written to the peer. It is released once the outstanding size falls to half the high watermark.

Nothing enforces the signal: encoding continues to succeed while it is set. Producers that are not driven by a dispatcher should check this before encoding more, or the write buffer grows without bound. See encode and write_ready.

Source

pub fn is_read_filter_paused(&self) -> bool

Checks whether transport reads are paused by the filter chain.

This is true while a filter is not ready for transport reads although the io state allows them, for example a filter that waits for its own resources. No input arrives during the pause, it is not caused by the peer, so read timeouts should not run. The dispatcher is notified when the pause starts and when it ends.

Source

pub fn is_write_filter_paused(&self) -> bool

Checks whether transport writes are paused by the filter chain.

This is true while buffered output is waiting for a filter that is not ready for transport writes. Output does not drain during the pause, it is not caused by the peer, so write timeouts should not run. The dispatcher is notified when the pause starts and when it ends.

Source

pub async fn write_ready(&self) -> Result<()>

Waits until the write buffer can accept more output.

Completes immediately unless write back-pressure is enabled. While it is, waits until the outstanding output falls to the release threshold (half of the high watermark), the level at which the dispatcher releases back-pressure. Any number of tasks can wait at once, the write task wakes them without involving the dispatcher.

Producers that are not driven by a dispatcher can await this before encoding more, so the write buffer does not grow without bound.

Fails once the connection is closing or closed, and with io::ErrorKind::TimedOut if the write timeout is set and expires first.

Source

pub fn close(&self)

Gracefully closes the connection.

Initiates the I/O stream shutdown process.

Source

pub fn terminate(&self)

Force-closes the connection.

The dispatcher does not wait for incomplete responses. The I/O stream is terminated without any graceful period, and whatever is still buffered is discarded.

The transport aborts the connection instead of closing it gracefully, so the peer most likely observes an RST rather than a clean end of stream, and output that has not been acknowledged yet is lost. That is what keeps a truncated response distinguishable from a complete one, but it also means this must not be used to end a connection normally. Use close for that.

Source

pub fn query<T: 'static>(&self) -> QueryItem<T>

Queries filter-specific data.

Source

pub fn encode<U>( &self, item: U::Item, codec: &U, ) -> Result<(), <U as Encoder>::Error>
where U: Encoder,

Encodes an item into the write buffer.

This method reports codec errors only. Any io::Error produced while buffering is discarded: if the connection is already closing or closed the item is not encoded, and a transport or filter error raised by an eager backend write is dropped. Such errors remain observable later through crate::Io::poll_flush or crate::Io::poll_recv. Use encode_slice or encode_bytes when they must be observed at the call site.

§Back-pressure is advisory

Encoding never blocks and never refuses. Once buffered output reaches the configured high watermark this arms write back-pressure and wakes the dispatch task, but the item is still buffered and Ok is still returned. A caller that keeps encoding without consulting is_wr_backpressure, or awaiting Io::poll_status_update or Io::poll_flush, will grow the write buffer without bound, because a slow peer cannot slow the producer down on its own. Honouring the signal is the caller’s responsibility.

Source

pub fn encode_slice(&self, src: &[u8]) -> Result<()>

Encodes the slice into the write buffer.

If this triggers an eager backend write, any transport or filter error from that write is returned immediately.

Write back-pressure is advisory here too; see encode.

Source

pub fn encode_bytes<B>(&self, src: B) -> Result<()>
where BytePage: From<B>,

Writes bytes to the write buffer.

If this triggers an eager backend write, any transport or filter error from that write is returned immediately.

Write back-pressure is advisory here too; see encode.

Source

pub fn decode<U>( &self, codec: &U, ) -> Result<Option<<U as Decoder>::Item>, <U as Decoder>::Error>
where U: Decoder,

Attempts to decode a frame from the read buffer.

Once the transport reached eof this uses Decoder::decode_eof instead of Decoder::decode.

This mutates the read state: it clears read readiness, and consuming enough bytes may release read backpressure. It also cancels a pause installed by Io::poll_read_pause and wakes the transport read task.

Decoded frames that share the read buffer’s allocation keep the whole buffer alive, see IoConfig::set_read_size.

Source

pub fn decode_item<U>( &self, codec: &U, ) -> Result<Decoded<<U as Decoder>::Item>, <U as Decoder>::Error>
where U: Decoder,

Attempts to decode a frame from the read buffer.

Decoded::consumed reports the bytes taken by this attempt and Decoded::remains the bytes left in the application-facing read buffer. Once the transport reached eof this uses Decoder::decode_eof instead of Decoder::decode.

Like decode, this mutates the read state: it clears read readiness, may release read backpressure, and cancels a pause installed by Io::poll_read_pause.

Source

pub fn send_buf(&self) -> Result<()>

Sends the write buffer to the I/O layer.

Requires the underlying runtime to implement .write(); otherwise, no action is taken.

Source

pub fn with_buf<F, R>(&self, f: F) -> Result<R>
where F: FnOnce(&mut FilterBuf<'_>) -> R,

Provides temporary access to the outermost filter buffers.

Filter callbacks run before and after f, and any produced write data is scheduled for delivery after the closure returns. Errors from an eager backend write are returned to the caller.

The destination exposed by FilterBuf::with_read_buffers is the application-facing read destination, so consuming enough of it releases read backpressure and cancels an installed read pause.

Buffer access is not reentrant. The closure must not use this connection to access an overlapping read or write buffer, or change the filter chain.

Source

pub fn with_read_dst<F, R>(&self, f: F) -> R
where F: FnOnce(&mut BytesMut) -> R,

Provides mutable access to the application-facing read destination.

This holds the decoded bytes the application consumes; see with_read_src for the transport-facing source.

This mutates the read state whether or not f consumes anything. While read back-pressure is active nothing is released until the buffer has fallen to at most half the high watermark, so until then read readiness and any installed read pause are left in place. Once it has, or when back-pressure was not active, read readiness is cleared and a pause installed by Io::poll_read_pause is cancelled, waking the transport read task.

Use crate::Io::poll_read_more rather than this method to check whether data is available.

The closure must not access this connection’s application-facing read destination again. Nested access is unsupported and may terminate the connection or lose nested buffer changes.

Source

pub fn read_dst_size(&self) -> usize

Returns the size of the application-facing read destination.

Unlike with_read_dst this does not change the read state, read readiness, read back-pressure and an installed read pause are left in place, and no buffer is allocated. Returns 0 when called from inside with_read_dst.

Source

pub fn with_write_src<F, R>(&self, f: F) -> Result<R>
where F: FnOnce(&mut BytePages) -> R,

Provides mutable access to the application-facing write source.

This holds the bytes the application produces; see with_write_dst for the transport-facing destination.

Returns an error without invoking f if the connection is closing or closed. Data appended by f is scheduled for delivery. If that starts an eager backend write, its transport or filter error is returned.

§Panics

Panics if the closure accesses the same application-facing write buffer again.

Source

pub fn with_read_src<F, R>(&self, f: F) -> R
where F: FnOnce(&mut BytesMut) -> R,

Provides mutable access to the transport-facing read source.

This is the buffer the transport fills; it is the counterpart of the application-facing destination exposed by with_read_dst. Primarily intended for transport and filter implementations.

Without a filter installed this is the same buffer as the application-facing destination, so consuming enough of it releases read backpressure and cancels an installed read pause. Unlike with_read_dst it never clears read readiness.

The closure must not access this connection’s transport-facing read source again. Nested access is unsupported and may terminate the connection or lose nested buffer changes.

Source

pub fn with_write_dst<F, R>(&self, f: F) -> R
where F: FnOnce(&mut BytePages) -> R,

Provides mutable access to the transport-facing write destination.

This is the buffer the transport drains; it is the counterpart of the application-facing source exposed by with_write_src. Primarily intended for transport and filter implementations.

§Panics

Panics if the closure accesses the same transport-facing write buffer again.

Source

pub fn notify_dispatcher(&self)

Wakeup dispatcher

Source

pub fn notify_timeout(&self)

Wakeup dispatcher and send Timeout error

Source

pub fn timer_handle(&self) -> TimerHandle

Returns the currently registered dispatcher timer handle.

TimerHandle::ZERO is returned when no timer is registered.

Source

pub fn start_timer(&self, timeout: Seconds) -> TimerHandle

Starts or updates the dispatcher timer.

The timer uses second-granularity deadlines. When it expires, poll_status_update reports IoStatusUpdate::Timeout.

A zero timeout cancels the current timer but does not consume a timeout notification that has already been delivered. Use stop_timer when leaving a protocol phase to also clear such a notification.

Source

pub fn stop_timer(&self)

Stops the timer and clears any pending timeout notification.

Source

pub fn on_disconnect(&self) -> Waiter<'static> ⓘ

Returns a future that resolves when the complete I/O stream disconnects.

A clean peer read EOF does not resolve this future because the write half remains usable. It resolves once the transport backend reports that teardown has finished, which happens after local shutdown or force termination. terminate requests that teardown but does not itself resolve the future.

Source

pub fn wake(&self, tag: usize)

Wakes all Waiter waiters of the tag.

Tags reserved for internal use are ignored.

Source

pub fn waiter(&self, tag: usize) -> Waiter<'_> ⓘ

Creates a Waiter for the tag.

The waiter registers on its first poll.

§Panics

Panics in debug builds if the tag is reserved for internal use, usize::MAX - 1 and usize::MAX - 2 are reserved.

Trait Implementations§

Source§

impl Debug for IoBoxed

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Deref for IoBoxed

Source§

type Target = Io<Sealed>

The resulting type after dereferencing.
Source§

fn deref(&self) -> &Self::Target

Dereferences the value.
Source§

impl<F: Filter> From<Io<F>> for IoBoxed

Source§

fn from(io: Io<F>) -> Self

Converts to this type from the input type.
Source§

impl From<IoBoxed> for Io<Sealed>

Source§

fn from(value: IoBoxed) -> Self

Converts to this type from the input type.
Source§

impl<F: Filter> RequestState<IoBoxed> for Io<F>

Source§

type State = ()

State extracted from the request.
Source§

fn unpack(self) -> ((), IoBoxed)

Splits this value into its state and request components.
Source§

impl RequestState<IoBoxed> for IoBoxed

Source§

type State = ()

State extracted from the request.
Source§

fn unpack(self) -> ((), IoBoxed)

Splits this value into its state and request components.
Source§

impl<F: Filter, St: 'static> RequestState<IoBoxed> for State<St, Io<F>>

Source§

type State = St

State extracted from the request.
Source§

fn unpack(self) -> (St, IoBoxed)

Splits this value into its state and request components.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<P, T> Receiver for P
where P: Deref<Target = T> + ?Sized, T: ?Sized,

Source§

type Target = T

🔬This is a nightly-only experimental API. (arbitrary_self_types)
The target type on which the method may be called.
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.