From c6904e0f5de91abc1cf569191a6c89ac2ec1a209 Mon Sep 17 00:00:00 2001 From: mdecimus Date: Wed, 16 Jul 2025 15:55:20 +0200 Subject: [PATCH] Fix queue event migration (closes #1833) --- crates/migration/src/queue.rs | 246 ++++++++++++++++------------ crates/smtp/src/queue/manager.rs | 2 +- crates/smtp/src/queue/spool.rs | 11 +- crates/trc/src/event/description.rs | 4 +- 4 files changed, 148 insertions(+), 115 deletions(-) diff --git a/crates/migration/src/queue.rs b/crates/migration/src/queue.rs index 88adced2..1775992e 100644 --- a/crates/migration/src/queue.rs +++ b/crates/migration/src/queue.rs @@ -7,7 +7,7 @@ use crate::LegacyBincode; use common::{ Server, - config::smtp::queue::{QueueExpiry, QueueName}, + config::smtp::queue::{DEFAULT_QUEUE_NAME, QueueExpiry, QueueName}, }; use smtp::queue::{ Error, ErrorDetails, HostResponse, Message, QueueId, QuotaKey, Recipient, Schedule, Status, @@ -15,64 +15,22 @@ use smtp::queue::{ }; use std::net::{IpAddr, Ipv4Addr}; use store::{ - IterateParams, Serialize, U64_LEN, ValueKey, - ahash::AHashSet, + IterateParams, SUBSPACE_QUEUE_EVENT, Serialize, U64_LEN, ValueKey, + ahash::AHashMap, write::{ - AlignedBytes, Archive, Archiver, BatchBuilder, QueueClass, ValueClass, - key::DeserializeBigEndian, now, + AlignedBytes, AnyClass, Archive, Archiver, BatchBuilder, QueueClass, ValueClass, + key::{DeserializeBigEndian, KeySerializer}, + now, }, }; use trc::AddContext; use utils::BlobHash; pub(crate) async fn migrate_queue_v011(server: &Server) -> trc::Result<()> { - let from_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( - store::write::QueueEvent { - due: 0, - queue_id: 0, - queue_name: [0; 8], - }, - ))); - let to_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( - store::write::QueueEvent { - due: u64::MAX, - queue_id: u64::MAX, - queue_name: [u8::MAX; 8], - }, - ))); - - let mut queue_ids = AHashSet::new(); - server - .store() - .iterate( - IterateParams::new(from_key, to_key).ascending().no_values(), - |key, _| { - queue_ids.insert(key.deserialize_be_u64(U64_LEN)?); - - Ok(true) - }, - ) - .await - .caused_by(trc::location!())?; - - let from_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(0))); - let to_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(u64::MAX))); - server - .store() - .iterate( - IterateParams::new(from_key, to_key).ascending().no_values(), - |key, _| { - queue_ids.insert(key.deserialize_be_u64(0)?); - - Ok(true) - }, - ) - .await - .caused_by(trc::location!())?; - let mut count = 0; + let now = now(); - for queue_id in queue_ids { + for (queue_id, due) in get_queue_events(server).await? { match server .store() .get_value::>(ValueKey::from(ValueClass::Queue( @@ -82,12 +40,28 @@ pub(crate) async fn migrate_queue_v011(server: &Server) -> trc::Result<()> { { Ok(Some(bincoded)) => { let mut batch = BatchBuilder::new(); - batch.set( - ValueClass::Queue(QueueClass::Message(queue_id)), - Archiver::new(Message::from(bincoded.inner)) - .serialize() - .caused_by(trc::location!())?, - ); + let message = Message::from(bincoded.inner); + if let Some(due) = due { + batch.clear(ValueClass::Any(AnyClass { + subspace: SUBSPACE_QUEUE_EVENT, + key: KeySerializer::new(16).write(due).write(queue_id).finalize(), + })); + } + batch + .set( + ValueClass::Queue(QueueClass::MessageEvent(store::write::QueueEvent { + due: due.unwrap_or(now), + queue_id, + queue_name: DEFAULT_QUEUE_NAME.into_inner(), + })), + vec![], + ) + .set( + ValueClass::Queue(QueueClass::Message(queue_id)), + Archiver::new(message) + .serialize() + .caused_by(trc::location!())?, + ); count += 1; server .store() @@ -95,7 +69,20 @@ pub(crate) async fn migrate_queue_v011(server: &Server) -> trc::Result<()> { .await .caused_by(trc::location!())?; } - Ok(None) => (), + Ok(None) => { + if let Some(due) = due { + let mut batch = BatchBuilder::new(); + batch.clear(ValueClass::Any(AnyClass { + subspace: SUBSPACE_QUEUE_EVENT, + key: KeySerializer::new(16).write(due).write(queue_id).finalize(), + })); + server + .store() + .write(batch.build_all()) + .await + .caused_by(trc::location!())?; + } + } Err(err) => { if server .store() @@ -124,53 +111,10 @@ pub(crate) async fn migrate_queue_v011(server: &Server) -> trc::Result<()> { } pub(crate) async fn migrate_queue_v012(server: &Server) -> trc::Result<()> { - let from_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( - store::write::QueueEvent { - due: 0, - queue_id: 0, - queue_name: [0; 8], - }, - ))); - let to_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( - store::write::QueueEvent { - due: u64::MAX, - queue_id: u64::MAX, - queue_name: [u8::MAX; 8], - }, - ))); - - let mut queue_ids = AHashSet::new(); - server - .store() - .iterate( - IterateParams::new(from_key, to_key).ascending().no_values(), - |key, _| { - queue_ids.insert(key.deserialize_be_u64(U64_LEN)?); - - Ok(true) - }, - ) - .await - .caused_by(trc::location!())?; - - let from_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(0))); - let to_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(u64::MAX))); - server - .store() - .iterate( - IterateParams::new(from_key, to_key).ascending().no_values(), - |key, _| { - queue_ids.insert(key.deserialize_be_u64(0)?); - - Ok(true) - }, - ) - .await - .caused_by(trc::location!())?; - let mut count = 0; + let now = now(); - for queue_id in queue_ids { + for (queue_id, due) in get_queue_events(server).await? { match server .store() .get_value::>(ValueKey::from(ValueClass::Queue( @@ -185,13 +129,29 @@ pub(crate) async fn migrate_queue_v012(server: &Server) -> trc::Result<()> { } }) { Ok(Some(archive)) => { + let message = Message::from(archive); let mut batch = BatchBuilder::new(); - batch.set( - ValueClass::Queue(QueueClass::Message(queue_id)), - Archiver::new(Message::from(archive)) - .serialize() - .caused_by(trc::location!())?, - ); + if let Some(due) = due { + batch.clear(ValueClass::Any(AnyClass { + subspace: SUBSPACE_QUEUE_EVENT, + key: KeySerializer::new(16).write(due).write(queue_id).finalize(), + })); + } + batch + .set( + ValueClass::Queue(QueueClass::MessageEvent(store::write::QueueEvent { + due: due.unwrap_or(now), + queue_id, + queue_name: DEFAULT_QUEUE_NAME.into_inner(), + })), + vec![], + ) + .set( + ValueClass::Queue(QueueClass::Message(queue_id)), + Archiver::new(message) + .serialize() + .caused_by(trc::location!())?, + ); count += 1; server .store() @@ -199,7 +159,20 @@ pub(crate) async fn migrate_queue_v012(server: &Server) -> trc::Result<()> { .await .caused_by(trc::location!())?; } - Ok(None) => (), + Ok(None) => { + if let Some(due) = due { + let mut batch = BatchBuilder::new(); + batch.clear(ValueClass::Any(AnyClass { + subspace: SUBSPACE_QUEUE_EVENT, + key: KeySerializer::new(16).write(due).write(queue_id).finalize(), + })); + server + .store() + .write(batch.build_all()) + .await + .caused_by(trc::location!())?; + } + } Err(err) => { if server .store() @@ -234,6 +207,61 @@ pub(crate) async fn migrate_queue_v012(server: &Server) -> trc::Result<()> { Ok(()) } +async fn get_queue_events(server: &Server) -> trc::Result>> { + let from_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( + store::write::QueueEvent { + due: 0, + queue_id: 0, + queue_name: [0; 8], + }, + ))); + let to_key = ValueKey::from(ValueClass::Queue(QueueClass::MessageEvent( + store::write::QueueEvent { + due: u64::MAX, + queue_id: u64::MAX, + queue_name: [u8::MAX; 8], + }, + ))); + + let mut queue_ids: AHashMap> = AHashMap::new(); + server + .store() + .iterate( + IterateParams::new(from_key, to_key).ascending().no_values(), + |key, _| { + queue_ids.insert( + key.deserialize_be_u64(U64_LEN)?, + Some(key.deserialize_be_u64(0)?), + ); + + Ok(true) + }, + ) + .await + .caused_by(trc::location!())?; + + let from_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(0))); + let to_key = ValueKey::from(ValueClass::Queue(QueueClass::Message(u64::MAX))); + server + .store() + .iterate( + IterateParams::new(from_key, to_key).ascending().no_values(), + |key, _| { + let queue_id = key.deserialize_be_u64(0)?; + + if !queue_ids.contains_key(&queue_id) { + queue_ids.insert(queue_id, None); + } + + Ok(true) + }, + ) + .await + .caused_by(trc::location!())?; + + Ok(queue_ids) +} + impl From> for Message where SIZE: AsU64, diff --git a/crates/smtp/src/queue/manager.rs b/crates/smtp/src/queue/manager.rs index 63374a3f..7cc5e42f 100644 --- a/crates/smtp/src/queue/manager.rs +++ b/crates/smtp/src/queue/manager.rs @@ -62,7 +62,7 @@ impl Queue { locked: AHashMap::with_capacity(128), locked_revision: 0, stats: AHashMap::new(), - next_refresh: Instant::now(), + next_refresh: Instant::now() + Duration::from_secs(1), is_paused: false, rx, } diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index 19cab484..33936445 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -126,9 +126,14 @@ impl SmtpSpool for Server { if due <= now { let queue_id = key.deserialize_be_u64(U64_LEN)?; - let queue_name = - QueueName::from_bytes(key.get(U64_LEN + U64_LEN..).unwrap_or_default()) - .unwrap_or_default(); + let queue_name = key + .get(U64_LEN + U64_LEN..) + .and_then(QueueName::from_bytes) + .ok_or_else(|| { + trc::StoreEvent::DataCorruption + .caused_by(trc::location!()) + .ctx(trc::Key::Key, key) + })?; let add_event = queue .stats diff --git a/crates/trc/src/event/description.rs b/crates/trc/src/event/description.rs index ce39841d..e8488605 100644 --- a/crates/trc/src/event/description.rs +++ b/crates/trc/src/event/description.rs @@ -396,7 +396,7 @@ impl SmtpEvent { pub fn description(&self) -> &'static str { match self { SmtpEvent::Error => "SMTP error occurred", - SmtpEvent::IdNotFound => "Remote host ID not found", + SmtpEvent::IdNotFound => "Strategy not found", SmtpEvent::ConcurrencyLimitExceeded => "Concurrency limit exceeded", SmtpEvent::TransferLimitExceeded => "Transfer limit exceeded", SmtpEvent::RateLimitExceeded => "Rate limit exceeded", @@ -483,7 +483,7 @@ impl SmtpEvent { pub fn explain(&self) -> &'static str { match self { SmtpEvent::Error => "An error occurred during an SMTP command", - SmtpEvent::IdNotFound => "The remote server ID was not found in the configuration", + SmtpEvent::IdNotFound => "The strategy ID was not found in the configuration", SmtpEvent::ConcurrencyLimitExceeded => "The concurrency limit was exceeded", SmtpEvent::TransferLimitExceeded => { "The remote host transferred more data than allowed"