diff --git a/crates/kumo-api-types/src/egress_path.rs b/crates/kumo-api-types/src/egress_path.rs index c7b66dd5..8f1766ba 100644 --- a/crates/kumo-api-types/src/egress_path.rs +++ b/crates/kumo-api-types/src/egress_path.rs @@ -131,6 +131,14 @@ pub fn find_rustls_cipher_suite(name: &str) -> Option { None } +#[derive(Deserialize, Serialize, Debug, Clone, Default, Copy, PartialEq, Eq)] +#[cfg_attr(feature = "lua", derive(FromLua))] +pub enum WakeupStrategy { + #[default] + Aggressive, + Relaxed, +} + #[derive(Deserialize, Serialize, Debug, Clone, Default, Copy, PartialEq, Eq)] #[cfg_attr(feature = "lua", derive(FromLua))] pub enum MemoryReductionPolicy { @@ -280,6 +288,11 @@ pub struct EgressPathConfig { #[serde(default)] pub refresh_strategy: ConfigRefreshStrategy, + #[serde(default)] + pub dispatcher_wakeup_strategy: WakeupStrategy, + #[serde(default)] + pub maintainer_wakeup_strategy: WakeupStrategy, + /// Specify an explicit provider name that should apply to this /// path. The provider name will be used when computing metrics /// rollups by provider. If omitted, then @@ -372,6 +385,8 @@ impl Default for EgressPathConfig { readyq_pool_name: None, low_memory_reduction_policy: MemoryReductionPolicy::default(), no_memory_reduction_policy: MemoryReductionPolicy::default(), + maintainer_wakeup_strategy: WakeupStrategy::default(), + dispatcher_wakeup_strategy: WakeupStrategy::default(), } } } diff --git a/crates/kumo-api-types/src/shaping.rs b/crates/kumo-api-types/src/shaping.rs index d86a0bf7..198ced86 100644 --- a/crates/kumo-api-types/src/shaping.rs +++ b/crates/kumo-api-types/src/shaping.rs @@ -1968,6 +1968,8 @@ MergedEntry { aggressive_connection_opening: false, refresh_interval: 60s, refresh_strategy: Ttl, + dispatcher_wakeup_strategy: Aggressive, + maintainer_wakeup_strategy: Aggressive, provider_name: None, remember_broken_tls: None, opportunistic_tls_reconnect_on_failed_handshake: false, @@ -2113,6 +2115,8 @@ MergedEntry { aggressive_connection_opening: false, refresh_interval: 60s, refresh_strategy: Ttl, + dispatcher_wakeup_strategy: Aggressive, + maintainer_wakeup_strategy: Aggressive, provider_name: None, remember_broken_tls: None, opportunistic_tls_reconnect_on_failed_handshake: false, @@ -2171,6 +2175,8 @@ MergedEntry { aggressive_connection_opening: false, refresh_interval: 60s, refresh_strategy: Ttl, + dispatcher_wakeup_strategy: Aggressive, + maintainer_wakeup_strategy: Aggressive, provider_name: None, remember_broken_tls: None, opportunistic_tls_reconnect_on_failed_handshake: false, @@ -2322,6 +2328,8 @@ MergedEntry { aggressive_connection_opening: false, refresh_interval: 60s, refresh_strategy: Ttl, + dispatcher_wakeup_strategy: Aggressive, + maintainer_wakeup_strategy: Aggressive, provider_name: None, remember_broken_tls: None, opportunistic_tls_reconnect_on_failed_handshake: false, diff --git a/crates/kumod/src/ready_queue.rs b/crates/kumod/src/ready_queue.rs index de0a280f..9c90fddf 100644 --- a/crates/kumod/src/ready_queue.rs +++ b/crates/kumod/src/ready_queue.rs @@ -24,7 +24,9 @@ use config::epoch::ConfigEpoch; use config::{declare_event, load_config}; use dashmap::DashMap; use dns_resolver::MailExchanger; -use kumo_api_types::egress_path::{ConfigRefreshStrategy, EgressPathConfig, MemoryReductionPolicy}; +use kumo_api_types::egress_path::{ + ConfigRefreshStrategy, EgressPathConfig, MemoryReductionPolicy, WakeupStrategy, +}; use kumo_server_common::config_handle::ConfigHandle; use kumo_server_lifecycle::{is_shutting_down, Activity, ShutdownSubcription}; use kumo_server_memory::{ @@ -465,6 +467,7 @@ impl ReadyQueueManager { notify_dispatcher, notify_maintainer, connections: FairMutex::new(vec![]), + num_connections: Arc::new(AtomicUsize::new(0)), path_config: ConfigHandle::new(path_config), queue_config: queue_config.clone(), egress_source, @@ -671,6 +674,7 @@ pub struct ReadyQueue { notify_maintainer: Arc, notify_dispatcher: Arc, connections: FairMutex>>, + num_connections: Arc, metrics: DeliveryMetrics, activity: Activity, consecutive_connection_failures: Arc, @@ -693,13 +697,47 @@ impl ReadyQueue { Some(res) => Some(res), None => { self.metrics.ready_full.inc(); - self.notify_maintainer.notify_one(); - self.notify_dispatcher.notify_waiters(); + self.wakeup_dispatcher_or_maintainer(); None } } } + pub fn wakeup_all_dispatchers(&self) { + self.notify_dispatcher.notify_waiters(); + } + + pub fn wakeup_dispatcher_or_maintainer(&self) { + let path_config = self.path_config.borrow(); + let num_connections = self.num_connections.load(Ordering::SeqCst); + + match path_config.dispatcher_wakeup_strategy { + WakeupStrategy::Aggressive => { + self.notify_dispatcher.notify_waiters(); + } + WakeupStrategy::Relaxed => { + if num_connections > 0 { + self.notify_dispatcher.notify_one(); + } + } + } + + match path_config.maintainer_wakeup_strategy { + WakeupStrategy::Aggressive => { + self.notify_maintainer.notify_one(); + } + WakeupStrategy::Relaxed => { + let approx_conn_goal = ideal_connection_count( + self.ready.len(), + path_config.connection_limit.limit as usize, + ); + if num_connections < approx_conn_goal { + self.notify_maintainer.notify_one(); + } + } + } + } + pub fn get_path_config(&self) -> &ConfigHandle { &self.path_config } @@ -720,8 +758,7 @@ impl ReadyQueue { } } reservation.redeem(msg); - self.notify_maintainer.notify_one(); - self.notify_dispatcher.notify_waiters(); + self.wakeup_dispatcher_or_maintainer(); } pub async fn insert(&self, msg: Message) -> Result<(), Message> { @@ -742,14 +779,12 @@ impl ReadyQueue { } match self.ready.push(msg) { Ok(()) => { - self.notify_maintainer.notify_one(); - self.notify_dispatcher.notify_waiters(); + self.wakeup_dispatcher_or_maintainer(); Ok(()) } Err(msg) => { self.metrics.ready_full.inc(); - self.notify_maintainer.notify_one(); - self.notify_dispatcher.notify_waiters(); + self.wakeup_dispatcher_or_maintainer(); Err(msg) } } @@ -909,7 +944,7 @@ impl ReadyQueue { ); self.reinsert_ready_queue("suspend", InsertReason::ReadyQueueWasSuspended.into()) .await; - self.notify_dispatcher.notify_waiters(); + self.wakeup_all_dispatchers(); return; } @@ -1011,6 +1046,7 @@ impl ReadyQueue { let egress_pool = self.egress_pool.clone(); let consecutive_connection_failures = self.consecutive_connection_failures.clone(); let states = self.states.clone(); + let num_connections = self.num_connections.clone(); tracing::trace!("spawning client for {name}"); if let Ok(handle) = self.readyq_spawn(format!("smtp client {name}"), async move { @@ -1028,6 +1064,7 @@ impl ReadyQueue { egress_pool, leases, states, + num_connections, ) .await { @@ -1100,7 +1137,7 @@ impl ReadyQueue { "{}: refreshed get_egress_path_config to generation {generation}", self.name ); - self.notify_dispatcher.notify_waiters(); + self.wakeup_all_dispatchers(); for msg in self.ready.update_capacity(max_ready) { if let Err(err) = Dispatcher::reinsert_message( msg, @@ -1190,6 +1227,7 @@ pub struct Dispatcher { batch_started: Option, pub states: Arc>, active_bounce: ArcSwap>>, + num_connections: Arc, } impl Drop for Dispatcher { @@ -1198,7 +1236,7 @@ impl Drop for Dispatcher { let msgs = std::mem::take(&mut self.msgs); let activity = self.activity.clone(); let name = self.name.to_string(); - let notify_dispatcher = self.notify_dispatcher.clone(); + self.num_connections.fetch_sub(1, Ordering::SeqCst); self.readyq_spawn("Dispatcher::drop".to_string(), async move { let had_msgs = !msgs.is_empty(); @@ -1236,9 +1274,8 @@ impl Drop for Dispatcher { tokio::time::sleep(Duration::from_secs(1)).await; let ready_queue = ReadyQueueManager::get_by_name(&name); if let Some(q) = ready_queue { - q.notify_maintainer.notify_one(); + q.wakeup_dispatcher_or_maintainer(); } - notify_dispatcher.notify_one(); } }) .ok(); @@ -1261,6 +1298,7 @@ impl Dispatcher { egress_pool: String, leases: Vec, states: Arc>, + num_connections: Arc, ) -> anyhow::Result<()> { let activity = Activity::get(format!("ready_queue Dispatcher {name}"))?; @@ -1295,7 +1333,9 @@ impl Dispatcher { session_id: Uuid::new_v4(), states, active_bounce: Arc::new(None).into(), + num_connections: num_connections.clone(), }; + dispatcher.num_connections.fetch_add(1, Ordering::SeqCst); let mut queue_dispatcher: Box = match &queue_config.borrow().protocol { DeliveryProto::Smtp { smtp } => { diff --git a/docs/changelog/main.md b/docs/changelog/main.md index 5b2c1147..ac7fa391 100644 --- a/docs/changelog/main.md +++ b/docs/changelog/main.md @@ -82,6 +82,11 @@ on systems with a large number of admin bounce entries. * Improved performance of TSA state storage, which in turn improves latency in tsa-daemon response times when many automation rules are triggering. +* New + [maintainer_wakeup_strategy](../reference/kumo/make_egress_path/maintainer_wakeup_strategy.md) + and + [dispatcher_wakeup_strategy](../reference/kumo/make_egress_path/dispatcher_wakeup_strategy.md) + options for fine tuning overall system performance. ## Fixes diff --git a/docs/reference/kumo/make_egress_path/dispatcher_wakeup_strategy.md b/docs/reference/kumo/make_egress_path/dispatcher_wakeup_strategy.md new file mode 100644 index 00000000..68490140 --- /dev/null +++ b/docs/reference/kumo/make_egress_path/dispatcher_wakeup_strategy.md @@ -0,0 +1,26 @@ +# maintainer_wakeup_strategy + +{{since('dev')}} + +Adjusts how aggressively the readyq Dispatch tasks will be awoken +as messages are placed into the readyq. + +Can have one of two values: + + * `"Aggressive"` - the default. Every attempt to place a message + into the ready queue will cause all idle outbound sessions + that are associated with the readyq to wakeup and attempt to + pull messages from it. + + * `"Relaxed"` - Each submission attempt will signal just a single + idle outbound session to wakeup and pull a message from the + queue. + +An `"Aggressive"` setting will cause more spurious wakeups in +the case that there are multiple sessions and just a single new +message being submitted, whereas a `"Relaxed"` setting will be +more optimal in terms of CPU usage, at the risk of potentially +missing a wakeup in some edge cases with bursty or low traffic. + +In earlier versions of KumoMTA this option did not exist, but the +behavior was equivalent to the `"Aggressive"` setting. diff --git a/docs/reference/kumo/make_egress_path/maintainer_wakeup_strategy.md b/docs/reference/kumo/make_egress_path/maintainer_wakeup_strategy.md new file mode 100644 index 00000000..e2c6f400 --- /dev/null +++ b/docs/reference/kumo/make_egress_path/maintainer_wakeup_strategy.md @@ -0,0 +1,30 @@ +# maintainer_wakeup_strategy + +{{since('dev')}} + +Adjusts how aggressively the readyq maintainer task will be awoken +as messages are placed into the readyq. + +Can have one of two values: + + * `"Aggressive"` - the default. Every attempt to place a message into the + ready queue will cause the associated maintainer task to wakeup to assess + whether more connections need to be established. + + * `"Relaxed"` - Each submission attempt will perform a quick approximation + and assessment of the current connection count to decide whether the + maintainer task needs to be signalled. If the number of connections + matches the ideal for the current queue size, then the maintainer will + not be signalled and it will wakeup periodically to reassess the + load. + +The primary purpose of the readyq maintainer is to establish new outbound +connections based on the queue size. A `"Relaxed"` setting will cause +a precise assessment of that state to occur less frequently, reducing +CPU overhead, but it may result in an increase in latency for +outbound traffic when conditions are bursty. + +In earlier versions of KumoMTA this option did not exist, but the +behavior was equivalent to the `"Aggressive"` setting. + +