From 21c541e505c13ff9bb47ea471a9da8c8fbe29877 Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Sat, 21 Mar 2026 19:51:25 +0100 Subject: [PATCH] Registry testing - part 11 --- crates/common/src/config/inner.rs | 2 + crates/common/src/config/mailstore/scripts.rs | 8 + crates/common/src/config/smtp/queue.rs | 22 +- crates/common/src/lib.rs | 1 + crates/common/src/storage/index.rs | 10 +- crates/common/src/telemetry/metrics/store.rs | 55 +- crates/email/src/message/copy.rs | 11 +- crates/email/src/message/ingest.rs | 30 +- crates/http/src/management/diagnose.rs | 505 +++++++++--------- crates/services/src/task_manager/index.rs | 7 +- .../src/task_manager/merge_threads.rs | 142 +++-- crates/smtp/src/outbound/lookup.rs | 2 +- crates/smtp/src/outbound/mod.rs | 15 +- crates/smtp/src/queue/spool.rs | 10 +- crates/store/src/write/batch.rs | 3 +- crates/utils/src/snowflake.rs | 8 +- tests/src/jmap/mail/query.rs | 7 +- tests/src/jmap/mail/sieve_script.rs | 18 + tests/src/jmap/mail/thread_merge.rs | 75 ++- tests/src/jmap/mod.rs | 26 +- tests/src/utils/server.rs | 24 +- tests/src/utils/storage.rs | 2 +- 22 files changed, 519 insertions(+), 464 deletions(-) diff --git a/crates/common/src/config/inner.rs b/crates/common/src/config/inner.rs index 481e1ec5..c88b2fae 100644 --- a/crates/common/src/config/inner.rs +++ b/crates/common/src/config/inner.rs @@ -80,6 +80,7 @@ impl Data { blocked_ips: RwLock::new(blocked_ips), jmap_id_gen: id_generator.clone(), queue_id_gen: id_generator.clone(), + registry_id_gen: id_generator.clone(), span_id_gen: id_generator, queue_status: true.into(), applications: Default::default(), @@ -220,6 +221,7 @@ impl Default for Data { jmap_id_gen: Default::default(), queue_id_gen: Default::default(), span_id_gen: Default::default(), + registry_id_gen: Default::default(), queue_status: true.into(), applications: Default::default(), logos: Default::default(), diff --git a/crates/common/src/config/mailstore/scripts.rs b/crates/common/src/config/mailstore/scripts.rs index d305db3c..2f75fd2a 100644 --- a/crates/common/src/config/mailstore/scripts.rs +++ b/crates/common/src/config/mailstore/scripts.rs @@ -135,6 +135,10 @@ impl Scripting { // Parse trusted scripts let mut trusted_scripts = AHashMap::new(); for script in bp.list_infallible::().await { + if !script.object.is_active { + continue; + } + match trusted_compiler.compile(script.object.contents.as_bytes()) { Ok(compiled) => { trusted_scripts.insert(script.object.name, compiled.into()); @@ -151,6 +155,10 @@ impl Scripting { // Parse untrusted scripts let mut untrusted_scripts = AHashMap::new(); for script in bp.list_infallible::().await { + if !script.object.is_active { + continue; + } + match untrusted_compiler.compile(script.object.contents.as_bytes()) { Ok(compiled) => { untrusted_scripts.insert(script.object.name, compiled.into()); diff --git a/crates/common/src/config/smtp/queue.rs b/crates/common/src/config/smtp/queue.rs index c2e4d7b4..8fd039bc 100644 --- a/crates/common/src/config/smtp/queue.rs +++ b/crates/common/src/config/smtp/queue.rs @@ -179,7 +179,7 @@ pub struct QueueQuota { #[derive(Clone, Hash, PartialEq, Eq)] pub struct RelayConfig { - pub address: HostOrIp, + pub address: HostOrIp, IpStr>, pub port: u16, pub protocol: ServerProtocol, pub auth: Option>, @@ -188,9 +188,15 @@ pub struct RelayConfig { } #[derive(Clone, Debug, Hash, PartialEq, Eq)] -pub enum HostOrIp { - Host(T), - Ip { ip: IpAddr, ip_str: String }, +pub enum HostOrIp { + Host(N), + Ip(I), +} + +#[derive(Clone, Debug, Hash, PartialEq, Eq)] +pub struct IpStr { + pub ip: IpAddr, + pub ip_str: Box, } #[derive(Debug, Clone, Copy, Default)] @@ -375,12 +381,12 @@ impl QueueConfig { route.name, RoutingStrategy::Relay(RelayConfig { address: if let Ok(ip) = route.address.parse() { - HostOrIp::Ip { + HostOrIp::Ip(IpStr { ip, - ip_str: route.address, - } + ip_str: route.address.into(), + }) } else { - HostOrIp::Host(route.address) + HostOrIp::Host(route.address.into()) }, port: route.port as u16, protocol: match route.protocol { diff --git a/crates/common/src/lib.rs b/crates/common/src/lib.rs index 9fb074b0..40c1e40b 100644 --- a/crates/common/src/lib.rs +++ b/crates/common/src/lib.rs @@ -154,6 +154,7 @@ pub struct Data { pub jmap_id_gen: SnowflakeIdGenerator, pub queue_id_gen: SnowflakeIdGenerator, pub span_id_gen: SnowflakeIdGenerator, + pub registry_id_gen: SnowflakeIdGenerator, pub queue_status: AtomicBool, pub applications: WebApplicationManager, diff --git a/crates/common/src/storage/index.rs b/crates/common/src/storage/index.rs index 95c95936..9b305caa 100644 --- a/crates/common/src/storage/index.rs +++ b/crates/common/src/storage/index.rs @@ -469,10 +469,7 @@ fn build_index( let object_account_id = batch.last_account_id().unwrap_or_default(); let object_type = batch.last_collection().unwrap_or(Collection::None); let object_id = batch.last_document_id().unwrap_or_default(); - let notification_id = SnowflakeIdGenerator::from_sequence_id( - object_type as u64 ^ object_account_id as u64, - ) - .unwrap_or_default(); + let notification_id = SnowflakeIdGenerator::global_id().unwrap_or_default(); for item in value.as_ref() { if set { @@ -630,10 +627,7 @@ fn merge_index( let object_account_id = batch.last_account_id().unwrap_or_default(); let object_type = batch.last_collection().unwrap_or(Collection::None); let object_id = batch.last_document_id().unwrap_or_default(); - let notification_id = SnowflakeIdGenerator::from_sequence_id( - object_type as u64 ^ object_account_id as u64, - ) - .unwrap_or_default(); + let notification_id = SnowflakeIdGenerator::global_id().unwrap_or_default(); match (has_old_acl, has_new_acl) { (true, true) => { diff --git a/crates/common/src/telemetry/metrics/store.rs b/crates/common/src/telemetry/metrics/store.rs index 03d372b8..30c3c3d3 100644 --- a/crates/common/src/telemetry/metrics/store.rs +++ b/crates/common/src/telemetry/metrics/store.rs @@ -35,6 +35,7 @@ pub trait MetricsStore: Sync + Send { pub struct MetricsHistory { events: AHashMap, histograms: AHashMap, + id_generator: SnowflakeIdGenerator, } #[derive(Default)] @@ -53,7 +54,7 @@ impl MetricsStore for Store { ) -> trc::Result<()> { let mut batch = BatchBuilder::new(); { - let mut history = history_.lock(); + let mut history_guard = history_.lock(); for event in [ MetricType::SmtpConnectionStart, MetricType::ImapConnectionStart, @@ -80,27 +81,20 @@ impl MetricsStore for Store { ] { let reading = Collector::read_metric_counter(event.event_id()); if reading > 0 { - let history = history.events.entry(event).or_insert(0); + let history = history_guard.events.entry(event).or_insert(0); let diff = reading - *history; + *history = reading; if diff > 0 { #[cfg(not(feature = "test_mode"))] - let metric_id = - SnowflakeIdGenerator::from_sequence_id(event.to_id() as u64) - .unwrap_or_default(); + let metric_id = history_guard.id_generator.generate(); #[cfg(feature = "test_mode")] let metric_id = _timestamp .map(|timestamp| { - SnowflakeIdGenerator::from_timestamp_and_sequence_id( - timestamp, - event.to_id() as u64, - ) + SnowflakeIdGenerator::global_id_from_timestamp(timestamp).unwrap() }) - .unwrap_or_else(|| { - SnowflakeIdGenerator::from_sequence_id(event.to_id() as u64) - }) - .unwrap_or_default(); + .unwrap_or_else(|| history_guard.id_generator.generate()); batch.set( ValueClass::Telemetry(TelemetryClass::Metric(metric_id)), @@ -111,7 +105,6 @@ impl MetricsStore for Store { .to_pickled_vec(), ); } - *history = reading; } } @@ -121,22 +114,14 @@ impl MetricsStore for Store { let value = gauge.get(); if value > 0 { #[cfg(not(feature = "test_mode"))] - let metric_id = - SnowflakeIdGenerator::from_sequence_id(metric.to_id() as u64) - .unwrap_or_default(); + let metric_id = history_guard.id_generator.generate(); #[cfg(feature = "test_mode")] let metric_id = _timestamp .map(|timestamp| { - SnowflakeIdGenerator::from_timestamp_and_sequence_id( - timestamp, - metric.to_id() as u64, - ) + SnowflakeIdGenerator::global_id_from_timestamp(timestamp).unwrap() }) - .unwrap_or_else(|| { - SnowflakeIdGenerator::from_sequence_id(metric.to_id() as u64) - }) - .unwrap_or_default(); + .unwrap_or_else(|| history_guard.id_generator.generate()); batch.set( ValueClass::Telemetry(TelemetryClass::Metric(metric_id)), @@ -160,37 +145,29 @@ impl MetricsStore for Store { | MetricType::DeliveryAttemptTime | MetricType::DnsLookupTime ) { - let history = history.histograms.entry(metric).or_default(); + let history = history_guard.histograms.entry(metric).or_default(); let sum = histogram.sum(); let count = histogram.count(); let diff_sum = sum - history.sum; let diff_count = count - history.count; + history.sum = sum; + history.count = count; if diff_sum > 0 || diff_count > 0 { #[cfg(not(feature = "test_mode"))] - let metric_id = - SnowflakeIdGenerator::from_sequence_id(metric.to_id() as u64) - .unwrap_or_default(); + let metric_id = history_guard.id_generator.generate(); #[cfg(feature = "test_mode")] let metric_id = _timestamp .map(|timestamp| { - SnowflakeIdGenerator::from_timestamp_and_sequence_id( - timestamp, - metric.to_id() as u64, - ) + SnowflakeIdGenerator::global_id_from_timestamp(timestamp).unwrap() }) - .unwrap_or_else(|| { - SnowflakeIdGenerator::from_sequence_id(metric.to_id() as u64) - }) - .unwrap_or_default(); + .unwrap_or_else(|| history_guard.id_generator.generate()); batch.set( ValueClass::Telemetry(TelemetryClass::Metric(metric_id)), Metric::Histogram(MetricSum { count, metric, sum }).to_pickled_vec(), ); } - history.sum = sum; - history.count = count; } } } diff --git a/crates/email/src/message/copy.rs b/crates/email/src/message/copy.rs index f64634bd..5418cf70 100644 --- a/crates/email/src/message/copy.rs +++ b/crates/email/src/message/copy.rs @@ -220,16 +220,9 @@ impl EmailCopy for Server { if !thread_result.merge_ids.is_empty() { batch.schedule_task(Task::MergeThreads(TaskMergeThreads { account_id: to_account_id.into(), - document_id: document_id.into(), status: TaskStatus::now(), - thread_ids: Map::new( - thread_result - .merge_ids - .into_iter() - .map(|id| id.into()) - .collect(), - ), - thread_hash: thread_result.thread_hash.to_string(), + thread_name: thread_result.thread_hash.to_string(), + message_ids: Map::new(message_ids.into_iter().map(|id| id.to_string()).collect()), })); } diff --git a/crates/email/src/message/ingest.rs b/crates/email/src/message/ingest.rs index 1dcb76f9..dc0afca6 100644 --- a/crates/email/src/message/ingest.rs +++ b/crates/email/src/message/ingest.rs @@ -34,6 +34,7 @@ use registry::{ }; use std::future::Future; use std::{borrow::Cow, cmp::Ordering, fmt::Write, time::Instant}; +use store::write::{AlignedBytes, Archive, RegistryClass}; use store::{ IndexKeyPrefix, IterateParams, SerializeInfallible, U32_LEN, ValueKey, ahash::AHashMap, @@ -42,10 +43,6 @@ use store::{ key::DeserializeBigEndian, now, }, }; -use store::{ - write::{AlignedBytes, Archive, RegistryClass}, - xxhash_rust, -}; use trc::{AddContext, MessageIngestEvent, SpamEvent}; use types::{ blob::{BlobClass, BlobId}, @@ -660,19 +657,17 @@ impl EmailIngest for Server { } // Merge threads if necessary - if !thread_result.merge_ids.is_empty() { + if !thread_result.merge_ids.is_empty() + || matches!( + params.source, + IngestSource::Jmap { .. } | IngestSource::Imap { .. } + ) + { batch.schedule_task(Task::MergeThreads(TaskMergeThreads { account_id: account_id.into(), - document_id: document_id.into(), status: TaskStatus::now(), - thread_ids: Map::new( - thread_result - .merge_ids - .into_iter() - .map(|id| id.into()) - .collect(), - ), - thread_hash: thread_result.thread_hash.to_string(), + thread_name: thread_result.thread_hash.to_string(), + message_ids: Map::new(message_ids.into_iter().map(|id| id.to_string()).collect()), })); } @@ -942,10 +937,7 @@ impl EmailIngest for Server { .to_pickled_vec(); let object_id = ObjectType::SpamTrainingSample.to_id(); - let item_id = SnowflakeIdGenerator::from_sequence_id(xxhash_rust::xxh3::xxh3_64( - sample.as_slice(), - )) - .unwrap_or_default(); + let item_id = SnowflakeIdGenerator::global_id().unwrap_or_default(); batch .set( BlobOp::Link { @@ -979,7 +971,7 @@ impl EmailIngest for Server { } } -fn has_message_id(a: &[CheekyHash], b: &[u8]) -> bool { +pub fn has_message_id(a: &[CheekyHash], b: &[u8]) -> bool { let mut i = 0; let mut j = 0; diff --git a/crates/http/src/management/diagnose.rs b/crates/http/src/management/diagnose.rs index 2a5014eb..a8b99e05 100644 --- a/crates/http/src/management/diagnose.rs +++ b/crates/http/src/management/diagnose.rs @@ -429,245 +429,255 @@ async fn delivery_diagnose( tx.send(DeliveryStage::IpLookupStart).await?; let now = Instant::now(); - let hostname = match host.fqdn_hostname() { - HostOrIp::Host(host) => host.into_owned(), - HostOrIp::Ip { ip_str, .. } => ip_str, + let remote_ips = match host.fqdn_hostname() { + HostOrIp::Host(hostname) => { + match server + .ip_lookup(&hostname, IpLookupStrategy::Ipv4thenIpv6, usize::MAX) + .await + { + Ok(remote_ips) if !remote_ips.is_empty() => remote_ips, + Ok(_) => { + tx.send(DeliveryStage::IpLookupError { + reason: "No IP addresses found for host".to_string(), + elapsed: now.elapsed_ms(), + }) + .await?; + continue; + } + Err(err) => { + tx.send(DeliveryStage::IpLookupError { + reason: err.to_string(), + elapsed: now.elapsed_ms(), + }) + .await?; + continue; + } + } + } + HostOrIp::Ip(ip) => vec![ip], }; - match server - .ip_lookup(&hostname, IpLookupStrategy::Ipv4thenIpv6, usize::MAX) - .await - { - Ok(remote_ips) if !remote_ips.is_empty() => { - tx.send(DeliveryStage::IpLookupSuccess { - remote_ips: remote_ips.clone(), - elapsed: now.elapsed_ms(), - }) + + tx.send(DeliveryStage::IpLookupSuccess { + remote_ips: remote_ips.clone(), + elapsed: now.elapsed_ms(), + }) + .await?; + + for remote_ip in remote_ips { + // Start connection + tx.send(DeliveryStage::ConnectionStart { remote_ip }) .await?; - for remote_ip in remote_ips { - // Start connection - tx.send(DeliveryStage::ConnectionStart { remote_ip }) - .await?; + let now = Instant::now(); + match SmtpClient::connect(SocketAddr::new(remote_ip, 25), timeout, 0).await { + Ok(mut client) => { + tx.send(DeliveryStage::ConnectionSuccess { + elapsed: now.elapsed_ms(), + }) + .await?; + + // Read greeting + tx.send(DeliveryStage::ReadGreetingStart).await?; let now = Instant::now(); - match SmtpClient::connect(SocketAddr::new(remote_ip, 25), timeout, 0).await { - Ok(mut client) => { - tx.send(DeliveryStage::ConnectionSuccess { + if let Err(status) = client.read_greeting(hostname).await { + tx.send(DeliveryStage::ReadGreetingError { + elapsed: now.elapsed_ms(), + reason: status.to_string(), + }) + .await?; + + continue; + } + tx.send(DeliveryStage::ReadGreetingSuccess { + elapsed: now.elapsed_ms(), + }) + .await?; + + // Say EHLO + tx.send(DeliveryStage::EhloStart).await?; + + let now = Instant::now(); + let capabilities = match tokio::time::timeout(timeout, async { + client + .stream + .write_all(format!("EHLO {local_host}\r\n",).as_bytes()) + .await?; + client.stream.flush().await?; + client.read_ehlo().await + }) + .await + { + Ok(Ok(capabilities)) => { + tx.send(DeliveryStage::EhloSuccess { elapsed: now.elapsed_ms(), }) .await?; - // Read greeting - tx.send(DeliveryStage::ReadGreetingStart).await?; + capabilities + } + Ok(Err(err)) => { + tx.send(DeliveryStage::EhloError { + elapsed: now.elapsed_ms(), + reason: err.to_string(), + }) + .await?; - let now = Instant::now(); - if let Err(status) = client.read_greeting(&hostname).await { - tx.send(DeliveryStage::ReadGreetingError { + continue; + } + Err(_) => { + tx.send(DeliveryStage::EhloError { + elapsed: now.elapsed_ms(), + reason: "Timed out reading response".to_string(), + }) + .await?; + + continue; + } + }; + + // Start TLS + tx.send(DeliveryStage::StartTlsStart).await?; + + let now = Instant::now(); + let mut client = match client + .try_start_tls( + &server.inner.data.smtp_connectors.pki_verify, + hostname, + &capabilities, + ) + .await + { + StartTlsResult::Success { smtp_client } => { + tx.send(DeliveryStage::StartTlsSuccess { + elapsed: now.elapsed_ms(), + }) + .await?; + + smtp_client + } + StartTlsResult::Error { error } => { + tx.send(DeliveryStage::StartTlsError { + elapsed: now.elapsed_ms(), + reason: error.to_string(), + }) + .await?; + + continue; + } + StartTlsResult::Unavailable { response, .. } => { + tx.send(DeliveryStage::StartTlsError { + elapsed: now.elapsed_ms(), + reason: response.map(|r| r.to_string()).unwrap_or_else(|| { + "STARTTLS not advertised by host".to_string() + }), + }) + .await?; + + continue; + } + }; + + // Verify DANE policy + if let Some(dane_policy) = &dane_policy { + if let Err(err) = dane_policy.verify( + 0, + hostname, + client.tls_connection().peer_certificates(), + ) { + tx.send(DeliveryStage::DaneVerifyError { + reason: err.to_string(), + }) + .await?; + } else { + tx.send(DeliveryStage::DaneVerifySuccess).await?; + } + } + + // Say EHLO again (some SMTP servers require this) + tx.send(DeliveryStage::EhloStart).await?; + + let now = Instant::now(); + match tokio::time::timeout(timeout, async { + client + .stream + .write_all(format!("EHLO {local_host}\r\n",).as_bytes()) + .await?; + client.stream.flush().await?; + client.read_ehlo().await + }) + .await + { + Ok(Ok(_)) => { + tx.send(DeliveryStage::EhloSuccess { + elapsed: now.elapsed_ms(), + }) + .await?; + } + Ok(Err(err)) => { + tx.send(DeliveryStage::EhloError { + elapsed: now.elapsed_ms(), + reason: err.to_string(), + }) + .await?; + + continue; + } + Err(_) => { + tx.send(DeliveryStage::EhloError { + elapsed: now.elapsed_ms(), + reason: "Timed out reading response".to_string(), + }) + .await?; + + continue; + } + } + + // Verify recipient + let mut is_success = email.is_none(); + if let Some(email) = &email { + // MAIL FROM + tx.send(DeliveryStage::MailFromStart).await?; + + let now = Instant::now(); + + match client.cmd(b"MAIL FROM:<>\r\n").await.and_then(|r| { + if r.is_positive_completion() { + Ok(r) + } else { + Err(mail_send::Error::UnexpectedReply(r)) + } + }) { + Ok(_) => { + tx.send(DeliveryStage::MailFromSuccess { elapsed: now.elapsed_ms(), - reason: status.to_string(), }) .await?; - continue; - } - tx.send(DeliveryStage::ReadGreetingSuccess { - elapsed: now.elapsed_ms(), - }) - .await?; - - // Say EHLO - tx.send(DeliveryStage::EhloStart).await?; - - let now = Instant::now(); - let capabilities = match tokio::time::timeout(timeout, async { - client - .stream - .write_all(format!("EHLO {local_host}\r\n",).as_bytes()) - .await?; - client.stream.flush().await?; - client.read_ehlo().await - }) - .await - { - Ok(Ok(capabilities)) => { - tx.send(DeliveryStage::EhloSuccess { - elapsed: now.elapsed_ms(), - }) - .await?; - - capabilities - } - Ok(Err(err)) => { - tx.send(DeliveryStage::EhloError { - elapsed: now.elapsed_ms(), - reason: err.to_string(), - }) - .await?; - - continue; - } - Err(_) => { - tx.send(DeliveryStage::EhloError { - elapsed: now.elapsed_ms(), - reason: "Timed out reading response".to_string(), - }) - .await?; - - continue; - } - }; - - // Start TLS - tx.send(DeliveryStage::StartTlsStart).await?; - - let now = Instant::now(); - let mut client = match client - .try_start_tls( - &server.inner.data.smtp_connectors.pki_verify, - &hostname, - &capabilities, - ) - .await - { - StartTlsResult::Success { smtp_client } => { - tx.send(DeliveryStage::StartTlsSuccess { - elapsed: now.elapsed_ms(), - }) - .await?; - - smtp_client - } - StartTlsResult::Error { error } => { - tx.send(DeliveryStage::StartTlsError { - elapsed: now.elapsed_ms(), - reason: error.to_string(), - }) - .await?; - - continue; - } - StartTlsResult::Unavailable { response, .. } => { - tx.send(DeliveryStage::StartTlsError { - elapsed: now.elapsed_ms(), - reason: response.map(|r| r.to_string()).unwrap_or_else( - || "STARTTLS not advertised by host".to_string(), - ), - }) - .await?; - - continue; - } - }; - - // Verify DANE policy - if let Some(dane_policy) = &dane_policy { - if let Err(err) = dane_policy.verify( - 0, - &hostname, - client.tls_connection().peer_certificates(), - ) { - tx.send(DeliveryStage::DaneVerifyError { - reason: err.to_string(), - }) - .await?; - } else { - tx.send(DeliveryStage::DaneVerifySuccess).await?; - } - } - - // Say EHLO again (some SMTP servers require this) - tx.send(DeliveryStage::EhloStart).await?; - - let now = Instant::now(); - match tokio::time::timeout(timeout, async { - client - .stream - .write_all(format!("EHLO {local_host}\r\n",).as_bytes()) - .await?; - client.stream.flush().await?; - client.read_ehlo().await - }) - .await - { - Ok(Ok(_)) => { - tx.send(DeliveryStage::EhloSuccess { - elapsed: now.elapsed_ms(), - }) - .await?; - } - Ok(Err(err)) => { - tx.send(DeliveryStage::EhloError { - elapsed: now.elapsed_ms(), - reason: err.to_string(), - }) - .await?; - - continue; - } - Err(_) => { - tx.send(DeliveryStage::EhloError { - elapsed: now.elapsed_ms(), - reason: "Timed out reading response".to_string(), - }) - .await?; - - continue; - } - } - - // Verify recipient - let mut is_success = email.is_none(); - if let Some(email) = &email { - // MAIL FROM - tx.send(DeliveryStage::MailFromStart).await?; + // RCPT TO + tx.send(DeliveryStage::RcptToStart).await?; let now = Instant::now(); - - match client.cmd(b"MAIL FROM:<>\r\n").await.and_then(|r| { - if r.is_positive_completion() { - Ok(r) - } else { - Err(mail_send::Error::UnexpectedReply(r)) - } - }) { + match client + .cmd(format!("RCPT TO:<{email}>\r\n").as_bytes()) + .await + .and_then(|r| { + if r.is_positive_completion() { + Ok(r) + } else { + Err(mail_send::Error::UnexpectedReply(r)) + } + }) { Ok(_) => { - tx.send(DeliveryStage::MailFromSuccess { + is_success = true; + tx.send(DeliveryStage::RcptToSuccess { elapsed: now.elapsed_ms(), }) .await?; - - // RCPT TO - tx.send(DeliveryStage::RcptToStart).await?; - - let now = Instant::now(); - match client - .cmd(format!("RCPT TO:<{email}>\r\n").as_bytes()) - .await - .and_then(|r| { - if r.is_positive_completion() { - Ok(r) - } else { - Err(mail_send::Error::UnexpectedReply(r)) - } - }) { - Ok(_) => { - is_success = true; - tx.send(DeliveryStage::RcptToSuccess { - elapsed: now.elapsed_ms(), - }) - .await?; - } - Err(err) => { - tx.send(DeliveryStage::RcptToError { - reason: err.to_string(), - elapsed: now.elapsed_ms(), - }) - .await?; - } - } } Err(err) => { - tx.send(DeliveryStage::MailFromError { + tx.send(DeliveryStage::RcptToError { reason: err.to_string(), elapsed: now.elapsed_ms(), }) @@ -675,44 +685,37 @@ async fn delivery_diagnose( } } } - - // QUIT - tx.send(DeliveryStage::QuitStart).await?; - - let now = Instant::now(); - client.quit().await; - tx.send(DeliveryStage::QuitCompleted { - elapsed: now.elapsed_ms(), - }) - .await?; - - if is_success { - break 'outer; + Err(err) => { + tx.send(DeliveryStage::MailFromError { + reason: err.to_string(), + elapsed: now.elapsed_ms(), + }) + .await?; } } - Err(err) => { - tx.send(DeliveryStage::ConnectionError { - elapsed: now.elapsed_ms(), - reason: err.to_string(), - }) - .await?; - } + } + + // QUIT + tx.send(DeliveryStage::QuitStart).await?; + + let now = Instant::now(); + client.quit().await; + tx.send(DeliveryStage::QuitCompleted { + elapsed: now.elapsed_ms(), + }) + .await?; + + if is_success { + break 'outer; } } - } - Ok(_) => { - tx.send(DeliveryStage::IpLookupError { - reason: "No IP addresses found for host".to_string(), - elapsed: now.elapsed_ms(), - }) - .await?; - } - Err(err) => { - tx.send(DeliveryStage::IpLookupError { - reason: err.to_string(), - elapsed: now.elapsed_ms(), - }) - .await?; + Err(err) => { + tx.send(DeliveryStage::ConnectionError { + elapsed: now.elapsed_ms(), + reason: err.to_string(), + }) + .await?; + } } } } diff --git a/crates/services/src/task_manager/index.rs b/crates/services/src/task_manager/index.rs index f8b58f08..cd78bf3e 100644 --- a/crates/services/src/task_manager/index.rs +++ b/crates/services/src/task_manager/index.rs @@ -585,10 +585,8 @@ async fn delete_email_metadata( use store::{ SerializeInfallible, write::{BlobLink, BlobOp, RegistryClass, now}, - xxhash_rust, }; use types::blob::BlobId; - use utils::snowflake::SnowflakeIdGenerator; let root_part = metadata.root_part(); let from: Option = root_part.headers.iter().find_map(|h| { @@ -632,10 +630,7 @@ async fn delete_email_metadata( }) .to_pickled_vec(); let object_id = ObjectType::ArchivedItem.to_id(); - let item_id = SnowflakeIdGenerator::from_sequence_id( - xxhash_rust::xxh3::xxh3_64(item.as_slice()), - ) - .unwrap_or_default(); + let item_id = server.inner.data.registry_id_gen.generate(); batch .set( diff --git a/crates/services/src/task_manager/merge_threads.rs b/crates/services/src/task_manager/merge_threads.rs index af8251ed..319cb212 100644 --- a/crates/services/src/task_manager/merge_threads.rs +++ b/crates/services/src/task_manager/merge_threads.rs @@ -6,15 +6,18 @@ use crate::task_manager::TaskResult; use common::{Server, storage::index::ObjectIndexBuilder}; -use email::message::{ingest::ThreadMerge, metadata::MessageData}; +use email::message::{ + ingest::{ThreadMerge, has_message_id}, + metadata::MessageData, +}; use registry::schema::structs::TaskMergeThreads; use std::{str::FromStr, time::Duration}; use store::{ - IndexKeyPrefix, IterateParams, U32_LEN, ValueKey, - ahash::{AHashMap, AHashSet}, + IterateParams, Key, U32_LEN, ValueKey, + ahash::AHashMap, rand::Rng, write::{ - AlignedBytes, Archive, BatchBuilder, IndexPropertyClass, ValueClass, + AlignedBytes, Archive, BatchBuilder, IndexPropertyClass, MergeResult, Params, ValueClass, key::DeserializeBigEndian, }, }; @@ -51,55 +54,70 @@ async fn merge_threads( server: &Server, task_merge_threads: &TaskMergeThreads, ) -> trc::Result { - let Ok(thread_hash) = CheekyHash::from_str(&task_merge_threads.thread_hash) else { + let Ok(thread_hash) = CheekyHash::from_str(&task_merge_threads.thread_name) else { return Ok(TaskResult::permanent("Invalid thread hash")); }; - let account_id = task_merge_threads.account_id.document_id(); - let key_len = IndexKeyPrefix::len() + thread_hash.len() + U32_LEN; - let document_id_pos = key_len - U32_LEN; - let merge_thread_ids = task_merge_threads - .thread_ids + let Ok(mut message_ids) = task_merge_threads + .message_ids .iter() - .map(|id| id.document_id()) - .collect::>(); - let mut thread_merge = ThreadMerge::new(); - let mut thread_index = AHashMap::new(); + .map(|id| CheekyHash::from_str(id)) + .collect::, _>>() + else { + return Ok(TaskResult::permanent("Invalid message ids")); + }; + message_ids.sort_unstable(); + + let account_id = task_merge_threads.account_id.document_id(); let mut try_count = 0; + let from_key = ValueKey { + account_id, + collection: Collection::Email.into(), + document_id: 0, + class: ValueClass::IndexProperty(IndexPropertyClass::Hash { + property: EmailField::Threading.into(), + hash: thread_hash, + }), + }; + let to_key = ValueKey { + account_id, + collection: Collection::Email.into(), + document_id: u32::MAX, + class: ValueClass::IndexProperty(IndexPropertyClass::Hash { + property: EmailField::Threading.into(), + hash: thread_hash, + }), + }; + let mut prefix = from_key.serialize(0); + let key_len = prefix.len(); + let document_id_pos = key_len - U32_LEN; + prefix.truncate(document_id_pos); + 'retry: loop { + // Merge threads + let mut thread_merge = ThreadMerge::new(); + let mut same_subject_messages: AHashMap> = AHashMap::new(); + // Find thread ids server .store() .iterate( - IterateParams::new( - ValueKey { - account_id, - collection: Collection::Email.into(), - document_id: 0, - class: ValueClass::IndexProperty(IndexPropertyClass::Hash { - property: EmailField::Threading.into(), - hash: thread_hash, - }), - }, - ValueKey { - account_id, - collection: Collection::Email.into(), - document_id: u32::MAX, - class: ValueClass::IndexProperty(IndexPropertyClass::Hash { - property: EmailField::Threading.into(), - hash: thread_hash, - }), - }, - ) - .ascending(), + IterateParams::new(from_key.clone(), to_key.clone()).ascending(), |key, value| { - if key.len() == key_len { + if key.len() == key_len && key.starts_with(&prefix) { + // Find matching references + let references = value.get(U32_LEN..).unwrap_or_default(); let thread_id = value.deserialize_be_u32(0)?; - if merge_thread_ids.contains(&thread_id) { - let document_id = key.deserialize_be_u32(document_id_pos)?; + let document_id = key.deserialize_be_u32(document_id_pos)?; + if has_message_id(&message_ids, references) { thread_merge.add(thread_id, document_id); - thread_index.insert(document_id, value.to_vec()); + } else { + // Keep track of messages with the same subject for potential future merges + same_subject_messages + .entry(thread_id) + .or_default() + .push(document_id); } } @@ -113,6 +131,17 @@ async fn merge_threads( // Another process merged the threads already? return Ok(TaskResult::Success); } + + // Add other messages with the same subject to the merge if they share a + // thread id with a message that has a matching message id + for thread_id in thread_merge.thread_ids().copied().collect::>() { + if let Some(document_ids) = same_subject_messages.get(&thread_id) { + for &document_id in document_ids { + thread_merge.add(thread_id, document_id); + } + } + } + let thread_id = thread_merge.merge_thread_id(); // Delete all but the most common threadId @@ -168,14 +197,41 @@ async fn merge_threads( .caused_by(trc::location!())?; // Update thread index property - let mut thread_index = thread_index.remove(&document_id).unwrap(); - thread_index[0..U32_LEN].copy_from_slice(&thread_id.to_be_bytes()); - batch.set( + batch.merge_fnc( ValueClass::IndexProperty(IndexPropertyClass::Hash { property: EmailField::Threading.into(), hash: thread_hash, }), - thread_index, + Params::with_capacity(3) + .with_u64(thread_id as u64) + .with_u64(group_thread_id as u64), + |params, _, bytes| { + let new_thread_id = params.u64(0) as u32; + let old_thread_id = params.u64(1) as u32; + + let mut thread_index = bytes + .filter(|v| v.len() > U32_LEN) + .ok_or_else(|| { + trc::StoreEvent::AssertValueFailed + .into_err() + .details("Message no longer exists.") + .caused_by(trc::location!()) + })? + .to_vec(); + + if thread_index.as_slice().deserialize_be_u32(0)? != old_thread_id { + return Err( + trc::StoreEvent::AssertValueFailed + .into_err() + .details("Thread id mismatch, likely due to concurrent modification.") + .caused_by(trc::location!()) + ); + } + + thread_index[0..U32_LEN].copy_from_slice(&new_thread_id.to_be_bytes()); + + Ok(MergeResult::Update(thread_index)) + }, ); } } diff --git a/crates/smtp/src/outbound/lookup.rs b/crates/smtp/src/outbound/lookup.rs index 411aeecf..62e4f739 100644 --- a/crates/smtp/src/outbound/lookup.rs +++ b/crates/smtp/src/outbound/lookup.rs @@ -147,7 +147,7 @@ impl DnsLookup for Server { }) } })?, - HostOrIp::Ip { ip, .. } => vec![ip], + HostOrIp::Ip(ip) => vec![ip], }; if !remote_ips.is_empty() { diff --git a/crates/smtp/src/outbound/mod.rs b/crates/smtp/src/outbound/mod.rs index 1bd98ad4..56859146 100644 --- a/crates/smtp/src/outbound/mod.rs +++ b/crates/smtp/src/outbound/mod.rs @@ -15,7 +15,7 @@ use common::config::{ use mail_auth::IpLookupStrategy; use mail_send::Credentials; use smtp_proto::{Response, Severity}; -use std::borrow::Cow; +use std::{borrow::Cow, net::IpAddr}; pub mod client; pub mod dane; @@ -247,14 +247,14 @@ impl NextHop<'_> { } } NextHop::Relay(host) => match &host.address { - HostOrIp::Host(host) => host.as_str(), - HostOrIp::Ip { ip_str, .. } => ip_str.as_str(), + HostOrIp::Host(host) => host.as_ref(), + HostOrIp::Ip(ip) => ip.ip_str.as_ref(), }, } } #[inline(always)] - pub fn fqdn_hostname(&self) -> HostOrIp> { + pub fn fqdn_hostname(&self) -> HostOrIp, IpAddr> { match self { NextHop::MX { host, .. } => { if !host.ends_with('.') { @@ -264,11 +264,8 @@ impl NextHop<'_> { } } NextHop::Relay(host) => match &host.address { - HostOrIp::Host(host) => HostOrIp::Host(host.as_str().into()), - HostOrIp::Ip { ip, ip_str } => HostOrIp::Ip { - ip: *ip, - ip_str: ip_str.as_str().into(), - }, + HostOrIp::Host(host) => HostOrIp::Host(host.as_ref().into()), + HostOrIp::Ip(ip) => HostOrIp::Ip(ip.ip), }, } } diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index ee569893..441c1c80 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -34,14 +34,11 @@ use store::write::{ AlignedBytes, Archive, Archiver, BatchBuilder, BlobLink, BlobOp, MergeResult, Params, QueueClass, RegistryClass, ValueClass, now, }; -use store::{ - Deserialize, IterateParams, Serialize, SerializeInfallible, U64_LEN, ValueKey, xxhash_rust, -}; +use store::{Deserialize, IterateParams, Serialize, SerializeInfallible, U64_LEN, ValueKey}; use trc::{AddContext, ServerEvent, SpamEvent}; use types::blob::BlobId; use types::blob_hash::BlobHash; use utils::DomainPart; -use utils::snowflake::SnowflakeIdGenerator; pub const LOCK_EXPIRY: u64 = 10 * 60; // 10 minutes pub const QUEUE_REFRESH: u64 = 5 * 60; // 5 minutes @@ -476,10 +473,7 @@ impl MessageWrapper { .to_pickled_vec(); let object_id = ObjectType::SpamTrainingSample.to_id(); - let item_id = SnowflakeIdGenerator::from_sequence_id(xxhash_rust::xxh3::xxh3_64( - sample.as_slice(), - )) - .unwrap_or_default(); + let item_id = server.inner.data.registry_id_gen.generate(); batch .set( BlobOp::Link { diff --git a/crates/store/src/write/batch.rs b/crates/store/src/write/batch.rs index 62ceace9..eaadf12a 100644 --- a/crates/store/src/write/batch.rs +++ b/crates/store/src/write/batch.rs @@ -474,8 +474,7 @@ impl BatchBuilder { let due = task.due_timestamp(); let class = task.object_type().to_id(); let task = task.to_pickled_vec(); - let id = SnowflakeIdGenerator::from_sequence_id(xxhash_rust::xxh3::xxh3_64(&task)) - .unwrap_or_default(); + let id = SnowflakeIdGenerator::global_id().unwrap_or_default(); self.set(ValueClass::TaskQueue(TaskQueueClass::Task { id }), task) .set( diff --git a/crates/utils/src/snowflake.rs b/crates/utils/src/snowflake.rs index d2996f3d..6fb63534 100644 --- a/crates/utils/src/snowflake.rs +++ b/crates/utils/src/snowflake.rs @@ -24,6 +24,7 @@ const NODE_ID_MASK: u64 = (1 << NODE_ID_LEN) - 1; const DEFAULT_EPOCH: u64 = 1632280000; // 52 years after UNIX_EPOCH static mut NODE_ID: u64 = 1; +static SEQUENCE_ID: AtomicU64 = AtomicU64::new(0); /* @@ -74,12 +75,13 @@ impl SnowflakeIdGenerator { .and_then(|diff| Self::from_duration(Duration::from_secs(diff))) } - pub fn from_timestamp_and_sequence_id(timestamp: u64, sequence: u64) -> Option { + pub fn global_id_from_timestamp(timestamp: u64) -> Option { + let sequence = SEQUENCE_ID.fetch_add(1, Ordering::Relaxed) & SEQUENCE_MASK; Self::from_timestamp(timestamp).map(|id| id | (sequence << NODE_ID_LEN) | node_id()) } - pub fn from_sequence_id(sequence: u64) -> Option { - let sequence = sequence & SEQUENCE_MASK; + pub fn global_id() -> Option { + let sequence = SEQUENCE_ID.fetch_add(1, Ordering::Relaxed) & SEQUENCE_MASK; (SystemTime::UNIX_EPOCH + Duration::from_secs(DEFAULT_EPOCH)) .elapsed() diff --git a/tests/src/jmap/mail/query.rs b/tests/src/jmap/mail/query.rs index 77e849b1..f211d5a1 100644 --- a/tests/src/jmap/mail/query.rs +++ b/tests/src/jmap/mail/query.rs @@ -161,8 +161,11 @@ pub async fn query(client: &Client, can_stem: bool) { ], if can_stem { vec![ - "T10330", "N01744", "N01743", "N04885", "N02688", "N02122", "A00059", "A00058", - "N02123", "T00651", "T09439", "N05001", "T05848", "T05508", + /*"T10330", "N01744", "N01743", "N04885", "N02688", "N02122", "A00059", "A00058", + "N02123", "T00651", "T09439", "N05001", "T05848", "T05508",*/ + "T09187", "T10330", "N01744", "N01743", "N04885", "N02688", "N02122", "A00059", + "A00057", "A00058", "N02123", "T00651", "T09439", "N05001", "A01072", "A01061", + "AR00050", "T02310", "T05848", "T05508", "P20078", "P20079", ] } else { vec!["T10330", "N02122", "N02123", "T09439"] diff --git a/tests/src/jmap/mail/sieve_script.rs b/tests/src/jmap/mail/sieve_script.rs index 65aed2a0..bb20681f 100644 --- a/tests/src/jmap/mail/sieve_script.rs +++ b/tests/src/jmap/mail/sieve_script.rs @@ -14,6 +14,7 @@ use jmap_client::{ email, mailbox, sieve::query::{Comparator, Filter}, }; +use registry::schema::{prelude::ObjectType, structs::SieveUserScript}; use std::{ fs, path::PathBuf, @@ -22,6 +23,20 @@ use std::{ pub async fn test(test: &TestServer) { println!("Running Sieve tests..."); + + // Create a global script + let admin = test.account("admin@example.com"); + admin + .registry_create_object(SieveUserScript { + contents: "require \"reject\";\nreject \"Rejected from a global script.\";\nstop;\n" + .into(), + description: None, + is_active: true, + name: "common".into(), + }) + .await; + admin.reload_settings().await; + let server = test.server.clone(); let account = test.account("jdoe@example.com"); let client = account.jmap_client().await; @@ -491,6 +506,9 @@ pub async fn test(test: &TestServer) { client.sieve_script_destroy(&id).await.unwrap(); } test.destroy_all_mailboxes(account).await; + admin + .registry_destroy_all(ObjectType::SieveUserScript) + .await; test.assert_is_empty().await; } diff --git a/tests/src/jmap/mail/thread_merge.rs b/tests/src/jmap/mail/thread_merge.rs index fd9a396e..72fb3113 100644 --- a/tests/src/jmap/mail/thread_merge.rs +++ b/tests/src/jmap/mail/thread_merge.rs @@ -4,7 +4,10 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -use crate::{store::deflate_test_resource, utils::server::TestServer}; +use crate::{ + store::deflate_test_resource, + utils::server::{DestroyAllMailboxes, TestServer}, +}; use ::email::{ cache::MessageCacheFetch, mailbox::INBOX_ID, @@ -15,7 +18,7 @@ use jmap_client::{email, mailbox::Role}; use mail_parser::{MessageParser, mailbox::mbox::MessageIterator}; use std::{io::Cursor, str::FromStr, time::Duration}; use store::{ - ahash::{AHashMap, AHashSet}, + ahash::AHashSet, rand::{self, Rng}, }; use types::id::Id; @@ -29,10 +32,20 @@ async fn test_single_thread(test_server: &TestServer) { println!("Running Email Merge Threads tests..."); let account = test_server.account("admin@example.com"); let mut client = account.jmap_client().await; - let mut all_mailboxes = AHashMap::default(); - for (base_test_num, test) in [test_1(), test_2(), test_3()].iter().enumerate() { - let base_test_num = ((base_test_num * 6) as u32) + 1; + let mut account_ids = Vec::new(); + for name in [ + "admin@example.com", + "jdoe@example.com", + "jane.smith@example.com", + "bill@example.com", + "robert@example.com", + "sales@example.com", + ] { + account_ids.push(test_server.account(name).id_string()); + } + + for (test_group_num, test) in [test_1(), test_2(), test_3()].iter().enumerate() { let mut messages = Vec::new(); let mut total_messages = 0; let mut messages_per_thread = @@ -41,10 +54,10 @@ async fn test_single_thread(test_server: &TestServer) { let mut mailbox_ids = Vec::with_capacity(6); - for test_num in 0..=5 { + for account_id in &account_ids { mailbox_ids.push( client - .set_default_account_id(Id::new((base_test_num + test_num) as u64).to_string()) + .set_default_account_id(*account_id) .mailbox_create("Thread nightmare", None::, Role::None) .await .unwrap() @@ -54,7 +67,7 @@ async fn test_single_thread(test_server: &TestServer) { for message in &messages { client - .set_default_account_id(Id::new(base_test_num as u64).to_string()) + .set_default_account_id(account_ids[0]) .email_import( message.to_string().into_bytes(), [mailbox_ids[0].clone()], @@ -67,7 +80,7 @@ async fn test_single_thread(test_server: &TestServer) { for message in messages.iter().rev() { client - .set_default_account_id(Id::new((base_test_num + 1) as u64).to_string()) + .set_default_account_id(account_ids[1]) .email_import( message.to_string().into_bytes(), [mailbox_ids[1].clone()], @@ -79,8 +92,7 @@ async fn test_single_thread(test_server: &TestServer) { } for chunk in messages.chunks(5) { - client.set_default_account_id(Id::new((base_test_num + 2) as u64).to_string()); - + client.set_default_account_id(account_ids[2]); for message in chunk { client .email_import( @@ -93,8 +105,7 @@ async fn test_single_thread(test_server: &TestServer) { .unwrap(); } - client.set_default_account_id(Id::new((base_test_num + 3) as u64).to_string()); - + client.set_default_account_id(account_ids[3]); for message in chunk.iter().rev() { client .email_import( @@ -109,8 +120,7 @@ async fn test_single_thread(test_server: &TestServer) { } for chunk in messages.chunks(5).rev() { - client.set_default_account_id(Id::new((base_test_num + 4) as u64).to_string()); - + client.set_default_account_id(account_ids[4]); for message in chunk { client .email_import( @@ -123,8 +133,7 @@ async fn test_single_thread(test_server: &TestServer) { .unwrap(); } - client.set_default_account_id(Id::new((base_test_num + 5) as u64).to_string()); - + client.set_default_account_id(account_ids[5]); for message in chunk.iter().rev() { client .email_import( @@ -137,14 +146,13 @@ async fn test_single_thread(test_server: &TestServer) { .unwrap(); } } - test_server.wait_for_tasks().await; for test_num in 0..=5 { let result = client - .set_default_account_id(Id::new((base_test_num + test_num) as u64).to_string()) + .set_default_account_id(account_ids[test_num]) .email_query( - email::query::Filter::in_mailbox(mailbox_ids[test_num as usize].clone()).into(), + email::query::Filter::in_mailbox(mailbox_ids[test_num].clone()).into(), None::>, ) .await @@ -154,7 +162,7 @@ async fn test_single_thread(test_server: &TestServer) { result.ids().len(), total_messages, "test# {}/{}", - base_test_num, + test_group_num, test_num ); @@ -164,15 +172,6 @@ async fn test_single_thread(test_server: &TestServer) { .map(|id| Id::from_str(id).unwrap().prefix_id()) .collect(); - assert_eq!( - thread_ids.len(), - messages_per_thread.len(), - "{:?}: test# {}/{}", - thread_ids, - base_test_num, - test_num - ); - let mut messages_per_thread_db = Vec::new(); for thread_id in thread_ids { @@ -189,19 +188,17 @@ async fn test_single_thread(test_server: &TestServer) { messages_per_thread_db.sort_unstable(); assert_eq!(messages_per_thread_db, messages_per_thread); - println!("passed test# {}/{}", base_test_num, test_num); + println!("passed test# {}/{}", test_group_num, test_num); } - all_mailboxes.insert(base_test_num as usize, mailbox_ids); - } - - // Delete all messages and make sure no keys are left in the store. - for (base_test_num, mailbox_ids) in all_mailboxes { - for (test_num, _) in mailbox_ids.into_iter().enumerate() { - account - .destroy_all_mailboxes_for_account((base_test_num + test_num) as u32) + for account_id in &account_ids { + client + .set_default_account_id(*account_id) + .destroy_all_mailboxes() .await; } + test_server.wait_for_tasks().await; + test_server.assert_is_empty().await; } test_server.assert_is_empty().await; diff --git a/tests/src/jmap/mod.rs b/tests/src/jmap/mod.rs index ac2b9658..369f9d67 100644 --- a/tests/src/jmap/mod.rs +++ b/tests/src/jmap/mod.rs @@ -9,8 +9,8 @@ use registry::{ schema::{ enums::{MtaProtocol, Permission}, structs::{ - CalendarAlarm, Expression, ExpressionMatch, Imap, Jmap, MtaOutboundStrategy, MtaRoute, - MtaRouteRelay, MtaStageAuth, Sharing, + CalendarAlarm, Expression, ExpressionMatch, Imap, Jmap, MtaExtensions, + MtaOutboundStrategy, MtaRoute, MtaRouteRelay, MtaStageAuth, Sharing, }, }, types::list::List, @@ -153,17 +153,29 @@ async fn jmap_tests() { address: "127.0.0.1".into(), port: 9999, allow_invalid_certs: true, - implicit_tls: true, + implicit_tls: false, name: "mock-smtp".into(), protocol: MtaProtocol::Smtp, ..Default::default() })) .await; + admin + .registry_create_object(MtaExtensions { + future_release: Expression { + match_: List::from_iter([ExpressionMatch { + if_: "!is_empty(authenticated_as)".into(), + then: "99999999d".into(), + }]), + else_: "false".to_string(), + }, + ..Default::default() + }) + .await; admin.reload_settings().await; test.insert_account(admin); - /*mail::get::test(&test).await; + mail::get::test(&test).await; mail::set::test(&test).await; mail::parse::test(&test).await; mail::query::test(&test).await; @@ -174,12 +186,12 @@ async fn jmap_tests() { mail::thread_get::test(&test).await; mail::thread_merge::test(&test).await; mail::mailbox::test(&test).await; - mail::acl::test(&test).await;*/ + mail::acl::test(&test).await; mail::sieve_script::test(&test).await; mail::vacation_response::test(&test).await; mail::submission::test(&test).await; - /*core::event_source::test(&test).await; + core::event_source::test(&test).await; core::websocket::test(&test).await; core::push_subscription::test(&test).await; core::blob::test(&test).await; @@ -200,7 +212,7 @@ async fn jmap_tests() { calendar::acl::test(&test).await; principal::get::test(&test).await; - principal::availability::test(&test).await;*/ + principal::availability::test(&test).await; if test.is_reset() { test.temp_dir.delete(); diff --git a/tests/src/utils/server.rs b/tests/src/utils/server.rs index c9bb0c42..6118e2a1 100644 --- a/tests/src/utils/server.rs +++ b/tests/src/utils/server.rs @@ -350,7 +350,7 @@ impl TestServer { pub async fn destroy_all_mailboxes(&self, account: &Account) { self.wait_for_tasks().await; - destroy_all_mailboxes_no_wait(&account.jmap_client().await).await; + account.jmap_client().await.destroy_all_mailboxes().await; } } @@ -358,16 +358,22 @@ impl Account { pub async fn destroy_all_mailboxes_for_account(&self, account_id: u32) { let mut client = self.jmap_client().await; client.set_default_account_id(Id::from(account_id)); - destroy_all_mailboxes_no_wait(&client).await; + client.destroy_all_mailboxes().await; } } -async fn destroy_all_mailboxes_no_wait(client: &Client) { - let mut request = client.build(); - request.query_mailbox().arguments().sort_as_tree(true); - let mut ids = request.send_query_mailbox().await.unwrap().take_ids(); - ids.reverse(); - for id in ids { - client.mailbox_destroy(&id, true).await.unwrap(); +pub trait DestroyAllMailboxes { + fn destroy_all_mailboxes(&self) -> impl Future; +} + +impl DestroyAllMailboxes for Client { + async fn destroy_all_mailboxes(&self) { + let mut request = self.build(); + request.query_mailbox().arguments().sort_as_tree(true); + let mut ids = request.send_query_mailbox().await.unwrap().take_ids(); + ids.reverse(); + for id in ids { + self.mailbox_destroy(&id, true).await.unwrap(); + } } } diff --git a/tests/src/utils/storage.rs b/tests/src/utils/storage.rs index 36ed7f4b..bd00d79b 100644 --- a/tests/src/utils/storage.rs +++ b/tests/src/utils/storage.rs @@ -186,7 +186,7 @@ pub async fn wait_for_tasks(server: &Server, skip_permanent_failures: bool) { if count % 10 == 0 { println!("Waiting for pending task {:?}...", task); } - tokio::time::sleep(std::time::Duration::from_millis(300)).await; + tokio::time::sleep(std::time::Duration::from_millis(200)).await; } else { break; }