use super::super::admission::{DemandId, DemandSnapshot, ProtocolRequirement, SnapshotVersion};
use super::super::partition::EligibilityGroup;
use super::{AcquisitionOutcome, AcquisitionStep, EstablishmentPermit};
use std::collections::{BTreeSet, HashMap};
use std::num::NonZeroUsize;
use std::task::{Context, Poll, Waker};
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub(in crate::client::pool) struct WaiterId(pub(in crate::client::pool) u64);
#[derive(Debug, Default)]
pub(super) struct AcquisitionQueue {
records: HashMap<WaiterId, WaiterRecord>,
waiting: WaitingQueueState,
h1_compatible_waiters: BTreeSet<WaiterId>,
h2_compatible_waiters: BTreeSet<WaiterId>,
next_waiter_id: u64,
next_demand_id: u64,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum EstablishmentPhase {
Submitted,
Started,
}
#[derive(Debug, Default)]
enum WaitingQueueState {
#[default]
Empty,
Active {
head: WaiterId,
tail: WaiterId,
len: NonZeroUsize,
demand: DemandTicket,
},
}
#[derive(Debug)]
struct WaiterRecord {
requirement: ProtocolRequirement,
state: WaiterState,
}
#[derive(Debug)]
enum WaiterState {
Waiting {
previous: Option<WaiterId>,
next: Option<WaiterId>,
waker: Option<Waker>,
},
DeliveryPending {
waker: Option<Waker>,
pending_result: Option<AcquisitionOutcome>,
},
DeliveryCancelled {
waker: Option<Waker>,
pending_result: Option<AcquisitionOutcome>,
},
ReadyToEstablish {
permit: EstablishmentPermit,
},
Launching {
phase: EstablishmentPhase,
waker: Option<Waker>,
},
Ready(AcquisitionOutcome),
}
#[derive(Debug)]
struct DemandTicket {
id: DemandId,
version: SnapshotVersion,
requirement: ProtocolRequirement,
}
impl DemandTicket {
fn snapshot(&self, eligibility_group: &EligibilityGroup) -> DemandSnapshot {
DemandSnapshot::active(
self.id,
self.version,
self.requirement,
eligibility_group.clone(),
)
}
}
impl AcquisitionQueue {
pub(super) fn register_waiter(
&mut self,
requirement: ProtocolRequirement,
eligibility_group: &EligibilityGroup,
bounded: bool,
) -> (WaiterId, Option<DemandSnapshot>) {
let waiter = self.take_waiter_id();
if !bounded {
self.records.insert(
waiter,
WaiterRecord {
requirement,
state: WaiterState::ReadyToEstablish {
permit: EstablishmentPermit::unbounded(),
},
},
);
self.add_protocol_waiters(waiter, requirement);
self.assert_consistent();
return (waiter, None);
}
let previous = match &self.waiting {
WaitingQueueState::Empty => None,
WaitingQueueState::Active { tail, .. } => Some(*tail),
};
let initial_demand = previous.is_none().then(|| self.new_demand(requirement));
if let Some(previous) = previous {
let previous = self
.records
.get_mut(&previous)
.expect("waiting tail disappeared");
let WaiterState::Waiting { next, .. } = &mut previous.state else {
unreachable!("waiting tail left the waiting state");
};
debug_assert!(next.is_none());
*next = Some(waiter);
}
let replaced = self.records.insert(
waiter,
WaiterRecord {
requirement,
state: WaiterState::Waiting {
previous,
next: None,
waker: None,
},
},
);
debug_assert!(replaced.is_none());
let snapshot = match (&mut self.waiting, initial_demand) {
(waiting @ WaitingQueueState::Empty, Some(demand)) => {
let snapshot = demand.snapshot(eligibility_group);
*waiting = WaitingQueueState::Active {
head: waiter,
tail: waiter,
len: NonZeroUsize::MIN,
demand,
};
Some(snapshot)
}
(WaitingQueueState::Active { tail, len, .. }, None) => {
*tail = waiter;
*len = len.checked_add(1).expect("waiter queue length exhausted");
None
}
_ => unreachable!("waiter queue occupancy changed during registration"),
};
self.assert_consistent();
(waiter, snapshot)
}
fn oldest_h1_compatible_waiter(&self) -> Option<WaiterId> {
let waiting = match &self.waiting {
WaitingQueueState::Active { head, .. }
if self.records[head].requirement.accepts_h1() =>
{
Some(*head)
}
WaitingQueueState::Empty | WaitingQueueState::Active { .. } => None,
};
match (waiting, self.h1_compatible_waiters.first().copied()) {
(Some(waiting), Some(launching)) => Some(waiting.min(launching)),
(Some(waiting), None) => Some(waiting),
(None, launching) => launching,
}
}
pub(super) fn has_h1_compatible_waiter(&self) -> bool {
self.oldest_h1_compatible_waiter().is_some()
}
pub(super) fn has_prior_h1_waiter(&self) -> bool {
!self.h1_compatible_waiters.is_empty()
}
pub(super) fn oldest_h2_compatible_waiter(&self) -> Option<WaiterId> {
let waiting = match &self.waiting {
WaitingQueueState::Active { head, .. }
if self.records[head].requirement.accepts_h2() =>
{
Some(*head)
}
WaitingQueueState::Active { .. } => None,
WaitingQueueState::Empty => None,
};
match (waiting, self.h2_compatible_waiters.first().copied()) {
(Some(waiting), Some(active)) => Some(waiting.min(active)),
(Some(waiting), None) => Some(waiting),
(None, active) => active,
}
}
pub(super) fn supersede_demand_snapshot(&mut self, demand: DemandId) -> bool {
let WaitingQueueState::Active {
demand: current, ..
} = &mut self.waiting
else {
return false;
};
if current.id != demand {
return false;
}
let acknowledged = current.version.next();
current.version = acknowledged.next();
self.assert_consistent();
true
}
pub(super) fn current_demand_snapshot(
&self,
eligibility_group: &EligibilityGroup,
) -> Option<DemandSnapshot> {
match &self.waiting {
WaitingQueueState::Empty => None,
WaitingQueueState::Active { demand, .. } => Some(demand.snapshot(eligibility_group)),
}
}
pub(super) fn is_oldest_h2_compatible_waiter(&self, waiter: WaiterId) -> bool {
self.oldest_h2_compatible_waiter() == Some(waiter)
}
pub(super) fn route_cutoff(&self) -> Option<WaiterId> {
(!self.records.is_empty()).then(|| {
WaiterId(
self.next_waiter_id
.checked_sub(1)
.expect("nonempty waiter set had no allocated identity"),
)
})
}
pub(super) fn has_h2_compatible_waiter_through(&self, cutoff: WaiterId) -> bool {
self.oldest_h2_compatible_waiter()
.is_some_and(|waiter| waiter <= cutoff)
}
pub(super) fn has_h2_compatible_waiter(&self) -> bool {
self.oldest_h2_compatible_waiter().is_some()
}
pub(super) fn is_launching_h2_waiter(&self, waiter: WaiterId) -> bool {
self.h2_compatible_waiters.contains(&waiter)
&& self.records.get(&waiter).is_some_and(|record| {
record.requirement.accepts_h2()
&& matches!(record.state, WaiterState::Launching { .. })
})
}
fn add_protocol_waiters(&mut self, waiter: WaiterId, requirement: ProtocolRequirement) {
if requirement.accepts_h1() {
self.h1_compatible_waiters.insert(waiter);
}
if requirement.accepts_h2() {
self.h2_compatible_waiters.insert(waiter);
}
}
fn remove_protocol_waiters(&mut self, waiter: WaiterId) {
self.h1_compatible_waiters.remove(&waiter);
self.h2_compatible_waiters.remove(&waiter);
}
pub(super) fn offer_returned_h1(
&mut self,
result: impl FnOnce() -> AcquisitionOutcome,
eligibility_group: &EligibilityGroup,
) -> (Option<WaiterId>, WaiterResolution) {
let Some(waiter) = self.oldest_h1_compatible_waiter() else {
return (None, WaiterResolution::refused(result()));
};
(
Some(waiter),
self.assign_protocol_outcome(waiter, result, eligibility_group),
)
}
fn assign_protocol_outcome(
&mut self,
waiter: WaiterId,
result: impl FnOnce() -> AcquisitionOutcome,
eligibility_group: &EligibilityGroup,
) -> WaiterResolution {
if matches!(
self.records.get(&waiter).map(|record| &record.state),
Some(WaiterState::Waiting { .. })
) {
let removed = self.pop_head(eligibility_group);
debug_assert_eq!(removed.waiter, waiter);
let record = self
.records
.get_mut(&waiter)
.expect("selected requesting waiter disappeared");
let WaiterState::Waiting { waker, .. } = &mut record.state else {
unreachable!("selected requesting waiter left the waiting state");
};
let waker = waker.take();
record.state = WaiterState::Ready(result());
let retired =
DemandSnapshot::inactive(removed.demand.id, removed.demand.version.next());
self.assert_consistent();
return WaiterResolution {
demand_updates: [Some(retired), removed.successor],
returned_step: None,
waker,
};
}
let record = self
.records
.get_mut(&waiter)
.expect("selected compatible waiter disappeared");
let (returned_step, waker) = match &mut record.state {
WaiterState::DeliveryPending {
waker: _,
pending_result,
} => {
debug_assert!(pending_result.is_none());
*pending_result = Some(result());
(None, None)
}
WaiterState::ReadyToEstablish { .. } | WaiterState::Launching { .. } => {
let previous = std::mem::replace(&mut record.state, WaiterState::Ready(result()));
match previous {
WaiterState::ReadyToEstablish { permit } => {
(Some(AcquisitionStep::StartEstablishment(permit)), None)
}
WaiterState::Launching { waker, .. } => (None, waker),
_ => {
unreachable!("selected compatible waiter changed state under the cell lock")
}
}
}
WaiterState::Waiting { .. }
| WaiterState::DeliveryCancelled { .. }
| WaiterState::Ready(_) => {
unreachable!("selected waiter was not eligible for a protocol result")
}
};
self.remove_protocol_waiters(waiter);
if returned_step.is_none() {
self.assert_consistent();
}
WaiterResolution {
demand_updates: [None, None],
returned_step,
waker,
}
}
pub(super) fn cancel_waiter(
&mut self,
waiter: WaiterId,
eligibility_group: &EligibilityGroup,
) -> Option<WaiterCancellation> {
let state = &self.records.get(&waiter)?.state;
let cancellation = if matches!(state, WaiterState::Waiting { .. }) {
let is_head = matches!(
self.waiting,
WaitingQueueState::Active { head, .. } if head == waiter
);
let demand_updates = if is_head {
let removed = self.pop_head(eligibility_group);
debug_assert_eq!(waiter, removed.waiter);
let record = self
.records
.remove(&waiter)
.expect("cancelled head waiter disappeared");
debug_assert!(matches!(record.state, WaiterState::Waiting { .. }));
let retired =
DemandSnapshot::inactive(removed.demand.id, removed.demand.version.next());
[Some(retired), removed.successor]
} else {
self.remove_non_head(waiter);
[None, None]
};
WaiterCancellation {
demand_updates,
returned_steps: [None, None],
}
} else if matches!(state, WaiterState::DeliveryPending { .. }) {
self.remove_protocol_waiters(waiter);
let record = self.records.get_mut(&waiter)?;
let WaiterState::DeliveryPending {
waker,
pending_result,
} = &mut record.state
else {
unreachable!("delivery-pending waiter changed state under the cell lock");
};
let waker = waker.take();
let pending_result = pending_result.take();
record.state = WaiterState::DeliveryCancelled {
waker,
pending_result,
};
WaiterCancellation {
demand_updates: [None, None],
returned_steps: [None, None],
}
} else if matches!(
state,
WaiterState::ReadyToEstablish { .. } | WaiterState::Ready(_)
) {
self.assert_consistent();
let record = self.records.remove(&waiter)?;
self.remove_protocol_waiters(waiter);
let event = match record.state {
WaiterState::ReadyToEstablish { permit } => {
AcquisitionStep::StartEstablishment(permit)
}
WaiterState::Ready(result) => AcquisitionStep::Resolved(result),
_ => unreachable!("ready waiter changed state under the cell lock"),
};
return Some(WaiterCancellation {
demand_updates: [None, None],
returned_steps: [Some(event), None],
});
} else if matches!(state, WaiterState::Launching { .. }) {
self.assert_consistent();
self.remove_protocol_waiters(waiter);
self.records.remove(&waiter)?;
return Some(WaiterCancellation {
demand_updates: [None, None],
returned_steps: [None, None],
});
} else {
debug_assert!(matches!(state, WaiterState::DeliveryCancelled { .. }));
return None;
};
self.assert_consistent();
Some(cancellation)
}
pub(super) fn poll_waiter(
&mut self,
waiter: WaiterId,
cx: &mut Context<'_>,
) -> Poll<AcquisitionStep> {
if matches!(
self.records.get(&waiter).map(|record| &record.state),
Some(WaiterState::ReadyToEstablish { .. })
) {
self.assert_consistent();
let record = self
.records
.get_mut(&waiter)
.expect("ready establishment waiter disappeared");
let previous = std::mem::replace(
&mut record.state,
WaiterState::Launching {
phase: EstablishmentPhase::Submitted,
waker: None,
},
);
let WaiterState::ReadyToEstablish { permit } = previous else {
unreachable!("ready establishment waiter changed state under the cell lock");
};
return Poll::Ready(AcquisitionStep::StartEstablishment(permit));
}
if matches!(
self.records.get(&waiter).map(|record| &record.state),
Some(WaiterState::Ready(_))
) {
return Poll::Ready(
self.take_ready_result(waiter)
.map(AcquisitionStep::Resolved)
.expect("ready waiter lost its acquisition outcome"),
);
}
let record = self
.records
.get_mut(&waiter)
.expect("polled a cancelled, consumed, or unknown acquisition waiter");
let waker = match &mut record.state {
WaiterState::Waiting { waker, .. }
| WaiterState::DeliveryPending { waker, .. }
| WaiterState::Launching { waker, .. } => waker,
WaiterState::DeliveryCancelled { .. } => {
panic!("polled a cancelled acquisition waiter")
}
WaiterState::ReadyToEstablish { .. } | WaiterState::Ready(_) => {
unreachable!("ready waiter changed state under the cell lock")
}
};
if waker
.as_ref()
.is_none_or(|waker| !waker.will_wake(cx.waker()))
{
*waker = Some(cx.waker().clone());
}
Poll::Pending
}
pub(super) fn start_establishment(&mut self, waiter: WaiterId) -> bool {
let Some(record) = self.records.get_mut(&waiter) else {
return false;
};
match &mut record.state {
WaiterState::Launching { phase, .. } if *phase == EstablishmentPhase::Submitted => {
*phase = EstablishmentPhase::Started;
self.assert_consistent();
true
}
WaiterState::Ready(_) => false,
WaiterState::Launching {
phase: EstablishmentPhase::Started,
..
} => {
panic!("establishment attempt was started more than once")
}
_ => {
panic!("establishment start did not name a launching waiter")
}
}
}
pub(super) fn reserve_delivery_waiter(
&mut self,
demand: DemandId,
eligibility_group: &EligibilityGroup,
) -> DeliveryReservation {
let current = matches!(
&self.waiting,
WaitingQueueState::Active {
demand: current,
..
} if current.id == demand
);
if !current {
return DeliveryReservation::Rejected;
}
let removed = self.pop_head(eligibility_group);
debug_assert_eq!(removed.demand.id, demand);
let record = self
.records
.get_mut(&removed.waiter)
.expect("reserved queue head disappeared");
let WaiterState::Waiting { waker, .. } = &mut record.state else {
unreachable!("reserved queue head left the waiting state");
};
let waker = waker.take();
let requirement = record.requirement;
record.state = WaiterState::DeliveryPending {
waker,
pending_result: None,
};
self.add_protocol_waiters(removed.waiter, requirement);
self.assert_consistent();
DeliveryReservation::Reserved {
waiter: removed.waiter,
successor: removed.successor,
}
}
pub(super) fn commit_capacity(
&mut self,
waiter: WaiterId,
permit: EstablishmentPermit,
) -> CellCommitOutcome {
let Some(record) = self.records.get_mut(&waiter) else {
return CellCommitOutcome::invalid(
AcquisitionStep::StartEstablishment(permit),
CellCommitError::MissingWaiter,
);
};
match &mut record.state {
WaiterState::DeliveryPending {
waker,
pending_result,
} => {
let waker = waker.take();
if let Some(result) = pending_result.take() {
record.state = WaiterState::Ready(result);
self.remove_protocol_waiters(waiter);
return CellCommitOutcome::Refused {
returned: [Some(AcquisitionStep::StartEstablishment(permit)), None],
waker,
};
}
record.state = WaiterState::ReadyToEstablish { permit };
self.assert_consistent();
CellCommitOutcome::Committed { waker }
}
WaiterState::DeliveryCancelled {
waker,
pending_result,
} => {
let waker = waker.take();
let pending_result = pending_result.take();
self.records.remove(&waiter);
CellCommitOutcome::Refused {
returned: [
pending_result.map(AcquisitionStep::Resolved),
Some(AcquisitionStep::StartEstablishment(permit)),
],
waker,
}
}
WaiterState::Waiting { .. }
| WaiterState::ReadyToEstablish { .. }
| WaiterState::Launching { .. }
| WaiterState::Ready(_) => CellCommitOutcome::invalid(
AcquisitionStep::StartEstablishment(permit),
CellCommitError::UnexpectedState,
),
}
}
pub(super) fn commit_borrowed_h1(
&mut self,
waiter: WaiterId,
result: AcquisitionOutcome,
) -> CellCommitOutcome {
let Some(record) = self.records.get_mut(&waiter) else {
return CellCommitOutcome::invalid(
AcquisitionStep::Resolved(result),
CellCommitError::MissingWaiter,
);
};
match &mut record.state {
WaiterState::DeliveryPending {
waker,
pending_result,
} => {
let waker = waker.take();
if let Some(local_result) = pending_result.take() {
record.state = WaiterState::Ready(local_result);
self.remove_protocol_waiters(waiter);
return CellCommitOutcome::Refused {
returned: [Some(AcquisitionStep::Resolved(result)), None],
waker,
};
}
record.state = WaiterState::Ready(result);
self.remove_protocol_waiters(waiter);
self.assert_consistent();
CellCommitOutcome::Committed { waker }
}
WaiterState::DeliveryCancelled {
waker,
pending_result,
} => {
let waker = waker.take();
let pending_result = pending_result.take();
self.records.remove(&waiter);
CellCommitOutcome::Refused {
returned: [
pending_result.map(AcquisitionStep::Resolved),
Some(AcquisitionStep::Resolved(result)),
],
waker,
}
}
WaiterState::Waiting { .. }
| WaiterState::ReadyToEstablish { .. }
| WaiterState::Launching { .. }
| WaiterState::Ready(_) => CellCommitOutcome::invalid(
AcquisitionStep::Resolved(result),
CellCommitError::UnexpectedState,
),
}
}
pub(super) fn commit_establishment(
&mut self,
waiter: WaiterId,
result: AcquisitionOutcome,
) -> CellCommitOutcome {
let Some(record) = self.records.get_mut(&waiter) else {
return CellCommitOutcome::refused(AcquisitionStep::Resolved(result));
};
if matches!(record.state, WaiterState::Ready(_)) {
return CellCommitOutcome::refused(AcquisitionStep::Resolved(result));
}
if !matches!(record.state, WaiterState::Launching { .. }) {
return CellCommitOutcome::invalid(
AcquisitionStep::Resolved(result),
CellCommitError::UnexpectedState,
);
}
let WaiterState::Launching { waker, .. } = &mut record.state else {
unreachable!("validated launching waiter changed state")
};
let waker = waker.take();
record.state = WaiterState::Ready(result);
self.remove_protocol_waiters(waiter);
self.assert_consistent();
CellCommitOutcome::Committed { waker }
}
pub(super) fn take_ready_result(&mut self, waiter: WaiterId) -> Option<AcquisitionOutcome> {
if !matches!(
self.records.get(&waiter).map(|record| &record.state),
Some(WaiterState::Ready(_))
) {
return None;
}
self.assert_consistent();
let record = self.records.remove(&waiter)?;
let WaiterState::Ready(result) = record.state else {
unreachable!("ready waiter changed state under the cell lock");
};
Some(result)
}
pub(super) fn offer_h2_activation(
&mut self,
cutoff: Option<WaiterId>,
result: impl FnOnce(WaiterId) -> AcquisitionOutcome,
eligibility_group: &EligibilityGroup,
) -> (Option<WaiterId>, WaiterResolution) {
let Some(waiter) = self.oldest_h2_compatible_waiter() else {
return (None, WaiterResolution::empty());
};
if cutoff.is_some_and(|cutoff| waiter > cutoff) {
return (None, WaiterResolution::empty());
}
(
Some(waiter),
self.assign_protocol_outcome(waiter, || result(waiter), eligibility_group),
)
}
fn pop_head(&mut self, eligibility_group: &EligibilityGroup) -> RemovedHead {
let waiting = std::mem::take(&mut self.waiting);
let WaitingQueueState::Active {
head,
tail,
len,
demand,
} = waiting
else {
unreachable!("removed a head from an empty waiter queue");
};
let next = match &self
.records
.get(&head)
.expect("waiting head disappeared")
.state
{
WaiterState::Waiting { previous, next, .. } => {
debug_assert!(previous.is_none());
*next
}
_ => unreachable!("waiting head left the waiting state"),
};
let successor = match next {
Some(next) => {
let requirement = {
let next_record = self
.records
.get_mut(&next)
.expect("next waiter disappeared");
let WaiterState::Waiting { previous, .. } = &mut next_record.state else {
unreachable!("next waiter left the waiting state");
};
debug_assert_eq!(*previous, Some(head));
*previous = None;
next_record.requirement
};
let next_demand = self.new_demand(requirement);
let snapshot = next_demand.snapshot(eligibility_group);
let len = NonZeroUsize::new(
len.get()
.checked_sub(1)
.expect("waiter queue length underflowed"),
)
.expect("nonempty waiter queue lost its length");
self.waiting = WaitingQueueState::Active {
head: next,
tail,
len,
demand: next_demand,
};
Some(snapshot)
}
None => {
debug_assert_eq!(head, tail);
debug_assert_eq!(len, NonZeroUsize::MIN);
None
}
};
RemovedHead {
waiter: head,
demand,
successor,
}
}
fn remove_non_head(&mut self, waiter: WaiterId) -> WaiterRecord {
let record = self
.records
.remove(&waiter)
.expect("removed waiter disappeared");
let (previous, next) = match &record.state {
WaiterState::Waiting { previous, next, .. } => (*previous, *next),
_ => unreachable!("removed waiter left the waiting state"),
};
let previous = previous.expect("non-head waiter had no predecessor");
let previous_record = self
.records
.get_mut(&previous)
.expect("previous waiter disappeared");
let WaiterState::Waiting {
next: previous_next,
..
} = &mut previous_record.state
else {
unreachable!("previous waiter left the waiting state");
};
debug_assert_eq!(*previous_next, Some(waiter));
*previous_next = next;
if let Some(next) = next {
let next_record = self
.records
.get_mut(&next)
.expect("next waiter disappeared");
let WaiterState::Waiting {
previous: next_previous,
..
} = &mut next_record.state
else {
unreachable!("next waiter left the waiting state");
};
debug_assert_eq!(*next_previous, Some(waiter));
*next_previous = Some(previous);
}
let WaitingQueueState::Active {
head, tail, len, ..
} = &mut self.waiting
else {
unreachable!("removed a waiter from an empty queue");
};
debug_assert_ne!(*head, waiter);
if next.is_none() {
debug_assert_eq!(*tail, waiter);
*tail = previous;
}
*len = NonZeroUsize::new(
len.get()
.checked_sub(1)
.expect("waiter queue length underflowed"),
)
.expect("removing a non-head waiter emptied the queue");
record
}
fn take_waiter_id(&mut self) -> WaiterId {
let value = self.next_waiter_id;
self.next_waiter_id = value.checked_add(1).expect("waiter identity exhausted");
WaiterId(value)
}
fn new_demand(&mut self, requirement: ProtocolRequirement) -> DemandTicket {
let value = self.next_demand_id;
self.next_demand_id = value.checked_add(1).expect("demand identity exhausted");
DemandTicket {
id: DemandId::from_u64(value),
version: SnapshotVersion::INITIAL,
requirement,
}
}
pub(super) fn assert_consistent(&self) {
#[cfg(any(debug_assertions, test))]
{
if std::thread::panicking() {
return;
}
self.assert_consistent_debug();
}
}
#[cfg(any(debug_assertions, test))]
fn assert_consistent_debug(&self) {
let waiting_records = self
.records
.values()
.filter(|record| matches!(record.state, WaiterState::Waiting { .. }))
.count();
match &self.waiting {
WaitingQueueState::Empty => {
assert_eq!(0, waiting_records, "empty queue retained waiting records");
}
WaitingQueueState::Active {
head,
tail,
len,
demand,
} => {
let head_record = self.records.get(head).expect("waiting head disappeared");
assert_eq!(
demand.requirement, head_record.requirement,
"aggregate demand did not describe the waiting head"
);
let mut current = Some(*head);
let mut previous = None;
let mut traversed = 0;
while let Some(waiter) = current {
assert!(
traversed < self.records.len(),
"waiter queue contains a cycle"
);
let record = self
.records
.get(&waiter)
.expect("linked waiter disappeared");
let WaiterState::Waiting {
previous: linked_previous,
next,
..
} = &record.state
else {
panic!("linked waiter left the waiting state");
};
assert_eq!(
previous, *linked_previous,
"waiter queue contains inconsistent backward links"
);
traversed += 1;
previous = Some(waiter);
current = *next;
}
assert_eq!(Some(*tail), previous, "waiting tail was not reachable");
assert_eq!(
len.get(),
traversed,
"waiter queue length did not match its links"
);
assert_eq!(
waiting_records, traversed,
"waiting record was not reachable from the queue head"
);
}
}
let expected_h1_waiters = self
.records
.iter()
.filter_map(|(waiter, record)| {
(record.requirement.accepts_h1()
&& matches!(
record.state,
WaiterState::DeliveryPending {
pending_result: None,
..
} | WaiterState::ReadyToEstablish { .. }
| WaiterState::Launching { .. }
))
.then_some(*waiter)
})
.collect::<BTreeSet<_>>();
assert_eq!(
expected_h1_waiters, self.h1_compatible_waiters,
"HTTP/1 compatible-waiter index did not match waiter state"
);
let expected_h2_waiters = self
.records
.iter()
.filter_map(|(waiter, record)| {
(record.requirement.accepts_h2()
&& matches!(
record.state,
WaiterState::DeliveryPending {
pending_result: None,
..
} | WaiterState::ReadyToEstablish { .. }
| WaiterState::Launching { .. }
))
.then_some(*waiter)
})
.collect::<BTreeSet<_>>();
assert_eq!(
expected_h2_waiters, self.h2_compatible_waiters,
"HTTP/2 compatible-waiter index did not match waiter state"
);
}
pub(super) fn pending_count(&self) -> usize {
self.records
.values()
.filter(|record| {
matches!(
record.state,
WaiterState::Waiting { .. }
| WaiterState::DeliveryPending { .. }
| WaiterState::ReadyToEstablish { .. }
| WaiterState::Launching { .. }
)
})
.count()
}
#[cfg(test)]
pub(super) fn probe(&self) -> AcquisitionProbe {
let (waiting, demand) = match &self.waiting {
WaitingQueueState::Empty => (0, None),
WaitingQueueState::Active { len, demand, .. } => (len.get(), Some(demand.id)),
};
AcquisitionProbe {
waiting,
retained: self.records.len(),
demand,
}
}
}
pub(super) enum DeliveryReservation {
Reserved {
waiter: WaiterId,
successor: Option<DemandSnapshot>,
},
Rejected,
}
pub(super) struct WaiterCancellation {
pub(super) demand_updates: [Option<DemandSnapshot>; 2],
pub(super) returned_steps: [Option<AcquisitionStep>; 2],
}
pub(super) struct WaiterResolution {
pub(super) demand_updates: [Option<DemandSnapshot>; 2],
pub(super) returned_step: Option<AcquisitionStep>,
pub(super) waker: Option<Waker>,
}
impl WaiterResolution {
fn refused(result: AcquisitionOutcome) -> Self {
Self {
demand_updates: [None, None],
returned_step: Some(AcquisitionStep::Resolved(result)),
waker: None,
}
}
fn empty() -> Self {
Self {
demand_updates: [None, None],
returned_step: None,
waker: None,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(super) enum CellCommitError {
MissingWaiter,
UnexpectedState,
}
pub(super) enum CellCommitOutcome {
Committed {
waker: Option<Waker>,
},
Refused {
returned: [Option<AcquisitionStep>; 2],
waker: Option<Waker>,
},
Invalid {
returned: [Option<AcquisitionStep>; 2],
waker: Option<Waker>,
error: CellCommitError,
},
}
impl CellCommitOutcome {
fn invalid(step: AcquisitionStep, error: CellCommitError) -> Self {
Self::Invalid {
returned: [Some(step), None],
waker: None,
error,
}
}
fn refused(step: AcquisitionStep) -> Self {
Self::Refused {
returned: [Some(step), None],
waker: None,
}
}
}
struct RemovedHead {
waiter: WaiterId,
demand: DemandTicket,
successor: Option<DemandSnapshot>,
}
#[cfg(test)]
#[derive(Debug)]
pub(super) struct AcquisitionProbe {
pub(super) waiting: usize,
pub(super) retained: usize,
pub(super) demand: Option<DemandId>,
}
#[cfg(all(test, not(smithy_http_client_loom)))]
mod tests {
use super::*;
#[test]
fn h1_required_waiter_is_not_h2_compatible() {
let mut queue = AcquisitionQueue::default();
let (waiter, snapshot) = queue.register_waiter(
ProtocolRequirement::H1Required,
&EligibilityGroup::Pool,
false,
);
assert!(snapshot.is_none());
assert_eq!(Some(waiter), queue.oldest_h1_compatible_waiter());
assert_eq!(None, queue.oldest_h2_compatible_waiter());
}
#[test]
fn head_cancellation_uses_the_successors_protocol_requirement() {
let mut queue = AcquisitionQueue::default();
let (head, initial) = queue.register_waiter(
ProtocolRequirement::H2Required,
&EligibilityGroup::Pool,
true,
);
let (_successor, no_new_demand) = queue.register_waiter(
ProtocolRequirement::H1Compatible,
&EligibilityGroup::Pool,
true,
);
assert!(initial.is_some());
assert!(no_new_demand.is_none());
let cancelled = queue
.cancel_waiter(head, &EligibilityGroup::Pool)
.expect("head waiter was not cancelled");
assert_eq!(
Some(DemandSnapshot::active(
DemandId::from_u64(1),
SnapshotVersion::INITIAL,
ProtocolRequirement::H1Compatible,
EligibilityGroup::Pool,
)),
cancelled.demand_updates[1]
);
}
#[test]
fn retired_demand_cannot_reserve_the_successor() {
let mut queue = AcquisitionQueue::default();
let (head, _initial) = queue.register_waiter(
ProtocolRequirement::H1Compatible,
&EligibilityGroup::Pool,
true,
);
queue.register_waiter(
ProtocolRequirement::H1Compatible,
&EligibilityGroup::Pool,
true,
);
queue
.cancel_waiter(head, &EligibilityGroup::Pool)
.expect("head waiter was not cancelled");
assert!(matches!(
queue.reserve_delivery_waiter(DemandId::from_u64(0), &EligibilityGroup::Pool),
DeliveryReservation::Rejected
));
}
}