use crate::Connection;
use cfg_if::cfg_if;
use s2n_quic_core::{
ensure,
event::{
api as events,
api::{ConnectionInfo, ConnectionMeta, MtuUpdated, Subscriber},
},
};
use std::time::Duration;
use tokio::sync::watch;
pub struct MtuConfirmComplete;
impl MtuConfirmComplete {
pub async fn wait_ready(conn: &mut Connection) {
let (mut receiver, peer_will_send) = conn
.query_event_context_mut(|context: &mut MtuConfirmContext| {
(
context.sender.subscribe(),
context.peer_will_send_completion,
)
})
.expect("connection context isn't properly set");
loop {
let ready = {
let state = receiver.borrow_and_update();
if peer_will_send {
state.is_ready()
} else {
state.local_ready
}
};
if ready {
if !peer_will_send {
cfg_if!(
if #[cfg(any(test, feature = "unstable-provider-io-testing"))] {
if tokio::runtime::Handle::try_current().is_err() {
crate::provider::io::testing::time::delay(Duration::from_secs(1)).await;
} else {
tokio::time::sleep(Duration::from_secs(1)).await;
}
} else {
tokio::time::sleep(Duration::from_secs(1)).await;
}
);
}
return;
}
if receiver.changed().await.is_err() {
return;
}
}
}
}
pub struct MtuConfirmContext {
sender: watch::Sender<MtuProbingState>,
peer_will_send_completion: bool,
}
impl Default for MtuConfirmContext {
fn default() -> Self {
let (sender, _receiver) = watch::channel(MtuProbingState::default());
Self {
sender,
peer_will_send_completion: false,
}
}
}
impl MtuConfirmContext {
fn update_and_check(&mut self, updater: impl FnOnce(&mut MtuProbingState)) {
self.sender.send_modify(|state| {
updater(state);
});
}
}
impl Drop for MtuConfirmContext {
fn drop(&mut self) {
self.sender.send_modify(|state| {
state.local_ready = true;
state.remote_ready = true;
});
}
}
#[derive(Debug, Clone, Copy, Default)]
struct MtuProbingState {
local_ready: bool,
remote_ready: bool,
}
impl MtuProbingState {
fn is_ready(&self) -> bool {
self.local_ready && self.remote_ready
}
}
impl Subscriber for MtuConfirmComplete {
type ConnectionContext = MtuConfirmContext;
#[inline]
fn create_connection_context(
&mut self,
_: &ConnectionMeta,
_info: &ConnectionInfo,
) -> Self::ConnectionContext {
MtuConfirmContext::default()
}
#[inline]
fn on_transport_parameters_received(
&mut self,
context: &mut Self::ConnectionContext,
_meta: &ConnectionMeta,
event: &events::TransportParametersReceived,
) {
context.peer_will_send_completion = event.transport_parameters.mtu_probing_complete_support;
}
#[inline]
fn on_connection_closed(
&mut self,
context: &mut Self::ConnectionContext,
_meta: &ConnectionMeta,
_event: &events::ConnectionClosed,
) {
let state = *context.sender.borrow();
ensure!(!state.is_ready());
#[cfg(feature = "provider-event-tracing")]
if context.peer_will_send_completion && !state.remote_ready {
tracing::warn!(
local_ready = state.local_ready,
"peer indicated MtuProbingComplete support but closed connection before sending it"
);
}
context.update_and_check(|state| {
state.local_ready = true;
state.remote_ready = true;
});
}
#[inline]
fn on_mtu_updated(
&mut self,
context: &mut Self::ConnectionContext,
_meta: &ConnectionMeta,
event: &MtuUpdated,
) {
if event.search_complete {
context.update_and_check(|state| {
state.local_ready = true;
});
}
}
#[inline]
fn on_mtu_probing_complete_received(
&mut self,
context: &mut Self::ConnectionContext,
_meta: &ConnectionMeta,
_event: &events::MtuProbingCompleteReceived,
) {
context.update_and_check(|state| {
state.remote_ready = true;
});
}
}