Fix MTA: Re-scheduling or editing a queued message reports success but persists nothing for recipients in a non-default virtual queue

This commit is contained in:
Maurus Decimus
2026-06-15 12:20:49 +02:00
parent e07fb0a7c1
commit e559cf5ab5
4 changed files with 187 additions and 13 deletions

View File

@@ -81,7 +81,6 @@ pub(crate) async fn queued_message_set(
}
// Process patches
let prev_event = archived_message.inner.next_event(None);
let mut message = map_message(archived_message.inner);
message.next_retry = None;
for (key, value) in value.into_expanded_object() {
@@ -101,7 +100,9 @@ pub(crate) async fn queued_message_set(
// Process changes
let mut has_changes = false;
let mut modified_rcpts = AHashSet::new();
let mut queued_message = archived_message.deserialize()?;
let prev_events = queued_message.next_events();
if queued_message.env_id.as_deref() != message.env_id.as_deref() {
queued_message.env_id = message.env_id.as_deref().map(|v| v.into());
has_changes = true;
@@ -110,7 +111,7 @@ pub(crate) async fn queued_message_set(
queued_message.priority = message.priority as i16;
has_changes = true;
}
for rcpt in queued_message.recipients.iter_mut() {
for (idx, rcpt) in queued_message.recipients.iter_mut().enumerate() {
if !message
.recipients
.iter()
@@ -121,13 +122,15 @@ pub(crate) async fn queued_message_set(
details: queue::Error::Io("Delivery canceled.".into()),
});
has_changes = true;
modified_rcpts.insert(idx);
}
}
for (address, rcpt) in message.recipients.into_iter() {
let Some(queued_rcpt) = queued_message
let Some((idx, queued_rcpt)) = queued_message
.recipients
.iter_mut()
.find(|r| r.address.as_ref() == address.as_str())
.enumerate()
.find(|(_, r)| r.address.as_ref() == address.as_str())
else {
set.response.not_updated.append(
id,
@@ -136,9 +139,10 @@ pub(crate) async fn queued_message_set(
);
continue 'outer;
};
let mut changed = false;
if rcpt.orcpt.as_deref() != queued_rcpt.orcpt.as_deref() {
queued_rcpt.orcpt = rcpt.orcpt.as_deref().map(|v| v.into());
has_changes = true;
changed = true;
}
let expiry = match rcpt.expires {
QueueExpiry::Ttl(ttl) => common::config::smtp::queue::QueueExpiry::Ttl(
@@ -152,7 +156,7 @@ pub(crate) async fn queued_message_set(
};
if expiry != queued_rcpt.expires {
queued_rcpt.expires = expiry;
has_changes = true;
changed = true;
}
for (due, count, field) in [
@@ -165,7 +169,7 @@ pub(crate) async fn queued_message_set(
};
if schedule != *field {
*field = schedule;
has_changes = true;
changed = true;
}
}
@@ -175,7 +179,7 @@ pub(crate) async fn queued_message_set(
let new_due = next_retry.timestamp() as u64;
if queued_rcpt.retry.due != new_due {
queued_rcpt.retry.due = new_due;
has_changes = true;
changed = true;
}
}
@@ -183,7 +187,12 @@ pub(crate) async fn queued_message_set(
&& !matches!(queued_rcpt.status, Status::Scheduled)
{
queued_rcpt.status = Status::Scheduled;
changed = true;
}
if changed {
has_changes = true;
modified_rcpts.insert(idx);
}
}
@@ -196,9 +205,11 @@ pub(crate) async fn queued_message_set(
Status::TemporaryFailure(_) | Status::Scheduled
)
}) {
message.save_changes(set.server, prev_event).await
message
.save_registry_changes(set.server, prev_events, modified_rcpts)
.await
} else {
message.remove(set.server, prev_event).await
message.remove_registry(set.server, prev_events).await
};
if !is_success {

View File

@@ -13,7 +13,7 @@ use crate::queue::{
FROM_AUTHENTICATED, FROM_AUTOGENERATED, FROM_DSN, FROM_REPORT, FROM_UNAUTHENTICATED,
FROM_UNAUTHENTICATED_DMARC, MessageWrapper,
};
use ahash::AHashSet;
use ahash::{AHashMap, AHashSet};
use common::config::smtp::queue::{ArchivedQueueExpiry, QueueName};
use common::ipc::QueueEvent;
use common::network::RcptResolution;
@@ -34,7 +34,7 @@ use store::write::{
AlignedBytes, Archive, Archiver, BatchBuilder, BlobLink, BlobOp, MergeResult, Params,
QueueClass, RegistryClass, ValueClass, now,
};
use store::{Deserialize, IterateParams, Serialize, SerializeInfallible, U64_LEN, ValueKey};
use store::{Deserialize, IterateParams, Serialize, SerializeInfallible, U32_LEN, U64_LEN, ValueKey};
use trc::{AddContext, ServerEvent, SpamEvent};
use types::blob::BlobId;
use types::blob_hash::BlobHash;
@@ -787,6 +787,168 @@ impl MessageWrapper {
}
}
pub async fn save_registry_changes(
mut self,
server: &Server,
prev_events: AHashMap<QueueName, u64>,
modified_rcpts: AHashSet<usize>,
) -> bool {
let mut batch = BatchBuilder::new();
self.release_quota(&mut batch);
for (queue_name, due) in prev_events {
batch.clear(ValueClass::Queue(QueueClass::MessageEvent(
store::write::QueueEvent {
due,
queue_id: self.queue_id,
queue_name: queue_name.into_inner(),
},
)));
}
for (queue_name, due) in self.message.next_events() {
batch.set(
ValueClass::Queue(QueueClass::MessageEvent(store::write::QueueEvent {
due,
queue_id: self.queue_id,
queue_name: queue_name.into_inner(),
})),
Vec::new(),
);
}
let message_bytes = match Archiver::new(self.message).serialize() {
Ok(data) => data,
Err(err) => {
trc::error!(
err.details("Failed to serialize message.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
return false;
}
};
let mut modified_bytes = Vec::with_capacity(modified_rcpts.len() * U32_LEN);
for idx in modified_rcpts {
modified_bytes.extend_from_slice(&(idx as u32).to_be_bytes());
}
batch.merge_fnc(
ValueClass::Queue(QueueClass::Message(self.queue_id)),
Params::with_capacity(3)
.with_u64(self.queue_id)
.with_bytes(modified_bytes)
.with_bytes(message_bytes),
|params, _, bytes| {
let mut cur_message = <Archive<AlignedBytes> as Deserialize>::deserialize(
bytes.ok_or_else(|| {
trc::StoreEvent::NotFound
.into_err()
.details("Message no longer exists.")
.caused_by(trc::location!())
.ctx(trc::Key::QueueId, params.u64(0))
})?,
)
.and_then(|archive| archive.deserialize::<Message>())
.caused_by(trc::location!())?;
let new_message_ =
<Archive<AlignedBytes> as Deserialize>::deserialize(params.bytes(2))
.caused_by(trc::location!())?;
let new_message = new_message_
.unarchive::<Message>()
.caused_by(trc::location!())?;
if cur_message.blob_hash.as_slice() == new_message.blob_hash.0.as_slice()
&& cur_message.recipients.len() == new_message.recipients.len()
{
cur_message.priority = new_message.priority.to_native();
cur_message.env_id = new_message.env_id.as_ref().map(|v| v.as_ref().into());
for idx in params.bytes(1).chunks_exact(U32_LEN) {
let rcpt_idx = u32::from_be_bytes(idx.try_into().unwrap()) as usize;
if let Some(rcpt) = new_message.recipients.get(rcpt_idx) {
cur_message.recipients[rcpt_idx] =
rkyv_deserialize(rcpt).caused_by(trc::location!())?;
}
}
Archiver::new(cur_message)
.serialize()
.caused_by(trc::location!())
.map(MergeResult::Update)
} else {
Err(trc::StoreEvent::UnexpectedError
.into_err()
.details("Message blob hash or recipient count mismatch.")
.caused_by(trc::location!())
.ctx(trc::Key::QueueId, params.u64(0)))
}
},
);
if let Err(err) = server.store().write(batch.build_all()).await {
trc::error!(
err.details("Failed to save changes.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
false
} else {
true
}
}
pub async fn remove_registry(
self,
server: &Server,
prev_events: AHashMap<QueueName, u64>,
) -> bool {
let mut batch = BatchBuilder::new();
for (queue_name, due) in prev_events {
batch.clear(ValueClass::Queue(QueueClass::MessageEvent(
store::write::QueueEvent {
due,
queue_id: self.queue_id,
queue_name: queue_name.into_inner(),
},
)));
}
for quota_key in self.message.quota_keys {
match quota_key {
QuotaKey::Count { key, .. } => {
batch.add(ValueClass::Queue(QueueClass::QuotaCount(key.to_vec())), -1);
}
QuotaKey::Size { key, .. } => {
batch.add(
ValueClass::Queue(QueueClass::QuotaSize(key.to_vec())),
-(self.message.size as i64),
);
}
}
}
batch
.clear(BlobOp::Link {
hash: self.message.blob_hash.clone(),
to: BlobLink::Id { id: self.queue_id },
})
.clear(ValueClass::Queue(QueueClass::Message(self.queue_id)));
if let Err(err) = server.store().write(batch.build_all()).await {
trc::error!(
err.details("Failed to write to update queue.")
.span_id(self.span_id)
.caused_by(trc::location!())
);
false
} else {
true
}
}
pub fn has_domain(&self, domains: &[String]) -> bool {
self.message.recipients.iter().any(|r| {
let domain = r.address.domain_part();