From 23e1700a13fd08709f0ae7bb34e35a81474a795a Mon Sep 17 00:00:00 2001 From: Wez Furlong Date: Thu, 8 Aug 2024 09:51:13 -0700 Subject: [PATCH] ready_queue: allow dynamically changing max_ready This was deliberately adjusted to only change when a queue was reapead in a prior release, in the interests of raw performance. While that might work for most domains, it's not ideal for busy domains that have continual traffic. This change adopts an ArcSwap to atomically swap out the underlying ArrayQueue. This introduces a small amount of overhead in the hot path to atomically acquire the current version of the queue, but it shouldn't be enough to be noticeable in practice. --- Cargo.lock | 1 + crates/kumod/Cargo.toml | 1 + crates/kumod/src/ready_queue.rs | 61 +++++++++++++++++++++++++++------ docs/changelog/main.md | 3 +- 4 files changed, 54 insertions(+), 12 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index b0343156..3f0511cc 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2762,6 +2762,7 @@ name = "kumod" version = "0.1.0" dependencies = [ "anyhow", + "arc-swap", "async-channel 2.3.1", "async-recursion", "async-trait", diff --git a/crates/kumod/Cargo.toml b/crates/kumod/Cargo.toml index 3e933549..9a59a5c8 100644 --- a/crates/kumod/Cargo.toml +++ b/crates/kumod/Cargo.toml @@ -7,6 +7,7 @@ edition = "2021" [dependencies] anyhow = "1.0" +arc-swap = "1.6" async-channel = "2.1" async-recursion = "1.1" async-trait = "0.1" diff --git a/crates/kumod/src/ready_queue.rs b/crates/kumod/src/ready_queue.rs index 9ab66df5..75c8217d 100644 --- a/crates/kumod/src/ready_queue.rs +++ b/crates/kumod/src/ready_queue.rs @@ -10,6 +10,7 @@ use crate::queue::{DeliveryProto, Queue, QueueConfig, QueueManager, QMAINT_RUNTI use crate::smtp_dispatcher::{MxListEntry, OpportunisticInsecureTlsHandshakeError, SmtpDispatcher}; use crate::spool::SpoolManager; use anyhow::Context; +use arc_swap::ArcSwap; use async_trait::async_trait; use config::{load_config, CallbackSignature}; use crossbeam_queue::ArrayQueue; @@ -49,41 +50,76 @@ pub fn set_readyq_threads(n: usize) { } pub struct Fifo { - queue: ArrayQueue, + queue: ArcSwap>, count: IntGauge, } impl Fifo { pub fn new(capacity: usize, count: IntGauge) -> Self { Self { - queue: ArrayQueue::new(capacity), + queue: Arc::new(ArrayQueue::new(capacity)).into(), count, } } pub fn push(&self, msg: Message) -> Result<(), Message> { - self.queue.push(msg)?; + self.queue.load().push(msg)?; self.count.inc(); Ok(()) } + #[must_use] pub fn pop(&self) -> Option { - let msg = self.queue.pop()?; + let msg = self.queue.load().pop()?; self.count.dec(); Some(msg) } + #[must_use] pub fn drain(&self) -> Vec { - let mut messages = Vec::with_capacity(self.queue.len()); - while let Some(msg) = self.queue.pop() { + let queue = self.queue.load(); + let mut messages = Vec::with_capacity(queue.len()); + while let Some(msg) = queue.pop() { messages.push(msg); } self.count.sub(messages.len() as i64); messages } + /// Adjust the capacity of the Fifo. + /// If the capacity is the same, nothing changes. + /// Otherwise, a new ArrayQueue is constructed and swapped in + /// to replace the existing queue. + /// The old queue is then drained into the new queue. + /// Any messages that won't fit into the new queue are + /// returned to the caller, who is responsible for re-inserting + /// those messages into the scheduled queue + #[must_use] + pub fn update_capacity(&self, capacity: usize) -> Vec { + let queue = self.queue.load(); + if queue.capacity() == capacity { + return vec![]; + } + + let queue = self.queue.swap(Arc::new(ArrayQueue::new(capacity)).into()); + let new_queue = self.queue.load(); + + let mut messages = Vec::with_capacity(queue.len()); + while let Some(msg) = queue.pop() { + // Note that we may race with other actors who are inserting + // into this queue, so even if the new capacity is greater + // than the prior capacity, there is still a chance that + // we'll have some overflow to deal with + if let Err(msg) = new_queue.push(msg) { + messages.push(msg); + } + } + self.count.sub(messages.len() as i64); + messages + } + pub fn len(&self) -> usize { - self.queue.len() + self.queue.load().len() } } @@ -359,13 +395,16 @@ impl ReadyQueueManager { { Ok(ReadyQueueConfig { path_config, .. }) => { if path_config != **queue.path_config.borrow() { + let max_ready = path_config.max_ready; + let generation = queue.path_config.update(path_config); - // Note that the Fifo type doesn't allow for dynamically - // changing the capacity of the ready queue, so you will - // need to allow the ready queue to be reaped before that - // change takes effect tracing::trace!("{name}: refreshed get_egress_path_config to generation {generation}"); queue.notify_dispatcher.notify_waiters(); + for msg in queue.ready.update_capacity(max_ready) { + if let Err(err) = Dispatcher::reinsert_message(msg).await { + tracing::error!("error reinserting message: {err:#}"); + } + } } } Err(err) => { diff --git a/docs/changelog/main.md b/docs/changelog/main.md index b2324659..d16a1ddc 100644 --- a/docs/changelog/main.md +++ b/docs/changelog/main.md @@ -95,4 +95,5 @@ classification with the same rule, the "A" classification would be the result, even if the "B" rule was the first one listed in your classification data file(s). - +* Changing the `max_ready` value for a ready queue no longer requires waiting for + the queue to be reaped before it will take effect.