//! Application-owned QUIC endpoint and connection lifecycles. //! //! s2n-quic's default Tokio IO adapter discards the endpoint task handle. This //! module keeps that handle and gives the application an explicit stop signal, //! so a peer runtime does not report itself stopped while its QUIC endpoint is //! still running. use std::{ future::Future as _, io, net::SocketAddr, ops::{Deref, DerefMut}, pin::Pin, sync::{Arc, Mutex}, task::{Context, Poll}, }; use lanspread_proto::{MAX_CONTROL_FRAME_BYTES, PeerEndpoint}; use s2n_quic::{ Client as QuicClient, Connection, client::Connect, provider::{ congestion_controller, io::{Provider as IoProvider, tokio::Provider as TokioIo}, limits::Limits, }, }; use s2n_quic_core::{ endpoint::{self, CloseError, Endpoint}, inet::SocketAddress, io::{rx, tx}, path::mtu, time::{Clock, Timestamp}, }; use tokio::task::{JoinError, JoinHandle}; use tokio_util::sync::{CancellationToken, WaitForCancellationFutureOwned}; use crate::{ config::{ QUIC_CONNECTION_DATA_WINDOW, QUIC_ENDPOINT_SHUTDOWN_GRACE, QUIC_HANDSHAKE_TIMEOUT, QUIC_IDLE_TIMEOUT, QUIC_INITIAL_CONGESTION_WINDOW, QUIC_MAX_SEND_BUFFER_SIZE, QUIC_SOCKET_BUFFER_SIZE, QUIC_STREAM_DATA_WINDOW, }, tls, }; /// Transport-level bound on simultaneously open bidirectional streams in /// either direction. The server independently enforces the same application /// control-stream task bound before spawning work. pub(crate) const MAX_OPEN_BIDIRECTIONAL_STREAMS: u64 = 32; /// One server connection exposes only enough aggregate receive credit for one /// maximal length-delimited control frame at a time. pub(crate) const SERVER_CONTROL_RECEIVE_WINDOW: u64 = MAX_CONTROL_FRAME_BYTES as u64 + 4; pub(crate) fn quic_client_limits() -> eyre::Result { Ok(Limits::default() .with_data_window(QUIC_CONNECTION_DATA_WINDOW)? .with_bidirectional_local_data_window(QUIC_STREAM_DATA_WINDOW)? .with_bidirectional_remote_data_window(0)? // The application protocol never opens or accepts unidirectional // streams, so expose neither stream IDs nor receive-window state. .with_unidirectional_data_window(0)? .with_max_open_local_bidirectional_streams(MAX_OPEN_BIDIRECTIONAL_STREAMS)? .with_max_open_remote_bidirectional_streams(0)? .with_max_open_local_unidirectional_streams(0)? .with_max_open_remote_unidirectional_streams(0)? .with_max_send_buffer_size(QUIC_MAX_SEND_BUFFER_SIZE)? .with_max_handshake_duration(QUIC_HANDSHAKE_TIMEOUT)? .with_max_idle_timeout(QUIC_IDLE_TIMEOUT)?) } /// Server-specific transport limits for public remote-initiated control /// streams. Large client receive windows remain available only on the /// separately configured outbound connector. pub(crate) fn quic_server_limits() -> eyre::Result { Ok(Limits::default() .with_data_window(SERVER_CONTROL_RECEIVE_WINDOW)? .with_bidirectional_local_data_window(0)? .with_bidirectional_remote_data_window(SERVER_CONTROL_RECEIVE_WINDOW)? .with_unidirectional_data_window(0)? .with_max_open_local_bidirectional_streams(0)? .with_max_open_remote_bidirectional_streams(MAX_OPEN_BIDIRECTIONAL_STREAMS)? .with_max_open_local_unidirectional_streams(0)? .with_max_open_remote_unidirectional_streams(0)? .with_max_send_buffer_size( u32::try_from(MAX_CONTROL_FRAME_BYTES) .map_err(|_| eyre::eyre!("control frame bound exceeds u32"))?, )? .with_max_handshake_duration(QUIC_HANDSHAKE_TIMEOUT)? .with_max_idle_timeout(QUIC_IDLE_TIMEOUT)?) } pub(crate) fn quic_congestion_controller() -> congestion_controller::Bbr { congestion_controller::bbr::Builder::default() .with_initial_congestion_window(QUIC_INITIAL_CONGESTION_WINDOW) .build() } /// Builds an IO provider together with the owner used to retrieve its endpoint /// task after s2n-quic has started it. pub(crate) fn tracked_quic_io( addr: SocketAddr, ) -> eyre::Result<(TrackedIoProvider, EndpointTaskControl)> { let inner = s2n_quic::provider::io::tokio::Builder::default() .with_receive_address(addr)? .with_send_buffer_size(QUIC_SOCKET_BUFFER_SIZE)? .with_recv_buffer_size(QUIC_SOCKET_BUFFER_SIZE)? .build()?; let control = EndpointTaskControl::new(); Ok(( TrackedIoProvider { inner, control: control.clone(), }, control, )) } /// Tokio IO provider that retains the s2n endpoint task instead of discarding /// its handle like the default adapter does. pub(crate) struct TrackedIoProvider { inner: TokioIo, control: EndpointTaskControl, } impl IoProvider for TrackedIoProvider { type PathHandle = ::PathHandle; type Error = io::Error; fn start>( self, endpoint: E, ) -> Result { // Lock before starting the endpoint. That way a poisoned/invalid slot // cannot cause a newly spawned endpoint task to escape unowned. let mut task_slot = self.control.lock_task_slot()?; if task_slot.is_some() { return Err(io::Error::other("QUIC endpoint task was already started")); } let endpoint = ControlledEndpoint::new(endpoint, self.control.stop.clone()); let (task, local_addr) = self.inner.start(endpoint)?; *task_slot = Some(task); Ok(local_addr) } } #[derive(Clone)] pub(crate) struct EndpointTaskControl { stop: CancellationToken, task: Arc, } struct EndpointTaskSlot(Mutex>>); impl Drop for EndpointTaskSlot { fn drop(&mut self) { let task = self .0 .get_mut() .unwrap_or_else(std::sync::PoisonError::into_inner) .take(); if let Some(task) = task { // A successful pinned s2n builder returns immediately after IO // start, so normal code always transfers this handle into // EndpointTask. This is only a synchronous panic-safety fallback; // the peer supervisor destroys and joins the isolated Tokio // runtime before its public owner may return. log::error!("QUIC endpoint started but ownership was not transferred; aborting it"); task.abort(); } } } impl EndpointTaskControl { fn new() -> Self { Self { stop: CancellationToken::new(), task: Arc::new(EndpointTaskSlot(Mutex::new(None))), } } fn lock_task_slot(&self) -> io::Result>>> { self.task .0 .lock() .map_err(|_| io::Error::other("QUIC endpoint task slot was poisoned")) } /// Takes ownership of the endpoint task after a successful s2n-quic start. pub(crate) fn take_started(&self) -> eyre::Result { let task = self .lock_task_slot()? .take() .ok_or_else(|| eyre::eyre!("QUIC IO provider did not start an endpoint task"))?; Ok(EndpointTask { stop: self.stop.clone(), task: Some(task), abort_requested: false, }) } } struct ControlledEndpoint { inner: E, stop: Pin>, } impl ControlledEndpoint { fn new(inner: E, stop: CancellationToken) -> Self { Self { inner, stop: Box::pin(stop.cancelled_owned()), } } } impl Endpoint for ControlledEndpoint { type PathHandle = E::PathHandle; type Subscriber = E::Subscriber; const ENDPOINT_TYPE: endpoint::Type = E::ENDPOINT_TYPE; fn receive(&mut self, queue: &mut Rx, clock: &C) where Rx: rx::Queue, C: Clock, { self.inner.receive(queue, clock); } fn transmit(&mut self, queue: &mut Tx, clock: &C) where Tx: tx::Queue, C: Clock, { self.inner.transmit(queue, clock); } fn poll_wakeups( &mut self, cx: &mut Context<'_>, clock: &C, ) -> Poll> { if self.stop.as_mut().poll(cx).is_ready() { return Poll::Ready(Err(CloseError)); } self.inner.poll_wakeups(cx, clock) } fn timeout(&self) -> Option { self.inner.timeout() } fn set_mtu_config(&mut self, config: mtu::Config) { self.inner.set_mtu_config(config); } fn subscriber(&mut self) -> &mut Self::Subscriber { self.inner.subscriber() } } /// Application-owned endpoint task. /// /// Every normal owner path consumes this value through [`Self::shutdown_and_join`]. /// Drop is only a panic-safety fallback: it prevents continued endpoint work, /// while the peer supervisor's final runtime teardown bounds and destroys the /// aborted task before `PeerRuntimeHandle` may finish joining. #[must_use = "the QUIC endpoint task must be explicitly stopped and joined"] pub(crate) struct EndpointTask { stop: CancellationToken, task: Option>, abort_requested: bool, } impl EndpointTask { pub(crate) async fn shutdown_and_join(mut self) -> eyre::Result<()> { self.shutdown_and_join_with_grace(QUIC_ENDPOINT_SHUTDOWN_GRACE) .await } async fn shutdown_and_join_with_grace( &mut self, grace: std::time::Duration, ) -> eyre::Result<()> { self.stop.cancel(); let Some(task) = self.task.as_mut() else { eyre::bail!("QUIC endpoint task was already joined"); }; // Keep the handle in `self` until every await has completed. If the // caller cancels this cleanup future, dropping `EndpointTask` can then // still abort the endpoint instead of silently detaching a handle that // was moved into the cancelled future's local state. let result = if let Ok(result) = tokio::time::timeout(grace, &mut *task).await { join_result(result, self.abort_requested) } else { log::warn!("QUIC endpoint did not stop within {grace:?}; aborting and joining it"); self.abort_requested = true; task.abort(); join_result(task.await, true) }; // A completed Tokio JoinHandle is inert. Remove it only after the // terminal await so `Drop` remains a cancellation-safety guard for the // whole shutdown operation. self.task = None; result } } impl Drop for EndpointTask { fn drop(&mut self) { if let Some(task) = self.task.take() { log::error!( "QUIC endpoint owner was dropped before join; aborting as panic-safety fallback" ); self.stop.cancel(); // The task remains owned by the isolated peer Tokio runtime. Its // supervisor thread cannot finish (and its public handle cannot // return from Drop/wait) until that runtime has been destroyed. task.abort(); } } } fn join_result(result: Result<(), JoinError>, expected_abort: bool) -> eyre::Result<()> { match result { Ok(()) => Ok(()), Err(error) if expected_abort && error.is_cancelled() => Ok(()), Err(error) => Err(eyre::eyre!("QUIC endpoint task failed to join: {error}")), } } /// Cloneable capability for opening connections on the runtime-owned client /// endpoint. #[derive(Clone, Debug)] pub(crate) struct QuicConnector { client: Option, } impl QuicConnector { pub(crate) async fn connect(&self, endpoint: &PeerEndpoint) -> eyre::Result { let client = self .client .as_ref() .ok_or_else(|| eyre::eyre!("QUIC connector is unavailable"))?; let server_name = tls::sni_for_peer(endpoint.peer_id)?; let connect = Connect::new(endpoint.addr).with_server_name(server_name); Ok(PeerConnection::new(client.connect(connect).await?)) } #[cfg(test)] pub(crate) fn unavailable() -> Self { Self { client: None } } } /// Runtime owner for the single outgoing client endpoint. pub(crate) struct QuicClientRuntime { client: QuicClient, endpoint: EndpointTask, } impl QuicClientRuntime { /// Waits for outstanding connections to settle, then always stops and /// joins the endpoint. A stuck transport can consume the grace period, but /// cannot make peer shutdown unbounded. pub(crate) async fn shutdown(mut self) -> eyre::Result<()> { let idle_result = match tokio::time::timeout(QUIC_ENDPOINT_SHUTDOWN_GRACE, self.client.wait_idle()).await { Ok(result) => result.map_err(eyre::Report::from), Err(_) => Err(eyre::eyre!( "QUIC client did not become idle within {QUIC_ENDPOINT_SHUTDOWN_GRACE:?}" )), }; // No application connector remains after the runtime's child scopes // have drained. Drop this final handle before forcing endpoint stop. drop(self.client); let endpoint_result = self.endpoint.shutdown_and_join().await; match (idle_result, endpoint_result) { (Ok(()), Ok(())) => Ok(()), (Err(error), Ok(())) | (Ok(()), Err(error)) => Err(error), (Err(idle_error), Err(endpoint_error)) => Err(eyre::eyre!( "QUIC client shutdown failed: {idle_error:#}; endpoint shutdown also failed: {endpoint_error:#}" )), } } /// Settles the test-only hostile-handshake fixture after the verifier has /// already proved that connection establishment was rejected. /// /// A failed TLS `CertificateVerify` can remain in s2n-quic's endpoint-owned /// closing state beyond the normal wait-idle grace. This seam is not a /// general timeout bypass: it drops the final client handle, then still /// explicitly stops and joins the endpoint task. #[cfg(test)] pub(crate) async fn shutdown_rejected_handshake_fixture(self) -> eyre::Result<()> { drop(self.client); self.endpoint.shutdown_and_join().await } } pub(crate) fn start_quic_client() -> eyre::Result<(QuicClientRuntime, QuicConnector)> { let (io, endpoint_control) = tracked_quic_io(SocketAddr::from(([0, 0, 0, 0], 0)))?; let client = QuicClient::builder() .with_tls(tls::client_provider()?)? .with_io(io)? .with_limits(quic_client_limits()?)? .with_congestion_controller(quic_congestion_controller())? .start()?; let endpoint = endpoint_control.take_started()?; let connector = QuicConnector { client: Some(client.clone()), }; Ok((QuicClientRuntime { client, endpoint }, connector)) } /// Connection whose lexical owner always initiates QUIC close before dropping /// the handle, including when a request future is cancelled or times out. pub(crate) struct PeerConnection { connection: Connection, } impl PeerConnection { fn new(connection: Connection) -> Self { Self { connection } } } impl Deref for PeerConnection { type Target = Connection; fn deref(&self) -> &Self::Target { &self.connection } } impl DerefMut for PeerConnection { fn deref_mut(&mut self) -> &mut Self::Target { &mut self.connection } } impl Drop for PeerConnection { fn drop(&mut self) { self.connection.close(0u32.into()); } } #[cfg(test)] mod tests { use std::{sync::Arc, time::Duration}; use tokio::sync::Notify; use tokio_util::sync::CancellationToken; use super::{EndpointTask, EndpointTaskControl, start_quic_client}; #[test] fn endpoint_control_rejects_take_before_provider_start() { let control = EndpointTaskControl::new(); assert!(control.take_started().is_err()); } #[tokio::test] async fn client_start_captures_an_endpoint_that_can_be_joined() { let (runtime, connector) = start_quic_client().expect("client endpoint should start"); drop(connector); tokio::time::timeout(Duration::from_secs(2), runtime.shutdown()) .await .expect("client endpoint shutdown should be bounded") .expect("client endpoint should join cleanly"); } #[tokio::test] async fn endpoint_owner_signals_and_joins_its_task() { let stop = CancellationToken::new(); let stopped = Arc::new(Notify::new()); let task_stop = stop.clone(); let task_stopped = stopped.clone(); let task = tokio::spawn(async move { task_stop.cancelled().await; task_stopped.notify_one(); }); let owner = EndpointTask { stop, task: Some(task), abort_requested: false, }; owner .shutdown_and_join() .await .expect("synthetic endpoint should join"); tokio::time::timeout(Duration::from_secs(1), stopped.notified()) .await .expect("endpoint cleanup should happen before owner returns"); } #[tokio::test] async fn endpoint_owner_aborts_but_still_joins_uncooperative_task() { let stop = CancellationToken::new(); let dropped = Arc::new(Notify::new()); let task_dropped = dropped.clone(); let task = tokio::spawn(async move { struct DropNotice(Arc); impl Drop for DropNotice { fn drop(&mut self) { self.0.notify_one(); } } let _drop_notice = DropNotice(task_dropped); std::future::pending::<()>().await; }); let owner = EndpointTask { stop, task: Some(task), abort_requested: false, }; let mut owner = owner; owner .shutdown_and_join_with_grace(Duration::from_millis(10)) .await .expect("aborted endpoint should still join"); tokio::time::timeout(Duration::from_secs(1), dropped.notified()) .await .expect("aborted endpoint should finish unwinding before owner returns"); } #[tokio::test] async fn cancelling_endpoint_cleanup_retains_task_for_retry() { let stop = CancellationToken::new(); let started = Arc::new(Notify::new()); let dropped = Arc::new(Notify::new()); let task_started = started.clone(); let task_dropped = dropped.clone(); let task = tokio::spawn(async move { struct DropNotice(Arc); impl Drop for DropNotice { fn drop(&mut self) { self.0.notify_one(); } } let _drop_notice = DropNotice(task_dropped); task_started.notify_one(); std::future::pending::<()>().await; }); let mut owner = EndpointTask { stop, task: Some(task), abort_requested: false, }; started.notified().await; let mut cleanup = Box::pin(owner.shutdown_and_join_with_grace(Duration::from_secs(10))); assert!( tokio::time::timeout(Duration::from_millis(10), cleanup.as_mut()) .await .is_err(), "synthetic endpoint should still be awaiting its grace period" ); drop(cleanup); assert!( owner.task.is_some(), "cancelling cleanup must leave the JoinHandle with its owner" ); owner .shutdown_and_join_with_grace(Duration::from_millis(10)) .await .expect("the owner should be able to retry and join after cancellation"); tokio::time::timeout(Duration::from_secs(1), dropped.notified()) .await .expect("retried cleanup must join the task retained by its owner"); } }