ready_queue: introduce dispatcher_wakeup_strategy and maintainer_wakeup_strategy

These allow optionally reducing how aggressively the dispatcher and
maintainers will be awoken when message(s) are submitted to the ready
queue.

The default behavior remains the same; the new thing here is the
ability to make it more relaxed, which should reduce some CPU
overheads for very busy systems with many queues.

Making things more relaxed does introduce a possibility for higher
outbound latency in some edge cases with low or bursty traffic.
This commit is contained in:
Wez Furlong
2025-04-16 15:19:27 -07:00
parent c22e98bcd5
commit 18e510ebe2
6 changed files with 138 additions and 14 deletions
+15
View File
@@ -131,6 +131,14 @@ pub fn find_rustls_cipher_suite(name: &str) -> Option<SupportedCipherSuite> {
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(),
}
}
}
+8
View File
@@ -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,
+54 -14
View File
@@ -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>,
notify_dispatcher: Arc<Notify>,
connections: FairMutex<Vec<JoinHandle<()>>>,
num_connections: Arc<AtomicUsize>,
metrics: DeliveryMetrics,
activity: Activity,
consecutive_connection_failures: Arc<AtomicUsize>,
@@ -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<EgressPathConfig> {
&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<tokio::time::Instant>,
pub states: Arc<FairMutex<ReadyQueueStates>>,
active_bounce: ArcSwap<Option<CachedEntry<AdminBounceEntry>>>,
num_connections: Arc<AtomicUsize>,
}
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<LimitLease>,
states: Arc<FairMutex<ReadyQueueStates>>,
num_connections: Arc<AtomicUsize>,
) -> 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<dyn QueueDispatcher> = match &queue_config.borrow().protocol {
DeliveryProto::Smtp { smtp } => {
+5
View File
@@ -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
@@ -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.
@@ -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.