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>
impl<C, U> Consumer<C, U>
Sourcepub async fn recv(&mut self) -> Option<Received<C, U>>
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).
Sourcepub async fn recv_control(&mut self) -> Option<C>
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.
Sourcepub fn drain(self) -> Drained<C, U>
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.