use super::super::super::connection::{ConnectionInfo, ConnectionState, DispatchGuard};
use super::super::super::partition::PartitionId;
use super::super::{OriginCell, WaiterId};
use super::{H2CloseHandle, H2GenerationId, H2RouteId, H2Sender};
use crate::sync::{Arc, Mutex, Weak};
pub(super) struct H2ActivationResources {
pub(super) sender: H2Sender,
pub(super) reused: bool,
pub(super) connection: Arc<ConnectionState>,
}
enum H2ActivationState {
Ready(H2ActivationResources),
Dispatched(H2RequestIdentity),
}
pub(in crate::client::pool) struct H2DispatchParts {
pub(in crate::client::pool) sender: H2Sender,
pub(in crate::client::pool) upload: H2UploadGuard,
pub(in crate::client::pool) response: H2ResponseGuard,
}
pub(in crate::client::pool) struct H2Activation {
connection_cell: Arc<OriginCell>,
generation: H2GenerationId,
request_partition: PartitionId,
state: Option<H2ActivationState>,
claim: Arc<H2RequestClaim>,
turn: Option<H2ActivationTurnGuard>,
active: bool,
}
impl std::fmt::Debug for H2Activation {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("H2Activation")
.field("generation", &self.generation)
.field(
"connection_id",
&self.state.as_ref().map(|state| match state {
H2ActivationState::Ready(resources) => resources.connection.id(),
H2ActivationState::Dispatched(identity) => identity.connection.id(),
}),
)
.field("active", &self.active)
.finish()
}
}
impl H2Activation {
pub(super) fn new(
connection_cell: Arc<OriginCell>,
generation: H2GenerationId,
resources: H2ActivationResources,
request_partition: PartitionId,
turn: Option<H2ActivationTurnGuard>,
) -> Self {
let claim = Arc::new(H2RequestClaim {
connection_cell: Weak::from_arc(&connection_cell),
generation,
state: Mutex::new(H2RequestClaimState::Prospective {
upload_finished: false,
response_finished: false,
}),
});
Self {
connection_cell,
generation,
request_partition,
state: Some(H2ActivationState::Ready(resources)),
claim,
turn,
active: true,
}
}
#[cfg(test)]
pub(in crate::client::pool) fn generation(&self) -> H2GenerationId {
self.generation
}
pub(in crate::client::pool) fn is_reused(&self) -> bool {
match self
.state
.as_ref()
.expect("HTTP/2 activation state missing")
{
H2ActivationState::Ready(resources) => resources.reused,
H2ActivationState::Dispatched(_) => {
panic!("HTTP/2 activation resources already taken")
}
}
}
pub(in crate::client::pool) fn connection(&self) -> &Arc<ConnectionState> {
match self
.state
.as_ref()
.expect("HTTP/2 activation state missing")
{
H2ActivationState::Ready(resources) => &resources.connection,
H2ActivationState::Dispatched(_) => {
panic!("HTTP/2 activation resources already taken")
}
}
}
pub(in crate::client::pool) fn take_dispatch_parts(&mut self) -> H2DispatchParts {
let state = self.state.take().expect("HTTP/2 activation state missing");
let H2ActivationState::Ready(resources) = state else {
self.state = Some(state);
panic!("HTTP/2 activation dispatch parts already taken");
};
let identity = H2RequestIdentity {
request_partition: self.request_partition,
connection: resources.connection.info().clone(),
generation: self.generation,
};
self.state = Some(H2ActivationState::Dispatched(identity.clone()));
H2DispatchParts {
sender: resources.sender,
upload: H2UploadGuard::new(self.claim.clone(), Some(identity.clone())),
response: H2ResponseGuard::new(self.claim.clone(), Some(identity)),
}
}
pub(super) fn attach_peer_turn(
&mut self,
requesting_cell: &Arc<OriginCell>,
route: H2RouteId,
waiter: WaiterId,
) {
debug_assert_eq!(
self.request_partition,
requesting_cell.id().partition(),
"peer HTTP/2 activation changed requesting partition"
);
assert!(
self.turn.is_none(),
"HTTP/2 activation acquired two requesting-cell turns"
);
self.turn = Some(H2ActivationTurnGuard::peer(requesting_cell, route, waiter));
}
pub(in crate::client::pool) fn close_handle(&self) -> H2CloseHandle {
H2CloseHandle::new(&self.connection_cell, self.generation)
}
pub(in crate::client::pool) fn accept(mut self, dispatch: DispatchGuard) {
let Some(H2ActivationState::Dispatched(identity)) = self.state.as_ref() else {
panic!("HTTP/2 activation accepted before dispatch parts were taken");
};
assert!(
OriginCell::accept_h2_activation(&self.connection_cell, self.generation),
"prospective HTTP/2 activation disappeared before acceptance"
);
self.active = false;
self.claim.accept(dispatch);
if let Some(turn) = self.turn.take() {
turn.release();
}
tracing::trace!(
connection_id = %identity.connection.id(),
request_partition = ?identity.request_partition,
connection_partition = ?identity.connection.owner_partition(),
origin_scheme = %identity.connection.origin().scheme(),
origin_host = identity.connection.origin().host(),
origin_port = ?identity.connection.origin().port(),
h2_generation = ?identity.generation,
"HTTP/2 request accepted"
);
}
}
impl Drop for H2Activation {
fn drop(&mut self) {
if self.active {
OriginCell::cancel_h2_activation(&self.connection_cell, self.generation);
self.claim.cancel();
if let Some(turn) = self.turn.take() {
turn.release();
}
match &self.state {
Some(H2ActivationState::Ready(resources)) => trace_activation_cancelled(
self.request_partition,
self.generation,
resources.connection.info(),
),
Some(H2ActivationState::Dispatched(identity)) => trace_activation_cancelled(
identity.request_partition,
identity.generation,
&identity.connection,
),
None => {}
}
}
}
}
pub(super) struct H2ActivationTurnGuard {
cell: Weak<OriginCell>,
waiter: WaiterId,
owner: H2ActivationGateOwner,
}
enum H2ActivationGateOwner {
Local { generation: H2GenerationId },
Peer { route: H2RouteId },
}
impl H2ActivationTurnGuard {
pub(super) fn local(
cell: &Arc<OriginCell>,
generation: H2GenerationId,
waiter: WaiterId,
) -> Self {
Self {
cell: Weak::from_arc(cell),
waiter,
owner: H2ActivationGateOwner::Local { generation },
}
}
fn peer(cell: &Arc<OriginCell>, route: H2RouteId, waiter: WaiterId) -> Self {
Self {
cell: Weak::from_arc(cell),
waiter,
owner: H2ActivationGateOwner::Peer { route },
}
}
fn release(self) {
if let Some(cell) = self.cell.upgrade() {
match self.owner {
H2ActivationGateOwner::Local { generation } => {
OriginCell::release_local_h2_turn(&cell, generation, self.waiter)
}
H2ActivationGateOwner::Peer { route } => {
OriginCell::release_peer_h2_turn(&cell, route, self.waiter)
}
}
}
}
}
struct H2RequestClaim {
connection_cell: Weak<OriginCell>,
generation: H2GenerationId,
state: Mutex<H2RequestClaimState>,
}
enum H2RequestClaimState {
Prospective {
upload_finished: bool,
response_finished: bool,
},
Accepted {
dispatch: DispatchGuard,
upload_finished: bool,
response_finished: bool,
},
Complete,
}
impl H2RequestClaim {
fn accept(&self, dispatch: DispatchGuard) {
let completed_dispatch = {
let mut state = self.state.lock();
let previous = std::mem::replace(&mut *state, H2RequestClaimState::Complete);
let (upload_finished, response_finished) = match previous {
H2RequestClaimState::Prospective {
upload_finished,
response_finished,
} => (upload_finished, response_finished),
other @ (H2RequestClaimState::Accepted { .. } | H2RequestClaimState::Complete) => {
*state = other;
drop(state);
panic!("HTTP/2 request claim accepted outside prospective state");
}
};
if upload_finished && response_finished {
Some(dispatch)
} else {
*state = H2RequestClaimState::Accepted {
dispatch,
upload_finished,
response_finished,
};
None
}
};
if let Some(dispatch) = completed_dispatch {
drop(dispatch);
self.release_generation();
}
}
fn cancel(&self) {
let mut state = self.state.lock();
if matches!(*state, H2RequestClaimState::Prospective { .. }) {
*state = H2RequestClaimState::Complete;
}
}
fn finish_upload(&self) -> bool {
self.finish_side(|upload_finished, _| *upload_finished = true)
}
fn finish_response(&self) -> bool {
self.finish_side(|_, response_finished| *response_finished = true)
}
fn finish_side(&self, finish: impl FnOnce(&mut bool, &mut bool)) -> bool {
let completed_dispatch = {
let mut state = self.state.lock();
let previous = std::mem::replace(&mut *state, H2RequestClaimState::Complete);
match previous {
H2RequestClaimState::Prospective {
mut upload_finished,
mut response_finished,
} => {
finish(&mut upload_finished, &mut response_finished);
*state = H2RequestClaimState::Prospective {
upload_finished,
response_finished,
};
None
}
H2RequestClaimState::Accepted {
dispatch,
mut upload_finished,
mut response_finished,
} => {
finish(&mut upload_finished, &mut response_finished);
if upload_finished && response_finished {
Some(dispatch)
} else {
*state = H2RequestClaimState::Accepted {
dispatch,
upload_finished,
response_finished,
};
None
}
}
H2RequestClaimState::Complete => None,
}
};
let request_complete = completed_dispatch.is_some();
if let Some(dispatch) = completed_dispatch {
drop(dispatch);
self.release_generation();
}
request_complete
}
fn release_generation(&self) {
if let Some(cell) = self.connection_cell.upgrade() {
OriginCell::release_h2_request(&cell, self.generation);
}
}
#[cfg(all(test, not(smithy_http_client_loom), feature = "rt-tokio"))]
fn for_test(cell: &Arc<OriginCell>) -> Arc<Self> {
Arc::new(Self {
connection_cell: Weak::from_arc(cell),
generation: H2GenerationId(0),
state: Mutex::new(H2RequestClaimState::Prospective {
upload_finished: false,
response_finished: false,
}),
})
}
}
#[derive(Clone)]
struct H2RequestIdentity {
request_partition: PartitionId,
connection: Arc<ConnectionInfo>,
generation: H2GenerationId,
}
fn trace_activation_cancelled(
request_partition: PartitionId,
generation: H2GenerationId,
connection: &ConnectionInfo,
) {
tracing::trace!(
connection_id = %connection.id(),
request_partition = ?request_partition,
connection_partition = ?connection.owner_partition(),
origin_scheme = %connection.origin().scheme(),
origin_host = connection.origin().host(),
origin_port = ?connection.origin().port(),
h2_generation = ?generation,
"HTTP/2 activation cancelled before request acceptance"
);
}
pub(in crate::client::pool) struct H2UploadGuard {
claim: Arc<H2RequestClaim>,
identity: Option<H2RequestIdentity>,
active: bool,
}
impl H2UploadGuard {
fn new(claim: Arc<H2RequestClaim>, identity: Option<H2RequestIdentity>) -> Self {
Self {
claim,
identity,
active: true,
}
}
#[cfg(all(test, not(smithy_http_client_loom), feature = "rt-tokio"))]
pub(in crate::client::pool) fn for_test(cell: &Arc<OriginCell>) -> (Self, H2RequestClaimProbe) {
let claim = H2RequestClaim::for_test(cell);
(
Self::new(claim.clone(), None),
H2RequestClaimProbe { claim },
)
}
pub(in crate::client::pool) fn finish(mut self) {
self.finish_once();
}
fn finish_once(&mut self) {
if self.active {
self.active = false;
let request_complete = self.claim.finish_upload();
trace_request_side(self.identity.take(), "upload", request_complete);
}
}
}
impl Drop for H2UploadGuard {
fn drop(&mut self) {
self.finish_once();
}
}
pub(in crate::client::pool) struct H2ResponseGuard {
claim: Arc<H2RequestClaim>,
identity: Option<H2RequestIdentity>,
active: bool,
}
impl H2ResponseGuard {
fn new(claim: Arc<H2RequestClaim>, identity: Option<H2RequestIdentity>) -> Self {
Self {
claim,
identity,
active: true,
}
}
#[cfg(all(test, not(smithy_http_client_loom), feature = "rt-tokio"))]
pub(in crate::client::pool) fn for_test(cell: &Arc<OriginCell>) -> (Self, H2RequestClaimProbe) {
let claim = H2RequestClaim::for_test(cell);
(
Self::new(claim.clone(), None),
H2RequestClaimProbe { claim },
)
}
pub(in crate::client::pool) fn finish(mut self) {
self.finish_once();
}
fn finish_once(&mut self) {
if self.active {
self.active = false;
let request_complete = self.claim.finish_response();
trace_request_side(self.identity.take(), "response", request_complete);
}
}
}
impl Drop for H2ResponseGuard {
fn drop(&mut self) {
self.finish_once();
}
}
fn trace_request_side(
identity: Option<H2RequestIdentity>,
request_side: &'static str,
request_complete: bool,
) {
if let Some(identity) = identity {
tracing::trace!(
connection_id = %identity.connection.id(),
request_partition = ?identity.request_partition,
connection_partition = ?identity.connection.owner_partition(),
origin_scheme = %identity.connection.origin().scheme(),
origin_host = identity.connection.origin().host(),
origin_port = ?identity.connection.origin().port(),
h2_generation = ?identity.generation,
request_side,
request_complete,
"HTTP/2 request side finished"
);
}
}
#[cfg(all(test, not(smithy_http_client_loom), feature = "rt-tokio"))]
pub(in crate::client::pool) struct H2RequestClaimProbe {
claim: Arc<H2RequestClaim>,
}
#[cfg(all(test, not(smithy_http_client_loom), feature = "rt-tokio"))]
impl H2RequestClaimProbe {
pub(in crate::client::pool) fn upload_finished(&self) -> bool {
matches!(
&*self.claim.state.lock(),
H2RequestClaimState::Prospective {
upload_finished: true,
..
} | H2RequestClaimState::Accepted {
upload_finished: true,
..
} | H2RequestClaimState::Complete
)
}
pub(in crate::client::pool) fn response_finished(&self) -> bool {
matches!(
&*self.claim.state.lock(),
H2RequestClaimState::Prospective {
response_finished: true,
..
} | H2RequestClaimState::Accepted {
response_finished: true,
..
} | H2RequestClaimState::Complete
)
}
}