Skip to main content

Consumer

Struct Consumer 

Source
pub struct Consumer<C, U> { /* private fields */ }
Expand description

The single consumer that merges the two lanes.

Implementations§

Source§

impl<C, U> Consumer<C, U>

Source

pub async fn recv(&mut self) -> Option<Received<C, U>>

Receive the next item.

Returns None once both lanes are closed and empty.

POLICY: strict control-first priority with an anti-starvation aging cap. Control is dequeued before any waiting user (P1, overtake); after aging_cap consecutive control dequeues one waiting user is forced through (P3 under flood). Users reset the streak; a cap of 0 disables aging entirely.

Once every UserSender is gone AND the ring is drained, the user stream’s one-shot end-marker (Received::UserLaneClosed) is delivered through the same path a user item takes — control keeps its priority over it, aging forces it through exactly as for users, and it resets the streak (it IS the user stream’s end-marker).

Source

pub async fn recv_control(&mut self) -> Option<C>

Receive the next CONTROL item only, never consuming the user lane. None once the control lane is closed and drained. Cancel-safe under the same registration protocol as recv; a user-lane push may cause a spurious wake (one extra loop turn), never a lost wakeup and never a consumed user item.

This path consumes no user items, so it must NOT call usr.release_one_waiter(): releasing a producer parked on a full ring would just re-park it — and the semantics WANT user producers to stay parked while the process is parked-for-restart (that is the backpressure). The aging streak (consec_control) is untouched.

Source

pub fn drain(self) -> Drained<C, U>

Consume the consumer and return everything still queued on both lanes, in FIFO order (P6 teardown seam).

§Teardown race

A send’s linearization point is the ring-slot publish inside UserLane::try_push. A UserSender::send parked on a full lane at this moment is RELEASED with its payload: it never linearized, so it resolves Err(UserClosed(item)) and the item comes back to the caller exactly once. A send whose publish completed before it observed teardown resolves Ok(()); its item is owned by the lane — collected here if still queued, or dropped with the lane otherwise. Pinned by tests/teardown_oracle.rs and edge_cases::drain_teardown_race_releases_blocked_sender_with_payload.

Trait Implementations§

Source§

impl<C, U> Drop for Consumer<C, U>

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

impl<C: Send, U: Send> Send for Consumer<C, U>

Auto Trait Implementations§

§

impl<C, U> Freeze for Consumer<C, U>

§

impl<C, U> !RefUnwindSafe for Consumer<C, U>

§

impl<C, U> !Sync for Consumer<C, U>

§

impl<C, U> Unpin for Consumer<C, U>

§

impl<C, U> UnsafeUnpin for Consumer<C, U>

§

impl<C, U> !UnwindSafe for Consumer<C, U>

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<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.