From 4dc645f4f1045e93e6cae93a7b9bc883f57bb17e Mon Sep 17 00:00:00 2001 From: Appelmans Date: Tue, 4 Aug 2026 11:55:44 -0700 Subject: [PATCH 1/9] Adds per connection offload runtime metrics --- dc/s2n-quic-dc/events/endpoint.rs | 9 ++ dc/s2n-quic-dc/src/event/generated.rs | 125 ++++++++++++++++++ .../src/event/generated/metrics/aggregate.rs | 78 ++++++++++- .../src/event/generated/metrics/probe.rs | 24 ++++ dc/s2n-quic-dc/src/path/secret/map.rs | 9 +- dc/s2n-quic-dc/src/path/secret/map/state.rs | 10 +- dc/s2n-quic-dc/src/path/secret/map/store.rs | 4 +- dc/s2n-quic-dc/src/psk/io.rs | 20 ++- .../src/bin/overload-server/mod.rs | 24 +++- 9 files changed, 292 insertions(+), 11 deletions(-) diff --git a/dc/s2n-quic-dc/events/endpoint.rs b/dc/s2n-quic-dc/events/endpoint.rs index c7757ff464..347da1cea6 100644 --- a/dc/s2n-quic-dc/events/endpoint.rs +++ b/dc/s2n-quic-dc/events/endpoint.rs @@ -21,3 +21,12 @@ struct DcConnectionTimeout<'a> { #[nominal_counter("peer_address.protocol")] peer_address: SocketAddress<'a>, } + +#[event("runtime:offload_metrics")] +#[subject(endpoint)] +struct OffloadRuntimeMetrics { + #[measure("global_queue_depth", Count)] + global_queue_depth: usize, + #[measure("num_alive_tasks", Count)] + num_alive_tasks: usize, +} diff --git a/dc/s2n-quic-dc/src/event/generated.rs b/dc/s2n-quic-dc/src/event/generated.rs index 06f83bad60..7a76f953c6 100644 --- a/dc/s2n-quic-dc/src/event/generated.rs +++ b/dc/s2n-quic-dc/src/event/generated.rs @@ -1892,6 +1892,24 @@ pub mod api { } #[derive(Clone, Debug)] #[non_exhaustive] + pub struct OffloadRuntimeMetrics { + pub global_queue_depth: usize, + pub num_alive_tasks: usize, + } + #[cfg(any(test, feature = "testing"))] + impl crate::event::snapshot::Fmt for OffloadRuntimeMetrics { + fn fmt(&self, fmt: &mut core::fmt::Formatter) -> core::fmt::Result { + let mut fmt = fmt.debug_struct("OffloadRuntimeMetrics"); + fmt.field("global_queue_depth", &self.global_queue_depth); + fmt.field("num_alive_tasks", &self.num_alive_tasks); + fmt.finish() + } + } + impl Event for OffloadRuntimeMetrics { + const NAME: &'static str = "runtime:offload_metrics"; + } + #[derive(Clone, Debug)] + #[non_exhaustive] pub struct PathSecretMapInitialized { /// The capacity of the path secret map pub capacity: usize, @@ -4012,6 +4030,24 @@ pub mod tracing { ); } #[inline] + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + let parent = self.parent(meta); + let api::OffloadRuntimeMetrics { + global_queue_depth, + num_alive_tasks, + } = event; + tracing::event!( + target : "offload_runtime_metrics", parent : parent, + tracing::Level::DEBUG, { global_queue_depth = + tracing::field::debug(global_queue_depth), num_alive_tasks = + tracing::field::debug(num_alive_tasks) } + ); + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -6418,6 +6454,24 @@ pub mod builder { } } #[derive(Clone, Debug)] + pub struct OffloadRuntimeMetrics { + pub global_queue_depth: usize, + pub num_alive_tasks: usize, + } + impl IntoEvent for OffloadRuntimeMetrics { + #[inline] + fn into_event(self) -> api::OffloadRuntimeMetrics { + let OffloadRuntimeMetrics { + global_queue_depth, + num_alive_tasks, + } = self; + api::OffloadRuntimeMetrics { + global_queue_depth: global_queue_depth.into_event(), + num_alive_tasks: num_alive_tasks.into_event(), + } + } + } + #[derive(Clone, Debug)] pub struct PathSecretMapInitialized { /// The capacity of the path secret map pub capacity: usize, @@ -8004,6 +8058,16 @@ mod traits { let _ = meta; let _ = event; } + ///Called when the `OffloadRuntimeMetrics` event is triggered + #[inline] + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + let _ = meta; + let _ = event; + } ///Called when the `PathSecretMapInitialized` event is triggered #[inline] fn on_path_secret_map_initialized( @@ -8956,6 +9020,14 @@ mod traits { self.as_ref().on_dc_connection_timeout(meta, event); } #[inline] + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + self.as_ref().on_offload_runtime_metrics(meta, event); + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -9888,6 +9960,15 @@ mod traits { (self.1).on_dc_connection_timeout(meta, event); } #[inline] + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + (self.0).on_offload_runtime_metrics(meta, event); + (self.1).on_offload_runtime_metrics(meta, event); + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -10305,6 +10386,8 @@ mod traits { fn on_endpoint_initialized(&self, event: builder::EndpointInitialized); ///Publishes a `DcConnectionTimeout` event to the publisher's subscriber fn on_dc_connection_timeout(&self, event: builder::DcConnectionTimeout); + ///Publishes a `OffloadRuntimeMetrics` event to the publisher's subscriber + fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics); ///Publishes a `PathSecretMapInitialized` event to the publisher's subscriber fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized); ///Publishes a `PathSecretMapUninitialized` event to the publisher's subscriber @@ -10664,6 +10747,13 @@ mod traits { self.subscriber.on_event(&self.meta, &event); } #[inline] + fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics) { + let event = event.into_event(); + self.subscriber + .on_offload_runtime_metrics(&self.meta, &event); + self.subscriber.on_event(&self.meta, &event); + } + #[inline] fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized) { let event = event.into_event(); self.subscriber @@ -11430,6 +11520,7 @@ pub mod testing { pub stream_connect_error: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, + pub offload_runtime_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -11528,6 +11619,7 @@ pub mod testing { stream_connect_error: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), + offload_runtime_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -11932,6 +12024,17 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + self.offload_runtime_metrics.fetch_add(1, Ordering::Relaxed); + let meta = crate::event::snapshot::Fmt::to_snapshot(meta); + let event = crate::event::snapshot::Fmt::to_snapshot(event); + let out = format!("{meta:?} {event:?}"); + self.output.lock().unwrap().push(out); + } fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -12428,6 +12531,7 @@ pub mod testing { pub connection_closed: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, + pub offload_runtime_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -12559,6 +12663,7 @@ pub mod testing { connection_closed: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), + offload_runtime_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -13432,6 +13537,17 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + self.offload_runtime_metrics.fetch_add(1, Ordering::Relaxed); + let meta = crate::event::snapshot::Fmt::to_snapshot(meta); + let event = crate::event::snapshot::Fmt::to_snapshot(event); + let out = format!("{meta:?} {event:?}"); + self.output.lock().unwrap().push(out); + } fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -13927,6 +14043,7 @@ pub mod testing { pub connection_closed: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, + pub offload_runtime_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -14048,6 +14165,7 @@ pub mod testing { connection_closed: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), + offload_runtime_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -14340,6 +14458,13 @@ pub mod testing { let out = format!("{event:?}"); self.output.lock().unwrap().push(out); } + fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics) { + self.offload_runtime_metrics.fetch_add(1, Ordering::Relaxed); + let event = event.into_event(); + let event = crate::event::snapshot::Fmt::to_snapshot(&event); + let out = format!("{event:?}"); + self.output.lock().unwrap().push(out); + } fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized) { self.path_secret_map_initialized .fetch_add(1, Ordering::Relaxed); diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs index 5756228fb0..129ed49ea7 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs @@ -248,6 +248,9 @@ mod id { ENDPOINT_INITIALIZED__UDP, DC_CONNECTION_TIMEOUT, DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL, + OFFLOAD_RUNTIME_METRICS, + OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, + OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -738,6 +741,11 @@ mod id { pub const DC_CONNECTION_TIMEOUT: usize = InfoId::DC_CONNECTION_TIMEOUT as usize; pub const DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL: usize = InfoId::DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL as usize; + pub const OFFLOAD_RUNTIME_METRICS: usize = InfoId::OFFLOAD_RUNTIME_METRICS as usize; + pub const OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH: usize = + InfoId::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; + pub const OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = + InfoId::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1022,6 +1030,7 @@ mod id { COUNTERS_CONNECTION_CLOSED, COUNTERS_ENDPOINT_INITIALIZED, COUNTERS_DC_CONNECTION_TIMEOUT, + COUNTERS_OFFLOAD_RUNTIME_METRICS, COUNTERS_PATH_SECRET_MAP_INITIALIZED, COUNTERS_PATH_SECRET_MAP_UNINITIALIZED, COUNTERS_PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED, @@ -1204,6 +1213,8 @@ mod id { Counters::COUNTERS_ENDPOINT_INITIALIZED as usize; pub const COUNTERS_DC_CONNECTION_TIMEOUT: usize = Counters::COUNTERS_DC_CONNECTION_TIMEOUT as usize; + pub const COUNTERS_OFFLOAD_RUNTIME_METRICS: usize = + Counters::COUNTERS_OFFLOAD_RUNTIME_METRICS as usize; pub const COUNTERS_PATH_SECRET_MAP_INITIALIZED: usize = Counters::COUNTERS_PATH_SECRET_MAP_INITIALIZED as usize; pub const COUNTERS_PATH_SECRET_MAP_UNINITIALIZED: usize = @@ -1586,6 +1597,8 @@ mod id { MEASURES_STREAM_CONTROL_PACKET_RECEIVED__PACKET_LEN, MEASURES_STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN, MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN, + MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, + MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__ENTRIES, @@ -1819,6 +1832,10 @@ mod id { Measures::MEASURES_STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN as usize; pub const MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN: usize = Measures::MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN as usize; + pub const MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH: usize = + Measures::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; + pub const MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = + Measures::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; pub const MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = Measures::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; pub const MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY: usize = @@ -3384,6 +3401,24 @@ static INFO: &[Info; 341usize] = &[ units: Units::None, } .build(), + info::Builder { + id: id::OFFLOAD_RUNTIME_METRICS, + name: Str::new("offload_runtime_metrics\0"), + units: Units::None, + } + .build(), + info::Builder { + id: id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, + name: Str::new("offload_runtime_metrics.global_queue_depth\0"), + units: Units::None, + } + .build(), + info::Builder { + id: id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, + name: Str::new("offload_runtime_metrics.num_alive_tasks\0"), + units: Units::None, + } + .build(), info::Builder { id: id::PATH_SECRET_MAP_INITIALIZED, name: Str::new("path_secret_map_initialized\0"), @@ -4090,7 +4125,7 @@ pub struct Subscriber { #[allow(dead_code)] nominal_counter_offsets: Box<[usize; 37usize]>, #[allow(dead_code)] - measures: Box<[R::Measure; 137usize]>, + measures: Box<[R::Measure; 139usize]>, #[allow(dead_code)] gauges: Box<[R::Gauge; 0usize]>, #[allow(dead_code)] @@ -4121,7 +4156,7 @@ impl Subscriber { let mut bool_counters = Vec::with_capacity(25usize); let mut nominal_counters = Vec::with_capacity(37usize); let mut nominal_counter_offsets = Vec::with_capacity(37usize); - let mut measures = Vec::with_capacity(137usize); + let mut measures = Vec::with_capacity(139usize); let mut gauges = Vec::with_capacity(0usize); let mut timers = Vec::with_capacity(28usize); let mut nominal_timers = Vec::with_capacity(0usize); @@ -4217,6 +4252,7 @@ impl Subscriber { counters.push(registry.register_counter(&INFO[id::CONNECTION_CLOSED])); counters.push(registry.register_counter(&INFO[id::ENDPOINT_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::DC_CONNECTION_TIMEOUT])); + counters.push(registry.register_counter(&INFO[id::OFFLOAD_RUNTIME_METRICS])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED])); counters.push( @@ -4997,6 +5033,11 @@ impl Subscriber { registry.register_measure(&INFO[id::STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN]), ); measures.push(registry.register_measure(&INFO[id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN])); + measures.push( + registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH]), + ); + measures + .push(registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS])); measures.push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY])); measures .push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED__CAPACITY])); @@ -5357,6 +5398,7 @@ impl Subscriber { id::COUNTERS_CONNECTION_CLOSED => (&INFO[id::CONNECTION_CLOSED], entry), id::COUNTERS_ENDPOINT_INITIALIZED => (&INFO[id::ENDPOINT_INITIALIZED], entry), id::COUNTERS_DC_CONNECTION_TIMEOUT => (&INFO[id::DC_CONNECTION_TIMEOUT], entry), + id::COUNTERS_OFFLOAD_RUNTIME_METRICS => (&INFO[id::OFFLOAD_RUNTIME_METRICS], entry), id::COUNTERS_PATH_SECRET_MAP_INITIALIZED => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED], entry) } @@ -6340,6 +6382,12 @@ impl Subscriber { id::MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN => { (&INFO[id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN], entry) } + id::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH => { + (&INFO[id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH], entry) + } + id::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS => { + (&INFO[id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS], entry) + } id::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY], entry) } @@ -8758,6 +8806,32 @@ impl event::Subscriber for Subscriber { let _ = meta; } #[inline] + fn on_offload_runtime_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeMetrics, + ) { + #[allow(unused_imports)] + use api::*; + self.count( + id::OFFLOAD_RUNTIME_METRICS, + id::COUNTERS_OFFLOAD_RUNTIME_METRICS, + 1usize, + ); + self.measure( + id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, + id::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, + event.global_queue_depth, + ); + self.measure( + id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, + id::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, + event.num_alive_tasks, + ); + let _ = event; + let _ = meta; + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs index 8aa34e4d81..0b937c4a34 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs @@ -244,6 +244,9 @@ mod id { ENDPOINT_INITIALIZED__UDP, DC_CONNECTION_TIMEOUT, DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL, + OFFLOAD_RUNTIME_METRICS, + OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, + OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -734,6 +737,11 @@ mod id { pub const DC_CONNECTION_TIMEOUT: usize = InfoId::DC_CONNECTION_TIMEOUT as usize; pub const DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL: usize = InfoId::DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL as usize; + pub const OFFLOAD_RUNTIME_METRICS: usize = InfoId::OFFLOAD_RUNTIME_METRICS as usize; + pub const OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH: usize = + InfoId::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; + pub const OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = + InfoId::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1054,6 +1062,7 @@ mod counter { id::CONNECTION_CLOSED => Self(connection_closed), id::ENDPOINT_INITIALIZED => Self(endpoint_initialized), id::DC_CONNECTION_TIMEOUT => Self(dc_connection_timeout), + id::OFFLOAD_RUNTIME_METRICS => Self(offload_runtime_metrics), id::PATH_SECRET_MAP_INITIALIZED => Self(path_secret_map_initialized), id::PATH_SECRET_MAP_UNINITIALIZED => Self(path_secret_map_uninitialized), id::PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED => { @@ -1344,6 +1353,9 @@ mod counter { s2n_quic_dc__event__counter__dc_connection_timeout] fn dc_connection_timeout(value: u64); #[link_name = + s2n_quic_dc__event__counter__offload_runtime_metrics] + fn offload_runtime_metrics(value: u64); + #[link_name = s2n_quic_dc__event__counter__path_secret_map_initialized] fn path_secret_map_initialized(value: u64); #[link_name = @@ -2220,6 +2232,12 @@ mod measure { id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN => { Self(stream_handshake_packet_rejected__conn) } + id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH => { + Self(offload_runtime_metrics__global_queue_depth) + } + id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS => { + Self(offload_runtime_metrics__num_alive_tasks) + } id::PATH_SECRET_MAP_INITIALIZED__CAPACITY => { Self(path_secret_map_initialized__capacity) } @@ -2639,6 +2657,12 @@ mod measure { s2n_quic_dc__event__measure__stream_handshake_packet_rejected__conn] fn stream_handshake_packet_rejected__conn(value: u64); #[link_name = + s2n_quic_dc__event__measure__offload_runtime_metrics__global_queue_depth] + fn offload_runtime_metrics__global_queue_depth(value: u64); + #[link_name = + s2n_quic_dc__event__measure__offload_runtime_metrics__num_alive_tasks] + fn offload_runtime_metrics__num_alive_tasks(value: u64); + #[link_name = s2n_quic_dc__event__measure__path_secret_map_initialized__capacity] fn path_secret_map_initialized__capacity(value: u64); #[link_name = diff --git a/dc/s2n-quic-dc/src/path/secret/map.rs b/dc/s2n-quic-dc/src/path/secret/map.rs index 7388de5cca..9cad0a0792 100644 --- a/dc/s2n-quic-dc/src/path/secret/map.rs +++ b/dc/s2n-quic-dc/src/path/secret/map.rs @@ -3,7 +3,7 @@ use crate::{ credentials::{Credentials, Id}, - event, + event::{self}, packet::{secret_control as control, Packet}, path::secret::{ open, @@ -16,7 +16,7 @@ use crate::{ use core::fmt; use s2n_quic_core::{dc, time, varint::VarInt}; use std::{net::SocketAddr, sync::Arc}; -use tokio::task::JoinHandle; +use tokio::{runtime::RuntimeMetrics, task::JoinHandle}; mod cleaner; mod disk; @@ -326,6 +326,11 @@ impl Map { self.store.on_dc_connection_timeout(peer_address); } + /// Emits a OffloadRuntimeMetrics event via the subscriber + pub fn on_offload_runtime_metrics(&self, metrics: &RuntimeMetrics) { + self.store.on_offload_runtime_metrics(metrics); + } + /// Emits a datagram encrypt event with the wire packet length pub(crate) fn on_datagram_encrypt(&self, packet_len: usize) { self.store.on_datagram_encrypt(packet_len); diff --git a/dc/s2n-quic-dc/src/path/secret/map/state.rs b/dc/s2n-quic-dc/src/path/secret/map/state.rs index c33644f308..aae77afba0 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/state.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/state.rs @@ -25,7 +25,7 @@ use std::{ sync::{Arc, Mutex, RwLock, Weak}, time::Duration, }; -use tokio::task::JoinHandle; +use tokio::{runtime::RuntimeMetrics, task::JoinHandle}; #[cfg(test)] mod tests; @@ -1352,6 +1352,14 @@ where event::builder::PathSecretMapDatagramDecrypt { packet_len }, ); } + + fn on_offload_runtime_metrics(&self, metrics: &RuntimeMetrics) { + self.subscriber() + .on_offload_runtime_metrics(event::builder::OffloadRuntimeMetrics { + global_queue_depth: metrics.global_queue_depth(), + num_alive_tasks: metrics.num_alive_tasks(), + }); + } } impl Drop for State diff --git a/dc/s2n-quic-dc/src/path/secret/map/store.rs b/dc/s2n-quic-dc/src/path/secret/map/store.rs index 8783efde80..4b558d2d41 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/store.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/store.rs @@ -12,7 +12,7 @@ use core::time::Duration; use s2n_codec::EncoderBuffer; use s2n_quic_core::varint::VarInt; use std::{net::SocketAddr, sync::Arc}; -use tokio::task::JoinHandle; +use tokio::{runtime::RuntimeMetrics, task::JoinHandle}; pub trait Store: 'static + Send + Sync { fn secrets_len(&self) -> usize; @@ -163,6 +163,8 @@ pub trait Store: 'static + Send + Sync { fn on_dc_connection_timeout(&self, peer_address: &SocketAddr); + fn on_offload_runtime_metrics(&self, metrics: &RuntimeMetrics); + fn on_datagram_encrypt(&self, packet_len: usize); fn on_datagram_decrypt(&self, packet_len: usize); diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index c3c33cebdb..ea84afd08c 100644 --- a/dc/s2n-quic-dc/src/psk/io.rs +++ b/dc/s2n-quic-dc/src/psk/io.rs @@ -100,6 +100,7 @@ impl s2n_quic::provider::tls::offload::ExporterHandler for DCExporter { pub struct Server { server: s2n_quic::Server, + offload_handle: Option, } impl Server { @@ -159,7 +160,7 @@ impl Server { }}; } - let server = if builder.thread_offload_count > 0 { + let (server, handle) = if builder.thread_offload_count > 0 { let runtime = tokio::runtime::Builder::new_multi_thread() // Hs=handshake, s=server, offload .thread_name("hs-s-offload") @@ -167,6 +168,8 @@ impl Server { .enable_all() .build()?; + let handle = runtime.handle().clone(); + let tls = s2n_quic::provider::tls::offload::OffloadBuilder::new() .with_endpoint(tls_materials_provider) .with_exporter(DCExporter { @@ -183,12 +186,18 @@ impl Server { let connection_limits = connection_limits.with_packet_buffer_size(DEFAULT_MTU as u32)?; - build_and_start!(tls, connection_limits) + (build_and_start!(tls, connection_limits), Some(handle)) } else { - build_and_start!(tls_materials_provider, connection_limits) + ( + build_and_start!(tls_materials_provider, connection_limits), + None, + ) }; - Ok(Self { server }) + Ok(Self { + server, + offload_handle: handle, + }) } #[allow(dead_code)] @@ -238,6 +247,9 @@ pub(super) async fn server< while let Some(mut connection) = server.server.accept().await { let map_clone = map.clone(); + if let Some(handle) = &server.offload_handle { + map_clone.on_offload_runtime_metrics(&handle.metrics()); + } tokio::spawn(async move { // The accepted connection must remain open until the client has finished inserting // the entry into its map. The client indicates this by sending a ConnectionClose diff --git a/quic/s2n-quic-bench/src/bin/overload-server/mod.rs b/quic/s2n-quic-bench/src/bin/overload-server/mod.rs index d53e70e133..81b3b156c0 100644 --- a/quic/s2n-quic-bench/src/bin/overload-server/mod.rs +++ b/quic/s2n-quic-bench/src/bin/overload-server/mod.rs @@ -35,7 +35,7 @@ fn new_map(capacity: usize) -> secret::Map { capacity, false, StdClock::default(), - s2n_quic_dc::event::disabled::Subscriber::default(), + QueueSubscriber, ) } @@ -93,6 +93,27 @@ mod mtls { } } +struct QueueSubscriber; +impl s2n_quic_dc::event::Subscriber for QueueSubscriber { + type ConnectionContext = (); + + fn create_connection_context( + &self, + _meta: &s2n_quic_dc::event::api::ConnectionMeta, + _info: &s2n_quic_dc::event::api::ConnectionInfo, + ) -> Self::ConnectionContext { + () + } + + fn on_offload_runtime_metrics( + &self, + meta: &s2n_quic_dc::event::api::EndpointMeta, + event: &s2n_quic_dc::event::api::OffloadRuntimeMetrics, + ) { + println!("{:?}", event); + } +} + pub fn main() -> Result<(), Box> { tracing_subscriber::registry() .with(tracing_subscriber::EnvFilter::from_env("S2N_LOG")) @@ -106,6 +127,7 @@ pub fn main() -> Result<(), Box> { let sub = s2n_quic::provider::event::disabled::Subscriber; let server = server::Provider::builder() + .with_thread_count(3) .start_blocking( "127.0.0.1:0".parse().unwrap(), mtls::build_server_mtls_provider(certificates::MTLS_CA_CERT)?, From 4d4e03e6c1477fb423ee32db524d1036b69db8d4 Mon Sep 17 00:00:00 2001 From: Appelmans Date: Wed, 5 Aug 2026 15:04:28 -0700 Subject: [PATCH 2/9] Adds per worker metrics to offload feature --- dc/s2n-quic-dc/events/endpoint.rs | 9 ++ dc/s2n-quic-dc/src/event/generated.rs | 127 ++++++++++++++++++ .../src/event/generated/metrics/aggregate.rs | 85 +++++++++++- .../src/event/generated/metrics/probe.rs | 25 ++++ dc/s2n-quic-dc/src/path/secret/map/state.rs | 9 ++ dc/s2n-quic-dc/src/psk/io.rs | 18 ++- 6 files changed, 266 insertions(+), 7 deletions(-) diff --git a/dc/s2n-quic-dc/events/endpoint.rs b/dc/s2n-quic-dc/events/endpoint.rs index 347da1cea6..2b49286bfc 100644 --- a/dc/s2n-quic-dc/events/endpoint.rs +++ b/dc/s2n-quic-dc/events/endpoint.rs @@ -30,3 +30,12 @@ struct OffloadRuntimeMetrics { #[measure("num_alive_tasks", Count)] num_alive_tasks: usize, } + +#[event("runtime:offload_worker_metrics")] +#[subject(endpoint)] +struct OffloadRuntimeWorkerMetrics { + #[measure("park_count", Count)] + park_count: u64, + #[measure("busy_duration", Duration)] + busy_duration: core::time::Duration, +} diff --git a/dc/s2n-quic-dc/src/event/generated.rs b/dc/s2n-quic-dc/src/event/generated.rs index 7a76f953c6..4a5c40f8da 100644 --- a/dc/s2n-quic-dc/src/event/generated.rs +++ b/dc/s2n-quic-dc/src/event/generated.rs @@ -1910,6 +1910,24 @@ pub mod api { } #[derive(Clone, Debug)] #[non_exhaustive] + pub struct OffloadRuntimeWorkerMetrics { + pub park_count: u64, + pub busy_duration: core::time::Duration, + } + #[cfg(any(test, feature = "testing"))] + impl crate::event::snapshot::Fmt for OffloadRuntimeWorkerMetrics { + fn fmt(&self, fmt: &mut core::fmt::Formatter) -> core::fmt::Result { + let mut fmt = fmt.debug_struct("OffloadRuntimeWorkerMetrics"); + fmt.field("park_count", &self.park_count); + fmt.field("busy_duration", &self.busy_duration); + fmt.finish() + } + } + impl Event for OffloadRuntimeWorkerMetrics { + const NAME: &'static str = "runtime:offload_worker_metrics"; + } + #[derive(Clone, Debug)] + #[non_exhaustive] pub struct PathSecretMapInitialized { /// The capacity of the path secret map pub capacity: usize, @@ -4048,6 +4066,23 @@ pub mod tracing { ); } #[inline] + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + let parent = self.parent(meta); + let api::OffloadRuntimeWorkerMetrics { + park_count, + busy_duration, + } = event; + tracing::event!( + target : "offload_runtime_worker_metrics", parent : parent, + tracing::Level::DEBUG, { park_count = tracing::field::debug(park_count), + busy_duration = tracing::field::debug(busy_duration) } + ); + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -6472,6 +6507,24 @@ pub mod builder { } } #[derive(Clone, Debug)] + pub struct OffloadRuntimeWorkerMetrics { + pub park_count: u64, + pub busy_duration: core::time::Duration, + } + impl IntoEvent for OffloadRuntimeWorkerMetrics { + #[inline] + fn into_event(self) -> api::OffloadRuntimeWorkerMetrics { + let OffloadRuntimeWorkerMetrics { + park_count, + busy_duration, + } = self; + api::OffloadRuntimeWorkerMetrics { + park_count: park_count.into_event(), + busy_duration: busy_duration.into_event(), + } + } + } + #[derive(Clone, Debug)] pub struct PathSecretMapInitialized { /// The capacity of the path secret map pub capacity: usize, @@ -8068,6 +8121,16 @@ mod traits { let _ = meta; let _ = event; } + ///Called when the `OffloadRuntimeWorkerMetrics` event is triggered + #[inline] + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + let _ = meta; + let _ = event; + } ///Called when the `PathSecretMapInitialized` event is triggered #[inline] fn on_path_secret_map_initialized( @@ -9028,6 +9091,14 @@ mod traits { self.as_ref().on_offload_runtime_metrics(meta, event); } #[inline] + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + self.as_ref().on_offload_runtime_worker_metrics(meta, event); + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -9969,6 +10040,15 @@ mod traits { (self.1).on_offload_runtime_metrics(meta, event); } #[inline] + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + (self.0).on_offload_runtime_worker_metrics(meta, event); + (self.1).on_offload_runtime_worker_metrics(meta, event); + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -10388,6 +10468,8 @@ mod traits { fn on_dc_connection_timeout(&self, event: builder::DcConnectionTimeout); ///Publishes a `OffloadRuntimeMetrics` event to the publisher's subscriber fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics); + ///Publishes a `OffloadRuntimeWorkerMetrics` event to the publisher's subscriber + fn on_offload_runtime_worker_metrics(&self, event: builder::OffloadRuntimeWorkerMetrics); ///Publishes a `PathSecretMapInitialized` event to the publisher's subscriber fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized); ///Publishes a `PathSecretMapUninitialized` event to the publisher's subscriber @@ -10754,6 +10836,13 @@ mod traits { self.subscriber.on_event(&self.meta, &event); } #[inline] + fn on_offload_runtime_worker_metrics(&self, event: builder::OffloadRuntimeWorkerMetrics) { + let event = event.into_event(); + self.subscriber + .on_offload_runtime_worker_metrics(&self.meta, &event); + self.subscriber.on_event(&self.meta, &event); + } + #[inline] fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized) { let event = event.into_event(); self.subscriber @@ -11521,6 +11610,7 @@ pub mod testing { pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, pub offload_runtime_metrics: AtomicU64, + pub offload_runtime_worker_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -11620,6 +11710,7 @@ pub mod testing { endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), offload_runtime_metrics: AtomicU64::new(0), + offload_runtime_worker_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -12035,6 +12126,18 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + self.offload_runtime_worker_metrics + .fetch_add(1, Ordering::Relaxed); + let meta = crate::event::snapshot::Fmt::to_snapshot(meta); + let event = crate::event::snapshot::Fmt::to_snapshot(event); + let out = format!("{meta:?} {event:?}"); + self.output.lock().unwrap().push(out); + } fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -12532,6 +12635,7 @@ pub mod testing { pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, pub offload_runtime_metrics: AtomicU64, + pub offload_runtime_worker_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -12664,6 +12768,7 @@ pub mod testing { endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), offload_runtime_metrics: AtomicU64::new(0), + offload_runtime_worker_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -13548,6 +13653,18 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + self.offload_runtime_worker_metrics + .fetch_add(1, Ordering::Relaxed); + let meta = crate::event::snapshot::Fmt::to_snapshot(meta); + let event = crate::event::snapshot::Fmt::to_snapshot(event); + let out = format!("{meta:?} {event:?}"); + self.output.lock().unwrap().push(out); + } fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -14044,6 +14161,7 @@ pub mod testing { pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, pub offload_runtime_metrics: AtomicU64, + pub offload_runtime_worker_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -14166,6 +14284,7 @@ pub mod testing { endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), offload_runtime_metrics: AtomicU64::new(0), + offload_runtime_worker_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -14465,6 +14584,14 @@ pub mod testing { let out = format!("{event:?}"); self.output.lock().unwrap().push(out); } + fn on_offload_runtime_worker_metrics(&self, event: builder::OffloadRuntimeWorkerMetrics) { + self.offload_runtime_worker_metrics + .fetch_add(1, Ordering::Relaxed); + let event = event.into_event(); + let event = crate::event::snapshot::Fmt::to_snapshot(&event); + let out = format!("{event:?}"); + self.output.lock().unwrap().push(out); + } fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized) { self.path_secret_map_initialized .fetch_add(1, Ordering::Relaxed); diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs index 129ed49ea7..4bd24c53f9 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs @@ -251,6 +251,9 @@ mod id { OFFLOAD_RUNTIME_METRICS, OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, + OFFLOAD_RUNTIME_WORKER_METRICS, + OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, + OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -746,6 +749,12 @@ mod id { InfoId::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; pub const OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = InfoId::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; + pub const OFFLOAD_RUNTIME_WORKER_METRICS: usize = + InfoId::OFFLOAD_RUNTIME_WORKER_METRICS as usize; + pub const OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT: usize = + InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT as usize; + pub const OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION: usize = + InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1031,6 +1040,7 @@ mod id { COUNTERS_ENDPOINT_INITIALIZED, COUNTERS_DC_CONNECTION_TIMEOUT, COUNTERS_OFFLOAD_RUNTIME_METRICS, + COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS, COUNTERS_PATH_SECRET_MAP_INITIALIZED, COUNTERS_PATH_SECRET_MAP_UNINITIALIZED, COUNTERS_PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED, @@ -1215,6 +1225,8 @@ mod id { Counters::COUNTERS_DC_CONNECTION_TIMEOUT as usize; pub const COUNTERS_OFFLOAD_RUNTIME_METRICS: usize = Counters::COUNTERS_OFFLOAD_RUNTIME_METRICS as usize; + pub const COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS: usize = + Counters::COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS as usize; pub const COUNTERS_PATH_SECRET_MAP_INITIALIZED: usize = Counters::COUNTERS_PATH_SECRET_MAP_INITIALIZED as usize; pub const COUNTERS_PATH_SECRET_MAP_UNINITIALIZED: usize = @@ -1599,6 +1611,8 @@ mod id { MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN, MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, + MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, + MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__ENTRIES, @@ -1836,6 +1850,10 @@ mod id { Measures::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; pub const MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = Measures::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; + pub const MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT: usize = + Measures::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT as usize; + pub const MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION: usize = + Measures::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION as usize; pub const MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = Measures::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; pub const MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY: usize = @@ -3419,6 +3437,24 @@ static INFO: &[Info; 341usize] = &[ units: Units::None, } .build(), + info::Builder { + id: id::OFFLOAD_RUNTIME_WORKER_METRICS, + name: Str::new("offload_runtime_worker_metrics\0"), + units: Units::None, + } + .build(), + info::Builder { + id: id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, + name: Str::new("offload_runtime_worker_metrics.park_count\0"), + units: Units::None, + } + .build(), + info::Builder { + id: id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, + name: Str::new("offload_runtime_worker_metrics.busy_duration\0"), + units: Units::Duration, + } + .build(), info::Builder { id: id::PATH_SECRET_MAP_INITIALIZED, name: Str::new("path_secret_map_initialized\0"), @@ -4117,7 +4153,7 @@ pub struct ConnectionContext { } pub struct Subscriber { #[allow(dead_code)] - counters: Box<[R::Counter; 114usize]>, + counters: Box<[R::Counter; 115usize]>, #[allow(dead_code)] bool_counters: Box<[R::BoolCounter; 25usize]>, #[allow(dead_code)] @@ -4125,7 +4161,7 @@ pub struct Subscriber { #[allow(dead_code)] nominal_counter_offsets: Box<[usize; 37usize]>, #[allow(dead_code)] - measures: Box<[R::Measure; 139usize]>, + measures: Box<[R::Measure; 141usize]>, #[allow(dead_code)] gauges: Box<[R::Gauge; 0usize]>, #[allow(dead_code)] @@ -4152,11 +4188,11 @@ impl Subscriber { #[allow(unused_mut)] #[inline] pub fn new(registry: R) -> Self { - let mut counters = Vec::with_capacity(114usize); + let mut counters = Vec::with_capacity(115usize); let mut bool_counters = Vec::with_capacity(25usize); let mut nominal_counters = Vec::with_capacity(37usize); let mut nominal_counter_offsets = Vec::with_capacity(37usize); - let mut measures = Vec::with_capacity(139usize); + let mut measures = Vec::with_capacity(141usize); let mut gauges = Vec::with_capacity(0usize); let mut timers = Vec::with_capacity(28usize); let mut nominal_timers = Vec::with_capacity(0usize); @@ -4253,6 +4289,7 @@ impl Subscriber { counters.push(registry.register_counter(&INFO[id::ENDPOINT_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::DC_CONNECTION_TIMEOUT])); counters.push(registry.register_counter(&INFO[id::OFFLOAD_RUNTIME_METRICS])); + counters.push(registry.register_counter(&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED])); counters.push( @@ -5038,6 +5075,11 @@ impl Subscriber { ); measures .push(registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS])); + measures + .push(registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT])); + measures.push( + registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION]), + ); measures.push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY])); measures .push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED__CAPACITY])); @@ -5399,6 +5441,9 @@ impl Subscriber { id::COUNTERS_ENDPOINT_INITIALIZED => (&INFO[id::ENDPOINT_INITIALIZED], entry), id::COUNTERS_DC_CONNECTION_TIMEOUT => (&INFO[id::DC_CONNECTION_TIMEOUT], entry), id::COUNTERS_OFFLOAD_RUNTIME_METRICS => (&INFO[id::OFFLOAD_RUNTIME_METRICS], entry), + id::COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS => { + (&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS], entry) + } id::COUNTERS_PATH_SECRET_MAP_INITIALIZED => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED], entry) } @@ -6388,6 +6433,12 @@ impl Subscriber { id::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS => { (&INFO[id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS], entry) } + id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT => { + (&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT], entry) + } + id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION => { + (&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION], entry) + } id::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY], entry) } @@ -8832,6 +8883,32 @@ impl event::Subscriber for Subscriber { let _ = meta; } #[inline] + fn on_offload_runtime_worker_metrics( + &self, + meta: &api::EndpointMeta, + event: &api::OffloadRuntimeWorkerMetrics, + ) { + #[allow(unused_imports)] + use api::*; + self.count( + id::OFFLOAD_RUNTIME_WORKER_METRICS, + id::COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS, + 1usize, + ); + self.measure( + id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, + id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, + event.park_count, + ); + self.measure( + id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, + id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, + event.busy_duration, + ); + let _ = event; + let _ = meta; + } + #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs index 0b937c4a34..91f3df4700 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs @@ -247,6 +247,9 @@ mod id { OFFLOAD_RUNTIME_METRICS, OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, + OFFLOAD_RUNTIME_WORKER_METRICS, + OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, + OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -742,6 +745,12 @@ mod id { InfoId::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; pub const OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = InfoId::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; + pub const OFFLOAD_RUNTIME_WORKER_METRICS: usize = + InfoId::OFFLOAD_RUNTIME_WORKER_METRICS as usize; + pub const OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT: usize = + InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT as usize; + pub const OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION: usize = + InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1063,6 +1072,7 @@ mod counter { id::ENDPOINT_INITIALIZED => Self(endpoint_initialized), id::DC_CONNECTION_TIMEOUT => Self(dc_connection_timeout), id::OFFLOAD_RUNTIME_METRICS => Self(offload_runtime_metrics), + id::OFFLOAD_RUNTIME_WORKER_METRICS => Self(offload_runtime_worker_metrics), id::PATH_SECRET_MAP_INITIALIZED => Self(path_secret_map_initialized), id::PATH_SECRET_MAP_UNINITIALIZED => Self(path_secret_map_uninitialized), id::PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED => { @@ -1356,6 +1366,9 @@ mod counter { s2n_quic_dc__event__counter__offload_runtime_metrics] fn offload_runtime_metrics(value: u64); #[link_name = + s2n_quic_dc__event__counter__offload_runtime_worker_metrics] + fn offload_runtime_worker_metrics(value: u64); + #[link_name = s2n_quic_dc__event__counter__path_secret_map_initialized] fn path_secret_map_initialized(value: u64); #[link_name = @@ -2238,6 +2251,12 @@ mod measure { id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS => { Self(offload_runtime_metrics__num_alive_tasks) } + id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT => { + Self(offload_runtime_worker_metrics__park_count) + } + id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION => { + Self(offload_runtime_worker_metrics__busy_duration) + } id::PATH_SECRET_MAP_INITIALIZED__CAPACITY => { Self(path_secret_map_initialized__capacity) } @@ -2663,6 +2682,12 @@ mod measure { s2n_quic_dc__event__measure__offload_runtime_metrics__num_alive_tasks] fn offload_runtime_metrics__num_alive_tasks(value: u64); #[link_name = + s2n_quic_dc__event__measure__offload_runtime_worker_metrics__park_count] + fn offload_runtime_worker_metrics__park_count(value: u64); + #[link_name = + s2n_quic_dc__event__measure__offload_runtime_worker_metrics__busy_duration] + fn offload_runtime_worker_metrics__busy_duration(value: u64); + #[link_name = s2n_quic_dc__event__measure__path_secret_map_initialized__capacity] fn path_secret_map_initialized__capacity(value: u64); #[link_name = diff --git a/dc/s2n-quic-dc/src/path/secret/map/state.rs b/dc/s2n-quic-dc/src/path/secret/map/state.rs index aae77afba0..7c8b8a74da 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/state.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/state.rs @@ -1359,6 +1359,15 @@ where global_queue_depth: metrics.global_queue_depth(), num_alive_tasks: metrics.num_alive_tasks(), }); + + for idx in 0..metrics.num_workers() { + self.subscriber().on_offload_runtime_worker_metrics( + event::builder::OffloadRuntimeWorkerMetrics { + park_count: metrics.worker_park_count(idx), + busy_duration: metrics.worker_total_busy_duration(idx), + }, + ); + } } } diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index ea84afd08c..7c4d88cbc2 100644 --- a/dc/s2n-quic-dc/src/psk/io.rs +++ b/dc/s2n-quic-dc/src/psk/io.rs @@ -245,11 +245,23 @@ pub(super) async fn server< } } + let map_clone2 = map.clone(); + if let Some(handle) = &server.offload_handle { + let h = handle.clone(); + tokio::spawn(async move { + let mut interval = tokio::time::interval(Duration::from_secs(1)); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + interval.tick().await; + let metrics = h.metrics(); + map_clone2.on_offload_runtime_metrics(&metrics); + } + }); + } + while let Some(mut connection) = server.server.accept().await { let map_clone = map.clone(); - if let Some(handle) = &server.offload_handle { - map_clone.on_offload_runtime_metrics(&handle.metrics()); - } + tokio::spawn(async move { // The accepted connection must remain open until the client has finished inserting // the entry into its map. The client indicates this by sending a ConnectionClose From 1fdc06f0e07fbcf5b064b3973337b16789e31858 Mon Sep 17 00:00:00 2001 From: Appelmans Date: Wed, 12 Aug 2026 16:22:50 -0700 Subject: [PATCH 3/9] Saving snapshot of instrument_runtime function --- .cargo/config.toml | 2 +- dc/s2n-quic-dc/src/psk/io.rs | 23 ++++++++++++++--------- 2 files changed, 15 insertions(+), 10 deletions(-) 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/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index 7c4d88cbc2..dba359941e 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::{Registry, TaskMonitor}; use std::{ any::Any, hash::BuildHasher, @@ -100,7 +101,7 @@ impl s2n_quic::provider::tls::offload::ExporterHandler for DCExporter { pub struct Server { server: s2n_quic::Server, - offload_handle: Option, + offload_metrics: Option, } impl Server { @@ -160,7 +161,7 @@ impl Server { }}; } - let (server, handle) = if builder.thread_offload_count > 0 { + let (server, registry) = if builder.thread_offload_count > 0 { let runtime = tokio::runtime::Builder::new_multi_thread() // Hs=handshake, s=server, offload .thread_name("hs-s-offload") @@ -169,7 +170,9 @@ impl Server { .build()?; let handle = runtime.handle().clone(); - + let registry = Registry::new(); + let interval_duration = Duration::new(1, 0); + registry.instrument_runtime("offload-runtime", &handle, interval_duration); let tls = s2n_quic::provider::tls::offload::OffloadBuilder::new() .with_endpoint(tls_materials_provider) .with_exporter(DCExporter { @@ -186,7 +189,7 @@ impl Server { let connection_limits = connection_limits.with_packet_buffer_size(DEFAULT_MTU as u32)?; - (build_and_start!(tls, connection_limits), Some(handle)) + (build_and_start!(tls, connection_limits), Some(registry)) } else { ( build_and_start!(tls_materials_provider, connection_limits), @@ -196,7 +199,7 @@ impl Server { Ok(Self { server, - offload_handle: handle, + offload_metrics: registry, }) } @@ -246,15 +249,17 @@ pub(super) async fn server< } let map_clone2 = map.clone(); - if let Some(handle) = &server.offload_handle { - let h = handle.clone(); + if let Some(registry) = server.offload_metrics.take() { + //let h = handle.clone(); tokio::spawn(async move { let mut interval = tokio::time::interval(Duration::from_secs(1)); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { interval.tick().await; - let metrics = h.metrics(); - map_clone2.on_offload_runtime_metrics(&metrics); + //let metrics = h.metrics(); + let metrics = registry.take_current_metrics_line(); + println!("metrics {:?}", metrics); + //map_clone2.on_offload_runtime_metrics(&metrics); } }); } From 4571bbb38cb24696e8372a481fa87c4fcd275eac Mon Sep 17 00:00:00 2001 From: Appelmans Date: Wed, 12 Aug 2026 21:45:21 -0700 Subject: [PATCH 4/9] snapshot of tokio-metrics iteration --- dc/s2n-quic-dc/Cargo.toml | 1 + dc/s2n-quic-dc/src/psk/io.rs | 44 +++++++++++++++++++++--------------- 2 files changed, 27 insertions(+), 18 deletions(-) diff --git a/dc/s2n-quic-dc/Cargo.toml b/dc/s2n-quic-dc/Cargo.toml index 9bb222fe5a..041fb3b1de 100644 --- a/dc/s2n-quic-dc/Cargo.toml +++ b/dc/s2n-quic-dc/Cargo.toml @@ -63,6 +63,7 @@ zeroize = "1" parking_lot = "0.12" bitvec = { version = "1.1.1", default-features = false } nix = {version = "0.31.1", features = ["socket", "uio", "fs", "time"] } +tokio-metrics = "0.5.1" [dev-dependencies] bach = { version = "0.1.0", features = ["net", "tokio-compat"] } diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index dba359941e..5c851e518c 100644 --- a/dc/s2n-quic-dc/src/psk/io.rs +++ b/dc/s2n-quic-dc/src/psk/io.rs @@ -13,7 +13,6 @@ use s2n_quic::{ server::Name, }; use s2n_quic_core::{endpoint::Type, inet::SocketAddress}; -use s2n_quic_dc_metrics::{Registry, TaskMonitor}; use std::{ any::Any, hash::BuildHasher, @@ -26,6 +25,7 @@ use std::{ time::Duration, }; use tokio::{runtime::Runtime, sync::Semaphore, time::Instant as TokioInstant}; +use tokio_metrics::TaskMonitor; pub use crate::stream::DEFAULT_IDLE_TIMEOUT; pub const DEFAULT_MAX_DATA: u64 = 1u64 << 25; @@ -49,10 +49,11 @@ pub type Result = core::result::Result; struct TokioExecutor { runtime: Runtime, + monitor: TaskMonitor, } impl s2n_quic::provider::tls::offload::Executor for TokioExecutor { fn spawn(&self, task: impl core::future::Future + Send + 'static) { - self.runtime.spawn(task); + self.runtime.spawn(self.monitor.instrument(task)); } } #[derive(Clone)] @@ -101,7 +102,7 @@ impl s2n_quic::provider::tls::offload::ExporterHandler for DCExporter { pub struct Server { server: s2n_quic::Server, - offload_metrics: Option, + offload_metrics: Option, } impl Server { @@ -161,7 +162,7 @@ impl Server { }}; } - let (server, registry) = if builder.thread_offload_count > 0 { + let (server, monitor) = if builder.thread_offload_count > 0 { let runtime = tokio::runtime::Builder::new_multi_thread() // Hs=handshake, s=server, offload .thread_name("hs-s-offload") @@ -170,9 +171,8 @@ impl Server { .build()?; let handle = runtime.handle().clone(); - let registry = Registry::new(); - let interval_duration = Duration::new(1, 0); - registry.instrument_runtime("offload-runtime", &handle, interval_duration); + let monitor = tokio_metrics::TaskMonitor::new(); + let tls = s2n_quic::provider::tls::offload::OffloadBuilder::new() .with_endpoint(tls_materials_provider) .with_exporter(DCExporter { @@ -180,7 +180,10 @@ impl Server { endpoint_type: Type::Server, map: map.clone(), }) - .with_executor(TokioExecutor { runtime }) + .with_executor(TokioExecutor { + runtime, + monitor: monitor.clone(), + }) .build(); // We need packet storage when offloading is turned on due to this issue: @@ -189,7 +192,10 @@ impl Server { let connection_limits = connection_limits.with_packet_buffer_size(DEFAULT_MTU as u32)?; - (build_and_start!(tls, connection_limits), Some(registry)) + ( + build_and_start!(tls, connection_limits), + Some(monitor.clone()), + ) } else { ( build_and_start!(tls_materials_provider, connection_limits), @@ -199,7 +205,7 @@ impl Server { Ok(Self { server, - offload_metrics: registry, + offload_metrics: monitor, }) } @@ -249,17 +255,19 @@ pub(super) async fn server< } let map_clone2 = map.clone(); - if let Some(registry) = server.offload_metrics.take() { + if let Some(monitor) = server.offload_metrics.take() { //let h = handle.clone(); + let frequency = std::time::Duration::from_millis(500); tokio::spawn(async move { - let mut interval = tokio::time::interval(Duration::from_secs(1)); - interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { - interval.tick().await; - //let metrics = h.metrics(); - let metrics = registry.take_current_metrics_line(); - println!("metrics {:?}", metrics); - //map_clone2.on_offload_runtime_metrics(&metrics); + for metrics in monitor.intervals() { + println!("mean poll duration {:?}", metrics.mean_poll_duration()); + println!( + "mean scheduled duration {:?}", + metrics.mean_scheduled_duration() + ); + tokio::time::sleep(frequency).await; + } } }); } From 69b4a8212374c6075ad0bcbc9f8e1081e43af789 Mon Sep 17 00:00:00 2001 From: Appelmans Date: Thu, 13 Aug 2026 10:50:14 -0700 Subject: [PATCH 5/9] Adds new events for scheduled/poll mean times in offload feature --- dc/s2n-quic-dc/events/endpoint.rs | 21 +- dc/s2n-quic-dc/src/event/generated.rs | 242 +++++------------- .../src/event/generated/metrics/aggregate.rs | 171 ++++--------- .../src/event/generated/metrics/probe.rs | 66 ++--- dc/s2n-quic-dc/src/path/secret/map.rs | 9 +- dc/s2n-quic-dc/src/path/secret/map/state.rs | 20 +- dc/s2n-quic-dc/src/path/secret/map/store.rs | 5 +- dc/s2n-quic-dc/src/psk/io.rs | 17 +- 8 files changed, 147 insertions(+), 404 deletions(-) diff --git a/dc/s2n-quic-dc/events/endpoint.rs b/dc/s2n-quic-dc/events/endpoint.rs index 2b49286bfc..980473da18 100644 --- a/dc/s2n-quic-dc/events/endpoint.rs +++ b/dc/s2n-quic-dc/events/endpoint.rs @@ -22,20 +22,11 @@ struct DcConnectionTimeout<'a> { peer_address: SocketAddress<'a>, } -#[event("runtime:offload_metrics")] +#[event("offload:task_metrics")] #[subject(endpoint)] -struct OffloadRuntimeMetrics { - #[measure("global_queue_depth", Count)] - global_queue_depth: usize, - #[measure("num_alive_tasks", Count)] - num_alive_tasks: usize, -} - -#[event("runtime:offload_worker_metrics")] -#[subject(endpoint)] -struct OffloadRuntimeWorkerMetrics { - #[measure("park_count", Count)] - park_count: u64, - #[measure("busy_duration", Duration)] - busy_duration: core::time::Duration, +struct OffloadTaskMetrics { + #[measure("mean_poll_duration", Duration)] + mean_poll_duration: core::time::Duration, + #[measure("mean_scheduled_duration", Duration)] + mean_scheduled_duration: core::time::Duration, } diff --git a/dc/s2n-quic-dc/src/event/generated.rs b/dc/s2n-quic-dc/src/event/generated.rs index 4a5c40f8da..afa7401e4d 100644 --- a/dc/s2n-quic-dc/src/event/generated.rs +++ b/dc/s2n-quic-dc/src/event/generated.rs @@ -1892,39 +1892,21 @@ pub mod api { } #[derive(Clone, Debug)] #[non_exhaustive] - pub struct OffloadRuntimeMetrics { - pub global_queue_depth: usize, - pub num_alive_tasks: usize, + pub struct OffloadTaskMetrics { + pub mean_poll_duration: core::time::Duration, + pub mean_scheduled_duration: core::time::Duration, } #[cfg(any(test, feature = "testing"))] - impl crate::event::snapshot::Fmt for OffloadRuntimeMetrics { + impl crate::event::snapshot::Fmt for OffloadTaskMetrics { fn fmt(&self, fmt: &mut core::fmt::Formatter) -> core::fmt::Result { - let mut fmt = fmt.debug_struct("OffloadRuntimeMetrics"); - fmt.field("global_queue_depth", &self.global_queue_depth); - fmt.field("num_alive_tasks", &self.num_alive_tasks); + let mut fmt = fmt.debug_struct("OffloadTaskMetrics"); + fmt.field("mean_poll_duration", &self.mean_poll_duration); + fmt.field("mean_scheduled_duration", &self.mean_scheduled_duration); fmt.finish() } } - impl Event for OffloadRuntimeMetrics { - const NAME: &'static str = "runtime:offload_metrics"; - } - #[derive(Clone, Debug)] - #[non_exhaustive] - pub struct OffloadRuntimeWorkerMetrics { - pub park_count: u64, - pub busy_duration: core::time::Duration, - } - #[cfg(any(test, feature = "testing"))] - impl crate::event::snapshot::Fmt for OffloadRuntimeWorkerMetrics { - fn fmt(&self, fmt: &mut core::fmt::Formatter) -> core::fmt::Result { - let mut fmt = fmt.debug_struct("OffloadRuntimeWorkerMetrics"); - fmt.field("park_count", &self.park_count); - fmt.field("busy_duration", &self.busy_duration); - fmt.finish() - } - } - impl Event for OffloadRuntimeWorkerMetrics { - const NAME: &'static str = "runtime:offload_worker_metrics"; + impl Event for OffloadTaskMetrics { + const NAME: &'static str = "offload:task_metrics"; } #[derive(Clone, Debug)] #[non_exhaustive] @@ -4048,38 +4030,21 @@ pub mod tracing { ); } #[inline] - fn on_offload_runtime_metrics( + fn on_offload_task_metrics( &self, meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, + event: &api::OffloadTaskMetrics, ) { let parent = self.parent(meta); - let api::OffloadRuntimeMetrics { - global_queue_depth, - num_alive_tasks, + let api::OffloadTaskMetrics { + mean_poll_duration, + mean_scheduled_duration, } = event; tracing::event!( - target : "offload_runtime_metrics", parent : parent, - tracing::Level::DEBUG, { global_queue_depth = - tracing::field::debug(global_queue_depth), num_alive_tasks = - tracing::field::debug(num_alive_tasks) } - ); - } - #[inline] - fn on_offload_runtime_worker_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, - ) { - let parent = self.parent(meta); - let api::OffloadRuntimeWorkerMetrics { - park_count, - busy_duration, - } = event; - tracing::event!( - target : "offload_runtime_worker_metrics", parent : parent, - tracing::Level::DEBUG, { park_count = tracing::field::debug(park_count), - busy_duration = tracing::field::debug(busy_duration) } + target : "offload_task_metrics", parent : parent, tracing::Level::DEBUG, + { mean_poll_duration = tracing::field::debug(mean_poll_duration), + mean_scheduled_duration = tracing::field::debug(mean_scheduled_duration) + } ); } #[inline] @@ -6489,38 +6454,20 @@ pub mod builder { } } #[derive(Clone, Debug)] - pub struct OffloadRuntimeMetrics { - pub global_queue_depth: usize, - pub num_alive_tasks: usize, + pub struct OffloadTaskMetrics { + pub mean_poll_duration: core::time::Duration, + pub mean_scheduled_duration: core::time::Duration, } - impl IntoEvent for OffloadRuntimeMetrics { + impl IntoEvent for OffloadTaskMetrics { #[inline] - fn into_event(self) -> api::OffloadRuntimeMetrics { - let OffloadRuntimeMetrics { - global_queue_depth, - num_alive_tasks, + fn into_event(self) -> api::OffloadTaskMetrics { + let OffloadTaskMetrics { + mean_poll_duration, + mean_scheduled_duration, } = self; - api::OffloadRuntimeMetrics { - global_queue_depth: global_queue_depth.into_event(), - num_alive_tasks: num_alive_tasks.into_event(), - } - } - } - #[derive(Clone, Debug)] - pub struct OffloadRuntimeWorkerMetrics { - pub park_count: u64, - pub busy_duration: core::time::Duration, - } - impl IntoEvent for OffloadRuntimeWorkerMetrics { - #[inline] - fn into_event(self) -> api::OffloadRuntimeWorkerMetrics { - let OffloadRuntimeWorkerMetrics { - park_count, - busy_duration, - } = self; - api::OffloadRuntimeWorkerMetrics { - park_count: park_count.into_event(), - busy_duration: busy_duration.into_event(), + api::OffloadTaskMetrics { + mean_poll_duration: mean_poll_duration.into_event(), + mean_scheduled_duration: mean_scheduled_duration.into_event(), } } } @@ -8111,22 +8058,12 @@ mod traits { let _ = meta; let _ = event; } - ///Called when the `OffloadRuntimeMetrics` event is triggered - #[inline] - fn on_offload_runtime_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, - ) { - let _ = meta; - let _ = event; - } - ///Called when the `OffloadRuntimeWorkerMetrics` event is triggered + ///Called when the `OffloadTaskMetrics` event is triggered #[inline] - fn on_offload_runtime_worker_metrics( + fn on_offload_task_metrics( &self, meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, + event: &api::OffloadTaskMetrics, ) { let _ = meta; let _ = event; @@ -9083,20 +9020,12 @@ mod traits { self.as_ref().on_dc_connection_timeout(meta, event); } #[inline] - fn on_offload_runtime_metrics( + fn on_offload_task_metrics( &self, meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, + event: &api::OffloadTaskMetrics, ) { - self.as_ref().on_offload_runtime_metrics(meta, event); - } - #[inline] - fn on_offload_runtime_worker_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, - ) { - self.as_ref().on_offload_runtime_worker_metrics(meta, event); + self.as_ref().on_offload_task_metrics(meta, event); } #[inline] fn on_path_secret_map_initialized( @@ -10031,22 +9960,13 @@ mod traits { (self.1).on_dc_connection_timeout(meta, event); } #[inline] - fn on_offload_runtime_metrics( + fn on_offload_task_metrics( &self, meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, + event: &api::OffloadTaskMetrics, ) { - (self.0).on_offload_runtime_metrics(meta, event); - (self.1).on_offload_runtime_metrics(meta, event); - } - #[inline] - fn on_offload_runtime_worker_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, - ) { - (self.0).on_offload_runtime_worker_metrics(meta, event); - (self.1).on_offload_runtime_worker_metrics(meta, event); + (self.0).on_offload_task_metrics(meta, event); + (self.1).on_offload_task_metrics(meta, event); } #[inline] fn on_path_secret_map_initialized( @@ -10466,10 +10386,8 @@ mod traits { fn on_endpoint_initialized(&self, event: builder::EndpointInitialized); ///Publishes a `DcConnectionTimeout` event to the publisher's subscriber fn on_dc_connection_timeout(&self, event: builder::DcConnectionTimeout); - ///Publishes a `OffloadRuntimeMetrics` event to the publisher's subscriber - fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics); - ///Publishes a `OffloadRuntimeWorkerMetrics` event to the publisher's subscriber - fn on_offload_runtime_worker_metrics(&self, event: builder::OffloadRuntimeWorkerMetrics); + ///Publishes a `OffloadTaskMetrics` event to the publisher's subscriber + fn on_offload_task_metrics(&self, event: builder::OffloadTaskMetrics); ///Publishes a `PathSecretMapInitialized` event to the publisher's subscriber fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized); ///Publishes a `PathSecretMapUninitialized` event to the publisher's subscriber @@ -10829,17 +10747,9 @@ mod traits { self.subscriber.on_event(&self.meta, &event); } #[inline] - fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics) { - let event = event.into_event(); - self.subscriber - .on_offload_runtime_metrics(&self.meta, &event); - self.subscriber.on_event(&self.meta, &event); - } - #[inline] - fn on_offload_runtime_worker_metrics(&self, event: builder::OffloadRuntimeWorkerMetrics) { + fn on_offload_task_metrics(&self, event: builder::OffloadTaskMetrics) { let event = event.into_event(); - self.subscriber - .on_offload_runtime_worker_metrics(&self.meta, &event); + self.subscriber.on_offload_task_metrics(&self.meta, &event); self.subscriber.on_event(&self.meta, &event); } #[inline] @@ -11609,8 +11519,7 @@ pub mod testing { pub stream_connect_error: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, - pub offload_runtime_metrics: AtomicU64, - pub offload_runtime_worker_metrics: AtomicU64, + pub offload_task_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -11709,8 +11618,7 @@ pub mod testing { stream_connect_error: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), - offload_runtime_metrics: AtomicU64::new(0), - offload_runtime_worker_metrics: AtomicU64::new(0), + offload_task_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -12115,24 +12023,12 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } - fn on_offload_runtime_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, - ) { - self.offload_runtime_metrics.fetch_add(1, Ordering::Relaxed); - let meta = crate::event::snapshot::Fmt::to_snapshot(meta); - let event = crate::event::snapshot::Fmt::to_snapshot(event); - let out = format!("{meta:?} {event:?}"); - self.output.lock().unwrap().push(out); - } - fn on_offload_runtime_worker_metrics( + fn on_offload_task_metrics( &self, meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, + event: &api::OffloadTaskMetrics, ) { - self.offload_runtime_worker_metrics - .fetch_add(1, Ordering::Relaxed); + self.offload_task_metrics.fetch_add(1, Ordering::Relaxed); let meta = crate::event::snapshot::Fmt::to_snapshot(meta); let event = crate::event::snapshot::Fmt::to_snapshot(event); let out = format!("{meta:?} {event:?}"); @@ -12634,8 +12530,7 @@ pub mod testing { pub connection_closed: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, - pub offload_runtime_metrics: AtomicU64, - pub offload_runtime_worker_metrics: AtomicU64, + pub offload_task_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -12767,8 +12662,7 @@ pub mod testing { connection_closed: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), - offload_runtime_metrics: AtomicU64::new(0), - offload_runtime_worker_metrics: AtomicU64::new(0), + offload_task_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -13642,24 +13536,12 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } - fn on_offload_runtime_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, - ) { - self.offload_runtime_metrics.fetch_add(1, Ordering::Relaxed); - let meta = crate::event::snapshot::Fmt::to_snapshot(meta); - let event = crate::event::snapshot::Fmt::to_snapshot(event); - let out = format!("{meta:?} {event:?}"); - self.output.lock().unwrap().push(out); - } - fn on_offload_runtime_worker_metrics( + fn on_offload_task_metrics( &self, meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, + event: &api::OffloadTaskMetrics, ) { - self.offload_runtime_worker_metrics - .fetch_add(1, Ordering::Relaxed); + self.offload_task_metrics.fetch_add(1, Ordering::Relaxed); let meta = crate::event::snapshot::Fmt::to_snapshot(meta); let event = crate::event::snapshot::Fmt::to_snapshot(event); let out = format!("{meta:?} {event:?}"); @@ -14160,8 +14042,7 @@ pub mod testing { pub connection_closed: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, - pub offload_runtime_metrics: AtomicU64, - pub offload_runtime_worker_metrics: AtomicU64, + pub offload_task_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -14283,8 +14164,7 @@ pub mod testing { connection_closed: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), - offload_runtime_metrics: AtomicU64::new(0), - offload_runtime_worker_metrics: AtomicU64::new(0), + offload_task_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -14577,16 +14457,8 @@ pub mod testing { let out = format!("{event:?}"); self.output.lock().unwrap().push(out); } - fn on_offload_runtime_metrics(&self, event: builder::OffloadRuntimeMetrics) { - self.offload_runtime_metrics.fetch_add(1, Ordering::Relaxed); - let event = event.into_event(); - let event = crate::event::snapshot::Fmt::to_snapshot(&event); - let out = format!("{event:?}"); - self.output.lock().unwrap().push(out); - } - fn on_offload_runtime_worker_metrics(&self, event: builder::OffloadRuntimeWorkerMetrics) { - self.offload_runtime_worker_metrics - .fetch_add(1, Ordering::Relaxed); + fn on_offload_task_metrics(&self, event: builder::OffloadTaskMetrics) { + self.offload_task_metrics.fetch_add(1, Ordering::Relaxed); let event = event.into_event(); let event = crate::event::snapshot::Fmt::to_snapshot(&event); let out = format!("{event:?}"); diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs index 4bd24c53f9..f2de4e388d 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs @@ -248,12 +248,9 @@ mod id { ENDPOINT_INITIALIZED__UDP, DC_CONNECTION_TIMEOUT, DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL, - OFFLOAD_RUNTIME_METRICS, - OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, - OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, - OFFLOAD_RUNTIME_WORKER_METRICS, - OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, - OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, + OFFLOAD_TASK_METRICS, + OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, + OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -744,17 +741,11 @@ mod id { pub const DC_CONNECTION_TIMEOUT: usize = InfoId::DC_CONNECTION_TIMEOUT as usize; pub const DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL: usize = InfoId::DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL as usize; - pub const OFFLOAD_RUNTIME_METRICS: usize = InfoId::OFFLOAD_RUNTIME_METRICS as usize; - pub const OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH: usize = - InfoId::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; - pub const OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = - InfoId::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; - pub const OFFLOAD_RUNTIME_WORKER_METRICS: usize = - InfoId::OFFLOAD_RUNTIME_WORKER_METRICS as usize; - pub const OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT: usize = - InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT as usize; - pub const OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION: usize = - InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION as usize; + pub const OFFLOAD_TASK_METRICS: usize = InfoId::OFFLOAD_TASK_METRICS as usize; + pub const OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION: usize = + InfoId::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION as usize; + pub const OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION: usize = + InfoId::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1039,8 +1030,7 @@ mod id { COUNTERS_CONNECTION_CLOSED, COUNTERS_ENDPOINT_INITIALIZED, COUNTERS_DC_CONNECTION_TIMEOUT, - COUNTERS_OFFLOAD_RUNTIME_METRICS, - COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS, + COUNTERS_OFFLOAD_TASK_METRICS, COUNTERS_PATH_SECRET_MAP_INITIALIZED, COUNTERS_PATH_SECRET_MAP_UNINITIALIZED, COUNTERS_PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED, @@ -1223,10 +1213,8 @@ mod id { Counters::COUNTERS_ENDPOINT_INITIALIZED as usize; pub const COUNTERS_DC_CONNECTION_TIMEOUT: usize = Counters::COUNTERS_DC_CONNECTION_TIMEOUT as usize; - pub const COUNTERS_OFFLOAD_RUNTIME_METRICS: usize = - Counters::COUNTERS_OFFLOAD_RUNTIME_METRICS as usize; - pub const COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS: usize = - Counters::COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS as usize; + pub const COUNTERS_OFFLOAD_TASK_METRICS: usize = + Counters::COUNTERS_OFFLOAD_TASK_METRICS as usize; pub const COUNTERS_PATH_SECRET_MAP_INITIALIZED: usize = Counters::COUNTERS_PATH_SECRET_MAP_INITIALIZED as usize; pub const COUNTERS_PATH_SECRET_MAP_UNINITIALIZED: usize = @@ -1609,10 +1597,8 @@ mod id { MEASURES_STREAM_CONTROL_PACKET_RECEIVED__PACKET_LEN, MEASURES_STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN, MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN, - MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, - MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, - MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, - MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, + MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, + MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__ENTRIES, @@ -1846,14 +1832,10 @@ mod id { Measures::MEASURES_STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN as usize; pub const MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN: usize = Measures::MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN as usize; - pub const MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH: usize = - Measures::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; - pub const MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = - Measures::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; - pub const MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT: usize = - Measures::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT as usize; - pub const MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION: usize = - Measures::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION as usize; + pub const MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION: usize = + Measures::MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION as usize; + pub const MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION: usize = + Measures::MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION as usize; pub const MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = Measures::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; pub const MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY: usize = @@ -3420,38 +3402,20 @@ static INFO: &[Info; 341usize] = &[ } .build(), info::Builder { - id: id::OFFLOAD_RUNTIME_METRICS, - name: Str::new("offload_runtime_metrics\0"), + id: id::OFFLOAD_TASK_METRICS, + name: Str::new("offload_task_metrics\0"), units: Units::None, } .build(), info::Builder { - id: id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, - name: Str::new("offload_runtime_metrics.global_queue_depth\0"), - units: Units::None, - } - .build(), - info::Builder { - id: id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, - name: Str::new("offload_runtime_metrics.num_alive_tasks\0"), - units: Units::None, - } - .build(), - info::Builder { - id: id::OFFLOAD_RUNTIME_WORKER_METRICS, - name: Str::new("offload_runtime_worker_metrics\0"), - units: Units::None, - } - .build(), - info::Builder { - id: id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, - name: Str::new("offload_runtime_worker_metrics.park_count\0"), - units: Units::None, + id: id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, + name: Str::new("offload_task_metrics.mean_poll_duration\0"), + units: Units::Duration, } .build(), info::Builder { - id: id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, - name: Str::new("offload_runtime_worker_metrics.busy_duration\0"), + id: id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, + name: Str::new("offload_task_metrics.mean_scheduled_duration\0"), units: Units::Duration, } .build(), @@ -4153,7 +4117,7 @@ pub struct ConnectionContext { } pub struct Subscriber { #[allow(dead_code)] - counters: Box<[R::Counter; 115usize]>, + counters: Box<[R::Counter; 114usize]>, #[allow(dead_code)] bool_counters: Box<[R::BoolCounter; 25usize]>, #[allow(dead_code)] @@ -4161,7 +4125,7 @@ pub struct Subscriber { #[allow(dead_code)] nominal_counter_offsets: Box<[usize; 37usize]>, #[allow(dead_code)] - measures: Box<[R::Measure; 141usize]>, + measures: Box<[R::Measure; 139usize]>, #[allow(dead_code)] gauges: Box<[R::Gauge; 0usize]>, #[allow(dead_code)] @@ -4188,11 +4152,11 @@ impl Subscriber { #[allow(unused_mut)] #[inline] pub fn new(registry: R) -> Self { - let mut counters = Vec::with_capacity(115usize); + let mut counters = Vec::with_capacity(114usize); let mut bool_counters = Vec::with_capacity(25usize); let mut nominal_counters = Vec::with_capacity(37usize); let mut nominal_counter_offsets = Vec::with_capacity(37usize); - let mut measures = Vec::with_capacity(141usize); + let mut measures = Vec::with_capacity(139usize); let mut gauges = Vec::with_capacity(0usize); let mut timers = Vec::with_capacity(28usize); let mut nominal_timers = Vec::with_capacity(0usize); @@ -4288,8 +4252,7 @@ impl Subscriber { counters.push(registry.register_counter(&INFO[id::CONNECTION_CLOSED])); counters.push(registry.register_counter(&INFO[id::ENDPOINT_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::DC_CONNECTION_TIMEOUT])); - counters.push(registry.register_counter(&INFO[id::OFFLOAD_RUNTIME_METRICS])); - counters.push(registry.register_counter(&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS])); + counters.push(registry.register_counter(&INFO[id::OFFLOAD_TASK_METRICS])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED])); counters.push( @@ -5070,15 +5033,10 @@ impl Subscriber { registry.register_measure(&INFO[id::STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN]), ); measures.push(registry.register_measure(&INFO[id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN])); - measures.push( - registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH]), - ); - measures - .push(registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS])); measures - .push(registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT])); + .push(registry.register_measure(&INFO[id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION])); measures.push( - registry.register_measure(&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION]), + registry.register_measure(&INFO[id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION]), ); measures.push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY])); measures @@ -5440,10 +5398,7 @@ impl Subscriber { id::COUNTERS_CONNECTION_CLOSED => (&INFO[id::CONNECTION_CLOSED], entry), id::COUNTERS_ENDPOINT_INITIALIZED => (&INFO[id::ENDPOINT_INITIALIZED], entry), id::COUNTERS_DC_CONNECTION_TIMEOUT => (&INFO[id::DC_CONNECTION_TIMEOUT], entry), - id::COUNTERS_OFFLOAD_RUNTIME_METRICS => (&INFO[id::OFFLOAD_RUNTIME_METRICS], entry), - id::COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS => { - (&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS], entry) - } + id::COUNTERS_OFFLOAD_TASK_METRICS => (&INFO[id::OFFLOAD_TASK_METRICS], entry), id::COUNTERS_PATH_SECRET_MAP_INITIALIZED => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED], entry) } @@ -6427,17 +6382,11 @@ impl Subscriber { id::MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN => { (&INFO[id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN], entry) } - id::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH => { - (&INFO[id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH], entry) - } - id::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS => { - (&INFO[id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS], entry) - } - id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT => { - (&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT], entry) + id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION => { + (&INFO[id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION], entry) } - id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION => { - (&INFO[id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION], entry) + id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION => { + (&INFO[id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION], entry) } id::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY], entry) @@ -8857,53 +8806,23 @@ impl event::Subscriber for Subscriber { let _ = meta; } #[inline] - fn on_offload_runtime_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeMetrics, - ) { - #[allow(unused_imports)] - use api::*; - self.count( - id::OFFLOAD_RUNTIME_METRICS, - id::COUNTERS_OFFLOAD_RUNTIME_METRICS, - 1usize, - ); - self.measure( - id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, - id::MEASURES_OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, - event.global_queue_depth, - ); - self.measure( - id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, - id::MEASURES_OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, - event.num_alive_tasks, - ); - let _ = event; - let _ = meta; - } - #[inline] - fn on_offload_runtime_worker_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadRuntimeWorkerMetrics, - ) { + fn on_offload_task_metrics(&self, meta: &api::EndpointMeta, event: &api::OffloadTaskMetrics) { #[allow(unused_imports)] use api::*; self.count( - id::OFFLOAD_RUNTIME_WORKER_METRICS, - id::COUNTERS_OFFLOAD_RUNTIME_WORKER_METRICS, + id::OFFLOAD_TASK_METRICS, + id::COUNTERS_OFFLOAD_TASK_METRICS, 1usize, ); self.measure( - id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, - id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, - event.park_count, + id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, + id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, + event.mean_poll_duration, ); self.measure( - id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, - id::MEASURES_OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, - event.busy_duration, + id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, + id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, + event.mean_scheduled_duration, ); let _ = event; let _ = meta; diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs index 91f3df4700..8fbedf914a 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs @@ -244,12 +244,9 @@ mod id { ENDPOINT_INITIALIZED__UDP, DC_CONNECTION_TIMEOUT, DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL, - OFFLOAD_RUNTIME_METRICS, - OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH, - OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS, - OFFLOAD_RUNTIME_WORKER_METRICS, - OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT, - OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION, + OFFLOAD_TASK_METRICS, + OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, + OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -740,17 +737,11 @@ mod id { pub const DC_CONNECTION_TIMEOUT: usize = InfoId::DC_CONNECTION_TIMEOUT as usize; pub const DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL: usize = InfoId::DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL as usize; - pub const OFFLOAD_RUNTIME_METRICS: usize = InfoId::OFFLOAD_RUNTIME_METRICS as usize; - pub const OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH: usize = - InfoId::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH as usize; - pub const OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS: usize = - InfoId::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS as usize; - pub const OFFLOAD_RUNTIME_WORKER_METRICS: usize = - InfoId::OFFLOAD_RUNTIME_WORKER_METRICS as usize; - pub const OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT: usize = - InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT as usize; - pub const OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION: usize = - InfoId::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION as usize; + pub const OFFLOAD_TASK_METRICS: usize = InfoId::OFFLOAD_TASK_METRICS as usize; + pub const OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION: usize = + InfoId::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION as usize; + pub const OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION: usize = + InfoId::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1071,8 +1062,7 @@ mod counter { id::CONNECTION_CLOSED => Self(connection_closed), id::ENDPOINT_INITIALIZED => Self(endpoint_initialized), id::DC_CONNECTION_TIMEOUT => Self(dc_connection_timeout), - id::OFFLOAD_RUNTIME_METRICS => Self(offload_runtime_metrics), - id::OFFLOAD_RUNTIME_WORKER_METRICS => Self(offload_runtime_worker_metrics), + id::OFFLOAD_TASK_METRICS => Self(offload_task_metrics), id::PATH_SECRET_MAP_INITIALIZED => Self(path_secret_map_initialized), id::PATH_SECRET_MAP_UNINITIALIZED => Self(path_secret_map_uninitialized), id::PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED => { @@ -1363,13 +1353,9 @@ mod counter { s2n_quic_dc__event__counter__dc_connection_timeout] fn dc_connection_timeout(value: u64); #[link_name = - s2n_quic_dc__event__counter__offload_runtime_metrics] - fn offload_runtime_metrics(value: u64); - #[link_name = - s2n_quic_dc__event__counter__offload_runtime_worker_metrics] - fn offload_runtime_worker_metrics(value: u64); - #[link_name = - s2n_quic_dc__event__counter__path_secret_map_initialized] + s2n_quic_dc__event__counter__offload_task_metrics] + fn offload_task_metrics(value: u64); + #[link_name = s2n_quic_dc__event__counter__path_secret_map_initialized] fn path_secret_map_initialized(value: u64); #[link_name = s2n_quic_dc__event__counter__path_secret_map_uninitialized] @@ -2245,17 +2231,11 @@ mod measure { id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN => { Self(stream_handshake_packet_rejected__conn) } - id::OFFLOAD_RUNTIME_METRICS__GLOBAL_QUEUE_DEPTH => { - Self(offload_runtime_metrics__global_queue_depth) - } - id::OFFLOAD_RUNTIME_METRICS__NUM_ALIVE_TASKS => { - Self(offload_runtime_metrics__num_alive_tasks) + id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION => { + Self(offload_task_metrics__mean_poll_duration) } - id::OFFLOAD_RUNTIME_WORKER_METRICS__PARK_COUNT => { - Self(offload_runtime_worker_metrics__park_count) - } - id::OFFLOAD_RUNTIME_WORKER_METRICS__BUSY_DURATION => { - Self(offload_runtime_worker_metrics__busy_duration) + id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION => { + Self(offload_task_metrics__mean_scheduled_duration) } id::PATH_SECRET_MAP_INITIALIZED__CAPACITY => { Self(path_secret_map_initialized__capacity) @@ -2676,17 +2656,11 @@ mod measure { s2n_quic_dc__event__measure__stream_handshake_packet_rejected__conn] fn stream_handshake_packet_rejected__conn(value: u64); #[link_name = - s2n_quic_dc__event__measure__offload_runtime_metrics__global_queue_depth] - fn offload_runtime_metrics__global_queue_depth(value: u64); - #[link_name = - s2n_quic_dc__event__measure__offload_runtime_metrics__num_alive_tasks] - fn offload_runtime_metrics__num_alive_tasks(value: u64); - #[link_name = - s2n_quic_dc__event__measure__offload_runtime_worker_metrics__park_count] - fn offload_runtime_worker_metrics__park_count(value: u64); + s2n_quic_dc__event__measure__offload_task_metrics__mean_poll_duration] + fn offload_task_metrics__mean_poll_duration(value: u64); #[link_name = - s2n_quic_dc__event__measure__offload_runtime_worker_metrics__busy_duration] - fn offload_runtime_worker_metrics__busy_duration(value: u64); + s2n_quic_dc__event__measure__offload_task_metrics__mean_scheduled_duration] + fn offload_task_metrics__mean_scheduled_duration(value: u64); #[link_name = s2n_quic_dc__event__measure__path_secret_map_initialized__capacity] fn path_secret_map_initialized__capacity(value: u64); diff --git a/dc/s2n-quic-dc/src/path/secret/map.rs b/dc/s2n-quic-dc/src/path/secret/map.rs index 9cad0a0792..e1365266fb 100644 --- a/dc/s2n-quic-dc/src/path/secret/map.rs +++ b/dc/s2n-quic-dc/src/path/secret/map.rs @@ -16,7 +16,8 @@ use crate::{ use core::fmt; use s2n_quic_core::{dc, time, varint::VarInt}; use std::{net::SocketAddr, sync::Arc}; -use tokio::{runtime::RuntimeMetrics, task::JoinHandle}; +use tokio::task::JoinHandle; +use tokio_metrics::TaskMetrics; mod cleaner; mod disk; @@ -326,9 +327,9 @@ impl Map { self.store.on_dc_connection_timeout(peer_address); } - /// Emits a OffloadRuntimeMetrics event via the subscriber - pub fn on_offload_runtime_metrics(&self, metrics: &RuntimeMetrics) { - self.store.on_offload_runtime_metrics(metrics); + /// Emits a OffloadTaskMetrics event via the subscriber + pub fn on_offload_task_metrics(&self, metrics: &TaskMetrics) { + self.store.on_offload_task_metrics(metrics); } /// Emits a datagram encrypt event with the wire packet length diff --git a/dc/s2n-quic-dc/src/path/secret/map/state.rs b/dc/s2n-quic-dc/src/path/secret/map/state.rs index 7c8b8a74da..b1d8bf73e4 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/state.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/state.rs @@ -25,7 +25,8 @@ use std::{ sync::{Arc, Mutex, RwLock, Weak}, time::Duration, }; -use tokio::{runtime::RuntimeMetrics, task::JoinHandle}; +use tokio::task::JoinHandle; +use tokio_metrics::TaskMetrics; #[cfg(test)] mod tests; @@ -1353,21 +1354,12 @@ where ); } - fn on_offload_runtime_metrics(&self, metrics: &RuntimeMetrics) { + fn on_offload_task_metrics(&self, metrics: &TaskMetrics) { self.subscriber() - .on_offload_runtime_metrics(event::builder::OffloadRuntimeMetrics { - global_queue_depth: metrics.global_queue_depth(), - num_alive_tasks: metrics.num_alive_tasks(), + .on_offload_task_metrics(event::builder::OffloadTaskMetrics { + mean_poll_duration: metrics.mean_poll_duration(), + mean_scheduled_duration: metrics.mean_scheduled_duration(), }); - - for idx in 0..metrics.num_workers() { - self.subscriber().on_offload_runtime_worker_metrics( - event::builder::OffloadRuntimeWorkerMetrics { - park_count: metrics.worker_park_count(idx), - busy_duration: metrics.worker_total_busy_duration(idx), - }, - ); - } } } diff --git a/dc/s2n-quic-dc/src/path/secret/map/store.rs b/dc/s2n-quic-dc/src/path/secret/map/store.rs index 4b558d2d41..061d35fd4f 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/store.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/store.rs @@ -12,7 +12,8 @@ use core::time::Duration; use s2n_codec::EncoderBuffer; use s2n_quic_core::varint::VarInt; use std::{net::SocketAddr, sync::Arc}; -use tokio::{runtime::RuntimeMetrics, task::JoinHandle}; +use tokio::task::JoinHandle; +use tokio_metrics::TaskMetrics; pub trait Store: 'static + Send + Sync { fn secrets_len(&self) -> usize; @@ -163,7 +164,7 @@ pub trait Store: 'static + Send + Sync { fn on_dc_connection_timeout(&self, peer_address: &SocketAddr); - fn on_offload_runtime_metrics(&self, metrics: &RuntimeMetrics); + fn on_offload_task_metrics(&self, metrics: &TaskMetrics); fn on_datagram_encrypt(&self, packet_len: usize); diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index 5c851e518c..507889e2d1 100644 --- a/dc/s2n-quic-dc/src/psk/io.rs +++ b/dc/s2n-quic-dc/src/psk/io.rs @@ -170,7 +170,6 @@ impl Server { .enable_all() .build()?; - let handle = runtime.handle().clone(); let monitor = tokio_metrics::TaskMonitor::new(); let tls = s2n_quic::provider::tls::offload::OffloadBuilder::new() @@ -254,20 +253,14 @@ pub(super) async fn server< } } - let map_clone2 = map.clone(); + let map_clone = map.clone(); if let Some(monitor) = server.offload_metrics.take() { - //let h = handle.clone(); - let frequency = std::time::Duration::from_millis(500); + let frequency = std::time::Duration::new(1, 0); tokio::spawn(async move { loop { - for metrics in monitor.intervals() { - println!("mean poll duration {:?}", metrics.mean_poll_duration()); - println!( - "mean scheduled duration {:?}", - metrics.mean_scheduled_duration() - ); - tokio::time::sleep(frequency).await; - } + let metrics = monitor.cumulative(); + map_clone.on_offload_task_metrics(&metrics); + tokio::time::sleep(frequency).await; } }); } From 2c7a05f59fceda0abd9c45c61efcb4738882859e Mon Sep 17 00:00:00 2001 From: Appelmans Date: Thu, 13 Aug 2026 10:53:01 -0700 Subject: [PATCH 6/9] Remove changes to overload-server --- .../src/bin/overload-server/mod.rs | 24 +------------------ 1 file changed, 1 insertion(+), 23 deletions(-) diff --git a/quic/s2n-quic-bench/src/bin/overload-server/mod.rs b/quic/s2n-quic-bench/src/bin/overload-server/mod.rs index 81b3b156c0..d53e70e133 100644 --- a/quic/s2n-quic-bench/src/bin/overload-server/mod.rs +++ b/quic/s2n-quic-bench/src/bin/overload-server/mod.rs @@ -35,7 +35,7 @@ fn new_map(capacity: usize) -> secret::Map { capacity, false, StdClock::default(), - QueueSubscriber, + s2n_quic_dc::event::disabled::Subscriber::default(), ) } @@ -93,27 +93,6 @@ mod mtls { } } -struct QueueSubscriber; -impl s2n_quic_dc::event::Subscriber for QueueSubscriber { - type ConnectionContext = (); - - fn create_connection_context( - &self, - _meta: &s2n_quic_dc::event::api::ConnectionMeta, - _info: &s2n_quic_dc::event::api::ConnectionInfo, - ) -> Self::ConnectionContext { - () - } - - fn on_offload_runtime_metrics( - &self, - meta: &s2n_quic_dc::event::api::EndpointMeta, - event: &s2n_quic_dc::event::api::OffloadRuntimeMetrics, - ) { - println!("{:?}", event); - } -} - pub fn main() -> Result<(), Box> { tracing_subscriber::registry() .with(tracing_subscriber::EnvFilter::from_env("S2N_LOG")) @@ -127,7 +106,6 @@ pub fn main() -> Result<(), Box> { let sub = s2n_quic::provider::event::disabled::Subscriber; let server = server::Provider::builder() - .with_thread_count(3) .start_blocking( "127.0.0.1:0".parse().unwrap(), mtls::build_server_mtls_provider(certificates::MTLS_CA_CERT)?, From 6ce6bbdbf9c266ee0c65e1a2347da910b97f5d7a Mon Sep 17 00:00:00 2001 From: Appelmans Date: Thu, 13 Aug 2026 11:29:04 -0700 Subject: [PATCH 7/9] Fix event code merge --- dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs index f2de4e388d..f0793c85ad 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs @@ -2014,7 +2014,7 @@ mod id { pub const TIMERS_STREAM_CONNECT_ERROR__LATENCY: usize = Timers::TIMERS_STREAM_CONNECT_ERROR__LATENCY as usize; } -static INFO: &[Info; 341usize] = &[ +static INFO: &[Info; 344usize] = &[ info::Builder { id: id::ACCEPTOR_TCP_STARTED, name: Str::new("acceptor_tcp_started\0"), @@ -4117,7 +4117,7 @@ pub struct ConnectionContext { } pub struct Subscriber { #[allow(dead_code)] - counters: Box<[R::Counter; 114usize]>, + counters: Box<[R::Counter; 115usize]>, #[allow(dead_code)] bool_counters: Box<[R::BoolCounter; 25usize]>, #[allow(dead_code)] @@ -4152,7 +4152,7 @@ impl Subscriber { #[allow(unused_mut)] #[inline] pub fn new(registry: R) -> Self { - let mut counters = Vec::with_capacity(114usize); + let mut counters = Vec::with_capacity(115usize); let mut bool_counters = Vec::with_capacity(25usize); let mut nominal_counters = Vec::with_capacity(37usize); let mut nominal_counter_offsets = Vec::with_capacity(37usize); From 77611f94da6b7b89d2a3963d35e7da975d68d3c3 Mon Sep 17 00:00:00 2001 From: Appelmans Date: Fri, 14 Aug 2026 14:46:25 -0700 Subject: [PATCH 8/9] Changing to Registry approach --- dc/s2n-quic-dc/Cargo.toml | 2 +- dc/s2n-quic-dc/src/path/secret/map.rs | 6 -- dc/s2n-quic-dc/src/path/secret/map/state.rs | 9 --- dc/s2n-quic-dc/src/path/secret/map/store.rs | 3 - dc/s2n-quic-dc/src/psk/io.rs | 64 ++++++++++----------- dc/s2n-quic-dc/src/psk/server/builder.rs | 9 +++ 6 files changed, 42 insertions(+), 51 deletions(-) diff --git a/dc/s2n-quic-dc/Cargo.toml b/dc/s2n-quic-dc/Cargo.toml index 041fb3b1de..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" @@ -63,7 +64,6 @@ zeroize = "1" parking_lot = "0.12" bitvec = { version = "1.1.1", default-features = false } nix = {version = "0.31.1", features = ["socket", "uio", "fs", "time"] } -tokio-metrics = "0.5.1" [dev-dependencies] bach = { version = "0.1.0", features = ["net", "tokio-compat"] } diff --git a/dc/s2n-quic-dc/src/path/secret/map.rs b/dc/s2n-quic-dc/src/path/secret/map.rs index e1365266fb..4356134e70 100644 --- a/dc/s2n-quic-dc/src/path/secret/map.rs +++ b/dc/s2n-quic-dc/src/path/secret/map.rs @@ -17,7 +17,6 @@ use core::fmt; use s2n_quic_core::{dc, time, varint::VarInt}; use std::{net::SocketAddr, sync::Arc}; use tokio::task::JoinHandle; -use tokio_metrics::TaskMetrics; mod cleaner; mod disk; @@ -327,11 +326,6 @@ impl Map { self.store.on_dc_connection_timeout(peer_address); } - /// Emits a OffloadTaskMetrics event via the subscriber - pub fn on_offload_task_metrics(&self, metrics: &TaskMetrics) { - self.store.on_offload_task_metrics(metrics); - } - /// Emits a datagram encrypt event with the wire packet length pub(crate) fn on_datagram_encrypt(&self, packet_len: usize) { self.store.on_datagram_encrypt(packet_len); diff --git a/dc/s2n-quic-dc/src/path/secret/map/state.rs b/dc/s2n-quic-dc/src/path/secret/map/state.rs index b1d8bf73e4..c33644f308 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/state.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/state.rs @@ -26,7 +26,6 @@ use std::{ time::Duration, }; use tokio::task::JoinHandle; -use tokio_metrics::TaskMetrics; #[cfg(test)] mod tests; @@ -1353,14 +1352,6 @@ where event::builder::PathSecretMapDatagramDecrypt { packet_len }, ); } - - fn on_offload_task_metrics(&self, metrics: &TaskMetrics) { - self.subscriber() - .on_offload_task_metrics(event::builder::OffloadTaskMetrics { - mean_poll_duration: metrics.mean_poll_duration(), - mean_scheduled_duration: metrics.mean_scheduled_duration(), - }); - } } impl Drop for State diff --git a/dc/s2n-quic-dc/src/path/secret/map/store.rs b/dc/s2n-quic-dc/src/path/secret/map/store.rs index 061d35fd4f..8783efde80 100644 --- a/dc/s2n-quic-dc/src/path/secret/map/store.rs +++ b/dc/s2n-quic-dc/src/path/secret/map/store.rs @@ -13,7 +13,6 @@ use s2n_codec::EncoderBuffer; use s2n_quic_core::varint::VarInt; use std::{net::SocketAddr, sync::Arc}; use tokio::task::JoinHandle; -use tokio_metrics::TaskMetrics; pub trait Store: 'static + Send + Sync { fn secrets_len(&self) -> usize; @@ -164,8 +163,6 @@ pub trait Store: 'static + Send + Sync { fn on_dc_connection_timeout(&self, peer_address: &SocketAddr); - fn on_offload_task_metrics(&self, metrics: &TaskMetrics); - fn on_datagram_encrypt(&self, packet_len: usize); fn on_datagram_decrypt(&self, packet_len: usize); diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index 507889e2d1..f7198b1acc 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, @@ -25,7 +26,6 @@ use std::{ time::Duration, }; use tokio::{runtime::Runtime, sync::Semaphore, time::Instant as TokioInstant}; -use tokio_metrics::TaskMonitor; pub use crate::stream::DEFAULT_IDLE_TIMEOUT; pub const DEFAULT_MAX_DATA: u64 = 1u64 << 25; @@ -49,11 +49,15 @@ pub type Result = core::result::Result; struct TokioExecutor { runtime: Runtime, - monitor: TaskMonitor, + monitor: Option, } impl s2n_quic::provider::tls::offload::Executor for TokioExecutor { fn spawn(&self, task: impl core::future::Future + Send + 'static) { - self.runtime.spawn(self.monitor.instrument(task)); + if let Some(monitor) = &self.monitor { + self.runtime.spawn(monitor.instrument(task)); + } else { + self.runtime.spawn(task); + } } } #[derive(Clone)] @@ -102,7 +106,6 @@ impl s2n_quic::provider::tls::offload::ExporterHandler for DCExporter { pub struct Server { server: s2n_quic::Server, - offload_metrics: Option, } impl Server { @@ -162,7 +165,7 @@ impl Server { }}; } - let (server, monitor) = if builder.thread_offload_count > 0 { + let server = if builder.thread_offload_count > 0 { let runtime = tokio::runtime::Builder::new_multi_thread() // Hs=handshake, s=server, offload .thread_name("hs-s-offload") @@ -170,7 +173,16 @@ impl Server { .enable_all() .build()?; - let monitor = tokio_metrics::TaskMonitor::new(); + 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) @@ -179,10 +191,7 @@ impl Server { endpoint_type: Type::Server, map: map.clone(), }) - .with_executor(TokioExecutor { - runtime, - monitor: monitor.clone(), - }) + .with_executor(TokioExecutor { runtime, monitor }) .build(); // We need packet storage when offloading is turned on due to this issue: @@ -191,21 +200,12 @@ impl Server { let connection_limits = connection_limits.with_packet_buffer_size(DEFAULT_MTU as u32)?; - ( - build_and_start!(tls, connection_limits), - Some(monitor.clone()), - ) + build_and_start!(tls, connection_limits) } else { - ( - build_and_start!(tls_materials_provider, connection_limits), - None, - ) + build_and_start!(tls_materials_provider, connection_limits) }; - Ok(Self { - server, - offload_metrics: monitor, - }) + Ok(Self { server }) } #[allow(dead_code)] @@ -254,16 +254,16 @@ pub(super) async fn server< } let map_clone = map.clone(); - if let Some(monitor) = server.offload_metrics.take() { - let frequency = std::time::Duration::new(1, 0); - tokio::spawn(async move { - loop { - let metrics = monitor.cumulative(); - map_clone.on_offload_task_metrics(&metrics); - tokio::time::sleep(frequency).await; - } - }); - } + // if let Some(monitor) = server.offload_metrics.take() { + // let frequency = std::time::Duration::new(1, 0); + // tokio::spawn(async move { + // loop { + // let metrics = monitor.cumulative(); + // map_clone.on_offload_task_metrics(&metrics); + // tokio::time::sleep(frequency).await; + // } + // }); + // } while let Some(mut connection) = server.server.accept().await { let map_clone = map.clone(); 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, From 27a2a978e1d0751649b93d8cda05e615a79a68a1 Mon Sep 17 00:00:00 2001 From: Appelmans Date: Fri, 14 Aug 2026 14:51:43 -0700 Subject: [PATCH 9/9] Removes vestigial code --- dc/s2n-quic-dc/events/endpoint.rs | 9 -- dc/s2n-quic-dc/src/event/generated.rs | 124 ------------------ .../src/event/generated/metrics/aggregate.rs | 80 +---------- .../src/event/generated/metrics/probe.rs | 25 +--- dc/s2n-quic-dc/src/path/secret/map.rs | 2 +- dc/s2n-quic-dc/src/psk/io.rs | 13 -- 6 files changed, 7 insertions(+), 246 deletions(-) diff --git a/dc/s2n-quic-dc/events/endpoint.rs b/dc/s2n-quic-dc/events/endpoint.rs index 980473da18..c7757ff464 100644 --- a/dc/s2n-quic-dc/events/endpoint.rs +++ b/dc/s2n-quic-dc/events/endpoint.rs @@ -21,12 +21,3 @@ struct DcConnectionTimeout<'a> { #[nominal_counter("peer_address.protocol")] peer_address: SocketAddress<'a>, } - -#[event("offload:task_metrics")] -#[subject(endpoint)] -struct OffloadTaskMetrics { - #[measure("mean_poll_duration", Duration)] - mean_poll_duration: core::time::Duration, - #[measure("mean_scheduled_duration", Duration)] - mean_scheduled_duration: core::time::Duration, -} diff --git a/dc/s2n-quic-dc/src/event/generated.rs b/dc/s2n-quic-dc/src/event/generated.rs index afa7401e4d..06f83bad60 100644 --- a/dc/s2n-quic-dc/src/event/generated.rs +++ b/dc/s2n-quic-dc/src/event/generated.rs @@ -1892,24 +1892,6 @@ pub mod api { } #[derive(Clone, Debug)] #[non_exhaustive] - pub struct OffloadTaskMetrics { - pub mean_poll_duration: core::time::Duration, - pub mean_scheduled_duration: core::time::Duration, - } - #[cfg(any(test, feature = "testing"))] - impl crate::event::snapshot::Fmt for OffloadTaskMetrics { - fn fmt(&self, fmt: &mut core::fmt::Formatter) -> core::fmt::Result { - let mut fmt = fmt.debug_struct("OffloadTaskMetrics"); - fmt.field("mean_poll_duration", &self.mean_poll_duration); - fmt.field("mean_scheduled_duration", &self.mean_scheduled_duration); - fmt.finish() - } - } - impl Event for OffloadTaskMetrics { - const NAME: &'static str = "offload:task_metrics"; - } - #[derive(Clone, Debug)] - #[non_exhaustive] pub struct PathSecretMapInitialized { /// The capacity of the path secret map pub capacity: usize, @@ -4030,24 +4012,6 @@ pub mod tracing { ); } #[inline] - fn on_offload_task_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadTaskMetrics, - ) { - let parent = self.parent(meta); - let api::OffloadTaskMetrics { - mean_poll_duration, - mean_scheduled_duration, - } = event; - tracing::event!( - target : "offload_task_metrics", parent : parent, tracing::Level::DEBUG, - { mean_poll_duration = tracing::field::debug(mean_poll_duration), - mean_scheduled_duration = tracing::field::debug(mean_scheduled_duration) - } - ); - } - #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -6454,24 +6418,6 @@ pub mod builder { } } #[derive(Clone, Debug)] - pub struct OffloadTaskMetrics { - pub mean_poll_duration: core::time::Duration, - pub mean_scheduled_duration: core::time::Duration, - } - impl IntoEvent for OffloadTaskMetrics { - #[inline] - fn into_event(self) -> api::OffloadTaskMetrics { - let OffloadTaskMetrics { - mean_poll_duration, - mean_scheduled_duration, - } = self; - api::OffloadTaskMetrics { - mean_poll_duration: mean_poll_duration.into_event(), - mean_scheduled_duration: mean_scheduled_duration.into_event(), - } - } - } - #[derive(Clone, Debug)] pub struct PathSecretMapInitialized { /// The capacity of the path secret map pub capacity: usize, @@ -8058,16 +8004,6 @@ mod traits { let _ = meta; let _ = event; } - ///Called when the `OffloadTaskMetrics` event is triggered - #[inline] - fn on_offload_task_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadTaskMetrics, - ) { - let _ = meta; - let _ = event; - } ///Called when the `PathSecretMapInitialized` event is triggered #[inline] fn on_path_secret_map_initialized( @@ -9020,14 +8956,6 @@ mod traits { self.as_ref().on_dc_connection_timeout(meta, event); } #[inline] - fn on_offload_task_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadTaskMetrics, - ) { - self.as_ref().on_offload_task_metrics(meta, event); - } - #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -9960,15 +9888,6 @@ mod traits { (self.1).on_dc_connection_timeout(meta, event); } #[inline] - fn on_offload_task_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadTaskMetrics, - ) { - (self.0).on_offload_task_metrics(meta, event); - (self.1).on_offload_task_metrics(meta, event); - } - #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -10386,8 +10305,6 @@ mod traits { fn on_endpoint_initialized(&self, event: builder::EndpointInitialized); ///Publishes a `DcConnectionTimeout` event to the publisher's subscriber fn on_dc_connection_timeout(&self, event: builder::DcConnectionTimeout); - ///Publishes a `OffloadTaskMetrics` event to the publisher's subscriber - fn on_offload_task_metrics(&self, event: builder::OffloadTaskMetrics); ///Publishes a `PathSecretMapInitialized` event to the publisher's subscriber fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized); ///Publishes a `PathSecretMapUninitialized` event to the publisher's subscriber @@ -10747,12 +10664,6 @@ mod traits { self.subscriber.on_event(&self.meta, &event); } #[inline] - fn on_offload_task_metrics(&self, event: builder::OffloadTaskMetrics) { - let event = event.into_event(); - self.subscriber.on_offload_task_metrics(&self.meta, &event); - self.subscriber.on_event(&self.meta, &event); - } - #[inline] fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized) { let event = event.into_event(); self.subscriber @@ -11519,7 +11430,6 @@ pub mod testing { pub stream_connect_error: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, - pub offload_task_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -11618,7 +11528,6 @@ pub mod testing { stream_connect_error: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), - offload_task_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -12023,17 +11932,6 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } - fn on_offload_task_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadTaskMetrics, - ) { - self.offload_task_metrics.fetch_add(1, Ordering::Relaxed); - let meta = crate::event::snapshot::Fmt::to_snapshot(meta); - let event = crate::event::snapshot::Fmt::to_snapshot(event); - let out = format!("{meta:?} {event:?}"); - self.output.lock().unwrap().push(out); - } fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -12530,7 +12428,6 @@ pub mod testing { pub connection_closed: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, - pub offload_task_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -12662,7 +12559,6 @@ pub mod testing { connection_closed: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), - offload_task_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -13536,17 +13432,6 @@ pub mod testing { let out = format!("{meta:?} {event:?}"); self.output.lock().unwrap().push(out); } - fn on_offload_task_metrics( - &self, - meta: &api::EndpointMeta, - event: &api::OffloadTaskMetrics, - ) { - self.offload_task_metrics.fetch_add(1, Ordering::Relaxed); - let meta = crate::event::snapshot::Fmt::to_snapshot(meta); - let event = crate::event::snapshot::Fmt::to_snapshot(event); - let out = format!("{meta:?} {event:?}"); - self.output.lock().unwrap().push(out); - } fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, @@ -14042,7 +13927,6 @@ pub mod testing { pub connection_closed: AtomicU64, pub endpoint_initialized: AtomicU64, pub dc_connection_timeout: AtomicU64, - pub offload_task_metrics: AtomicU64, pub path_secret_map_initialized: AtomicU64, pub path_secret_map_uninitialized: AtomicU64, pub path_secret_map_background_handshake_requested: AtomicU64, @@ -14164,7 +14048,6 @@ pub mod testing { connection_closed: AtomicU64::new(0), endpoint_initialized: AtomicU64::new(0), dc_connection_timeout: AtomicU64::new(0), - offload_task_metrics: AtomicU64::new(0), path_secret_map_initialized: AtomicU64::new(0), path_secret_map_uninitialized: AtomicU64::new(0), path_secret_map_background_handshake_requested: AtomicU64::new(0), @@ -14457,13 +14340,6 @@ pub mod testing { let out = format!("{event:?}"); self.output.lock().unwrap().push(out); } - fn on_offload_task_metrics(&self, event: builder::OffloadTaskMetrics) { - self.offload_task_metrics.fetch_add(1, Ordering::Relaxed); - let event = event.into_event(); - let event = crate::event::snapshot::Fmt::to_snapshot(&event); - let out = format!("{event:?}"); - self.output.lock().unwrap().push(out); - } fn on_path_secret_map_initialized(&self, event: builder::PathSecretMapInitialized) { self.path_secret_map_initialized .fetch_add(1, Ordering::Relaxed); diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs index f0793c85ad..5756228fb0 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/aggregate.rs @@ -248,9 +248,6 @@ mod id { ENDPOINT_INITIALIZED__UDP, DC_CONNECTION_TIMEOUT, DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL, - OFFLOAD_TASK_METRICS, - OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, - OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -741,11 +738,6 @@ mod id { pub const DC_CONNECTION_TIMEOUT: usize = InfoId::DC_CONNECTION_TIMEOUT as usize; pub const DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL: usize = InfoId::DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL as usize; - pub const OFFLOAD_TASK_METRICS: usize = InfoId::OFFLOAD_TASK_METRICS as usize; - pub const OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION: usize = - InfoId::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION as usize; - pub const OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION: usize = - InfoId::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1030,7 +1022,6 @@ mod id { COUNTERS_CONNECTION_CLOSED, COUNTERS_ENDPOINT_INITIALIZED, COUNTERS_DC_CONNECTION_TIMEOUT, - COUNTERS_OFFLOAD_TASK_METRICS, COUNTERS_PATH_SECRET_MAP_INITIALIZED, COUNTERS_PATH_SECRET_MAP_UNINITIALIZED, COUNTERS_PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED, @@ -1213,8 +1204,6 @@ mod id { Counters::COUNTERS_ENDPOINT_INITIALIZED as usize; pub const COUNTERS_DC_CONNECTION_TIMEOUT: usize = Counters::COUNTERS_DC_CONNECTION_TIMEOUT as usize; - pub const COUNTERS_OFFLOAD_TASK_METRICS: usize = - Counters::COUNTERS_OFFLOAD_TASK_METRICS as usize; pub const COUNTERS_PATH_SECRET_MAP_INITIALIZED: usize = Counters::COUNTERS_PATH_SECRET_MAP_INITIALIZED as usize; pub const COUNTERS_PATH_SECRET_MAP_UNINITIALIZED: usize = @@ -1597,8 +1586,6 @@ mod id { MEASURES_STREAM_CONTROL_PACKET_RECEIVED__PACKET_LEN, MEASURES_STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN, MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN, - MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, - MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY, MEASURES_PATH_SECRET_MAP_UNINITIALIZED__ENTRIES, @@ -1832,10 +1819,6 @@ mod id { Measures::MEASURES_STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN as usize; pub const MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN: usize = Measures::MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN as usize; - pub const MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION: usize = - Measures::MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION as usize; - pub const MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION: usize = - Measures::MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION as usize; pub const MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = Measures::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; pub const MEASURES_PATH_SECRET_MAP_UNINITIALIZED__CAPACITY: usize = @@ -2014,7 +1997,7 @@ mod id { pub const TIMERS_STREAM_CONNECT_ERROR__LATENCY: usize = Timers::TIMERS_STREAM_CONNECT_ERROR__LATENCY as usize; } -static INFO: &[Info; 344usize] = &[ +static INFO: &[Info; 341usize] = &[ info::Builder { id: id::ACCEPTOR_TCP_STARTED, name: Str::new("acceptor_tcp_started\0"), @@ -3401,24 +3384,6 @@ static INFO: &[Info; 344usize] = &[ units: Units::None, } .build(), - info::Builder { - id: id::OFFLOAD_TASK_METRICS, - name: Str::new("offload_task_metrics\0"), - units: Units::None, - } - .build(), - info::Builder { - id: id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, - name: Str::new("offload_task_metrics.mean_poll_duration\0"), - units: Units::Duration, - } - .build(), - info::Builder { - id: id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, - name: Str::new("offload_task_metrics.mean_scheduled_duration\0"), - units: Units::Duration, - } - .build(), info::Builder { id: id::PATH_SECRET_MAP_INITIALIZED, name: Str::new("path_secret_map_initialized\0"), @@ -4117,7 +4082,7 @@ pub struct ConnectionContext { } pub struct Subscriber { #[allow(dead_code)] - counters: Box<[R::Counter; 115usize]>, + counters: Box<[R::Counter; 114usize]>, #[allow(dead_code)] bool_counters: Box<[R::BoolCounter; 25usize]>, #[allow(dead_code)] @@ -4125,7 +4090,7 @@ pub struct Subscriber { #[allow(dead_code)] nominal_counter_offsets: Box<[usize; 37usize]>, #[allow(dead_code)] - measures: Box<[R::Measure; 139usize]>, + measures: Box<[R::Measure; 137usize]>, #[allow(dead_code)] gauges: Box<[R::Gauge; 0usize]>, #[allow(dead_code)] @@ -4152,11 +4117,11 @@ impl Subscriber { #[allow(unused_mut)] #[inline] pub fn new(registry: R) -> Self { - let mut counters = Vec::with_capacity(115usize); + let mut counters = Vec::with_capacity(114usize); let mut bool_counters = Vec::with_capacity(25usize); let mut nominal_counters = Vec::with_capacity(37usize); let mut nominal_counter_offsets = Vec::with_capacity(37usize); - let mut measures = Vec::with_capacity(139usize); + let mut measures = Vec::with_capacity(137usize); let mut gauges = Vec::with_capacity(0usize); let mut timers = Vec::with_capacity(28usize); let mut nominal_timers = Vec::with_capacity(0usize); @@ -4252,7 +4217,6 @@ impl Subscriber { counters.push(registry.register_counter(&INFO[id::CONNECTION_CLOSED])); counters.push(registry.register_counter(&INFO[id::ENDPOINT_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::DC_CONNECTION_TIMEOUT])); - counters.push(registry.register_counter(&INFO[id::OFFLOAD_TASK_METRICS])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_INITIALIZED])); counters.push(registry.register_counter(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED])); counters.push( @@ -5033,11 +4997,6 @@ impl Subscriber { registry.register_measure(&INFO[id::STREAM_CONTROL_PACKET_RECEIVED__CONTROL_DATA_LEN]), ); measures.push(registry.register_measure(&INFO[id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN])); - measures - .push(registry.register_measure(&INFO[id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION])); - measures.push( - registry.register_measure(&INFO[id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION]), - ); measures.push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY])); measures .push(registry.register_measure(&INFO[id::PATH_SECRET_MAP_UNINITIALIZED__CAPACITY])); @@ -5398,7 +5357,6 @@ impl Subscriber { id::COUNTERS_CONNECTION_CLOSED => (&INFO[id::CONNECTION_CLOSED], entry), id::COUNTERS_ENDPOINT_INITIALIZED => (&INFO[id::ENDPOINT_INITIALIZED], entry), id::COUNTERS_DC_CONNECTION_TIMEOUT => (&INFO[id::DC_CONNECTION_TIMEOUT], entry), - id::COUNTERS_OFFLOAD_TASK_METRICS => (&INFO[id::OFFLOAD_TASK_METRICS], entry), id::COUNTERS_PATH_SECRET_MAP_INITIALIZED => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED], entry) } @@ -6382,12 +6340,6 @@ impl Subscriber { id::MEASURES_STREAM_HANDSHAKE_PACKET_REJECTED__CONN => { (&INFO[id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN], entry) } - id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION => { - (&INFO[id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION], entry) - } - id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION => { - (&INFO[id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION], entry) - } id::MEASURES_PATH_SECRET_MAP_INITIALIZED__CAPACITY => { (&INFO[id::PATH_SECRET_MAP_INITIALIZED__CAPACITY], entry) } @@ -8806,28 +8758,6 @@ impl event::Subscriber for Subscriber { let _ = meta; } #[inline] - fn on_offload_task_metrics(&self, meta: &api::EndpointMeta, event: &api::OffloadTaskMetrics) { - #[allow(unused_imports)] - use api::*; - self.count( - id::OFFLOAD_TASK_METRICS, - id::COUNTERS_OFFLOAD_TASK_METRICS, - 1usize, - ); - self.measure( - id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, - id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, - event.mean_poll_duration, - ); - self.measure( - id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, - id::MEASURES_OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, - event.mean_scheduled_duration, - ); - let _ = event; - let _ = meta; - } - #[inline] fn on_path_secret_map_initialized( &self, meta: &api::EndpointMeta, diff --git a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs index 8fbedf914a..8aa34e4d81 100644 --- a/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs +++ b/dc/s2n-quic-dc/src/event/generated/metrics/probe.rs @@ -244,9 +244,6 @@ mod id { ENDPOINT_INITIALIZED__UDP, DC_CONNECTION_TIMEOUT, DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL, - OFFLOAD_TASK_METRICS, - OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION, - OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION, PATH_SECRET_MAP_INITIALIZED, PATH_SECRET_MAP_INITIALIZED__CAPACITY, PATH_SECRET_MAP_UNINITIALIZED, @@ -737,11 +734,6 @@ mod id { pub const DC_CONNECTION_TIMEOUT: usize = InfoId::DC_CONNECTION_TIMEOUT as usize; pub const DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL: usize = InfoId::DC_CONNECTION_TIMEOUT__PEER_ADDRESS__PROTOCOL as usize; - pub const OFFLOAD_TASK_METRICS: usize = InfoId::OFFLOAD_TASK_METRICS as usize; - pub const OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION: usize = - InfoId::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION as usize; - pub const OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION: usize = - InfoId::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION as usize; pub const PATH_SECRET_MAP_INITIALIZED: usize = InfoId::PATH_SECRET_MAP_INITIALIZED as usize; pub const PATH_SECRET_MAP_INITIALIZED__CAPACITY: usize = InfoId::PATH_SECRET_MAP_INITIALIZED__CAPACITY as usize; @@ -1062,7 +1054,6 @@ mod counter { id::CONNECTION_CLOSED => Self(connection_closed), id::ENDPOINT_INITIALIZED => Self(endpoint_initialized), id::DC_CONNECTION_TIMEOUT => Self(dc_connection_timeout), - id::OFFLOAD_TASK_METRICS => Self(offload_task_metrics), id::PATH_SECRET_MAP_INITIALIZED => Self(path_secret_map_initialized), id::PATH_SECRET_MAP_UNINITIALIZED => Self(path_secret_map_uninitialized), id::PATH_SECRET_MAP_BACKGROUND_HANDSHAKE_REQUESTED => { @@ -1353,9 +1344,7 @@ mod counter { s2n_quic_dc__event__counter__dc_connection_timeout] fn dc_connection_timeout(value: u64); #[link_name = - s2n_quic_dc__event__counter__offload_task_metrics] - fn offload_task_metrics(value: u64); - #[link_name = s2n_quic_dc__event__counter__path_secret_map_initialized] + s2n_quic_dc__event__counter__path_secret_map_initialized] fn path_secret_map_initialized(value: u64); #[link_name = s2n_quic_dc__event__counter__path_secret_map_uninitialized] @@ -2231,12 +2220,6 @@ mod measure { id::STREAM_HANDSHAKE_PACKET_REJECTED__CONN => { Self(stream_handshake_packet_rejected__conn) } - id::OFFLOAD_TASK_METRICS__MEAN_POLL_DURATION => { - Self(offload_task_metrics__mean_poll_duration) - } - id::OFFLOAD_TASK_METRICS__MEAN_SCHEDULED_DURATION => { - Self(offload_task_metrics__mean_scheduled_duration) - } id::PATH_SECRET_MAP_INITIALIZED__CAPACITY => { Self(path_secret_map_initialized__capacity) } @@ -2656,12 +2639,6 @@ mod measure { s2n_quic_dc__event__measure__stream_handshake_packet_rejected__conn] fn stream_handshake_packet_rejected__conn(value: u64); #[link_name = - s2n_quic_dc__event__measure__offload_task_metrics__mean_poll_duration] - fn offload_task_metrics__mean_poll_duration(value: u64); - #[link_name = - s2n_quic_dc__event__measure__offload_task_metrics__mean_scheduled_duration] - fn offload_task_metrics__mean_scheduled_duration(value: u64); - #[link_name = s2n_quic_dc__event__measure__path_secret_map_initialized__capacity] fn path_secret_map_initialized__capacity(value: u64); #[link_name = diff --git a/dc/s2n-quic-dc/src/path/secret/map.rs b/dc/s2n-quic-dc/src/path/secret/map.rs index 4356134e70..7388de5cca 100644 --- a/dc/s2n-quic-dc/src/path/secret/map.rs +++ b/dc/s2n-quic-dc/src/path/secret/map.rs @@ -3,7 +3,7 @@ use crate::{ credentials::{Credentials, Id}, - event::{self}, + event, packet::{secret_control as control, Packet}, path::secret::{ open, diff --git a/dc/s2n-quic-dc/src/psk/io.rs b/dc/s2n-quic-dc/src/psk/io.rs index f7198b1acc..9f6867ef35 100644 --- a/dc/s2n-quic-dc/src/psk/io.rs +++ b/dc/s2n-quic-dc/src/psk/io.rs @@ -253,21 +253,8 @@ pub(super) async fn server< } } - let map_clone = map.clone(); - // if let Some(monitor) = server.offload_metrics.take() { - // let frequency = std::time::Duration::new(1, 0); - // tokio::spawn(async move { - // loop { - // let metrics = monitor.cumulative(); - // map_clone.on_offload_task_metrics(&metrics); - // tokio::time::sleep(frequency).await; - // } - // }); - // } - while let Some(mut connection) = server.server.accept().await { let map_clone = map.clone(); - tokio::spawn(async move { // The accepted connection must remain open until the client has finished inserting // the entry into its map. The client indicates this by sending a ConnectionClose