use super::cell::{H1IdleProbeDecision, H1ReservationDecision, OriginCell};
use super::origin::OriginKey;
use super::partition::{EligibilityGroup, PartitionId};
use super::registry::AdmissionPolicy;
use super::stats::ConnectionCapacityStats;
use crate::sync::{Arc, Mutex, Weak};
use std::collections::HashMap;
use std::fmt;
use std::num::NonZeroUsize;
mod delivery;
mod demand;
mod h1;
mod h2;
mod order;
use self::demand::{
DemandAssignment, DemandAssignmentId, DemandAssignmentOutcome, DemandSchedule,
PreparedCapacityDelivery,
};
pub(in crate::client::pool) use self::demand::{
DemandId, DemandSnapshot, ProtocolRequirement, SnapshotVersion,
};
use self::h1::{
H1CancellationAction, H1CapacityReclaim, H1IdleProbeAction, H1ReservationAction,
H1SupplierSettlement, H1Supply, H1SupplyOutcome, PreparedH1Match,
};
use self::h2::{H2CapacityReclaim, H2RouteGuard, H2Supply, PreparedH2Reclaim, PreparedH2Route};
use self::order::{IntrusiveLinks, IntrusiveOrder};
pub(in crate::client::pool) use delivery::DeliveryGuard;
pub(in crate::client::pool) use h1::{
H1Candidate, H1MatchId, H1SupplyStatus, PreparedH1IdleProbe, PreparedH1Reservation,
};
pub(in crate::client::pool) use h2::H2SupplyStatus;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::client::pool) struct SupplyRevision<T> {
pub(in crate::client::pool) revision: u64,
pub(in crate::client::pool) status: T,
}
impl<T> SupplyRevision<T> {
pub(in crate::client::pool) fn new(revision: u64, status: T) -> Self {
Self { revision, status }
}
}
pub(super) struct CapacityPermit(u64);
impl fmt::Debug for CapacityPermit {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_tuple("CapacityPermit").field(&self.0).finish()
}
}
pub(in crate::client::pool) struct OriginAdmission {
origin: OriginKey,
can_reclaim_h2_for_h1: bool,
state: Mutex<AdmissionState>,
}
impl OriginAdmission {
pub(in crate::client::pool) fn new(origin: OriginKey, policy: AdmissionPolicy) -> Arc<Self> {
Arc::new(Self {
origin,
can_reclaim_h2_for_h1: policy.can_reclaim_h2_for_h1(),
state: Mutex::new(AdmissionState::new(policy.connection_limit())),
})
}
#[cfg(test)]
pub(super) fn for_test(limit: NonZeroUsize) -> Arc<Self> {
Self::for_test_with_h2_reclaim(limit, true)
}
#[cfg(test)]
pub(super) fn for_test_with_h2_reclaim(
limit: NonZeroUsize,
can_reclaim_h2_for_h1: bool,
) -> Arc<Self> {
Self::new(
OriginKey::from_parts(http_1x::uri::Scheme::HTTPS, "example.com", None)
.expect("test origin is valid"),
AdmissionPolicy::new(limit, can_reclaim_h2_for_h1),
)
}
pub(in crate::client::pool) fn connection_capacity_stats(&self) -> ConnectionCapacityStats {
let state = self.state.lock();
ConnectionCapacityStats::new(
state.capacity.limit,
state.capacity.limit - state.capacity.available,
)
}
fn origin(&self) -> &OriginKey {
&self.origin
}
pub(in crate::client::pool) fn register_cell(
origin: &Arc<Self>,
candidate: Arc<OriginCell>,
) -> Arc<OriginCell> {
let partition = candidate.id().partition();
assert_eq!(
origin.origin(),
candidate.id().origin(),
"cell origin did not match its admission authority"
);
let existing = {
let mut state = origin.state.lock();
match state.cells.get(&partition).and_then(Weak::upgrade) {
Some(existing) => Some(existing),
None => {
state.cells.insert(partition, Weak::from_arc(&candidate));
None
}
}
};
if let Some(existing) = existing {
drop(candidate);
existing
} else {
candidate
}
}
pub(in crate::client::pool) fn submit_demand_snapshot(
admission: &Arc<Self>,
requester: PartitionId,
snapshot: DemandSnapshot,
) {
let action = {
let mut state = admission.state.lock();
state.apply_demand_snapshot(requester, snapshot);
Self::prepare_action(admission, &mut state)
};
Self::run_action_chain(action);
}
fn prepare_action(origin: &Arc<Self>, state: &mut AdmissionState) -> Option<AdmissionAction> {
if let Some(cancellation) = state.h1_supply.prepare_cancellation() {
return Some(AdmissionAction::CancelH1Reservation(
H1CancellationAction::new(origin.clone(), cancellation),
));
}
if let Some(prepared) = state.prepare_capacity_delivery() {
return Some(AdmissionAction::Deliver(DeliveryGuard::capacity(
origin.clone(),
prepared.assignment,
prepared.permit,
)));
}
if let Some(route) = state.prepare_h2_route() {
return Some(AdmissionAction::AttachH2Route(H2RouteGuard::new(
origin.clone(),
route,
)));
}
if let Some(prepared) = state.h1_supply.prepare_match(&state.demand) {
return Some(match prepared {
PreparedH1Match::ProbeIdle(probe) => {
AdmissionAction::ProbeH1Supplier(H1IdleProbeAction::new(origin.clone(), probe))
}
PreparedH1Match::Reserve(reservation) => AdmissionAction::ReserveH1Supplier(
H1ReservationAction::new(origin.clone(), reservation),
),
});
}
if !origin.can_reclaim_h2_for_h1 {
return None;
}
state
.h2_supply
.prepare_reclaim(&state.demand)
.map(|reclaim| {
AdmissionAction::ReclaimCapacity(CapacityReclaim::FromH2(H2CapacityReclaim::new(
origin.clone(),
reclaim,
)))
})
}
fn cell(&self, id: &PartitionId) -> Option<Arc<OriginCell>> {
let cell = {
let state = self.state.lock();
state.cells.get(id).cloned()
};
cell.and_then(|cell| cell.upgrade())
}
pub(in crate::client::pool) fn run_action_chain(mut action: Option<AdmissionAction>) {
while let Some(current) = action {
action = match current {
AdmissionAction::Deliver(delivery) => delivery.deliver(),
AdmissionAction::ProbeH1Supplier(probe) => probe.probe_supplier(),
AdmissionAction::ReserveH1Supplier(reservation) => reservation.reserve_supplier(),
AdmissionAction::CancelH1Reservation(cancellation) => {
cancellation.cancel_reservation()
}
AdmissionAction::SettleH1Supplier(settlement) => settlement.settle_supplier(),
AdmissionAction::AttachH2Route(route) => route.attach_route(),
AdmissionAction::ReclaimCapacity(reclaim) => reclaim.reclaim_capacity(),
};
}
}
#[cfg(test)]
fn assignment_is_current(&self, assignment: &DemandAssignment) -> bool {
self.state.lock().demand.assignment_is_current(assignment)
}
fn return_permit(origin: &Arc<Self>, permit: CapacityPermit) {
let action = {
let mut state = origin.state.lock();
state.return_permit(permit);
Self::prepare_action(origin, &mut state)
};
Self::run_action_chain(action);
}
fn settle_delivery(
admission: &Arc<Self>,
assignment: &DemandAssignment,
permit: Option<CapacityPermit>,
outcome: DemandAssignmentOutcome,
) -> Option<AdmissionAction> {
let mut state = admission.state.lock();
if let Some(permit) = permit {
state.return_permit(permit);
}
state.settle_assignment(assignment, outcome);
Self::prepare_action(admission, &mut state)
}
pub(in crate::client::pool) fn apply_h1_supply_revision(
admission: &Arc<Self>,
supplier: PartitionId,
eligibility_group: EligibilityGroup,
revision: SupplyRevision<H1SupplyStatus>,
) {
h1::apply_supply_revision(admission, supplier, eligibility_group, revision);
}
pub(in crate::client::pool) fn reject_returned_h1_match(
admission: &Arc<Self>,
match_id: H1MatchId,
supplier: PartitionId,
revision: SupplyRevision<H1SupplyStatus>,
) {
h1::reject_returned_match(admission, match_id, supplier, revision);
}
pub(in crate::client::pool) fn settle_h1_idle_probe(
admission: &Arc<Self>,
match_id: H1MatchId,
supplier: PartitionId,
decision: H1IdleProbeDecision<H1Candidate>,
) -> Option<AdmissionAction> {
h1::settle_idle_probe(admission, match_id, supplier, decision)
}
pub(in crate::client::pool) fn settle_h1_reservation(
admission: &Arc<Self>,
match_id: H1MatchId,
supplier: PartitionId,
decision: H1ReservationDecision<H1Candidate>,
) -> Option<AdmissionAction> {
h1::settle_reservation(admission, match_id, supplier, decision)
}
pub(in crate::client::pool) fn resolve_h1_match(
admission: &Arc<Self>,
match_id: H1MatchId,
candidate: H1Candidate,
) -> Option<AdmissionAction> {
h1::resolve_match(admission, match_id, candidate)
}
fn settle_h1_match(
admission: &Arc<Self>,
match_id: H1MatchId,
outcome: H1SupplyOutcome,
) -> Option<AdmissionAction> {
h1::settle_match(admission, match_id, outcome)
}
fn settle_borrow_delivery(
admission: &Arc<Self>,
match_id: H1MatchId,
assignment: &DemandAssignment,
outcome: DemandAssignmentOutcome,
transferred_supplier: Option<PartitionId>,
refused_outcome: Option<H1SupplyOutcome>,
) -> Option<AdmissionAction> {
h1::settle_borrow_delivery(
admission,
match_id,
assignment,
outcome,
transferred_supplier,
refused_outcome,
)
}
pub(in crate::client::pool) fn apply_h2_supply_revision(
admission: &Arc<Self>,
supplier: PartitionId,
eligibility_group: EligibilityGroup,
revision: SupplyRevision<H2SupplyStatus>,
) {
h2::apply_supply_revision(admission, supplier, eligibility_group, revision);
}
fn settle_h2_route(
admission: &Arc<Self>,
prepared: &PreparedH2Route,
stale_generation: Option<super::cell::h2::H2GenerationId>,
outcome: DemandAssignmentOutcome,
) -> Option<AdmissionAction> {
h2::settle_route(admission, prepared, stale_generation, outcome)
}
fn settle_h2_reclaim(
admission: &Arc<Self>,
prepared: &PreparedH2Reclaim,
revision: Option<SupplyRevision<H2SupplyStatus>>,
) -> Option<AdmissionAction> {
h2::settle_reclaim(admission, prepared, revision)
}
#[cfg(test)]
pub(super) fn submit_action_without_running(
admission: &Arc<Self>,
requester: PartitionId,
snapshot: DemandSnapshot,
) -> Option<AdmissionAction> {
let mut state = admission.state.lock();
state.apply_demand_snapshot(requester, snapshot);
Self::prepare_action(admission, &mut state)
}
#[cfg(test)]
pub(super) fn submit_without_running(
admission: &Arc<Self>,
requester: PartitionId,
snapshot: DemandSnapshot,
) -> Option<DeliveryGuard> {
match Self::submit_action_without_running(admission, requester, snapshot) {
Some(AdmissionAction::Deliver(delivery)) => Some(delivery),
Some(_) => panic!("capacity-only test unexpectedly prepared another action"),
None => None,
}
}
#[cfg(test)]
fn probe(&self) -> AdmissionProbe {
let state = self.state.lock();
AdmissionProbe {
limit: state.capacity.limit,
available: state.available_capacity(),
ordered: state.demand.len(),
queued: state.demand.queued_len(),
assigned: state.demand.pending_assignment_count(),
}
}
#[cfg(test)]
pub(super) fn lease_for_test(origin: &Arc<Self>) -> CapacityLease {
let permit = origin
.state
.lock()
.take_permit()
.expect("test origin had no available capacity");
CapacityLease::new(origin.clone(), permit)
}
#[cfg(test)]
pub(super) fn available_capacity_for_test(&self) -> usize {
self.state.lock().available_capacity()
}
#[cfg(test)]
pub(super) fn ordered_demand_count_for_test(&self) -> usize {
self.state.lock().demand.len()
}
#[cfg(all(test, smithy_http_client_loom))]
pub(super) fn clear_modeled_cells_for_test(&self) {
self.state.lock().cells.clear();
}
}
pub(super) enum AdmissionAction {
Deliver(DeliveryGuard),
ProbeH1Supplier(H1IdleProbeAction),
ReserveH1Supplier(H1ReservationAction),
CancelH1Reservation(H1CancellationAction),
SettleH1Supplier(H1SupplierSettlement),
AttachH2Route(H2RouteGuard),
ReclaimCapacity(CapacityReclaim),
}
impl AdmissionAction {
#[cfg(all(test, smithy_http_client_loom))]
pub(super) fn run_once_for_test(self) -> Option<Self> {
match self {
Self::Deliver(delivery) => delivery.deliver(),
Self::ProbeH1Supplier(probe) => probe.probe_supplier(),
Self::ReserveH1Supplier(reservation) => reservation.reserve_supplier(),
Self::CancelH1Reservation(cancellation) => cancellation.cancel_reservation(),
Self::SettleH1Supplier(settlement) => settlement.settle_supplier(),
Self::AttachH2Route(route) => route.attach_route(),
Self::ReclaimCapacity(reclaim) => reclaim.reclaim_capacity(),
}
}
}
pub(super) enum CapacityReclaim {
FromH1(H1CapacityReclaim),
FromH2(H2CapacityReclaim),
}
impl CapacityReclaim {
fn reclaim_capacity(self) -> Option<AdmissionAction> {
match self {
Self::FromH1(reclaim) => reclaim.reclaim_capacity(),
Self::FromH2(reclaim) => reclaim.reclaim_capacity(),
}
}
}
impl fmt::Debug for OriginAdmission {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("OriginAdmission")
.field("state", &self.state)
.finish()
}
}
pub(in crate::client::pool) struct CapacityLease {
admission: Arc<OriginAdmission>,
permit: Option<CapacityPermit>,
}
impl CapacityLease {
fn new(admission: Arc<OriginAdmission>, permit: CapacityPermit) -> Self {
Self {
admission,
permit: Some(permit),
}
}
}
impl fmt::Debug for CapacityLease {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("CapacityLease")
.field("permit", &self.permit)
.finish_non_exhaustive()
}
}
impl Drop for CapacityLease {
fn drop(&mut self) {
if let Some(permit) = self.permit.take() {
OriginAdmission::return_permit(&self.admission, permit);
}
}
}
#[derive(Debug)]
struct CapacityBudget {
limit: usize,
available: usize,
next_permit_id: u64,
}
impl CapacityBudget {
fn new(limit: NonZeroUsize) -> Self {
let limit = limit.get();
Self {
limit,
available: limit,
next_permit_id: 0,
}
}
fn take_permit(&mut self) -> Option<CapacityPermit> {
if self.available == 0 {
return None;
}
let id = self.next_permit_id;
self.next_permit_id = id.checked_add(1).expect("permit identity exhausted");
self.available -= 1;
Some(CapacityPermit(id))
}
fn return_permit(&mut self, _permit: CapacityPermit) {
let available = self
.available
.checked_add(1)
.expect("available capacity count overflowed");
assert!(
available <= self.limit,
"available capacity exceeded the configured limit"
);
self.available = available;
}
}
#[derive(Debug)]
struct AdmissionState {
cells: HashMap<PartitionId, Weak<OriginCell>>,
capacity: CapacityBudget,
demand: DemandSchedule,
h1_supply: H1Supply,
h2_supply: H2Supply,
next_assignment_id: u64,
}
impl AdmissionState {
fn new(limit: NonZeroUsize) -> Self {
Self {
cells: HashMap::new(),
capacity: CapacityBudget::new(limit),
demand: DemandSchedule::default(),
h1_supply: H1Supply::default(),
h2_supply: H2Supply::default(),
next_assignment_id: 0,
}
}
fn take_permit(&mut self) -> Option<CapacityPermit> {
self.capacity.take_permit()
}
fn return_permit(&mut self, permit: CapacityPermit) {
self.capacity.return_permit(permit);
}
#[cfg(test)]
fn available_capacity(&self) -> usize {
self.capacity.available
}
fn apply_demand_snapshot(&mut self, requester: PartitionId, snapshot: DemandSnapshot) {
let old_group = self.demand.group_for(&requester);
self.demand.apply_snapshot(requester, snapshot);
self.reconcile_demand_indexes(&requester, old_group);
}
fn reconcile_demand_indexes(
&mut self,
requesting_partition: &PartitionId,
old_group: Option<EligibilityGroup>,
) {
self.h1_supply
.reconcile_requester(requesting_partition, &self.demand);
if let Some(old_group) = old_group {
self.h2_supply.reconcile_group(&old_group, &self.demand);
}
if let Some(group) = self.demand.group_for(requesting_partition) {
self.h2_supply.reconcile_group(&group, &self.demand);
}
}
fn prepare_capacity_delivery(&mut self) -> Option<PreparedCapacityDelivery> {
if !self.demand.head_is_queued() {
return None;
}
let permit = self.take_permit()?;
let assignment_id = self.take_assignment_id();
let old_group = self
.demand
.queued_head()
.and_then(|head| self.demand.group_for(&head.requester));
let assignment = self
.demand
.prepare_origin_assignment(assignment_id)
.expect("queued demand head disappeared");
self.reconcile_demand_indexes(&assignment.requester, old_group);
Some(PreparedCapacityDelivery { permit, assignment })
}
fn settle_assignment(
&mut self,
assignment: &DemandAssignment,
outcome: DemandAssignmentOutcome,
) {
let old_group = self.demand.group_for(&assignment.requester);
self.demand.settle_assignment(assignment, outcome);
self.reconcile_demand_indexes(&assignment.requester, old_group);
}
fn prepare_h2_route(&mut self) -> Option<PreparedH2Route> {
if !self.h2_supply.has_route_ready_group() {
return None;
}
let assignment_id = self.take_assignment_id();
self.h2_supply
.prepare_route(&mut self.demand, assignment_id)
}
fn take_assignment_id(&mut self) -> DemandAssignmentId {
let value = self.next_assignment_id;
self.next_assignment_id = value
.checked_add(1)
.expect("demand assignment identity exhausted");
DemandAssignmentId(value)
}
}
#[cfg(test)]
#[derive(Debug, Eq, PartialEq)]
struct AdmissionProbe {
limit: usize,
available: usize,
ordered: usize,
queued: usize,
assigned: usize,
}
#[cfg(all(test, not(smithy_http_client_loom)))]
mod tests {
use super::*;
use crate::client::pool::partition::PartitionId;
fn cell(origin: &Arc<OriginAdmission>, partition: usize) -> Arc<OriginCell> {
let cell = Arc::new(OriginCell::new(
PartitionId::from_index(partition),
OriginKey::from_parts(http_1x::uri::Scheme::HTTPS, "example.com", None).unwrap(),
EligibilityGroup::Pool,
Some(origin.clone()),
None,
));
OriginAdmission::register_cell(origin, cell)
}
fn demand(id: u64) -> DemandSnapshot {
DemandSnapshot::active(
DemandId::from_u64(id),
SnapshotVersion::INITIAL,
ProtocolRequirement::H1Compatible,
EligibilityGroup::Pool,
)
}
#[test]
#[should_panic(expected = "cell origin did not match its admission authority")]
fn registration_rejects_a_cell_from_another_origin() {
let admission = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let candidate = Arc::new(OriginCell::new(
PartitionId::from_index(1),
OriginKey::from_parts(http_1x::uri::Scheme::HTTPS, "other.example.com", None).unwrap(),
EligibilityGroup::Pool,
Some(admission.clone()),
None,
));
OriginAdmission::register_cell(&admission, candidate);
}
#[test]
fn permits_are_linear_and_never_reused() {
let mut state = AdmissionState::new(NonZeroUsize::new(2).unwrap());
assert_eq!(2, state.available_capacity());
let first = state.take_permit().unwrap();
let first_id = first.0;
let second = state.take_permit().unwrap();
assert_eq!(0, state.available_capacity());
assert!(state.take_permit().is_none());
state.return_permit(first);
let third = state.take_permit().unwrap();
assert_ne!(first_id, third.0);
state.return_permit(second);
state.return_permit(third);
assert_eq!(2, state.available_capacity());
}
#[test]
fn equal_or_older_snapshots_do_not_replace_current_demand() {
let mut state = AdmissionState::new(NonZeroUsize::new(1).unwrap());
let requesting_partition = PartitionId::from_index(1);
let current = demand(2);
state.apply_demand_snapshot(requesting_partition, current.clone());
state.apply_demand_snapshot(
requesting_partition,
DemandSnapshot::inactive(DemandId::from_u64(2), SnapshotVersion::INITIAL),
);
state.apply_demand_snapshot(requesting_partition, demand(1));
assert_eq!(1, state.demand.len());
assert_eq!(
Some(¤t),
state.demand.latest_for_test(&requesting_partition)
);
}
#[test]
fn cancellation_churn_does_not_retain_order_entries() {
let mut state = AdmissionState::new(NonZeroUsize::new(1).unwrap());
let held = state.take_permit().unwrap();
let requesting_partition = PartitionId::from_index(1);
for id in 0..2_000 {
let id = DemandId::from_u64(id);
state.apply_demand_snapshot(
requesting_partition,
DemandSnapshot::active(
id,
SnapshotVersion::INITIAL,
ProtocolRequirement::H1Compatible,
EligibilityGroup::Pool,
),
);
state.apply_demand_snapshot(
requesting_partition,
DemandSnapshot::inactive(id, SnapshotVersion::INITIAL.next()),
);
}
assert_eq!(0, state.demand.len());
state.return_permit(held);
}
#[test]
fn removing_middle_and_tail_demands_repairs_order() {
let mut state = AdmissionState::new(NonZeroUsize::new(1).unwrap());
let held = state.take_permit().unwrap();
let targets: Vec<_> = (1..=5).map(PartitionId::from_index).collect();
for (index, requesting_partition) in targets[..4].iter().enumerate() {
state.apply_demand_snapshot(*requesting_partition, demand(index as u64 + 1));
}
state.apply_demand_snapshot(
targets[1],
DemandSnapshot::inactive(DemandId::from_u64(2), SnapshotVersion::INITIAL.next()),
);
state.apply_demand_snapshot(
targets[3],
DemandSnapshot::inactive(DemandId::from_u64(4), SnapshotVersion::INITIAL.next()),
);
state.apply_demand_snapshot(targets[4], demand(5));
assert_eq!(3, state.demand.len());
state.return_permit(held);
for expected in [&targets[0], &targets[2], &targets[4]] {
let pending = state.prepare_capacity_delivery().unwrap();
assert_eq!(expected, &pending.assignment.requester);
state.settle_assignment(
&pending.assignment,
DemandAssignmentOutcome::Accepted { successor: None },
);
state.return_permit(pending.permit);
}
assert_eq!(0, state.demand.len());
}
#[test]
fn new_demand_moves_queued_cell_to_the_tail() {
let mut state = AdmissionState::new(NonZeroUsize::new(1).unwrap());
let held = state.take_permit().unwrap();
let first = PartitionId::from_index(1);
let second = PartitionId::from_index(2);
state.apply_demand_snapshot(first, demand(1));
state.apply_demand_snapshot(second, demand(2));
state.apply_demand_snapshot(first, demand(3));
state.return_permit(held);
assert_eq!(
second,
state
.prepare_capacity_delivery()
.unwrap()
.assignment
.requester
);
}
#[test]
fn losing_registered_cell_returns_capacity_after_admission_unlocks() {
let origin = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let retained = cell(&origin, 1);
let candidate = Arc::new(OriginCell::new(
retained.id().partition(),
retained.id().origin().clone(),
EligibilityGroup::Pool,
Some(origin.clone()),
None,
));
assert_eq!(retained.id(), candidate.id());
let (_waiter, snapshot) =
candidate.register_waiter_without_publish(ProtocolRequirement::H1Compatible);
let mut delivery =
OriginAdmission::submit_without_running(&origin, candidate.id().partition(), snapshot)
.expect("candidate demand did not reserve capacity");
assert!(delivery.resolve_payload_for_test());
assert!(OriginCell::receive_delivery(&candidate, delivery).is_none());
assert_eq!(0, origin.available_capacity_for_test());
let winner = OriginAdmission::register_cell(&origin, candidate);
assert!(Arc::ptr_eq(&retained, &winner));
assert_eq!(
1,
origin.available_capacity_for_test(),
"losing candidate did not return capacity after registration"
);
}
#[test]
fn separate_origins_conserve_capacity_independently() {
let first = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let second = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let first_cell = cell(&first, 1);
let second_cell = cell(&second, 1);
let first_delivery =
OriginAdmission::submit_without_running(&first, first_cell.id().partition(), demand(1))
.unwrap();
let second_delivery = OriginAdmission::submit_without_running(
&second,
second_cell.id().partition(),
demand(1),
)
.unwrap();
assert_eq!(0, first.available_capacity_for_test());
assert_eq!(0, second.available_capacity_for_test());
drop(first_delivery);
assert_eq!(1, first.available_capacity_for_test());
assert_eq!(0, second.available_capacity_for_test());
drop(second_delivery);
assert_eq!(1, second.available_capacity_for_test());
}
}
#[cfg(all(test, smithy_http_client_loom))]
mod loom_tests {
use super::*;
use crate::client::pool::partition::PartitionId;
fn id() -> PartitionId {
PartitionId::from_index(1)
}
fn demand() -> DemandSnapshot {
DemandSnapshot::active(
DemandId::from_u64(1),
SnapshotVersion::INITIAL,
ProtocolRequirement::H1Compatible,
EligibilityGroup::Pool,
)
}
#[test]
fn release_and_demand_submission_conserve_one_permit() {
loom::model(|| {
let origin = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let permit = origin.state.lock().take_permit().unwrap();
let lease = CapacityLease::new(origin.clone(), permit);
let release = loom::thread::spawn(move || drop(lease));
let publish_origin = origin.clone();
let publish = loom::thread::spawn(move || {
let delivery =
OriginAdmission::submit_without_running(&publish_origin, id(), demand());
drop(delivery);
});
release.join().unwrap();
publish.join().unwrap();
let probe = origin.probe();
assert_eq!(1, probe.limit);
assert_eq!(1, probe.available);
assert_eq!(0, probe.assigned);
assert!(
matches!((probe.ordered, probe.queued), (0, 0) | (1, 1)),
"demand publication left duplicate or unschedulable demand: {probe:?}"
);
});
}
#[test]
fn cancellation_preserves_an_outstanding_demand_assignment() {
loom::model(|| {
let origin = OriginAdmission::for_test(NonZeroUsize::new(2).unwrap());
let requesting_partition = id();
let delivery =
OriginAdmission::submit_without_running(&origin, requesting_partition, demand())
.unwrap();
let dropping = loom::thread::spawn(move || drop(delivery));
let replacement_origin = origin.clone();
let replacing = loom::thread::spawn(move || {
replacement_origin.state.lock().apply_demand_snapshot(
requesting_partition,
DemandSnapshot::inactive(
DemandId::from_u64(1),
SnapshotVersion::INITIAL.next(),
),
);
OriginAdmission::submit_without_running(
&replacement_origin,
requesting_partition,
DemandSnapshot::active(
DemandId::from_u64(2),
SnapshotVersion::INITIAL,
ProtocolRequirement::H1Compatible,
EligibilityGroup::Pool,
),
)
});
dropping.join().unwrap();
if let Some(replacement) = replacing.join().unwrap() {
assert!(
replacement.is_current(),
"replacement delivery did not own the current demand assignment"
);
drop(replacement);
}
let probe = origin.probe();
assert_eq!(2, probe.available);
assert_eq!(0, probe.assigned);
assert_eq!(0, probe.ordered);
});
}
}