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.