diff --git a/.cargo/config.toml b/.cargo/config.toml index 6c02d18a06..70ff9d2763 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -1,3 +1,3 @@ [build] -rustflags=['--cfg', 's2n_internal_dev'] +rustflags=['--cfg', 's2n_internal_dev', '--cfg', 'tokio_unstable'] rustdocflags=['--cfg', 's2n_internal_dev'] diff --git a/dc/s2n-quic-dc/Cargo.toml b/dc/s2n-quic-dc/Cargo.toml index 9bb222fe5a..879945215c 100644 --- a/dc/s2n-quic-dc/Cargo.toml +++ b/dc/s2n-quic-dc/Cargo.toml @@ -48,6 +48,7 @@ s2n-codec = { version = "=0.86.0", path = "../../common/s2n-codec", default-feat s2n-quic = { version = "=1.86.0", path = "../../quic/s2n-quic", features = ["unstable-provider-connection-close-formatter", "unstable-provider-dc", "unstable-provider-random", "unstable-offload-tls"] } s2n-quic-core = { version = "=0.86.0", path = "../../quic/s2n-quic-core", default-features = false } s2n-quic-platform = { version = "=0.86.0", path = "../../quic/s2n-quic-platform" } +s2n-quic-dc-metrics = { version= "=0.86.0", path = "../s2n-quic-dc-metrics" } s2n-tls = "0.3.38" slotmap = "1" hashbrown = "0.17" diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index c3c33cebdb..9f6867ef35 100644 --- a/dc/s2n-quic-dc/src/psk/io.rs +++ b/dc/s2n-quic-dc/src/psk/io.rs @@ -13,6 +13,7 @@ use s2n_quic::{ server::Name, }; use s2n_quic_core::{endpoint::Type, inet::SocketAddress}; +use s2n_quic_dc_metrics::TaskMonitor; use std::{ any::Any, hash::BuildHasher, @@ -48,10 +49,15 @@ pub type Result = core::result::Result; struct TokioExecutor { runtime: Runtime, + monitor: Option, } impl s2n_quic::provider::tls::offload::Executor for TokioExecutor { fn spawn(&self, task: impl core::future::Future + Send + 'static) { - self.runtime.spawn(task); + if let Some(monitor) = &self.monitor { + self.runtime.spawn(monitor.instrument(task)); + } else { + self.runtime.spawn(task); + } } } #[derive(Clone)] @@ -167,6 +173,17 @@ impl Server { .enable_all() .build()?; + let monitor = if let Some(registry) = builder.registry { + registry.instrument_runtime( + "offload runtime stats", + runtime.handle(), + Duration::new(1, 0), + ); + Some(registry.register_task_monitor("offload task stats")) + } else { + None + }; + let tls = s2n_quic::provider::tls::offload::OffloadBuilder::new() .with_endpoint(tls_materials_provider) .with_exporter(DCExporter { @@ -174,7 +191,7 @@ impl Server { endpoint_type: Type::Server, map: map.clone(), }) - .with_executor(TokioExecutor { runtime }) + .with_executor(TokioExecutor { runtime, monitor }) .build(); // We need packet storage when offloading is turned on due to this issue: diff --git a/dc/s2n-quic-dc/src/psk/server/builder.rs b/dc/s2n-quic-dc/src/psk/server/builder.rs index 2697d5a85c..a2be195715 100644 --- a/dc/s2n-quic-dc/src/psk/server/builder.rs +++ b/dc/s2n-quic-dc/src/psk/server/builder.rs @@ -9,6 +9,7 @@ use crate::{ }, }; use s2n_quic::provider::{event::Subscriber as Sub, tls::Provider as Prov}; +use s2n_quic_dc_metrics::Registry; use std::{net::SocketAddr, time::Duration}; use super::Provider; @@ -25,6 +26,7 @@ pub struct Builder< #[cfg(any(test, feature = "testing"))] pub(crate) endpoint_limits: Option, pub(crate) thread_offload_count: usize, + pub(crate) registry: Option, } /// A wrapper type for test endpoint limiters @@ -55,6 +57,7 @@ impl Default for Builder { #[cfg(any(test, feature = "testing"))] endpoint_limits: None, thread_offload_count: DEFAULT_THREAD_COUNT, + registry: None, } } } @@ -75,6 +78,7 @@ impl Builder { #[cfg(any(test, feature = "testing"))] endpoint_limits: self.endpoint_limits, thread_offload_count: self.thread_offload_count, + registry: self.registry, } } @@ -147,6 +151,11 @@ impl Builder { self } + pub fn with_registry(mut self, registry: Registry) -> Self { + self.registry = Some(registry); + self + } + /// Starts the server listening to the given address. pub async fn start< TlsProvider: Prov + Send + Sync + 'static,