From fb993f0cc121558fbc03a997e67205328984d2f0 Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Sun, 5 Jul 2026 12:05:46 +0200 Subject: [PATCH] Cluster: Broadcast MTA queue refresh events to all nodes --- CHANGELOG.md | 1 + crates/common/src/ipc.rs | 1 + crates/services/src/broadcast/mod.rs | 4 ++++ crates/services/src/broadcast/subscriber.rs | 10 ++++++++++ crates/smtp/src/queue/spool.rs | 4 +++- tests/src/system/oidc.rs | 5 ++++- 6 files changed, 23 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7dd82866..eae8ddd8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -30,6 +30,7 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If - Snowflake past id generation fails when the provided duration is longer than 4 years. - Calendar scheduling: Wrong RSVP base URL is used. - Network listener: Accept loop spins all CPU cores with no back-off when the process hits `EMFILE` (too many open files). +- Cluster: Broadcast MTA queue refresh events to all nodes. ## [0.16.11] - 2026-06-25 diff --git a/crates/common/src/ipc.rs b/crates/common/src/ipc.rs index 8371dd7a..dc8e4f54 100644 --- a/crates/common/src/ipc.rs +++ b/crates/common/src/ipc.rs @@ -78,6 +78,7 @@ pub enum BroadcastEvent { CacheInvalidateAll, CacheInvalidateNegative, MtaQueueStatus { is_running: bool }, + QueueRefresh, } #[derive(Debug, Clone, Copy)] diff --git a/crates/services/src/broadcast/mod.rs b/crates/services/src/broadcast/mod.rs index 8955e2ca..d804ea93 100644 --- a/crates/services/src/broadcast/mod.rs +++ b/crates/services/src/broadcast/mod.rs @@ -135,6 +135,9 @@ impl BroadcastBatch> { serialized.push(11u8); } } + BroadcastEvent::QueueRefresh => { + serialized.push(12u8); + } } } serialized @@ -266,6 +269,7 @@ where 9 => Ok(Some(BroadcastEvent::CacheInvalidateNegative)), 10 => Ok(Some(BroadcastEvent::MtaQueueStatus { is_running: true })), 11 => Ok(Some(BroadcastEvent::MtaQueueStatus { is_running: false })), + 12 => Ok(Some(BroadcastEvent::QueueRefresh)), _ => Err(()), } } else { diff --git a/crates/services/src/broadcast/subscriber.rs b/crates/services/src/broadcast/subscriber.rs index f5107b79..b62cab31 100644 --- a/crates/services/src/broadcast/subscriber.rs +++ b/crates/services/src/broadcast/subscriber.rs @@ -153,6 +153,15 @@ pub fn spawn_broadcast_subscriber(inner: Arc, mut shutdown_rx: watch::Rec .send(QueueEvent::Paused(!is_running)) .await; } + BroadcastEvent::QueueRefresh => { + if inner.shared_core.load().network.roles.outbound_mta { + let _ = inner + .ipc + .queue_tx + .send(QueueEvent::Refresh) + .await; + } + } BroadcastEvent::RegistryChange(change) => { match Box::pin(inner.build_server().reload_registry(change)).await { Ok(result) => { @@ -255,5 +264,6 @@ fn log_event(event: &BroadcastEvent) -> trc::Value { "MtaQueuePaused".into() } } + BroadcastEvent::QueueRefresh => "QueueRefresh".into(), } } diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index f6a9f301..23ccd9a9 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -17,7 +17,7 @@ use crate::queue::{ use ahash::{AHashMap, AHashSet}; use common::config::smtp::auth::DkimSigners; use common::config::smtp::queue::{ArchivedQueueExpiry, QueueName}; -use common::ipc::QueueEvent; +use common::ipc::{BroadcastEvent, QueueEvent}; use common::network::RcptResolution; use common::{KV_LOCK_QUEUE_MESSAGE, Server}; use mail_auth::AuthenticatedMessage; @@ -601,6 +601,8 @@ impl MessageWrapper { ); } + server.cluster_broadcast(BroadcastEvent::QueueRefresh).await; + true } diff --git a/tests/src/system/oidc.rs b/tests/src/system/oidc.rs index 336c366f..a8caf925 100644 --- a/tests/src/system/oidc.rs +++ b/tests/src/system/oidc.rs @@ -256,7 +256,10 @@ pub async fn test(test: &mut TestServer) { }, ) .await; - assert_eq!(status, 201, "registration should return 201 for {good_uri}: {body}"); + assert_eq!( + status, 201, + "registration should return 201 for {good_uri}: {body}" + ); } // Register the client used for the flow with a private-use scheme redirect URI