diff --git a/CHANGELOG.md b/CHANGELOG.md index 12e0ae1d..63f24316 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,7 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If - JMAP: `*/changes` methods leak ids of non-shared objects (reported by @5ud0er). - Sieve: Do not allow invalid certs in `http_header` function. - FoundationDB: Fix read version cache expiration logic. +- MTA: Re-scheduling or editing a queued message reports success but persists nothing for recipients in a non-`default` virtual queue (reported by @DorianCoding). ## [0.16.8] - 2026-06-06 diff --git a/crates/jmap/src/registry/mapping/queued_message.rs b/crates/jmap/src/registry/mapping/queued_message.rs index 35cb7494..ff980685 100644 --- a/crates/jmap/src/registry/mapping/queued_message.rs +++ b/crates/jmap/src/registry/mapping/queued_message.rs @@ -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 { diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index aaae4921..fe1f798d 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -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, + modified_rcpts: AHashSet, + ) -> 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 = 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::()) + .caused_by(trc::location!())?; + + let new_message_ = + as Deserialize>::deserialize(params.bytes(2)) + .caused_by(trc::location!())?; + let new_message = new_message_ + .unarchive::() + .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, + ) -> 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(); diff --git a/tests/src/smtp/management/queue.rs b/tests/src/smtp/management/queue.rs index fa4d1d3a..cda46abd 100644 --- a/tests/src/smtp/management/queue.rs +++ b/tests/src/smtp/management/queue.rs @@ -87,7 +87,7 @@ async fn manage_queue() { .await; let queue_id = admin .registry_create_object(MtaVirtualQueue { - name: "default".into(), + name: "myqueue".into(), threads_per_node: 25, description: None, })