smtp_dispatcher: improve handling of unilateral disconnects

The goal is to treat peer-initiated-unilateral-disconnects as being
somewhat equivalent to the way that we would idle out the connection if
we had no messages for it, on the supposition that the most likely
cause of a disconnect in between command verbs that we send is that
the peer decided that we were idle too long.  That isn't the only
possible reason, but it is the motivating example.

For the idled-out case we should close the current session and
have our message(s) go out on a separate session.  We don't want
to try the next host in the connection plan in case we're talking
to some kind of honey pit configuration where only the first
MX in the plan is valid, and even talking to a secondary can
harm your ability to send to the first in the future.

Handling this is a little tricky because we need to take care of
the distinction between getting a unilateral before or after
we've attempted to send a message.

In the before-case we don't want to blindly assume that we should start
a new session because that would mean that a persistent issue on the
first MX would have us spinning our wheels trying new connections only
to the first MX over and over.

This commit adds a couple of integration tests that contrive situations
where we exercise some pertinent cases.

refs: https://github.com/KumoCorp/kumomta/pull/482
This commit is contained in:
Wez Furlong
2026-02-25 13:26:17 +00:00
parent b018f8f128
commit 5bb60b06d0
12 changed files with 371 additions and 84 deletions
@@ -17,6 +17,9 @@ kumo.on('init', function()
relay_hosts = { '0.0.0.0/0' },
batch_handling = 'BatchByDomain',
max_recipients_per_message = 4,
-- This client_timeout value is coupled with assumptions
-- in disconnect_peer_idle_out!
client_timeout = '3s',
}
local client_ca = os.getenv 'KUMOD_CLIENT_REQUIRED_CA'
if client_ca then
+15
View File
@@ -296,6 +296,21 @@ kumo.on('get_queue_config', function(domain, tenant, campaign, routing_domain)
}
end
-- This domain is coupled with the broken_first_choice_mx test
if domain == 'broken-first-mx.example.com' then
protocol = {
-- Redirect traffic to the sink, with two hosts in the connection plan,
-- but the first one is broken/not routable
smtp = {
mx_list = {
-- 255.255.255.255 should not be reachable/routable on any system
{ name = 'brokenhost', addr = '255.255.255.255:' .. SINK_PORT },
{ name = 'workinghost', addr = '127.0.0.1:' .. SINK_PORT },
},
},
}
end
if domain == 'nxdomain' then
-- this nxdomain domain is a special domain that is assumed not
-- to resolve. It is generated by the retry_schedule integration
@@ -0,0 +1,70 @@
use crate::kumod::{DaemonWithMaildir, MailGenParams};
use kumo_log_types::RecordType::Delivery;
use std::time::Duration;
/// We're sending to broken-first-mx which is defined in source.lua
/// to have a non-routable first MX in its connection plan, followed
/// by the regular sink. This should cause a connection failure
/// for the first candidate, but we should then successfully
/// deliver to the second candidate.
#[tokio::test]
async fn broken_first_choice_mx() -> anyhow::Result<()> {
let mut daemon = DaemonWithMaildir::start().await?;
let mut client = daemon.smtp_client().await?;
let response = MailGenParams {
recip: Some("first@broken-first-mx.example.com"),
..Default::default()
}
.send(&mut client)
.await?;
eprintln!("{response:?}");
anyhow::ensure!(response.code == 250);
daemon
.wait_for_source_summary(
|summary| summary.get(&Delivery).copied().unwrap_or(0) >= 1,
Duration::from_secs(50),
)
.await;
daemon
.wait_for_maildir_count(1, Duration::from_secs(10))
.await;
daemon.stop_both().await?;
let delivery_summary = daemon.dump_logs().await?;
k9::snapshot!(
delivery_summary,
"
DeliverySummary {
source_counts: {
Reception: 1,
Delivery: 1,
},
sink_counts: {
Reception: 1,
Delivery: 1,
},
}
"
);
k9::snapshot!(
daemon.source.accounting_stats()?,
"
AccountingStats {
received: 1,
delivered: 1,
}
"
);
let logs = daemon.source.collect_logs().await?;
let delivered = logs.iter().find(|record| record.kind == Delivery).unwrap();
k9::assert_equal!(delivered.peer_address.as_ref().unwrap().name, "workinghost");
Ok(())
}
@@ -0,0 +1,87 @@
use crate::kumod::{DaemonWithMaildir, MailGenParams};
use kumo_log_types::RecordType::Delivery;
use std::time::Duration;
#[tokio::test]
async fn disconnect_peer_idle_out() -> anyhow::Result<()> {
let mut daemon = DaemonWithMaildir::start().await?;
let mut client = daemon.smtp_client().await?;
let response = MailGenParams {
recip: Some("first@example.com"),
..Default::default()
}
.send(&mut client)
.await?;
eprintln!("{response:?}");
anyhow::ensure!(response.code == 250);
// Wait for the sink to idle out its connection.
// We're assuming that it is set to 3s, so we wait
// 4s to make things a little less racy
tokio::time::sleep(Duration::from_secs(4)).await;
let response = MailGenParams {
recip: Some("second@example.com"),
..Default::default()
}
.send(&mut client)
.await?;
eprintln!("{response:?}");
anyhow::ensure!(response.code == 250);
daemon
.wait_for_source_summary(
|summary| summary.get(&Delivery).copied().unwrap_or(0) >= 2,
Duration::from_secs(50),
)
.await;
daemon
.wait_for_maildir_count(2, Duration::from_secs(10))
.await;
daemon.stop_both().await?;
let delivery_summary = daemon.dump_logs().await?;
k9::snapshot!(
delivery_summary,
"
DeliverySummary {
source_counts: {
Reception: 2,
Delivery: 2,
},
sink_counts: {
Reception: 2,
Delivery: 2,
},
}
"
);
k9::snapshot!(
daemon.source.accounting_stats()?,
"
AccountingStats {
received: 2,
delivered: 2,
}
"
);
let logs = daemon.source.collect_logs().await?;
let mut delivered_sessions = logs
.iter()
.filter_map(|record| match record.kind {
kumo_log_types::RecordType::Delivery => Some(record.session_id),
_ => None,
})
.collect::<Vec<_>>();
delivered_sessions.sort();
delivered_sessions.dedup();
assert_eq!(delivered_sessions.len(), 2, "two separate sessions");
Ok(())
}
+2
View File
@@ -2,8 +2,10 @@
mod arc;
mod auth_deliver;
mod auth_deliver_invalid_password;
mod broken_first_choice_mx;
mod disconnect_in_data;
mod disconnect_in_mail_from;
mod disconnect_peer_idle_out;
mod disconnect_reconnect_same_host;
mod disconnect_terminate_ok;
mod eightbitmime;
+8 -3
View File
@@ -1,7 +1,7 @@
use crate::delivery_metrics::MetricsWrappedConnection;
use crate::logging::disposition::{log_disposition, LogDisposition, RecordType};
use crate::queue::{DeliveryProto, QueueConfig, QueueManager};
use crate::ready_queue::{Dispatcher, QueueDispatcher};
use crate::ready_queue::{AttemptConnectionDisposition, Dispatcher, QueueDispatcher};
use crate::smtp_server::{default_hostname, TraceHeaders};
use crate::spool::SpoolManager;
use anyhow::Context;
@@ -1324,12 +1324,17 @@ impl QueueDispatcher for HttpInjectionGeneratorDispatcher {
None => Ok(false),
}
}
async fn attempt_connection(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<()> {
async fn attempt_connection(
&mut self,
dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition> {
if self.connection.is_none() {
self.connection
.replace(dispatcher.metrics.wrap_connection(()));
Ok(AttemptConnectionDisposition::ConnectedNew)
} else {
Ok(AttemptConnectionDisposition::ReusedExisting)
}
Ok(())
}
async fn have_more_connection_candidates(&mut self, _dispatcher: &mut Dispatcher) -> bool {
false
+9 -4
View File
@@ -1,7 +1,7 @@
use crate::delivery_metrics::MetricsWrappedConnection;
use crate::logging::disposition::{log_disposition, LogDisposition};
use crate::queue::{IncrementAttempts, InsertReason, QueueManager};
use crate::ready_queue::{Dispatcher, QueueDispatcher};
use crate::ready_queue::{AttemptConnectionDisposition, Dispatcher, QueueDispatcher};
use crate::smtp_server::RejectError;
use crate::spool::SpoolManager;
use anyhow::Context;
@@ -133,9 +133,14 @@ impl QueueDispatcher for LuaQueueDispatcher {
self.proto_config.max_batch_latency
}
async fn attempt_connection(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<()> {
async fn attempt_connection(
&mut self,
dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition> {
match &self.connection {
ConnectionState::Connected(_) => return Ok(()),
ConnectionState::Connected(_) => {
return Ok(AttemptConnectionDisposition::ReusedExisting)
}
ConnectionState::Disconnected => {
anyhow::bail!("only one connection attempt per session");
}
@@ -162,7 +167,7 @@ impl QueueDispatcher for LuaQueueDispatcher {
self.connection = ConnectionState::Connected(connection_wrapper.map_connection(connection));
dispatcher.delivered_this_connection = 0;
Ok(())
Ok(AttemptConnectionDisposition::ConnectedNew)
}
async fn have_more_connection_candidates(&mut self, _dispatcher: &mut Dispatcher) -> bool {
+3
View File
@@ -79,6 +79,9 @@ pub enum InsertReason {
/// The safey net in Dispatcher::Drop re-queued the message.
/// This shouldn't happen; if you see this in a log, please report it!
DispatcherDrop,
/// The peer unilaterally closed the connection before we started
/// delivery, so we want to try a new connection plan
PeerClosedConnection,
}
#[cfg(test)]
+106 -64
View File
@@ -1180,6 +1180,14 @@ impl Drop for ReadyQueue {
}
}
#[derive(Copy, Clone, Debug)]
pub enum AttemptConnectionDisposition {
ConnectedNew,
ReusedExisting,
PeerClosedConnectionNeedNewSession,
PeerClosedConnectionContinueSession,
}
#[async_trait]
pub trait QueueDispatcher: Debug + Send {
async fn deliver_message(
@@ -1188,7 +1196,10 @@ pub trait QueueDispatcher: Debug + Send {
dispatcher: &mut Dispatcher,
) -> anyhow::Result<()>;
async fn attempt_connection(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<()>;
async fn attempt_connection(
&mut self,
dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition>;
async fn have_more_connection_candidates(&mut self, dispatcher: &mut Dispatcher) -> bool;
async fn close_connection(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<bool>;
@@ -1431,78 +1442,109 @@ impl Dispatcher {
}
};
if let Err(err) = result {
if OpportunisticInsecureTlsHandshakeError::is_match_anyhow(&err) {
num_opportunistic_tls_failures += 1;
}
connection_failures.push(format!("{err:#}"));
if !queue_dispatcher
.have_more_connection_candidates(&mut dispatcher)
.await
{
for msg in dispatcher.msgs.drain(..) {
let summary = if num_opportunistic_tls_failures == connection_failures.len()
{
"All failures are related to OpportunisticInsecure STARTTLS. \
match result {
Err(err) => {
if OpportunisticInsecureTlsHandshakeError::is_match_anyhow(&err) {
num_opportunistic_tls_failures += 1;
}
connection_failures.push(format!("{err:#}"));
if !queue_dispatcher
.have_more_connection_candidates(&mut dispatcher)
.await
{
for msg in dispatcher.msgs.drain(..) {
let summary =
if num_opportunistic_tls_failures == connection_failures.len() {
"All failures are related to OpportunisticInsecure STARTTLS. \
Consider setting enable_tls=Disabled for this site. "
} else {
""
};
} else {
""
};
let response = Response {
code: 400,
enhanced_code: None,
content: format!(
"KumoMTA internal: \
let response = Response {
code: 400,
enhanced_code: None,
content: format!(
"KumoMTA internal: \
failed to connect to any candidate \
hosts: {summary}{}",
connection_failures.join(", ")
),
command: None,
};
connection_failures.join(", ")
),
command: None,
};
log_disposition(LogDisposition {
kind: RecordType::TransientFailure,
msg: msg.clone(),
site: &dispatcher.name,
peer_address: None,
response: response.clone(),
egress_pool: Some(&dispatcher.egress_pool),
egress_source: Some(&dispatcher.egress_source.name),
relay_disposition: None,
delivery_protocol: Some(&dispatcher.delivery_protocol),
tls_info: None,
source_address: None,
provider: dispatcher.path_config.borrow().provider_name.as_deref(),
session_id: Some(dispatcher.session_id),
recipient_list: None,
})
.await;
QueueManager::requeue_message(
msg,
IncrementAttempts::Yes,
None,
response,
InsertReason::LoggedTransientFailure.into(),
)
.await?;
dispatcher.metrics.inc_transfail();
}
log_disposition(LogDisposition {
kind: RecordType::TransientFailure,
msg: msg.clone(),
site: &dispatcher.name,
peer_address: None,
response: response.clone(),
egress_pool: Some(&dispatcher.egress_pool),
egress_source: Some(&dispatcher.egress_source.name),
relay_disposition: None,
delivery_protocol: Some(&dispatcher.delivery_protocol),
tls_info: None,
source_address: None,
provider: dispatcher.path_config.borrow().provider_name.as_deref(),
session_id: Some(dispatcher.session_id),
recipient_list: None,
})
.await;
QueueManager::requeue_message(
msg,
IncrementAttempts::Yes,
None,
response,
InsertReason::LoggedTransientFailure.into(),
)
.await?;
dispatcher.metrics.inc_transfail();
}
if consecutive_connection_failures.fetch_add(1, Ordering::SeqCst)
> dispatcher
.path_config
.borrow()
.consecutive_connection_failures_before_delay
{
dispatcher.delay_ready_queue().await;
if consecutive_connection_failures.fetch_add(1, Ordering::SeqCst)
> dispatcher
.path_config
.borrow()
.consecutive_connection_failures_before_delay
{
dispatcher.delay_ready_queue().await;
}
dispatcher.release_leases().await;
return Err(err);
}
tracing::debug!("{err:#}");
// Try the next candidate MX address
continue;
}
Ok(
AttemptConnectionDisposition::ReusedExisting
| AttemptConnectionDisposition::ConnectedNew,
) => {
// fall through to below logic to do the send
}
Ok(AttemptConnectionDisposition::PeerClosedConnectionNeedNewSession) => {
tracing::debug!(
"{} Peer closed the connection, will make new session",
dispatcher.name
);
dispatcher.release_leases().await;
return Err(err);
queue_dispatcher.close_connection(&mut dispatcher).await?;
// Push the message(s) batch into the ready queue so that they
// can be immediately tried in another session
dispatcher
.reinsert_ready_queue(InsertReason::PeerClosedConnection.into())
.await;
return Ok(());
}
Ok(AttemptConnectionDisposition::PeerClosedConnectionContinueSession) => {
tracing::debug!(
"{} Peer closed the connection, continue with session",
dispatcher.name
);
// We don't have a current connection, so continue around
// the loop to try to open a new one
continue;
}
tracing::debug!("{err:#}");
// Try the next candidate MX address
continue;
}
connection_failures.clear();
+54 -7
View File
@@ -4,7 +4,7 @@ use crate::http_server::admin_trace_smtp_client_v1::{
};
use crate::logging::disposition::{log_disposition, LogDisposition, RecordType};
use crate::queue::{IncrementAttempts, InsertReason, QueueManager, QueueState};
use crate::ready_queue::{Dispatcher, QueueDispatcher};
use crate::ready_queue::{AttemptConnectionDisposition, Dispatcher, QueueDispatcher};
use crate::smtp_server::ShuttingDownError;
use crate::spool::SpoolManager;
use anyhow::Context;
@@ -294,7 +294,10 @@ impl SmtpDispatcher {
}))
}
async fn attempt_connection_impl(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<()> {
async fn attempt_connection_impl(
&mut self,
dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition> {
if let Some(client) = &mut self.client {
if client.is_connected() {
// If we get a unilateral response here now it can either be:
@@ -311,23 +314,64 @@ impl SmtpDispatcher {
match client.check_unilateral_response().await {
Ok(None) => {
// Still connected
return Ok(());
return Ok(AttemptConnectionDisposition::ReusedExisting);
}
Ok(Some(response)) => {
// We got an explicit signal that the connection is closing
// out. We don't consider this to be a transport error
// so we'll close the current connection.
// If we've managed to send messages so far in this session,
// then we're consider it to be a successful plan so we should
// not advance to the next candidate in the plan and make a
// new session.
// However, if we haven't managed to send anything so far,
// closing and restarting the plan now would likely lead to
// a persistent recurrence of this same state, so in that
// situation we need to ensure that the reconnect_strategy
// is applied to decide what to do.
tracing::debug!(
"{} sent a unilateral response: \
{response:?}, treating the connection as closed",
dispatcher.name
);
let _ = self.close_connection(dispatcher).await;
// Decide whether this is a successful plan
self.update_state_for_reconnect(dispatcher);
// if so, we can/should close the current session and start
// a new one with a fresh plan
if self.terminated_ok {
let _ = self.close_connection(dispatcher).await;
return Ok(
AttemptConnectionDisposition::PeerClosedConnectionNeedNewSession,
);
}
// otherwise, continue to next candidate host if that is what
// the reconnect_strategy indicates.
// While this is a return that leaves this function,
// our caller will typically continue its loop and call
// back in, but on the next call we won't show as connected
// and will instead reach the logic below to make a new
// connection in the current session, if that is what
// the reconnect_strategy indicated.
return Ok(
AttemptConnectionDisposition::PeerClosedConnectionContinueSession,
);
}
Err(err) => {
// We got a transport error of some kind.
// update_state_for_reconnect applies any reconnect_strategy
// which will decide whether we continue with the connection
// plan or whether we need to go to a new session.
tracing::debug!(
"{} had error: {err:#} \
while checking for liveness, treating it as closed",
dispatcher.name
);
self.client.take();
self.update_state_for_reconnect(dispatcher);
return Ok(
AttemptConnectionDisposition::PeerClosedConnectionContinueSession,
);
}
};
}
@@ -845,7 +889,7 @@ impl SmtpDispatcher {
.replace(connection_wrapper.map_connection(client));
self.client_address.replace(address);
dispatcher.delivered_this_connection = 0;
Ok(())
Ok(AttemptConnectionDisposition::ConnectedNew)
}
async fn resolve_cached_client_cert(
@@ -971,7 +1015,10 @@ impl QueueDispatcher for SmtpDispatcher {
self.terminated_ok
}
async fn attempt_connection(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<()> {
async fn attempt_connection(
&mut self,
dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition> {
self.attempt_connection_impl(dispatcher)
.await
.map_err(|err| {
+8 -3
View File
@@ -6,7 +6,7 @@ use crate::logging::disposition::{log_disposition, LogDisposition, RecordType};
use crate::logging::rejection::{log_rejection, LogRejection};
use crate::metrics_helper::smtp_rejected_for_service;
use crate::queue::{DeliveryProto, IncrementAttempts, InsertReason, QueueConfig, QueueManager};
use crate::ready_queue::{Dispatcher, QueueDispatcher};
use crate::ready_queue::{AttemptConnectionDisposition, Dispatcher, QueueDispatcher};
use crate::spool::SpoolManager;
use anyhow::{anyhow, Context};
use async_trait::async_trait;
@@ -3432,12 +3432,17 @@ impl QueueDispatcher for DeferredSmtpInjectionDispatcher {
}
}
async fn attempt_connection(&mut self, dispatcher: &mut Dispatcher) -> anyhow::Result<()> {
async fn attempt_connection(
&mut self,
dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition> {
if self.connection.is_none() {
self.connection
.replace(dispatcher.metrics.wrap_connection(()));
Ok(AttemptConnectionDisposition::ConnectedNew)
} else {
Ok(AttemptConnectionDisposition::ReusedExisting)
}
Ok(())
}
async fn have_more_connection_candidates(&mut self, _dispatcher: &mut Dispatcher) -> bool {
+6 -3
View File
@@ -1,7 +1,7 @@
use crate::http_server::inject_v1::activity_for_peer;
use crate::logging::disposition::{log_disposition, LogDisposition};
use crate::queue::{DeliveryProto, QueueConfig, QueueManager};
use crate::ready_queue::{Dispatcher, QueueDispatcher};
use crate::ready_queue::{AttemptConnectionDisposition, Dispatcher, QueueDispatcher};
use crate::spool::SpoolManager;
use anyhow::Context;
use async_trait::async_trait;
@@ -149,8 +149,11 @@ impl QueueDispatcher for XferDispatcher {
Ok(true)
}
async fn attempt_connection(&mut self, _dispatcher: &mut Dispatcher) -> anyhow::Result<()> {
Ok(())
async fn attempt_connection(
&mut self,
_dispatcher: &mut Dispatcher,
) -> anyhow::Result<AttemptConnectionDisposition> {
Ok(AttemptConnectionDisposition::ReusedExisting)
}
async fn have_more_connection_candidates(&mut self, _dispatcher: &mut Dispatcher) -> bool {