use super::h1::{H1Candidate, H1MatchId, H1SupplyOutcome};
use super::{
AdmissionAction, CapacityLease, CapacityPermit, DemandAssignment, DemandAssignmentOutcome,
DemandId, DemandSnapshot, OriginAdmission,
};
use crate::client::pool::cell::{
AcquisitionOutcome, AcquisitionStep, EstablishmentPermit, OriginCell,
};
use crate::client::pool::partition::PartitionId;
use crate::sync::Arc;
use aws_smithy_runtime_api::client::connection::ConnectionId;
use std::fmt;
enum DeliveryPayload {
Capacity(CapacityPermit),
BorrowedH1 {
match_id: H1MatchId,
supplier: PartitionId,
candidate: H1Candidate,
},
}
pub(in crate::client::pool) struct DeliveryGuard {
admission: Arc<OriginAdmission>,
assignment: DemandAssignment,
state: DeliveryState,
}
enum DeliveryState {
Pending(DeliveryPayload),
Ready {
step: AcquisitionStep,
settlement: DeliverySettlementKind,
},
Disarmed,
}
impl DeliveryGuard {
pub(super) fn capacity(
admission: Arc<OriginAdmission>,
assignment: DemandAssignment,
permit: CapacityPermit,
) -> Self {
Self {
admission,
assignment,
state: DeliveryState::Pending(DeliveryPayload::Capacity(permit)),
}
}
pub(super) fn borrowed_h1(
admission: Arc<OriginAdmission>,
assignment: DemandAssignment,
match_id: H1MatchId,
supplier: PartitionId,
candidate: H1Candidate,
) -> Self {
Self {
admission,
assignment,
state: DeliveryState::Pending(DeliveryPayload::BorrowedH1 {
match_id,
supplier,
candidate,
}),
}
}
pub(in crate::client::pool) fn demand(&self) -> DemandId {
self.assignment.demand
}
#[cfg(test)]
pub(in crate::client::pool) fn is_current(&self) -> bool {
self.admission.assignment_is_current(&self.assignment)
}
pub(super) fn deliver(mut self) -> Option<AdmissionAction> {
if !self.resolve_payload() {
return None;
}
match self.admission.cell(&self.assignment.requester) {
Some(requesting_cell) => OriginCell::receive_delivery(&requesting_cell, self),
None => self.refuse(None),
}
}
fn resolve_payload(&mut self) -> bool {
let state = std::mem::replace(&mut self.state, DeliveryState::Disarmed);
let DeliveryState::Pending(payload) = state else {
unreachable!("delivery payload resolved more than once");
};
let (step, settlement) = match payload {
DeliveryPayload::Capacity(permit) => (
AcquisitionStep::StartEstablishment(EstablishmentPermit::bounded(
CapacityLease::new(self.admission.clone(), permit),
)),
DeliverySettlementKind::Capacity,
),
DeliveryPayload::BorrowedH1 {
match_id,
supplier,
candidate,
} => match candidate.commit() {
Ok(selection) => {
let connection_id = selection.connection_id();
(
AcquisitionStep::Resolved(AcquisitionOutcome::H1(selection)),
DeliverySettlementKind::BorrowedH1 {
connection_id,
match_id,
supplier,
},
)
}
Err(candidate) => {
let outcome = candidate.reject();
let next = OriginAdmission::settle_borrow_delivery(
&self.admission,
match_id,
&self.assignment,
DemandAssignmentOutcome::RetrySamePosition,
None,
Some(outcome),
);
OriginAdmission::run_action_chain(next);
return false;
}
},
};
self.state = DeliveryState::Ready { step, settlement };
true
}
#[cfg(test)]
pub(in crate::client::pool) fn resolve_payload_for_test(&mut self) -> bool {
self.resolve_payload()
}
pub(in crate::client::pool) fn into_step(
mut self,
successor: Option<DemandSnapshot>,
) -> (AcquisitionStep, DeliverySettlement) {
let state = std::mem::replace(&mut self.state, DeliveryState::Disarmed);
let DeliveryState::Ready { step, settlement } = state else {
unreachable!("delivery committed before payload resolution");
};
(
step,
DeliverySettlement {
admission: self.admission.clone(),
assignment: self.assignment.clone(),
successor,
kind: Some(settlement),
},
)
}
pub(in crate::client::pool) fn refuse(
mut self,
successor: Option<DemandSnapshot>,
) -> Option<AdmissionAction> {
let state = std::mem::replace(&mut self.state, DeliveryState::Disarmed);
self.settle_state(state, DemandAssignmentOutcome::Refused { successor })
}
fn settle_state(
&self,
state: DeliveryState,
outcome: DemandAssignmentOutcome,
) -> Option<AdmissionAction> {
match state {
DeliveryState::Pending(DeliveryPayload::Capacity(permit)) => {
OriginAdmission::settle_delivery(
&self.admission,
&self.assignment,
Some(permit),
outcome,
)
}
DeliveryState::Pending(DeliveryPayload::BorrowedH1 {
match_id,
candidate,
..
}) => {
let supply = candidate.reject();
OriginAdmission::settle_borrow_delivery(
&self.admission,
match_id,
&self.assignment,
outcome,
None,
Some(supply),
)
}
DeliveryState::Ready { step, settlement } => match settlement {
DeliverySettlementKind::Capacity => {
drop(step);
OriginAdmission::settle_delivery(
&self.admission,
&self.assignment,
None,
outcome,
)
}
DeliverySettlementKind::BorrowedH1 {
match_id, supplier, ..
} => {
let supplier_cell = self.admission.cell(&supplier);
drop(step);
let supply = match supplier_cell {
Some(supplier_cell) => H1SupplyOutcome::supplier_live(
supplier,
supplier_cell.cancel_h1_reservation(match_id),
),
None => H1SupplyOutcome::supplier_expired(supplier),
};
OriginAdmission::settle_borrow_delivery(
&self.admission,
match_id,
&self.assignment,
outcome,
None,
Some(supply),
)
}
},
DeliveryState::Disarmed => None,
}
}
}
impl fmt::Debug for DeliveryGuard {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("DeliveryGuard")
.field("assignment", &self.assignment)
.finish_non_exhaustive()
}
}
impl Drop for DeliveryGuard {
fn drop(&mut self) {
let state = std::mem::replace(&mut self.state, DeliveryState::Disarmed);
if matches!(state, DeliveryState::Disarmed) {
return;
}
let next = self.settle_state(state, DemandAssignmentOutcome::RetrySamePosition);
OriginAdmission::run_action_chain(next);
}
}
pub(in crate::client::pool) struct DeliverySettlement {
admission: Arc<OriginAdmission>,
assignment: DemandAssignment,
successor: Option<DemandSnapshot>,
kind: Option<DeliverySettlementKind>,
}
enum DeliverySettlementKind {
Capacity,
BorrowedH1 {
connection_id: ConnectionId,
match_id: H1MatchId,
supplier: PartitionId,
},
}
impl DeliverySettlement {
pub(in crate::client::pool) fn suppress_h2_successor(&mut self) {
if self
.successor
.as_ref()
.is_some_and(DemandSnapshot::accepts_h2)
{
self.successor = None;
}
}
pub(in crate::client::pool) fn accept(mut self) -> Option<AdmissionAction> {
let kind = self
.kind
.take()
.expect("delivery settlement completed more than once");
let successor = self.successor.take();
self.settle(kind, DemandAssignmentOutcome::Accepted { successor }, None)
}
pub(in crate::client::pool) fn refuse(
mut self,
returned_steps: [Option<AcquisitionStep>; 2],
) -> Option<AdmissionAction> {
let kind = self
.kind
.take()
.expect("delivery settlement completed more than once");
let successor = self.successor.take();
self.settle(
kind,
DemandAssignmentOutcome::Refused { successor },
Some(returned_steps),
)
}
fn settle(
&self,
kind: DeliverySettlementKind,
outcome: DemandAssignmentOutcome,
returned_steps: Option<[Option<AcquisitionStep>; 2]>,
) -> Option<AdmissionAction> {
match kind {
DeliverySettlementKind::Capacity => {
drop(returned_steps);
OriginAdmission::settle_delivery(&self.admission, &self.assignment, None, outcome)
}
DeliverySettlementKind::BorrowedH1 {
connection_id,
match_id,
supplier,
} => {
let refused = returned_steps.is_some();
let supplier_cell = refused.then(|| self.admission.cell(&supplier)).flatten();
drop(returned_steps);
let refused_outcome = refused.then(|| match supplier_cell {
Some(supplier_cell) => H1SupplyOutcome::supplier_live(
supplier,
supplier_cell.cancel_h1_reservation(match_id),
),
None => H1SupplyOutcome::supplier_expired(supplier),
});
let transferred_supplier = (!refused).then_some(supplier);
let action = OriginAdmission::settle_borrow_delivery(
&self.admission,
match_id,
&self.assignment,
outcome,
transferred_supplier,
refused_outcome,
);
if !refused {
tracing::trace!(
connection_id = %connection_id,
request_partition = ?self.assignment.requester,
connection_partition = ?supplier,
origin_scheme = %self.admission.origin().scheme(),
origin_host = self.admission.origin().host(),
origin_port = ?self.admission.origin().port(),
"HTTP/1 connection borrowed for peer demand"
);
}
action
}
}
}
}
impl fmt::Debug for DeliverySettlement {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("DeliverySettlement")
.field("assignment", &self.assignment)
.finish_non_exhaustive()
}
}
impl Drop for DeliverySettlement {
fn drop(&mut self) {
let Some(kind) = self.kind.take() else {
return;
};
let successor = self.successor.take();
let next = self.settle(kind, DemandAssignmentOutcome::Accepted { successor }, None);
OriginAdmission::run_action_chain(next);
}
}
#[cfg(all(test, not(smithy_http_client_loom)))]
mod tests {
use super::super::{ProtocolRequirement, SnapshotVersion};
use super::*;
use crate::client::pool::origin::OriginKey;
use crate::client::pool::partition::EligibilityGroup;
use std::num::NonZeroUsize;
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]
fn assignment_currency_includes_the_demand_id() {
let origin = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let requesting_partition = cell(&origin, 1);
let delivery = OriginAdmission::submit_without_running(
&origin,
requesting_partition.id().partition(),
demand(1),
)
.unwrap();
assert!(delivery.is_current());
origin
.state
.lock()
.apply_demand_snapshot(requesting_partition.id().partition(), demand(2));
assert!(!delivery.is_current());
delivery.refuse(None);
}
#[test]
fn stale_successor_cannot_leave_active_demand_idle() {
let origin = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let requesting_partition = cell(&origin, 1);
let delivery = OriginAdmission::submit_without_running(
&origin,
requesting_partition.id().partition(),
demand(1),
)
.unwrap();
origin.state.lock().apply_demand_snapshot(
requesting_partition.id().partition(),
DemandSnapshot::active(
DemandId::from_u64(1),
SnapshotVersion::INITIAL.next(),
ProtocolRequirement::H1Compatible,
EligibilityGroup::Pool,
),
);
delivery.refuse(Some(demand(1)));
assert_eq!(1, origin.probe().available);
assert_eq!(0, origin.probe().ordered);
}
#[test]
fn dropped_delivery_refunnels_capacity_and_preserves_order() {
let origin = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let first = cell(&origin, 1);
let second = cell(&origin, 2);
let (first_waiter, first_demand) =
first.register_waiter_without_publish(ProtocolRequirement::H1Compatible);
let (second_waiter, second_demand) =
second.register_waiter_without_publish(ProtocolRequirement::H1Compatible);
let delivery =
OriginAdmission::submit_without_running(&origin, first.id().partition(), first_demand)
.unwrap();
{
let mut state = origin.state.lock();
state.apply_demand_snapshot(second.id().partition(), second_demand);
}
drop(delivery);
let first_lease = OriginCell::take_ready_lease(&first, first_waiter)
.expect("dropped delivery did not retry the original head");
assert!(OriginCell::take_ready_lease(&second, second_waiter).is_none());
assert_eq!(1, origin.probe().ordered);
drop(first_lease);
let second_lease = OriginCell::take_ready_lease(&second, second_waiter)
.expect("younger demand did not run after the original head");
drop(second_lease);
}
#[test]
fn expired_requesting_cell_refunnels_capacity() {
let origin = OriginAdmission::for_test(NonZeroUsize::new(1).unwrap());
let requesting_partition = cell(&origin, 1);
let requesting_cell_id = requesting_partition.id().partition();
let (_waiter, snapshot) =
requesting_partition.register_waiter_without_publish(ProtocolRequirement::H1Compatible);
let delivery =
OriginAdmission::submit_without_running(&origin, requesting_cell_id, snapshot).unwrap();
drop(requesting_partition);
assert!(origin.cell(&requesting_cell_id).is_none());
OriginAdmission::run_action_chain(Some(AdmissionAction::Deliver(delivery)));
let probe = origin.probe();
assert_eq!(1, probe.available);
assert_eq!(0, probe.assigned);
assert_eq!(0, probe.ordered);
}
}