From 7564f516427b2d5c208714670ccd77f8748078e7 Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Sun, 12 Apr 2026 18:52:18 +0200 Subject: [PATCH] Telemetry fixes --- Cargo.lock | 3 + crates/common/Cargo.toml | 1 + crates/common/src/auth/authentication.rs | 4 +- crates/common/src/auth/oauth/mod.rs | 8 +- crates/common/src/auth/oauth/token.rs | 21 +- crates/common/src/auth/rate_limit.rs | 27 +- crates/common/src/config/network.rs | 10 +- crates/common/src/config/telemetry.rs | 29 +- crates/common/src/manager/defaults.rs | 26 +- crates/common/src/network/acme/http.rs | 2 +- crates/common/src/telemetry/metrics/mod.rs | 3 + crates/common/src/telemetry/metrics/store.rs | 16 +- .../common/src/telemetry/metrics/test_data.rs | 91 ++ crates/common/src/telemetry/tracers/mod.rs | 100 ++ crates/common/src/telemetry/tracers/store.rs | 84 +- crates/dav/Cargo.toml | 1 + crates/dav/src/calendar/copy_move.rs | 2 +- crates/dav/src/calendar/delete.rs | 2 +- crates/http-proto/Cargo.toml | 1 + crates/http-proto/src/context.rs | 21 +- crates/http/Cargo.toml | 1 + crates/http/src/api/diagnose.rs | 2 + crates/http/src/api/mod.rs | 62 +- crates/http/src/api/telemetry.rs | 17 +- crates/http/src/auth/authenticate.rs | 4 +- crates/http/src/auth/oauth/auth.rs | 13 +- crates/http/src/auth/oauth/openid.rs | 2 +- crates/http/src/auth/oauth/token.rs | 2 +- crates/jmap-proto/src/request/method.rs | 4 + crates/jmap-proto/src/request/mod.rs | 2 + crates/jmap-proto/src/request/parser.rs | 17 +- crates/jmap/src/api/auth.rs | 2 + crates/jmap/src/api/request.rs | 12 + crates/jmap/src/calendar_event/query.rs | 44 +- crates/jmap/src/contact/query.rs | 44 +- crates/jmap/src/registry/mapping/account.rs | 8 +- crates/jmap/src/registry/mapping/telemetry.rs | 60 +- crates/jmap/src/vacation/get.rs | 144 +- crates/main/Cargo.toml | 8 + crates/main/src/main.rs | 10 + crates/main/src/test_data.rs | 1403 +++++++++++++++++ crates/trc/Cargo.toml | 1 + crates/trc/src/ipc/collector.rs | 7 +- crates/trc/src/ipc/metrics.rs | 4 + resources/html-templates/login.html | 10 +- resources/html-templates/login.html.min | 2 +- tests/src/telemetry/metrics.rs | 77 +- 47 files changed, 2070 insertions(+), 344 deletions(-) create mode 100644 crates/common/src/telemetry/metrics/test_data.rs create mode 100644 crates/main/src/test_data.rs diff --git a/Cargo.lock b/Cargo.lock index 318718b6..0b2ee0e3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7338,14 +7338,17 @@ dependencies = [ "email", "groupware", "http 0.16.0", + "http_proto", "imap", "jemallocator", "jmap", "managesieve", "migration", "pop3", + "registry", "services", "smtp", + "smtp-proto", "spam-filter", "store", "tokio", diff --git a/crates/common/Cargo.toml b/crates/common/Cargo.toml index 13e22783..e5140730 100644 --- a/crates/common/Cargo.toml +++ b/crates/common/Cargo.toml @@ -93,6 +93,7 @@ libc = "0.2.126" [features] test_mode = [] +dev_mode = [] enterprise = [] foundation = [] diff --git a/crates/common/src/auth/authentication.rs b/crates/common/src/auth/authentication.rs index 21de2509..10c8c2b9 100644 --- a/crates/common/src/auth/authentication.rs +++ b/crates/common/src/auth/authentication.rs @@ -288,8 +288,8 @@ impl Server { .await; } - let todo = "fix"; - if token == "TEST_MODE_BYPASS" { + #[cfg(feature = "dev_mode")] + if std::env::var("API_TOKEN_ADMIN").is_ok_and(|admin_token| &admin_token == token) { return Ok(AccessToken::new_admin()); } diff --git a/crates/common/src/auth/oauth/mod.rs b/crates/common/src/auth/oauth/mod.rs index 3e5170fa..7f3fa161 100644 --- a/crates/common/src/auth/oauth/mod.rs +++ b/crates/common/src/auth/oauth/mod.rs @@ -24,7 +24,7 @@ pub enum GrantType { RefreshToken, LiveTracing, LiveMetrics, - Diagnose, + LiveDelivery, Rsvp, } @@ -35,7 +35,7 @@ impl GrantType { GrantType::RefreshToken => "refresh_token", GrantType::LiveTracing => "live_tracing", GrantType::LiveMetrics => "live_metrics", - GrantType::Diagnose => "diagnose", + GrantType::LiveDelivery => "live_delivery", GrantType::Rsvp => "rsvp", } } @@ -46,7 +46,7 @@ impl GrantType { GrantType::RefreshToken => 1, GrantType::LiveTracing => 2, GrantType::LiveMetrics => 3, - GrantType::Diagnose => 4, + GrantType::LiveDelivery => 4, GrantType::Rsvp => 5, } } @@ -57,7 +57,7 @@ impl GrantType { 1 => Some(GrantType::RefreshToken), 2 => Some(GrantType::LiveTracing), 3 => Some(GrantType::LiveMetrics), - 4 => Some(GrantType::Diagnose), + 4 => Some(GrantType::LiveDelivery), 5 => Some(GrantType::Rsvp), _ => None, } diff --git a/crates/common/src/auth/oauth/token.rs b/crates/common/src/auth/oauth/token.rs index ff74f9de..adca3622 100644 --- a/crates/common/src/auth/oauth/token.rs +++ b/crates/common/src/auth/oauth/token.rs @@ -6,8 +6,7 @@ use super::{CLIENT_ID_MAX_LEN, GrantType, RANDOM_CODE_LEN, crypto::SymmetricEncrypt}; use crate::Server; -use mail_builder::encoders::base64::base64_encode; -use mail_parser::decoders::base64::base64_decode; +use base64::{Engine, engine::general_purpose}; use registry::schema::structs::Account; use std::time::SystemTime; use store::{ @@ -102,7 +101,7 @@ impl Server { token.push_leb128(expiry); token.extend_from_slice(client_id.as_bytes()); - Ok(String::from_utf8(base64_encode(&token).unwrap_or_default()).unwrap()) + Ok(general_purpose::URL_SAFE_NO_PAD.encode(&token)) } pub async fn validate_access_token( @@ -111,13 +110,15 @@ impl Server { token_: &str, ) -> trc::Result { // Base64 decode token - let token = base64_decode(token_.as_bytes()).ok_or_else(|| { - trc::AuthEvent::Error - .into_err() - .ctx(trc::Key::Reason, "Failed to decode token") - .caused_by(trc::location!()) - .details(token_.to_string()) - })?; + let token = general_purpose::URL_SAFE_NO_PAD + .decode(token_.as_bytes()) + .map_err(|_| { + trc::AuthEvent::Error + .into_err() + .ctx(trc::Key::Reason, "Failed to decode token") + .caused_by(trc::location!()) + .details(token_.to_string()) + })?; let (account_id, grant_type, issued_at, expiry, client_id) = token .get((RANDOM_CODE_LEN + SymmetricEncrypt::ENCRYPT_TAG_LEN)..) .and_then(|bytes| { diff --git a/crates/common/src/auth/rate_limit.rs b/crates/common/src/auth/rate_limit.rs index 302f1e4d..56258d3b 100644 --- a/crates/common/src/auth/rate_limit.rs +++ b/crates/common/src/auth/rate_limit.rs @@ -16,20 +16,23 @@ impl Server { pub async fn is_http_authenticated_request_allowed( &self, access_token: &AccessToken, + addr: IpAddr, ) -> trc::Result> { let is_rate_allowed = if let Some(rate) = &self.core.network.http.rate_authenticated { - self.core - .storage - .memory - .is_rate_allowed( - KV_RATE_LIMIT_HTTP_AUTHENTICATED, - &access_token.account_id().to_be_bytes(), - rate, - false, - ) - .await - .caused_by(trc::location!())? - .is_none() + self.is_ip_allowed(addr) + || self + .core + .storage + .memory + .is_rate_allowed( + KV_RATE_LIMIT_HTTP_AUTHENTICATED, + &access_token.account_id().to_be_bytes(), + rate, + false, + ) + .await + .caused_by(trc::location!())? + .is_none() } else { true }; diff --git a/crates/common/src/config/network.rs b/crates/common/src/config/network.rs index 5aff2e6d..fa896fca 100644 --- a/crates/common/src/config/network.rs +++ b/crates/common/src/config/network.rs @@ -381,9 +381,13 @@ impl Http { .unwrap_or_default(); // Add permissive CORS headers - let todo = "fix"; - if true { - // http.use_permissive_cors { + #[cfg(feature = "dev_mode")] + let use_permissive_cors = true; + + #[cfg(not(feature = "dev_mode"))] + let use_permissive_cors = http.use_permissive_cors; + + if use_permissive_cors { http_headers.push(( hyper::header::ACCESS_CONTROL_ALLOW_ORIGIN, hyper::header::HeaderValue::from_static("*"), diff --git a/crates/common/src/config/telemetry.rs b/crates/common/src/config/telemetry.rs index 53808015..cbffeeb1 100644 --- a/crates/common/src/config/telemetry.rs +++ b/crates/common/src/config/telemetry.rs @@ -520,7 +520,7 @@ impl Tracers { if tracers.is_empty() { for event_type in EventType::variants() { let event_level = custom_levels - .get(&event_type) + .get(event_type) .copied() .unwrap_or(event_type.level()); if Level::Info.is_contained(event_level) { @@ -540,6 +540,33 @@ impl Tracers { }); } + #[cfg(feature = "dev_mode")] + if let Ok(level) = std::env::var("LOG") { + use std::str::FromStr; + + let level = Level::from_str(&level).expect("Invalid LOG level"); + for event_type in EventType::variants() { + let event_level = custom_levels + .get(event_type) + .copied() + .unwrap_or(event_type.level()); + if level.is_contained(event_level) { + global_interests.set(event_type.to_id() as usize); + } + } + + tracers.push(TelemetrySubscriber { + id: "default".to_string(), + interests: global_interests.clone(), + typ: TelemetrySubscriberType::ConsoleTracer(ConsoleTracer { + ansi: true, + multiline: false, + buffered: true, + }), + lossy: false, + }); + } + Tracers { subscribers: tracers, interests: global_interests, diff --git a/crates/common/src/manager/defaults.rs b/crates/common/src/manager/defaults.rs index 28373fcc..94b7de09 100644 --- a/crates/common/src/manager/defaults.rs +++ b/crates/common/src/manager/defaults.rs @@ -5,6 +5,10 @@ */ use crate::auth::permissions::DefaultPermissions; +use aws_lc_rs::{ + rand::SystemRandom, + signature::{ECDSA_P256_SHA256_FIXED_SIGNING, EcdsaKeyPair}, +}; use registry::{ schema::{ enums::*, @@ -13,10 +17,6 @@ use registry::{ }, types::{duration::Duration, error::Error, list::List, map::Map}, }; -use aws_lc_rs::{ - rand::SystemRandom, - signature::{ECDSA_P256_SHA256_FIXED_SIGNING, EcdsaKeyPair}, -}; use std::str::FromStr; use store::{ rand::{Rng, distr::Alphanumeric, rng}, @@ -71,7 +71,7 @@ async fn insert_safe_defaults(bp: &mut Bootstrap) -> trc::Result<()> { { for object in [ MtaInboundThrottle { - description: "Sender IP throttle".to_string().into(), + description: "Sender IP throttle".to_string(), enable: true, key: Map::new(vec![MtaInboundThrottleKey::RemoteIp]), rate: Rate { @@ -81,7 +81,7 @@ async fn insert_safe_defaults(bp: &mut Bootstrap) -> trc::Result<()> { ..Default::default() }, MtaInboundThrottle { - description: "Sender address to recipient throttle".to_string().into(), + description: "Sender address to recipient throttle".to_string(), enable: true, key: Map::new(vec![ MtaInboundThrottleKey::SenderDomain, @@ -503,7 +503,7 @@ async fn insert_safe_defaults(bp: &mut Bootstrap) -> trc::Result<()> { } } - #[cfg(not(feature = "test_mode"))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] if bp.registry.count_object(ObjectType::Asn).await? == 0 { bp.registry .write(RegistryWrite::insert( @@ -520,6 +520,18 @@ async fn insert_safe_defaults(bp: &mut Bootstrap) -> trc::Result<()> { .await?; } + if bp.registry.count_object(ObjectType::TracingStore).await? == 0 { + bp.registry + .write(RegistryWrite::insert(&TracingStore::Default.into())) + .await?; + } + + if bp.registry.count_object(ObjectType::MetricsStore).await? == 0 { + bp.registry + .write(RegistryWrite::insert(&MetricsStore::Default.into())) + .await?; + } + if bp.registry.count_object(ObjectType::Tracer).await? == 0 { bp.registry .write(RegistryWrite::insert( diff --git a/crates/common/src/network/acme/http.rs b/crates/common/src/network/acme/http.rs index d9bac583..518692f9 100644 --- a/crates/common/src/network/acme/http.rs +++ b/crates/common/src/network/acme/http.rs @@ -25,7 +25,7 @@ pub(crate) async fn https( .timeout(Duration::from_secs(30)) .http1_only(); - #[cfg(debug_assertions)] + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] { builder = builder.danger_accept_invalid_certs( url.starts_with("https://localhost") || url.starts_with("https://127.0.0.1"), diff --git a/crates/common/src/telemetry/metrics/mod.rs b/crates/common/src/telemetry/metrics/mod.rs index 38289598..0c0addf6 100644 --- a/crates/common/src/telemetry/metrics/mod.rs +++ b/crates/common/src/telemetry/metrics/mod.rs @@ -13,3 +13,6 @@ pub mod prometheus; #[cfg(feature = "enterprise")] pub mod store; // SPDX-SnippetEnd + +#[cfg(any(feature = "dev_mode", feature = "test_mode"))] +pub mod test_data; diff --git a/crates/common/src/telemetry/metrics/store.rs b/crates/common/src/telemetry/metrics/store.rs index 30c3c3d3..70055241 100644 --- a/crates/common/src/telemetry/metrics/store.rs +++ b/crates/common/src/telemetry/metrics/store.rs @@ -86,10 +86,10 @@ impl MetricsStore for Store { *history = reading; if diff > 0 { - #[cfg(not(feature = "test_mode"))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] let metric_id = history_guard.id_generator.generate(); - #[cfg(feature = "test_mode")] + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] let metric_id = _timestamp .map(|timestamp| { SnowflakeIdGenerator::global_id_from_timestamp(timestamp).unwrap() @@ -113,10 +113,10 @@ impl MetricsStore for Store { if matches!(metric, MetricType::QueueCount | MetricType::ServerMemory) { let value = gauge.get(); if value > 0 { - #[cfg(not(feature = "test_mode"))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] let metric_id = history_guard.id_generator.generate(); - #[cfg(feature = "test_mode")] + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] let metric_id = _timestamp .map(|timestamp| { SnowflakeIdGenerator::global_id_from_timestamp(timestamp).unwrap() @@ -144,6 +144,10 @@ impl MetricsStore for Store { | MetricType::DeliveryTotalTime | MetricType::DeliveryAttemptTime | MetricType::DnsLookupTime + | MetricType::StoreDataReadTime + | MetricType::StoreDataWriteTime + | MetricType::StoreBlobReadTime + | MetricType::StoreBlobWriteTime ) { let history = history_guard.histograms.entry(metric).or_default(); let sum = histogram.sum(); @@ -153,10 +157,10 @@ impl MetricsStore for Store { history.sum = sum; history.count = count; if diff_sum > 0 || diff_count > 0 { - #[cfg(not(feature = "test_mode"))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] let metric_id = history_guard.id_generator.generate(); - #[cfg(feature = "test_mode")] + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] let metric_id = _timestamp .map(|timestamp| { SnowflakeIdGenerator::global_id_from_timestamp(timestamp).unwrap() diff --git a/crates/common/src/telemetry/metrics/test_data.rs b/crates/common/src/telemetry/metrics/test_data.rs new file mode 100644 index 00000000..a5ebeaa9 --- /dev/null +++ b/crates/common/src/telemetry/metrics/test_data.rs @@ -0,0 +1,91 @@ +/* + * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * + * SPDX-License-Identifier: LicenseRef-SEL + * + * This file is subject to the Stalwart Enterprise License Agreement (SEL) and + * is NOT open source software. + * + */ + +use crate::{ + Server, + telemetry::metrics::store::{MetricsStore, SharedMetricHistory}, +}; +use std::time::Duration; +use store::rand::{self, Rng}; +use trc::*; + +impl Server { + pub async fn insert_test_metrics(&self) { + self.metrics_store() + .purge_metrics(Duration::from_secs(0)) + .await + .unwrap(); + let mut start_time = store::write::now() - (90 * 24 * 60 * 60); + let timestamp = store::write::now(); + let history = SharedMetricHistory::default(); + + while start_time <= timestamp { + for event_type in [ + EventType::Smtp(SmtpEvent::ConnectionStart), + EventType::Imap(ImapEvent::ConnectionStart), + EventType::Pop3(Pop3Event::ConnectionStart), + EventType::ManageSieve(ManageSieveEvent::ConnectionStart), + EventType::Http(HttpEvent::ConnectionStart), + EventType::Delivery(DeliveryEvent::AttemptStart), + EventType::Queue(QueueEvent::MessageQueued), + EventType::Queue(QueueEvent::AuthenticatedMessageQueued), + EventType::Queue(QueueEvent::DsnQueued), + EventType::Queue(QueueEvent::ReportQueued), + EventType::MessageIngest(MessageIngestEvent::Ham), + EventType::MessageIngest(MessageIngestEvent::Spam), + EventType::Auth(AuthEvent::Failed), + EventType::Security(SecurityEvent::AuthenticationBan), + EventType::Security(SecurityEvent::ScanBan), + EventType::Security(SecurityEvent::AbuseBan), + EventType::Security(SecurityEvent::LoiterBan), + EventType::Security(SecurityEvent::IpBlocked), + EventType::IncomingReport(IncomingReportEvent::DmarcReport), + EventType::IncomingReport(IncomingReportEvent::DmarcReportWithWarnings), + EventType::IncomingReport(IncomingReportEvent::TlsReport), + EventType::IncomingReport(IncomingReportEvent::TlsReportWithWarnings), + ] { + // Generate a random value between 0 and 100 + Collector::update_event_counter(event_type, rand::rng().random_range(0..=100)) + } + + Collector::update_gauge(MetricType::QueueCount, rand::rng().random_range(0..=1000)); + Collector::update_gauge( + MetricType::ServerMemory, + rand::rng().random_range(100 * 1024 * 1024..=300 * 1024 * 1024), + ); + Collector::update_gauge(MetricType::UserCount, rand::rng().random_range(100..=500)); + Collector::update_gauge(MetricType::DomainCount, rand::rng().random_range(10..=50)); + + for metric_type in [ + MetricType::MessageIngestTime, + MetricType::MessageIngestIndexTime, + MetricType::DeliveryTotalTime, + MetricType::DeliveryAttemptTime, + MetricType::DnsLookupTime, + MetricType::StoreDataReadTime, + MetricType::StoreDataWriteTime, + MetricType::StoreBlobReadTime, + MetricType::StoreBlobWriteTime, + ] { + Collector::update_histogram(metric_type, rand::rng().random_range(2..=1000)) + } + Collector::update_histogram( + MetricType::DeliveryTotalTime, + rand::rng().random_range(1000..=5000), + ); + + self.metrics_store() + .write_metrics(start_time.into(), history.clone()) + .await + .unwrap(); + start_time += 60 * 60; + } + } +} diff --git a/crates/common/src/telemetry/tracers/mod.rs b/crates/common/src/telemetry/tracers/mod.rs index 02d20bfa..6e5da60a 100644 --- a/crates/common/src/telemetry/tracers/mod.rs +++ b/crates/common/src/telemetry/tracers/mod.rs @@ -16,3 +16,103 @@ pub mod stdout; #[cfg(feature = "enterprise")] pub mod store; // SPDX-SnippetEnd + +use registry::{ + schema::structs::{ + Trace, TraceEvent, TraceKeyValue, TraceValue, TraceValueBoolean, TraceValueDuration, + TraceValueEvent, TraceValueFloat, TraceValueInteger, TraceValueIpAddr, TraceValueList, + TraceValueString, TraceValueUTCDateTime, TraceValueUnsignedInt, + }, + types::{datetime::UTCDateTime, ipaddr::IpAddr, list::List}, +}; +use trc::{Event, EventDetails, Value}; + +pub trait TraceEvents { + fn build_trace_events<'x>( + span_events: impl IntoIterator>, + num_events: usize, + ) -> Vec; + + fn from_events<'x>( + span_events: impl IntoIterator>, + num_events: usize, + ) -> Self; +} + +impl TraceEvents for Trace { + fn build_trace_events<'x>( + span_events: impl IntoIterator>, + num_events: usize, + ) -> Vec { + let mut events = Vec::with_capacity(num_events); + + for event in span_events { + let mut key_values = Vec::with_capacity(event.keys.len()); + for (key, value) in &event.keys { + key_values.push(TraceKeyValue { + key: *key, + value: map_value(value), + }); + } + + events.push(TraceEvent { + event: event.inner.typ, + timestamp: UTCDateTime::from_timestamp(event.inner.timestamp as i64), + key_values: key_values.into(), + }); + } + + events + } + + fn from_events<'x>( + span_events: impl IntoIterator>, + num_events: usize, + ) -> Self { + Trace { + events: Self::build_trace_events(span_events, num_events).into(), + } + } +} + +fn map_value(value: &Value) -> TraceValue { + match value { + Value::String(value) => TraceValue::String(TraceValueString { + value: value.to_string(), + }), + Value::UInt(value) => TraceValue::UnsignedInt(TraceValueUnsignedInt { value: *value }), + Value::Int(value) => TraceValue::Integer(TraceValueInteger { value: *value }), + Value::Float(value) => TraceValue::Float(TraceValueFloat { + value: (*value).into(), + }), + Value::Timestamp(value) => TraceValue::UTCDateTime(TraceValueUTCDateTime { + value: UTCDateTime::from_timestamp(*value as i64), + }), + Value::Duration(value) => TraceValue::Duration(TraceValueDuration { value: *value }), + Value::Bytes(items) => TraceValue::String(TraceValueString { + value: String::from_utf8_lossy(items).to_string(), + }), + Value::Bool(value) => TraceValue::Boolean(TraceValueBoolean { value: *value }), + Value::Ipv4(ipv4_addr) => TraceValue::IpAddr(TraceValueIpAddr { + value: IpAddr((*ipv4_addr).into()), + }), + Value::Ipv6(ipv6_addr) => TraceValue::IpAddr(TraceValueIpAddr { + value: IpAddr((*ipv6_addr).into()), + }), + Value::Event(event) => TraceValue::Event(TraceValueEvent { + value: event + .keys() + .iter() + .map(|(k, v)| TraceKeyValue { + key: *k, + value: map_value(v), + }) + .collect(), + event: event.event_type(), + }), + Value::Array(values) => TraceValue::List(TraceValueList { + value: List::from_iter(values.iter().map(map_value)), + }), + Value::None => TraceValue::Null, + } +} diff --git a/crates/common/src/telemetry/tracers/store.rs b/crates/common/src/telemetry/tracers/store.rs index 92db1625..01f37e4c 100644 --- a/crates/common/src/telemetry/tracers/store.rs +++ b/crates/common/src/telemetry/tracers/store.rs @@ -8,17 +8,14 @@ * */ -use crate::config::telemetry::StoreTracer; +use crate::{config::telemetry::StoreTracer, telemetry::tracers::TraceEvents}; use ahash::AHashMap; use registry::{ pickle::Pickle, schema::structs::{ - Task, TaskIndexTrace, TaskStatus, Trace, TraceEvent, TraceKeyValue, TraceValue, - TraceValueBoolean, TraceValueDuration, TraceValueEvent, TraceValueFloat, TraceValueInteger, - TraceValueIpAddr, TraceValueList, TraceValueString, TraceValueUTCDateTime, - TraceValueUnsignedInt, + Task, TaskIndexTrace, TaskStatus, Trace, TraceKeyValue, TraceValue, TraceValueIpAddr, + TraceValueList, TraceValueString, TraceValueUnsignedInt, }, - types::{datetime::UTCDateTime, ipaddr::IpAddr, list::List}, }; use std::{collections::HashSet, future::Future, time::Duration}; use store::{ @@ -27,8 +24,8 @@ use store::{ write::{BatchBuilder, SearchIndex, TelemetryClass, ValueClass}, }; use trc::{ - AddContext, AuthEvent, Event, EventDetails, EventType, Key, MessageIngestEvent, - OutgoingReportEvent, QueueEvent, Value, ipc::subscriber::SubscriberBuilder, + AddContext, AuthEvent, EventType, Key, MessageIngestEvent, OutgoingReportEvent, QueueEvent, + Value, ipc::subscriber::SubscriberBuilder, }; use utils::snowflake::SnowflakeIdGenerator; @@ -61,7 +58,7 @@ pub(crate) fn spawn_store_tracer(builder: SubscriberBuilder, settings: StoreTrac batch .set( ValueClass::Telemetry(TelemetryClass::Span(span_id)), - map_events( + Trace::from_events( [span.as_ref()] .into_iter() .chain(events.iter().map(|event| event.as_ref())) @@ -88,75 +85,6 @@ pub(crate) fn spawn_store_tracer(builder: SubscriberBuilder, settings: StoreTrac }); } -fn map_events<'x>( - span_events: impl IntoIterator>, - num_events: usize, -) -> Trace { - let mut events = Vec::with_capacity(num_events); - - for event in span_events { - let mut key_values = Vec::with_capacity(event.keys.len()); - for (key, value) in &event.keys { - key_values.push(TraceKeyValue { - key: *key, - value: map_value(value), - }); - } - - events.push(TraceEvent { - event: event.inner.typ, - timestamp: UTCDateTime::from_timestamp(event.inner.timestamp as i64), - key_values: key_values.into(), - }); - } - - Trace { - events: events.into(), - } -} - -fn map_value(value: &Value) -> TraceValue { - match value { - Value::String(value) => TraceValue::String(TraceValueString { - value: value.to_string(), - }), - Value::UInt(value) => TraceValue::UnsignedInt(TraceValueUnsignedInt { value: *value }), - Value::Int(value) => TraceValue::Integer(TraceValueInteger { value: *value }), - Value::Float(value) => TraceValue::Float(TraceValueFloat { - value: (*value).into(), - }), - Value::Timestamp(value) => TraceValue::UTCDateTime(TraceValueUTCDateTime { - value: UTCDateTime::from_timestamp(*value as i64), - }), - Value::Duration(value) => TraceValue::Duration(TraceValueDuration { value: *value }), - Value::Bytes(items) => TraceValue::String(TraceValueString { - value: String::from_utf8_lossy(items).to_string(), - }), - Value::Bool(value) => TraceValue::Boolean(TraceValueBoolean { value: *value }), - Value::Ipv4(ipv4_addr) => TraceValue::IpAddr(TraceValueIpAddr { - value: IpAddr((*ipv4_addr).into()), - }), - Value::Ipv6(ipv6_addr) => TraceValue::IpAddr(TraceValueIpAddr { - value: IpAddr((*ipv6_addr).into()), - }), - Value::Event(event) => TraceValue::Event(TraceValueEvent { - value: event - .keys() - .iter() - .map(|(k, v)| TraceKeyValue { - key: *k, - value: map_value(v), - }) - .collect(), - event: event.event_type(), - }), - Value::Array(values) => TraceValue::List(TraceValueList { - value: List::from_iter(values.iter().map(map_value)), - }), - Value::None => TraceValue::Null, - } -} - pub trait TracingStore: Sync + Send { fn purge_spans( &self, diff --git a/crates/dav/Cargo.toml b/crates/dav/Cargo.toml index d350a76d..514fc568 100644 --- a/crates/dav/Cargo.toml +++ b/crates/dav/Cargo.toml @@ -26,4 +26,5 @@ chrono = "0.4.40" [features] test_mode = [] +dev_mode = [] enterprise = [] diff --git a/crates/dav/src/calendar/copy_move.rs b/crates/dav/src/calendar/copy_move.rs index 16c821ae..98cc424e 100644 --- a/crates/dav/src/calendar/copy_move.rs +++ b/crates/dav/src/calendar/copy_move.rs @@ -70,7 +70,7 @@ impl CalendarCopyMoveRequestHandler for Server { let from_resource = from_resources .by_path(from_resource_name) .ok_or(DavError::Code(StatusCode::NOT_FOUND))?; - #[cfg(not(debug_assertions))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] if is_move && from_resource.is_container() && self diff --git a/crates/dav/src/calendar/delete.rs b/crates/dav/src/calendar/delete.rs index 33a21b3e..4c209924 100644 --- a/crates/dav/src/calendar/delete.rs +++ b/crates/dav/src/calendar/delete.rs @@ -82,7 +82,7 @@ impl CalendarDeleteRequestHandler for Server { let mut batch = BatchBuilder::new(); if delete_resource.is_container() { // Deleting the default calendar is not allowed - #[cfg(not(debug_assertions))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] if self .core .groupware diff --git a/crates/http-proto/Cargo.toml b/crates/http-proto/Cargo.toml index 312de3b2..3d0dcd0b 100644 --- a/crates/http-proto/Cargo.toml +++ b/crates/http-proto/Cargo.toml @@ -20,4 +20,5 @@ compact_str = "0.9.0" [features] test_mode = [] +dev_mode = [] enterprise = [] diff --git a/crates/http-proto/src/context.rs b/crates/http-proto/src/context.rs index d358c93c..a79f18ee 100644 --- a/crates/http-proto/src/context.rs +++ b/crates/http-proto/src/context.rs @@ -18,22 +18,31 @@ impl<'x> HttpContext<'x> { Self { session, req } } + #[allow(unused_variables)] pub fn resolve_response_url(&self, server: &Server) -> String { if self.session.is_tls { - #[cfg(not(feature = "test_mode"))] + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] { server.core.network.http.url_https.clone() } - #[cfg(feature = "test_mode")] + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] { format!("https://127.0.0.1:{}", self.session.local_port) } } else { - format!( - "{}:{}", - server.core.network.http.url_http, self.session.local_port - ) + #[cfg(not(any(feature = "dev_mode", feature = "test_mode")))] + { + format!( + "{}:{}", + server.core.network.http.url_http, self.session.local_port + ) + } + + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] + { + format!("http://127.0.0.1:{}", self.session.local_port) + } } } diff --git a/crates/http/Cargo.toml b/crates/http/Cargo.toml index fb279cf3..73240089 100644 --- a/crates/http/Cargo.toml +++ b/crates/http/Cargo.toml @@ -49,4 +49,5 @@ hashify = { version = "0.2" } [features] test_mode = [] +dev_mode = [] enterprise = [] diff --git a/crates/http/src/api/diagnose.rs b/crates/http/src/api/diagnose.rs index e98401cd..5c6f3109 100644 --- a/crates/http/src/api/diagnose.rs +++ b/crates/http/src/api/diagnose.rs @@ -87,6 +87,7 @@ pub(crate) enum DeliveryStage { reason: String, }, IpLookupStart, + #[serde(rename_all = "camelCase")] IpLookupSuccess { remote_ips: Vec, elapsed: u64, @@ -95,6 +96,7 @@ pub(crate) enum DeliveryStage { reason: String, elapsed: u64, }, + #[serde(rename_all = "camelCase")] ConnectionStart { remote_ip: IpAddr, }, diff --git a/crates/http/src/api/mod.rs b/crates/http/src/api/mod.rs index c71daf0b..11be5a6d 100644 --- a/crates/http/src/api/mod.rs +++ b/crates/http/src/api/mod.rs @@ -25,14 +25,13 @@ use common::{ }; use http_body_util::{StreamBody, combinators::BoxBody}; use http_proto::{ - HttpRequest, HttpResponse, HttpSessionData, JsonResponse, ToHttpResponse, + HttpRequest, HttpResponse, HttpSessionData, ToHttpResponse, request::{decode_path_element, fetch_body}, }; use hyper::{Method, StatusCode, header}; use jmap::api::{ToJmapHttpResponse, ToRequestError}; use jmap_proto::error::request::RequestError; use registry::schema::enums::Permission; -use serde_json::json; use std::time::Duration; use utils::url_params::UrlParams; @@ -76,7 +75,7 @@ impl ManagementApi for Server { .await } "discover" => { - if let Some(email) = path.get(2).copied() { + if let Some(email) = path.get(1).copied() { self.is_http_anonymous_request_allowed(session.remote_ip) .await?; self.handle_discover_request(req, session, decode_path_element(email).as_ref()) @@ -112,10 +111,17 @@ impl ManagementApi for Server { access_token.enforce_permission(Permission::LiveTracing)?; // Issue a live telemetry token valid for 60 seconds - Ok(JsonResponse::new(json!({ - "data": self.encode_access_token(GrantType::LiveTracing, account_id, "web", 60).await?, - })) - .into_http_response()) + Ok(HttpResponse::new(StatusCode::OK) + .with_no_cache() + .with_text_body( + self.encode_access_token( + GrantType::LiveTracing, + account_id, + "web", + 60, + ) + .await?, + )) } #[cfg(feature = "enterprise")] Some("metrics") if self.core.is_enterprise_edition() => { @@ -123,10 +129,17 @@ impl ManagementApi for Server { access_token.enforce_permission(Permission::LiveMetrics)?; // Issue a live telemetry token valid for 60 seconds - Ok(JsonResponse::new(json!({ - "data": self.encode_access_token(GrantType::LiveMetrics, account_id, "web", 60).await?, - })) - .into_http_response()) + Ok(HttpResponse::new(StatusCode::OK) + .with_no_cache() + .with_text_body( + self.encode_access_token( + GrantType::LiveMetrics, + account_id, + "web", + 60, + ) + .await?, + )) } // SPDX-SnippetEnd Some("delivery") => { @@ -134,10 +147,17 @@ impl ManagementApi for Server { access_token.enforce_permission(Permission::LiveDeliveryTest)?; // Issue a live telemetry token valid for 60 seconds - Ok(JsonResponse::new(json!({ - "data": self.encode_access_token(GrantType::Diagnose, account_id, "web", 60).await?, - })) - .into_http_response()) + Ok(HttpResponse::new(StatusCode::OK) + .with_no_cache() + .with_text_body( + self.encode_access_token( + GrantType::LiveDelivery, + account_id, + "web", + 60, + ) + .await?, + )) } Some("tracing") | Some("metrics") => { Err(trc::ResourceEvent::NotFound @@ -189,7 +209,7 @@ impl ManagementApi for Server { // SPDX-FileCopyrightText: 2020 Stalwart Labs LLC // SPDX-License-Identifier: LicenseRef-SEL #[cfg(feature = "enterprise")] - ("traces", _, &Method::GET) if self.core.is_enterprise_edition() => { + ("tracing", _, &Method::GET) if self.core.is_enterprise_edition() => { use crate::api::telemetry::TelemetryApi; self.handle_telemetry_api_request(req, true, &access_token) @@ -203,7 +223,7 @@ impl ManagementApi for Server { .await } // SPDX-SnippetEnd - ("traces" | "metrics", _, &Method::GET) => { + ("tracing" | "metrics", _, &Method::GET) => { Err(trc::ResourceEvent::NotFound .ctx(trc::Key::Details, "Enterprise feature")) } @@ -227,12 +247,12 @@ impl ManagementApi for Server { #[cfg(feature = "enterprise")] if self.core.is_enterprise_edition() { let path = req.uri().path(); - let (grant_type, permissions) = if path.starts_with("/api/telemetry/traces") { + let (grant_type, permissions) = if path.starts_with("/api/live/tracing") { (GrantType::LiveTracing, Permission::LiveTracing) - } else if path.starts_with("/api/telemetry/metrics") { + } else if path.starts_with("/api/live/metrics") { (GrantType::LiveMetrics, Permission::LiveMetrics) - } else if path.starts_with("/api/diagnose") { - (GrantType::Diagnose, Permission::LiveDeliveryTest) + } else if path.starts_with("/api/live/delivery") { + (GrantType::LiveDelivery, Permission::LiveDeliveryTest) } else { return Err(trc::ResourceEvent::NotFound.into_err()); }; diff --git a/crates/http/src/api/telemetry.rs b/crates/http/src/api/telemetry.rs index 8f080bd4..45c1417a 100644 --- a/crates/http/src/api/telemetry.rs +++ b/crates/http/src/api/telemetry.rs @@ -8,7 +8,7 @@ * */ -use common::{Server, auth::AccessToken}; +use common::{Server, auth::AccessToken, telemetry::tracers::TraceEvents}; use http_body_util::{StreamBody, combinators::BoxBody}; use http_proto::*; use hyper::{ @@ -16,7 +16,7 @@ use hyper::{ body::{Bytes, Frame}, }; use mail_parser::DateTime; -use registry::schema::enums::Permission; +use registry::schema::{enums::Permission, structs::Trace}; use std::future::Future; use std::{ fmt::Write, @@ -26,7 +26,6 @@ use store::ahash::{AHashMap, AHashSet}; use trc::{ Collector, EventType, Key, MetricType, Value, ipc::{bitset::Bitset, subscriber::SubscriberBuilder}, - serializers::json::JsonEventSerializer, }; use utils::url_params::UrlParams; @@ -84,6 +83,8 @@ impl TelemetryApi for Server { let mut last_message = Instant::now() - throttle; let mut timeout = ping_interval; + yield Ok(Frame::data(ping_payload.clone())); + loop { match tokio::time::timeout(timeout, rx.recv()).await { Ok(Some(event_batch)) => { @@ -148,11 +149,11 @@ impl TelemetryApi for Server { let elapsed = last_message.elapsed(); if elapsed >= throttle { last_message = Instant::now(); + let num_events = events.len(); yield Ok(Frame::data(Bytes::from(format!( "event: trace\ndata: {}\n\n", serde_json::to_string( - &JsonEventSerializer::new(std::mem::take(&mut events)) - .with_description()).unwrap_or_default() + &Trace::build_trace_events(events.iter().map(|e| e.as_ref()), num_events)).unwrap_or_default() )))); ping_interval @@ -231,7 +232,7 @@ impl TelemetryApi for Server { } let _ = write!( &mut metrics, - "{{\"id\":\"{}\",\"type\":\"counter\",\"value\":{}}}", + "{{\"metric\":\"{}\",\"@type\":\"Counter\",\"count\":{}}}", counter.id().as_str(), counter.value() ); @@ -246,7 +247,7 @@ impl TelemetryApi for Server { } let _ = write!( &mut metrics, - "{{\"id\":\"{}\",\"type\":\"gauge\",\"value\":{}}}", + "{{\"metric\":\"{}\",\"@type\":\"Gauge\",\"count\":{}}}", gauge.id().as_str(), gauge.get() ); @@ -261,7 +262,7 @@ impl TelemetryApi for Server { } let _ = write!( &mut metrics, - "{{\"id\":\"{}\",\"type\":\"histogram\",\"count\":{},\"sum\":{}}}", + "{{\"metric\":\"{}\",\"@type\":\"Histogram\",\"count\":{},\"sum\":{}}}", histogram.id().as_str(), histogram.count(), histogram.sum() diff --git a/crates/http/src/auth/authenticate.rs b/crates/http/src/auth/authenticate.rs index 800120d5..5d9f2714 100644 --- a/crates/http/src/auth/authenticate.rs +++ b/crates/http/src/auth/authenticate.rs @@ -41,7 +41,7 @@ impl Authenticator for Server { if access_token.revision() == http_cache.revision { // Enforce authenticated rate limit return self - .is_http_authenticated_request_allowed(&access_token) + .is_http_authenticated_request_allowed(&access_token, session.remote_ip) .await .map(|in_flight| (in_flight, access_token)); } @@ -103,7 +103,7 @@ impl Authenticator for Server { ); // Enforce authenticated rate limit - self.is_http_authenticated_request_allowed(&access_token) + self.is_http_authenticated_request_allowed(&access_token, session.remote_ip) .await .map(|in_flight| (in_flight, access_token)) } else { diff --git a/crates/http/src/auth/oauth/auth.rs b/crates/http/src/auth/oauth/auth.rs index be0304f9..7cf6c77e 100644 --- a/crates/http/src/auth/oauth/auth.rs +++ b/crates/http/src/auth/oauth/auth.rs @@ -77,6 +77,7 @@ pub trait OAuthApiHandler: Sync + Send { #[serde(tag = "type")] #[serde(rename_all = "camelCase")] pub enum LoginRequest { + #[serde(rename_all = "camelCase")] AuthCode { account_name: String, account_secret: String, @@ -97,6 +98,7 @@ pub enum LoginRequest { #[serde(default)] state: Option, }, + #[serde(rename_all = "camelCase")] AuthDevice { account_name: String, account_secret: String, @@ -167,6 +169,7 @@ impl OAuthApiHandler for Server { .as_ref() .is_some_and(|uri| uri.starts_with("http://")) { + #[cfg(not(feature = "dev_mode"))] return Err(trc::AuthEvent::Error .into_err() .details("Redirect URI must be HTTPS.")); @@ -175,16 +178,6 @@ impl OAuthApiHandler for Server { // Parse and validate PKCE challenge (RFC 7636). let pkce_challenge = match code_challenge { Some(challenge) => { - if !(43..=128).contains(&challenge.len()) - && challenge.bytes().all(|b| { - b.is_ascii_alphanumeric() || matches!(b, b'-' | b'.' | b'_' | b'~') - }) - { - return Err(trc::AuthEvent::Error - .into_err() - .details("Invalid PKCE code_challenge.")); - } - // Default to "plain" when the method is omitted, per RFC 7636 4.3. match code_challenge_method.as_deref().unwrap_or("plain") { "S256" => PkceCodeChallenge::S256(challenge), diff --git a/crates/http/src/auth/oauth/openid.rs b/crates/http/src/auth/oauth/openid.rs index b4493e34..747846d0 100644 --- a/crates/http/src/auth/oauth/openid.rs +++ b/crates/http/src/auth/oauth/openid.rs @@ -64,7 +64,7 @@ impl OpenIdHandler for Server { let base_url = HttpContext::new(session, req).resolve_response_url(self); Ok(JsonResponse::new(OpenIdMetadata { - authorization_endpoint: format!("{base_url}/authorize/code",), + authorization_endpoint: format!("{base_url}/login",), token_endpoint: format!("{base_url}/auth/token"), userinfo_endpoint: format!("{base_url}/auth/userinfo"), jwks_uri: format!("{base_url}/auth/jwks.json"), diff --git a/crates/http/src/auth/oauth/token.rs b/crates/http/src/auth/oauth/token.rs index ea48c441..99f2b1d5 100644 --- a/crates/http/src/auth/oauth/token.rs +++ b/crates/http/src/auth/oauth/token.rs @@ -339,7 +339,7 @@ impl TokenHandler for Server { fn verify_pkce(stored: &ArchivedPkceCodeChallenge, verifier: Option<&str>) -> bool { let is_valid_pkce_challenge = |challenge: &str| { - !(43..=128).contains(&challenge.len()) + (43..=128).contains(&challenge.len()) && challenge .bytes() .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'-' | b'.' | b'_' | b'~')) diff --git a/crates/jmap-proto/src/request/method.rs b/crates/jmap-proto/src/request/method.rs index 52ada7e7..a1193453 100644 --- a/crates/jmap-proto/src/request/method.rs +++ b/crates/jmap-proto/src/request/method.rs @@ -144,6 +144,7 @@ impl MethodName { (MethodFunction::Get, MethodObject::AddressBook) => "AddressBook/get", (MethodFunction::Changes, MethodObject::AddressBook) => "AddressBook/changes", (MethodFunction::Set, MethodObject::AddressBook) => "AddressBook/set", + (MethodFunction::Query, MethodObject::AddressBook) => "AddressBook/query", (MethodFunction::Get, MethodObject::ContactCard) => "ContactCard/get", (MethodFunction::Changes, MethodObject::ContactCard) => "ContactCard/changes", @@ -172,6 +173,7 @@ impl MethodName { (MethodFunction::Get, MethodObject::Calendar) => "Calendar/get", (MethodFunction::Changes, MethodObject::Calendar) => "Calendar/changes", (MethodFunction::Set, MethodObject::Calendar) => "Calendar/set", + (MethodFunction::Query, MethodObject::Calendar) => "Calendar/query", (MethodFunction::Get, MethodObject::CalendarEvent) => "CalendarEvent/get", (MethodFunction::Changes, MethodObject::CalendarEvent) => "CalendarEvent/changes", @@ -277,6 +279,7 @@ impl MethodName { "AddressBook/get" => (MethodObject::AddressBook, MethodFunction::Get), "AddressBook/changes" => (MethodObject::AddressBook, MethodFunction::Changes), "AddressBook/set" => (MethodObject::AddressBook, MethodFunction::Set), + "AddressBook/query" => (MethodObject::AddressBook, MethodFunction::Query), "ContactCard/get" => (MethodObject::ContactCard, MethodFunction::Get), "ContactCard/changes" => (MethodObject::ContactCard, MethodFunction::Changes), @@ -301,6 +304,7 @@ impl MethodName { "Calendar/get" => (MethodObject::Calendar, MethodFunction::Get), "Calendar/changes" => (MethodObject::Calendar, MethodFunction::Changes), "Calendar/set" => (MethodObject::Calendar, MethodFunction::Set), + "Calendar/query" => (MethodObject::Calendar, MethodFunction::Query), "CalendarEvent/get" => (MethodObject::CalendarEvent, MethodFunction::Get), "CalendarEvent/changes" => (MethodObject::CalendarEvent, MethodFunction::Changes), diff --git a/crates/jmap-proto/src/request/mod.rs b/crates/jmap-proto/src/request/mod.rs index bda8ef27..22c246cd 100644 --- a/crates/jmap-proto/src/request/mod.rs +++ b/crates/jmap-proto/src/request/mod.rs @@ -136,8 +136,10 @@ pub enum QueryRequestMethod { Sieve(QueryRequest), Principal(QueryRequest), Quota(QueryRequest), + AddressBook(QueryRequest), ContactCard(QueryRequest), FileNode(QueryRequest), + Calendar(QueryRequest), CalendarEvent(QueryRequest), CalendarEventNotification(QueryRequest), ShareNotification(QueryRequest), diff --git a/crates/jmap-proto/src/request/parser.rs b/crates/jmap-proto/src/request/parser.rs index 3e029519..49987b8f 100644 --- a/crates/jmap-proto/src/request/parser.rs +++ b/crates/jmap-proto/src/request/parser.rs @@ -32,7 +32,8 @@ impl<'x> Request<'x> { } Err(err) => Err(trc::JmapEvent::NotRequest .into_err() - .details(err.to_string())), + .reason(err.to_string()) + .details(String::from_utf8_lossy(json).into_owned())), } } else { Err(trc::LimitEvent::SizeRequest.into_err()) @@ -413,6 +414,13 @@ impl<'de> Visitor<'de> for CallVisitor { return Err(de::Error::invalid_length(1, &self)); } }, + (MethodFunction::Query, MethodObject::Calendar) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Query(QueryRequestMethod::Calendar(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, (MethodFunction::Query, MethodObject::CalendarEvent) => match seq.next_element() { Ok(Some(value)) => RequestMethod::Query(QueryRequestMethod::CalendarEvent(value)), Err(err) => RequestMethod::invalid(err), @@ -431,6 +439,13 @@ impl<'de> Visitor<'de> for CallVisitor { } } } + (MethodFunction::Query, MethodObject::AddressBook) => match seq.next_element() { + Ok(Some(value)) => RequestMethod::Query(QueryRequestMethod::AddressBook(value)), + Err(err) => RequestMethod::invalid(err), + Ok(None) => { + return Err(de::Error::invalid_length(1, &self)); + } + }, (MethodFunction::Query, MethodObject::ContactCard) => match seq.next_element() { Ok(Some(value)) => RequestMethod::Query(QueryRequestMethod::ContactCard(value)), Err(err) => RequestMethod::invalid(err), diff --git a/crates/jmap/src/api/auth.rs b/crates/jmap/src/api/auth.rs index 3b945eaa..f86a5157 100644 --- a/crates/jmap/src/api/auth.rs +++ b/crates/jmap/src/api/auth.rs @@ -281,8 +281,10 @@ impl JmapAuthorization for AccessToken { QueryRequestMethod::Sieve(_) => Permission::JmapSieveScriptQuery, QueryRequestMethod::Principal(_) => Permission::JmapPrincipalQuery, QueryRequestMethod::Quota(_) => Permission::JmapQuotaQuery, + QueryRequestMethod::AddressBook(_) => Permission::JmapAddressBookGet, QueryRequestMethod::ContactCard(_) => Permission::JmapContactCardQuery, QueryRequestMethod::FileNode(_) => Permission::JmapFileNodeQuery, + QueryRequestMethod::Calendar(_) => Permission::JmapCalendarGet, QueryRequestMethod::CalendarEvent(_) => Permission::JmapCalendarEventQuery, QueryRequestMethod::CalendarEventNotification(_) => { Permission::JmapCalendarEventNotificationQuery diff --git a/crates/jmap/src/api/request.rs b/crates/jmap/src/api/request.rs index 9dbe1411..a654c53e 100644 --- a/crates/jmap/src/api/request.rs +++ b/crates/jmap/src/api/request.rs @@ -394,6 +394,12 @@ impl RequestHandler for Server { self.quota_query(req, access_token).await?.into() } + QueryRequestMethod::AddressBook(mut req) => { + set_account_id_if_missing(&mut req.account_id, access_token); + access_token.assert_has_access(req.account_id, Collection::AddressBook)?; + + self.address_book_query(req, access_token).await?.into() + } QueryRequestMethod::ContactCard(mut req) => { set_account_id_if_missing(&mut req.account_id, access_token); access_token.assert_has_access(req.account_id, Collection::ContactCard)?; @@ -406,6 +412,12 @@ impl RequestHandler for Server { self.file_node_query(req, access_token).await?.into() } + QueryRequestMethod::Calendar(mut req) => { + set_account_id_if_missing(&mut req.account_id, access_token); + access_token.assert_has_access(req.account_id, Collection::Calendar)?; + + self.calendar_query(req, access_token).await?.into() + } QueryRequestMethod::CalendarEvent(mut req) => { set_account_id_if_missing(&mut req.account_id, access_token); access_token.assert_has_access(req.account_id, Collection::CalendarEvent)?; diff --git a/crates/jmap/src/calendar_event/query.rs b/crates/jmap/src/calendar_event/query.rs index 9bbc2a7d..6893f71a 100644 --- a/crates/jmap/src/calendar_event/query.rs +++ b/crates/jmap/src/calendar_event/query.rs @@ -11,8 +11,12 @@ use common::{Server, auth::AccessToken}; use groupware::{cache::GroupwareCache, calendar::CalendarEvent}; use jmap_proto::{ method::query::{Filter, QueryRequest, QueryResponse}, - object::calendar_event::{self, CalendarEventComparator, CalendarEventFilter}, + object::{ + calendar, + calendar_event::{self, CalendarEventComparator, CalendarEventFilter}, + }, request::MaybeInvalid, + types::state::State, }; use nlp::language::Language; use std::{cmp::Ordering, sync::Arc}; @@ -35,6 +39,12 @@ pub trait CalendarEventQuery: Sync + Send { request: QueryRequest, access_token: &AccessToken, ) -> impl Future> + Send; + + fn calendar_query( + &self, + request: QueryRequest, + access_token: &AccessToken, + ) -> impl Future> + Send; } impl CalendarEventQuery for Server { @@ -363,6 +373,38 @@ impl CalendarEventQuery for Server { response.build() } } + + async fn calendar_query( + &self, + request: QueryRequest, + access_token: &AccessToken, + ) -> trc::Result { + let account_id = request.account_id.document_id(); + let cache = self + .fetch_dav_resources( + access_token.account_id(), + account_id, + SyncCollection::Calendar, + ) + .await?; + + let results = cache.document_ids(true).collect::>(); + + let mut response = QueryResponseBuilder::new( + results.len() as usize, + self.core.jmap.query_max_results, + State::Initial, + &request, + ); + + for document_id in results { + if !response.add(0, document_id) { + break; + } + } + + response.build() + } } fn local_timestamp(dt: &JSCalendarDateTime, tz: Tz) -> Option { diff --git a/crates/jmap/src/contact/query.rs b/crates/jmap/src/contact/query.rs index 19cf5939..acf48b31 100644 --- a/crates/jmap/src/contact/query.rs +++ b/crates/jmap/src/contact/query.rs @@ -9,8 +9,12 @@ use common::{Server, auth::AccessToken}; use groupware::cache::GroupwareCache; use jmap_proto::{ method::query::{Filter, QueryRequest, QueryResponse}, - object::contact::{ContactCard, ContactCardComparator, ContactCardFilter}, + object::{ + addressbook::AddressBook, + contact::{ContactCard, ContactCardComparator, ContactCardFilter}, + }, request::MaybeInvalid, + types::state::State, }; use store::{ IterateParams, U32_LEN, U64_LEN, ValueKey, @@ -32,6 +36,12 @@ pub trait ContactCardQuery: Sync + Send { request: QueryRequest, access_token: &AccessToken, ) -> impl Future> + Send; + + fn address_book_query( + &self, + request: QueryRequest, + access_token: &AccessToken, + ) -> impl Future> + Send; } #[derive(Clone)] @@ -326,4 +336,36 @@ impl ContactCardQuery for Server { response.build() } + + async fn address_book_query( + &self, + request: QueryRequest, + access_token: &AccessToken, + ) -> trc::Result { + let account_id = request.account_id.document_id(); + let cache = self + .fetch_dav_resources( + access_token.account_id(), + account_id, + SyncCollection::Calendar, + ) + .await?; + + let results = cache.document_ids(true).collect::>(); + + let mut response = QueryResponseBuilder::new( + results.len() as usize, + self.core.jmap.query_max_results, + State::Initial, + &request, + ); + + for document_id in results { + if !response.add(0, document_id) { + break; + } + } + + response.build() + } } diff --git a/crates/jmap/src/registry/mapping/account.rs b/crates/jmap/src/registry/mapping/account.rs index e62af935..90973e78 100644 --- a/crates/jmap/src/registry/mapping/account.rs +++ b/crates/jmap/src/registry/mapping/account.rs @@ -747,10 +747,10 @@ pub(crate) async fn account_get( for credential in account.credentials { match (credential, get.object_type) { - ( - Credential::AppPassword(pass) | Credential::ApiKey(pass), - ObjectType::AppPassword, - ) if ids.contains(&pass.credential_id) => { + (Credential::AppPassword(pass), ObjectType::AppPassword) + | (Credential::ApiKey(pass), ObjectType::ApiKey) + if ids.contains(&pass.credential_id) => + { let id = pass.credential_id; let mut credential = pass.into_value(); credential diff --git a/crates/jmap/src/registry/mapping/telemetry.rs b/crates/jmap/src/registry/mapping/telemetry.rs index 0b24dc72..9753e6ac 100644 --- a/crates/jmap/src/registry/mapping/telemetry.rs +++ b/crates/jmap/src/registry/mapping/telemetry.rs @@ -27,7 +27,8 @@ use registry::{ }; use std::str::FromStr; use store::{ - IterateParams, ValueKey, + Deserialize, IterateParams, ValueKey, + ahash::AHashSet, registry::RegistryFilterOp, search::{ SearchComparator, SearchField, SearchFilter, SearchOperator, SearchQuery, @@ -35,7 +36,7 @@ use store::{ }, write::{SearchIndex, TelemetryClass, ValueClass, key::DeserializeBigEndian, now}, }; -use trc::{AddContext, EventType}; +use trc::{AddContext, EventType, MetricType}; use types::id::Id; use utils::snowflake::SnowflakeIdGenerator; @@ -283,6 +284,7 @@ pub(crate) async fn metric_query( ) -> trc::Result { let mut ts_from = 0u64; let mut ts_to = u64::MAX; + let mut metric_type = None; req.request .extract_filters(|property, op, value| match property { @@ -307,6 +309,22 @@ pub(crate) async fn metric_query( false } } + Property::Metric => { + if let Some(mt) = value + .as_array() + .map(|v| { + v.iter() + .filter_map(|s| s.as_str().and_then(MetricType::parse)) + .collect::>() + }) + .filter(|v| !v.is_empty()) + { + metric_type = Some(mt); + true + } else { + false + } + } _ => false, })?; @@ -317,6 +335,12 @@ pub(crate) async fn metric_query( if ts_from != 0 { ts_from = SnowflakeIdGenerator::from_timestamp(ts_from).unwrap_or(0); } + if let Some(anchor) = req.request.anchor { + let anchor = anchor.id(); + if anchor > ts_from { + ts_from = anchor; + } + } if ts_to != u64::MAX { ts_to = SnowflakeIdGenerator::from_timestamp(ts_to).unwrap_or(u64::MAX); @@ -327,26 +351,36 @@ pub(crate) async fn metric_query( // Build response let mut response = QueryResponseBuilder::new( - req.server.core.jmap.query_max_results, + req.server.core.jmap.query_max_results + 1, req.server.core.jmap.query_max_results, State::Initial, &req.request, ); - if response.response.total.is_some() { - response.response.total = Some(0); - } + let mut total = 0; req.server .metrics_store() .iterate( IterateParams::new(from_key, to_key) .set_ascending(params.sort_ascending) - .no_values(), - |key, _| { + .set_values(metric_type.is_some()), + |key, value| { let id = key.deserialize_be_u64(0)?; - if let Some(total) = response.response.total.as_mut() { - *total += 1; + + if let Some(ref types) = metric_type { + let mt = match Metric::deserialize(value)? { + Metric::Counter(metric_count) => metric_count.metric, + Metric::Gauge(metric_count) => metric_count.metric, + Metric::Histogram(metric_sum) => metric_sum.metric, + }; + if !types.contains(&mt) { + return Ok(true); + } + } + + total += 1; + if response.response.total.is_some() { if !response.is_full() { response.add_id(id.into()); } @@ -359,7 +393,11 @@ pub(crate) async fn metric_query( .await .caused_by(trc::location!())?; - if let (Some(total), Some(limit)) = (response.response.total, response.response.limit) + if response.response.total.is_none() { + response.response.total = Some(total); + } + + if let Some(limit) = response.response.limit && total < limit { response.response.limit = None; diff --git a/crates/jmap/src/vacation/get.rs b/crates/jmap/src/vacation/get.rs index 9f25526c..45140103 100644 --- a/crates/jmap/src/vacation/get.rs +++ b/crates/jmap/src/vacation/get.rs @@ -83,8 +83,9 @@ impl VacationResponseGet for Server { true }; if do_get { - if let Some(document_id) = self.get_vacation_sieve_script_id(account_id).await? { - if let Some(sieve_) = self + let mut result = Map::with_capacity(properties.len()); + if let Some(document_id) = self.get_vacation_sieve_script_id(account_id).await? + && let Some(sieve_) = self .store() .get_value::>(ValueKey::archive( account_id, @@ -92,78 +93,83 @@ impl VacationResponseGet for Server { document_id, )) .await? - { - let active_script_id = self.sieve_script_get_active_id(account_id).await?; - let sieve = sieve_ - .unarchive::() - .caused_by(trc::location!())?; - let vacation = sieve.vacation_response.as_ref(); - let mut result = Map::with_capacity(properties.len()); - for property in &properties { - match property { - VacationResponseProperty::Id => { - result.insert_unchecked( - VacationResponseProperty::Id, - Id::singleton(), - ); - } - VacationResponseProperty::IsEnabled => { - result.insert_unchecked( - VacationResponseProperty::IsEnabled, - active_script_id == Some(document_id), - ); - } - VacationResponseProperty::FromDate => { - result.insert_unchecked( - VacationResponseProperty::FromDate, - vacation.and_then(|r| { - r.from_date - .as_ref() - .map(u64::from) - .map(UTCDate::from) - .map(|v| Value::Element(VacationResponseValue::Date(v))) - }), - ); - } - VacationResponseProperty::ToDate => { - result.insert_unchecked( - VacationResponseProperty::ToDate, - vacation.and_then(|r| { - r.to_date - .as_ref() - .map(u64::from) - .map(UTCDate::from) - .map(|v| Value::Element(VacationResponseValue::Date(v))) - }), - ); - } - VacationResponseProperty::Subject => { - result.insert_unchecked( - VacationResponseProperty::Subject, - vacation.and_then(|r| r.subject.as_ref()), - ); - } - VacationResponseProperty::TextBody => { - result.insert_unchecked( - VacationResponseProperty::TextBody, - vacation.and_then(|r| r.text_body.as_ref()), - ); - } - VacationResponseProperty::HtmlBody => { - result.insert_unchecked( - VacationResponseProperty::HtmlBody, - vacation.and_then(|r| r.html_body.as_ref()), - ); - } + { + let active_script_id = self.sieve_script_get_active_id(account_id).await?; + let sieve = sieve_ + .unarchive::() + .caused_by(trc::location!())?; + let vacation = sieve.vacation_response.as_ref(); + for property in &properties { + match property { + VacationResponseProperty::Id => { + result.insert_unchecked(VacationResponseProperty::Id, Id::singleton()); + } + VacationResponseProperty::IsEnabled => { + result.insert_unchecked( + VacationResponseProperty::IsEnabled, + active_script_id == Some(document_id), + ); + } + VacationResponseProperty::FromDate => { + result.insert_unchecked( + VacationResponseProperty::FromDate, + vacation.and_then(|r| { + r.from_date + .as_ref() + .map(u64::from) + .map(UTCDate::from) + .map(|v| Value::Element(VacationResponseValue::Date(v))) + }), + ); + } + VacationResponseProperty::ToDate => { + result.insert_unchecked( + VacationResponseProperty::ToDate, + vacation.and_then(|r| { + r.to_date + .as_ref() + .map(u64::from) + .map(UTCDate::from) + .map(|v| Value::Element(VacationResponseValue::Date(v))) + }), + ); + } + VacationResponseProperty::Subject => { + result.insert_unchecked( + VacationResponseProperty::Subject, + vacation.and_then(|r| r.subject.as_ref()), + ); + } + VacationResponseProperty::TextBody => { + result.insert_unchecked( + VacationResponseProperty::TextBody, + vacation.and_then(|r| r.text_body.as_ref()), + ); + } + VacationResponseProperty::HtmlBody => { + result.insert_unchecked( + VacationResponseProperty::HtmlBody, + vacation.and_then(|r| r.html_body.as_ref()), + ); } } - response.list.push(result.into()); - } else { - response.not_found.push(Id::singleton()); } } else { - response.not_found.push(Id::singleton()); + for property in &properties { + match property { + VacationResponseProperty::Id => { + result.insert_unchecked(VacationResponseProperty::Id, Id::singleton()); + } + VacationResponseProperty::IsEnabled => { + result.insert_unchecked(VacationResponseProperty::IsEnabled, false); + } + _ => { + result.insert_unchecked(property.clone(), Value::Null); + } + } + } } + response.list.push(result.into()); } Ok(response) diff --git a/crates/main/Cargo.toml b/crates/main/Cargo.toml index b0142524..b7c0ae84 100644 --- a/crates/main/Cargo.toml +++ b/crates/main/Cargo.toml @@ -20,6 +20,7 @@ coordinator = { path = "../coordinator" } jmap = { path = "../jmap" } types = { path = "../types" } smtp = { path = "../smtp" } +smtp-proto = { version = "0.2", features = ["rkyv", "serde"] } imap = { path = "../imap" } pop3 = { path = "../pop3" } spam-filter = { path = "../spam-filter" } @@ -33,6 +34,8 @@ groupware = { path = "../groupware" } services = { path = "../services" } trc = { path = "../trc" } utils = { path = "../utils" } +registry = { path = "../registry" } +http_proto = { path = "../http-proto" } migration = { path = "../migration" } tokio = { version = "1.47", features = ["full"] } @@ -68,3 +71,8 @@ enterprise = [ "jmap/enterprise", "trc/enterprise", "services/enterprise", "migration/enterprise" ] +dev_mode = [ "common/dev_mode", + "trc/dev_mode", + "dav/dev_mode", + "http/dev_mode", + "http_proto/dev_mode" ] diff --git a/crates/main/src/main.rs b/crates/main/src/main.rs index 83b41e73..c9cce9c1 100644 --- a/crates/main/src/main.rs +++ b/crates/main/src/main.rs @@ -20,6 +20,9 @@ use std::time::Duration; use trc::Collector; use utils::wait_for_shutdown; +#[cfg(feature = "dev_mode")] +pub mod test_data; + #[cfg(not(target_env = "msvc"))] use jemallocator::Jemalloc; @@ -59,6 +62,13 @@ async fn main() -> std::io::Result<()> { init.inner.build_server().log_license_details(); // SPDX-SnippetEnd + #[cfg(feature = "dev_mode")] + if std::env::var("INSERT_TEST_DATA").is_ok() { + let server = init.inner.build_server(); + //test_data::insert_test_data(&server).await; + server.insert_test_metrics().await; + } + // Spawn servers let (shutdown_tx, shutdown_rx) = init.servers.spawn(|server, acceptor, shutdown_rx| { match &server.protocol { diff --git a/crates/main/src/test_data.rs b/crates/main/src/test_data.rs new file mode 100644 index 00000000..568cc4ef --- /dev/null +++ b/crates/main/src/test_data.rs @@ -0,0 +1,1403 @@ +/* + * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * + * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + */ + +use common::{ + Server, + config::smtp::queue::{QueueExpiry, QueueName}, +}; +use registry::{ + pickle::Pickle, + schema::{ + enums::{ + ArfAuthFailureType, ArfDeliveryResult, ArfFeedbackType, ArfIdentityAlignment, + DkimAuthResult, DmarcAlignment, DmarcDisposition, DmarcResult, SpfAuthResult, + SpfDomainScope, TlsPolicyType, TlsResultType, + }, + prelude::ObjectType, + structs::{ + ArfExternalReport, ArfFeedbackReport, DmarcDkimResult, DmarcExternalReport, + DmarcInternalReport, DmarcReport, DmarcReportRecord, DmarcSpfResult, TlsExternalReport, + TlsFailureDetails, TlsInternalReport, TlsReport, TlsReportPolicy, + }, + }, + types::{EnumImpl, datetime::UTCDateTime, float::Float, ipaddr::IpAddr, list::List, map::Map}, +}; +use smtp::{ + queue::{ + Error, ErrorDetails, FROM_AUTHENTICATED, FROM_DSN, FROM_REPORT, HostResponse, Message, + MessageWrapper, RCPT_DSN_SENT, RCPT_SPAM_PAYLOAD, Recipient, Schedule, Status, + UnexpectedResponse, + }, + reporting::index::{ExternalReportIndex, InternalReportIndex}, +}; +use smtp_proto::Response; +use std::net::Ipv4Addr; +use store::write::{BatchBuilder, RegistryClass, ValueClass}; +use types::blob_hash::BlobHash; + +pub async fn insert_test_data(server: &Server) { + let mut hashes = Vec::new(); + let future = one_year_from_now_u64(); + for message in sample_raw_messages() { + let (hash, _) = server + .put_temporary_blob(u32::MAX, message.as_bytes(), future) + .await + .unwrap(); + hashes.push(hash); + } + + for message in sample_queued_messages(hashes) { + let qm = MessageWrapper::new( + message, + server.inner.data.queue_id_gen.generate(), + QueueName::default(), + ); + assert!(qm.save_changes(server, None).await); + } + + for report in sample_tls_internal_reports() { + let object_id = ObjectType::TlsInternalReport.to_id(); + let item_id = server.inner.data.queue_id_gen.generate(); + let mut batch = BatchBuilder::new(); + report.write_ops(&mut batch, item_id, true); + let report_bytes = report.to_pickled_vec(); + batch.set( + ValueClass::Registry(RegistryClass::Item { object_id, item_id }), + report_bytes, + ); + server.store().write(batch.build_all()).await.unwrap(); + } + + for report in sample_dmarc_internal_reports() { + let object_id = ObjectType::DmarcInternalReport.to_id(); + let item_id = server.inner.data.queue_id_gen.generate(); + let mut batch = BatchBuilder::new(); + report.write_ops(&mut batch, item_id, true); + let report_bytes = report.to_pickled_vec(); + batch.set( + ValueClass::Registry(RegistryClass::Item { object_id, item_id }), + report_bytes, + ); + server.store().write(batch.build_all()).await.unwrap(); + } + + for report in sample_tls_external_reports() { + let mut batch = BatchBuilder::new(); + let item_id = server.inner.data.queue_id_gen.generate(); + report.write_ops(&mut batch, item_id, true); + server.store().write(batch.build_all()).await.unwrap(); + } + + for report in sample_dmarc_external_reports() { + let mut batch = BatchBuilder::new(); + let item_id = server.inner.data.queue_id_gen.generate(); + report.write_ops(&mut batch, item_id, true); + server.store().write(batch.build_all()).await.unwrap(); + } + + for report in sample_arf_external_reports() { + let mut batch = BatchBuilder::new(); + let item_id = server.inner.data.queue_id_gen.generate(); + report.write_ops(&mut batch, item_id, true); + server.store().write(batch.build_all()).await.unwrap(); + } +} + +fn sample_queued_messages(blob_hashes: Vec) -> Vec { + assert!( + blob_hashes.len() >= 3, + "Need at least 3 blob hashes for sample queued messages" + ); + let future = one_year_from_now_u64(); + let now = store::write::now(); + let raw_messages = sample_raw_messages(); + + vec![ + // Message 1: Normal outbound message with two recipients, one scheduled and one completed + Message { + created: now, + blob_hash: blob_hashes[0].clone(), + return_path: "sender@myserver.com".into(), + recipients: vec![ + Recipient { + address: "alice@example.com".into(), + retry: Schedule { + due: future, + inner: 0, + }, + notify: Schedule { + due: future, + inner: 0, + }, + expires: QueueExpiry::Ttl(365 * 24 * 3600), + queue: Default::default(), + status: Status::Scheduled, + flags: 0, + orcpt: None, + }, + Recipient { + address: "bob@example.org".into(), + retry: Schedule { + due: future, + inner: 2, + }, + notify: Schedule { + due: future, + inner: 1, + }, + expires: QueueExpiry::Ttl(365 * 24 * 3600), + queue: Default::default(), + status: Status::Completed(HostResponse { + hostname: "mx.example.org".into(), + response: Response { + code: 250, + esc: [2, 1, 5], + message: "OK".into(), + }, + }), + flags: RCPT_DSN_SENT, + orcpt: Some("rfc822;bob@example.org".into()), + }, + ], + received_from_ip: std::net::IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10)), + received_via_port: 25, + flags: FROM_AUTHENTICATED, + env_id: Some("env-001".into()), + priority: 0, + size: raw_messages[0].len() as u64, + quota_keys: Box::new([]), + }, + // Message 2: DSN bounce message with a temporary failure recipient + Message { + created: now, + blob_hash: blob_hashes[1].clone(), + return_path: "".into(), + recipients: vec![Recipient { + address: "postmaster@remote.net".into(), + retry: Schedule { + due: future, + inner: 3, + }, + notify: Schedule { + due: future, + inner: 0, + }, + expires: QueueExpiry::Ttl(365 * 24 * 3600), + queue: Default::default(), + status: Status::TemporaryFailure(ErrorDetails { + entity: "mx.remote.net".into(), + details: Error::UnexpectedResponse(UnexpectedResponse { + command: "RCPT TO".into(), + response: Response { + code: 450, + esc: [4, 2, 1], + message: "Mailbox temporarily unavailable".into(), + }, + }), + }), + flags: 0, + orcpt: None, + }], + received_from_ip: std::net::IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)), + received_via_port: 587, + flags: FROM_DSN, + env_id: None, + priority: -5, + size: raw_messages[1].len() as u64, + quota_keys: Box::new([]), + }, + // Message 3: Report message with a permanent failure recipient + Message { + created: now, + blob_hash: blob_hashes[2].clone(), + return_path: "reports@myserver.com".into(), + recipients: vec![Recipient { + address: "abuse@bigcorp.com".into(), + retry: Schedule { + due: future, + inner: 0, + }, + notify: Schedule { + due: future, + inner: 2, + }, + expires: QueueExpiry::Ttl(365 * 24 * 3600), + queue: Default::default(), + status: Status::PermanentFailure(ErrorDetails { + entity: "mx.bigcorp.com".into(), + details: Error::ConnectionError("Rejected by policy".into()), + }), + flags: RCPT_SPAM_PAYLOAD, + orcpt: None, + }], + received_from_ip: std::net::IpAddr::V4(Ipv4Addr::new(172, 16, 0, 5)), + received_via_port: 465, + flags: FROM_REPORT | FROM_AUTHENTICATED, + env_id: Some("env-report-99".into()), + priority: 10, + size: raw_messages[2].len() as u64, + quota_keys: Box::new([]), + }, + ] +} + +fn sample_raw_messages() -> Vec { + vec![ + // Raw message 1: Normal outbound email (matches Message 1) + concat!( + "From: sender@myserver.com\r\n", + "To: alice@example.com, bob@example.org\r\n", + "Subject: Quarterly Report Q1 2027\r\n", + "Date: Sat, 12 Apr 2027 09:00:00 +0000\r\n", + "Message-ID: \r\n", + "MIME-Version: 1.0\r\n", + "Content-Type: text/plain; charset=utf-8\r\n", + "\r\n", + "Hi team,\r\n", + "\r\n", + "Please find attached the quarterly report for Q1 2027.\r\n", + "Let me know if you have any questions.\r\n", + "\r\n", + "Best regards,\r\n", + "The Sender\r\n", + ) + .to_string(), + // Raw message 2: DSN bounce (matches Message 2) + concat!( + "From: <>\r\n", + "To: postmaster@remote.net\r\n", + "Subject: Delivery Status Notification (Failure)\r\n", + "Date: Sat, 12 Apr 2027 09:05:00 +0000\r\n", + "Message-ID: \r\n", + "MIME-Version: 1.0\r\n", + "Content-Type: multipart/report; report-type=delivery-status;\r\n", + " boundary=\"boundary-dsn-002\"\r\n", + "\r\n", + "--boundary-dsn-002\r\n", + "Content-Type: text/plain; charset=utf-8\r\n", + "\r\n", + "This is an automatically generated Delivery Status Notification.\r\n", + "Delivery to the following recipient failed temporarily:\r\n", + "\r\n", + " postmaster@remote.net\r\n", + "\r\n", + "The server will retry delivery.\r\n", + "\r\n", + "--boundary-dsn-002\r\n", + "Content-Type: message/delivery-status\r\n", + "\r\n", + "Reporting-MTA: dns; myserver.com\r\n", + "Arrival-Date: Sat, 12 Apr 2027 09:00:00 +0000\r\n", + "\r\n", + "Final-Recipient: rfc822; postmaster@remote.net\r\n", + "Action: delayed\r\n", + "Status: 4.2.1\r\n", + "Remote-MTA: dns; mx.remote.net\r\n", + "Diagnostic-Code: smtp; 450 Mailbox temporarily unavailable\r\n", + "\r\n", + "--boundary-dsn-002--\r\n", + ) + .to_string(), + // Raw message 3: Report message (matches Message 3) + concat!( + "From: reports@myserver.com\r\n", + "To: abuse@bigcorp.com\r\n", + "Subject: DMARC Aggregate Report for bigcorp.com\r\n", + "Date: Sat, 12 Apr 2027 09:10:00 +0000\r\n", + "Message-ID: \r\n", + "MIME-Version: 1.0\r\n", + "Content-Type: multipart/mixed;\r\n", + " boundary=\"boundary-report-003\"\r\n", + "\r\n", + "--boundary-report-003\r\n", + "Content-Type: text/plain; charset=utf-8\r\n", + "\r\n", + "This is a DMARC aggregate report for the domain bigcorp.com\r\n", + "generated by myserver.com.\r\n", + "\r\n", + "Report period: 2027-04-11T00:00:00Z to 2027-04-12T00:00:00Z\r\n", + "\r\n", + "--boundary-report-003\r\n", + "Content-Type: application/gzip\r\n", + "Content-Disposition: attachment;\r\n", + " filename=\"myserver.com!bigcorp.com!1744329600!1744416000.xml.gz\"\r\n", + "Content-Transfer-Encoding: base64\r\n", + "\r\n", + "H4sIAAAAAAAAA2NgGAWjYBSMglEwCkbBKBgFo2AUDAIAAP//\r\n", + "\r\n", + "--boundary-report-003--\r\n", + ) + .to_string(), + ] +} + +fn sample_tls_internal_reports() -> Vec { + let future = one_year_from_now(); + + vec![ + // Report 1: Successful TLS sessions, STS policy + TlsInternalReport { + created_at: now(), + deliver_at: future, + domain: "example.com".to_string(), + http_rua: Map::new(vec!["https://example.com/tlsrpt".to_string()]), + mail_rua: Map::new(vec!["mailto:tls-reports@example.com".to_string()]), + policy_identifiers: Map::new(vec![1]), + report: TlsReport { + contact_info: Some("admin@myserver.com".to_string()), + date_range_end: now(), + date_range_start: days_ago(1), + organization_name: Some("My Mail Server".to_string()), + policies: List::from(vec![TlsReportPolicy { + failure_details: List::from(vec![]), + mx_hosts: Map::new(vec!["mx1.example.com".to_string()]), + policy_domain: "example.com".to_string(), + policy_strings: Map::new(vec![ + "mode: enforce".to_string(), + "max_age: 86400".to_string(), + ]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 0, + total_successful_sessions: 150, + }]), + report_id: "tls-int-report-001".to_string(), + }, + }, + // Report 2: Mixed results with certificate mismatch failures + TlsInternalReport { + created_at: now(), + deliver_at: future, + domain: "secure-mail.org".to_string(), + http_rua: Map::new(vec![]), + mail_rua: Map::new(vec![ + "mailto:tlsrpt@secure-mail.org".to_string(), + "mailto:security@secure-mail.org".to_string(), + ]), + policy_identifiers: Map::new(vec![2, 3]), + report: TlsReport { + contact_info: Some("postmaster@myserver.com".to_string()), + date_range_end: now(), + date_range_start: days_ago(1), + organization_name: Some("My Mail Server".to_string()), + policies: List::from(vec![TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "Certificate CN does not match hostname".to_string(), + ), + failed_session_count: 5, + failure_reason_code: Some("certificate-host-mismatch".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 10, + )))), + receiving_mx_helo: Some("mx.secure-mail.org".to_string()), + receiving_mx_hostname: Some("mx.secure-mail.org".to_string()), + result_type: TlsResultType::CertificateHostMismatch, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 192, 0, 2, 1, + )))), + }]), + mx_hosts: Map::new(vec!["mx.secure-mail.org".to_string()]), + policy_domain: "secure-mail.org".to_string(), + policy_strings: Map::new(vec!["mode: testing".to_string()]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 5, + total_successful_sessions: 95, + }]), + report_id: "tls-int-report-002".to_string(), + }, + }, + // Report 3: DANE/TLSA policy with validation failure + TlsInternalReport { + created_at: now(), + deliver_at: future, + domain: "dane-enabled.net".to_string(), + http_rua: Map::new(vec!["https://dane-enabled.net/tlsrpt".to_string()]), + mail_rua: Map::new(vec![]), + policy_identifiers: Map::new(vec![4]), + report: TlsReport { + contact_info: None, + date_range_end: now(), + date_range_start: days_ago(1), + organization_name: Some("My Mail Server".to_string()), + policies: List::from(vec![TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: None, + failed_session_count: 2, + failure_reason_code: Some("tlsa-invalid".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 25, + )))), + receiving_mx_helo: Some("mail.dane-enabled.net".to_string()), + receiving_mx_hostname: Some("mail.dane-enabled.net".to_string()), + result_type: TlsResultType::TlsaInvalid, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 192, 0, 2, 1, + )))), + }]), + mx_hosts: Map::new(vec!["mail.dane-enabled.net".to_string()]), + policy_domain: "dane-enabled.net".to_string(), + policy_strings: Map::new(vec![]), + policy_type: TlsPolicyType::Tlsa, + total_failed_sessions: 2, + total_successful_sessions: 48, + }]), + report_id: "tls-int-report-003".to_string(), + }, + }, + ] +} + +fn sample_dmarc_internal_reports() -> Vec { + let future = one_year_from_now(); + + vec![ + // Report 1: Clean domain with all-pass records + DmarcInternalReport { + created_at: now(), + deliver_at: future, + domain: "trusted-sender.com".to_string(), + policy_identifier: 100, + report: DmarcReport { + date_range_begin: days_ago(1), + date_range_end: now(), + email: "dmarc@myserver.com".to_string(), + errors: Map::new(vec![]), + extensions: List::from(vec![]), + extra_contact_info: Some("https://myserver.com/dmarc".to_string()), + org_name: "My Mail Server".to_string(), + policy_adkim: DmarcAlignment::Relaxed, + policy_aspf: DmarcAlignment::Relaxed, + policy_disposition: DmarcDisposition::None, + policy_domain: "trusted-sender.com".to_string(), + policy_failure_reporting_options: Map::new(vec![]), + policy_subdomain_disposition: DmarcDisposition::Quarantine, + policy_testing_mode: false, + policy_version: Some("DMARC1".to_string()), + records: List::from(vec![DmarcReportRecord { + count: 500, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "trusted-sender.com".to_string(), + human_result: None, + result: DkimAuthResult::Pass, + selector: "selector1".to_string(), + }]), + envelope_from: "trusted-sender.com".to_string(), + envelope_to: Some("myserver.com".to_string()), + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Pass, + extensions: List::from(vec![]), + header_from: "trusted-sender.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 93, 184, 216, 34, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "trusted-sender.com".to_string(), + human_result: None, + result: SpfAuthResult::Pass, + scope: SpfDomainScope::MailFrom, + }]), + }]), + report_id: "dmarc-int-001".to_string(), + version: Float::from(1.0), + }, + rua: Map::new(vec!["mailto:dmarc-rua@trusted-sender.com".to_string()]), + }, + // Report 2: Domain with DKIM failure and strict alignment + DmarcInternalReport { + created_at: now(), + deliver_at: future, + domain: "strict-domain.org".to_string(), + policy_identifier: 200, + report: DmarcReport { + date_range_begin: days_ago(1), + date_range_end: now(), + email: "dmarc@myserver.com".to_string(), + errors: Map::new(vec!["DKIM signature verification failed".to_string()]), + extensions: List::from(vec![]), + extra_contact_info: None, + org_name: "My Mail Server".to_string(), + policy_adkim: DmarcAlignment::Strict, + policy_aspf: DmarcAlignment::Strict, + policy_disposition: DmarcDisposition::Reject, + policy_domain: "strict-domain.org".to_string(), + policy_failure_reporting_options: Map::new(vec![]), + policy_subdomain_disposition: DmarcDisposition::Reject, + policy_testing_mode: false, + policy_version: Some("DMARC1".to_string()), + records: List::from(vec![DmarcReportRecord { + count: 12, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "strict-domain.org".to_string(), + human_result: Some("signature verification failed".to_string()), + result: DkimAuthResult::Fail, + selector: "dkim2024".to_string(), + }]), + envelope_from: "strict-domain.org".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Fail, + evaluated_spf: DmarcResult::Pass, + extensions: List::from(vec![]), + header_from: "strict-domain.org".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 50, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "strict-domain.org".to_string(), + human_result: None, + result: SpfAuthResult::Pass, + scope: SpfDomainScope::MailFrom, + }]), + }]), + report_id: "dmarc-int-002".to_string(), + version: Float::from(1.0), + }, + rua: Map::new(vec!["mailto:dmarc@strict-domain.org".to_string()]), + }, + // Report 3: Domain in testing mode with multiple record types + DmarcInternalReport { + created_at: now(), + deliver_at: future, + domain: "new-policy.io".to_string(), + policy_identifier: 300, + report: DmarcReport { + date_range_begin: days_ago(1), + date_range_end: now(), + email: "dmarc@myserver.com".to_string(), + errors: Map::new(vec![]), + extensions: List::from(vec![]), + extra_contact_info: None, + org_name: "My Mail Server".to_string(), + policy_adkim: DmarcAlignment::Relaxed, + policy_aspf: DmarcAlignment::Strict, + policy_disposition: DmarcDisposition::Quarantine, + policy_domain: "new-policy.io".to_string(), + policy_failure_reporting_options: Map::new(vec![]), + policy_subdomain_disposition: DmarcDisposition::None, + policy_testing_mode: true, + policy_version: Some("DMARC1".to_string()), + records: List::from(vec![ + DmarcReportRecord { + count: 200, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "new-policy.io".to_string(), + human_result: None, + result: DkimAuthResult::Pass, + selector: "sel1".to_string(), + }]), + envelope_from: "new-policy.io".to_string(), + envelope_to: Some("myserver.com".to_string()), + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Pass, + extensions: List::from(vec![]), + header_from: "new-policy.io".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 5, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "new-policy.io".to_string(), + human_result: None, + result: SpfAuthResult::Pass, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 3, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "new-policy.io".to_string(), + human_result: Some("no signature found".to_string()), + result: DkimAuthResult::None, + selector: "".to_string(), + }]), + envelope_from: "spoofed.example".to_string(), + envelope_to: Some("myserver.com".to_string()), + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Fail, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "new-policy.io".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 99)))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "spoofed.example".to_string(), + human_result: Some("SPF record not found".to_string()), + result: SpfAuthResult::None, + scope: SpfDomainScope::MailFrom, + }]), + }, + ]), + report_id: "dmarc-int-003".to_string(), + version: Float::from(1.0), + }, + rua: Map::new(vec![ + "mailto:dmarc@new-policy.io".to_string(), + "mailto:dmarc-backup@new-policy.io".to_string(), + ]), + }, + ] +} + +fn sample_tls_external_reports() -> Vec { + let future = one_year_from_now(); + + vec![ + // Report 1: Clean report from a large provider + TlsExternalReport { + expires_at: future, + from: "tls-reports@bigprovider.com".to_string(), + member_tenant_id: None, + received_at: now(), + report: TlsReport { + contact_info: Some("postmaster@bigprovider.com".to_string()), + date_range_end: now(), + date_range_start: days_ago(1), + organization_name: Some("Big Provider Inc.".to_string()), + policies: List::from(vec![ + TlsReportPolicy { + failure_details: List::from(vec![]), + mx_hosts: Map::new(vec![ + "mx1.myserver.com".to_string(), + "mx2.myserver.com".to_string(), + ]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec![ + "mode: enforce".to_string(), + "max_age: 604800".to_string(), + ]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 0, + total_successful_sessions: 12500, + }, + TlsReportPolicy { + failure_details: List::from(vec![]), + mx_hosts: Map::new(vec!["mx3.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec![]), + policy_type: TlsPolicyType::Tlsa, + total_failed_sessions: 0, + total_successful_sessions: 3200, + }, + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "Fallback to plaintext after STARTTLS failure".to_string(), + ), + failed_session_count: 2, + failure_reason_code: Some("starttls-not-supported".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 192, 0, 2, 50, + )))), + receiving_mx_helo: Some("backup-mx.myserver.com".to_string()), + receiving_mx_hostname: Some("backup-mx.myserver.com".to_string()), + result_type: TlsResultType::StartTlsNotSupported, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 5, + )))), + }]), + mx_hosts: Map::new(vec!["backup-mx.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec!["mode: testing".to_string()]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 2, + total_successful_sessions: 50, + }, + ]), + report_id: "tls-ext-001-bigprovider".to_string(), + }, + subject: "TLS-RPT report for myserver.com".to_string(), + to: Map::new(vec!["tls-rpt@myserver.com".to_string()]), + }, + // Report 2: Report with expired certificate failures + TlsExternalReport { + expires_at: future, + from: "noreply@securemail.org".to_string(), + member_tenant_id: None, + received_at: now(), + report: TlsReport { + contact_info: None, + date_range_end: now(), + date_range_start: days_ago(1), + organization_name: Some("SecureMail".to_string()), + policies: List::from(vec![ + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "Certificate expired 2 days ago".to_string(), + ), + failed_session_count: 30, + failure_reason_code: Some("certificate-expired".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 192, 0, 2, 10, + )))), + receiving_mx_helo: Some("mx1.myserver.com".to_string()), + receiving_mx_hostname: Some("mx1.myserver.com".to_string()), + result_type: TlsResultType::CertificateExpired, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 1, + )))), + }]), + mx_hosts: Map::new(vec!["mx1.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec!["mode: enforce".to_string()]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 30, + total_successful_sessions: 70, + }, + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "CN=old.myserver.com does not match mx2.myserver.com".to_string(), + ), + failed_session_count: 15, + failure_reason_code: Some("certificate-host-mismatch".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 192, 0, 2, 11, + )))), + receiving_mx_helo: Some("mx2.myserver.com".to_string()), + receiving_mx_hostname: Some("mx2.myserver.com".to_string()), + result_type: TlsResultType::CertificateHostMismatch, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 2, + )))), + }]), + mx_hosts: Map::new(vec!["mx2.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec!["mode: enforce".to_string()]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 15, + total_successful_sessions: 85, + }, + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "Untrusted CA in certificate chain".to_string(), + ), + failed_session_count: 8, + failure_reason_code: Some("certificate-not-trusted".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 192, 0, 2, 12, + )))), + receiving_mx_helo: Some("mx3.myserver.com".to_string()), + receiving_mx_hostname: Some("mx3.myserver.com".to_string()), + result_type: TlsResultType::CertificateNotTrusted, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 3, + )))), + }]), + mx_hosts: Map::new(vec!["mx3.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec!["mode: testing".to_string()]), + policy_type: TlsPolicyType::Sts, + total_failed_sessions: 8, + total_successful_sessions: 120, + }, + ]), + report_id: "tls-ext-002-securemail".to_string(), + }, + subject: "SMTP TLS Reporting for myserver.com".to_string(), + to: Map::new(vec!["tls-rpt@myserver.com".to_string()]), + }, + // Report 3: DANE report with no policy found + TlsExternalReport { + expires_at: future, + from: "reports@mailhoster.net".to_string(), + member_tenant_id: None, + received_at: now(), + report: TlsReport { + contact_info: Some("abuse@mailhoster.net".to_string()), + date_range_end: now(), + date_range_start: days_ago(7), + organization_name: Some("MailHoster".to_string()), + policies: List::from(vec![ + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: None, + failed_session_count: 10, + failure_reason_code: None, + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 25, + )))), + receiving_mx_helo: Some("mx2.myserver.com".to_string()), + receiving_mx_hostname: Some("mx2.myserver.com".to_string()), + result_type: TlsResultType::StartTlsNotSupported, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 50, + )))), + }]), + mx_hosts: Map::new(vec!["mx2.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec![]), + policy_type: TlsPolicyType::NoPolicyFound, + total_failed_sessions: 10, + total_successful_sessions: 0, + }, + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "DANE TLSA record invalid for mx1.myserver.com".to_string(), + ), + failed_session_count: 20, + failure_reason_code: Some("tlsa-invalid".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 26, + )))), + receiving_mx_helo: Some("mx1.myserver.com".to_string()), + receiving_mx_hostname: Some("mx1.myserver.com".to_string()), + result_type: TlsResultType::TlsaInvalid, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 51, + )))), + }]), + mx_hosts: Map::new(vec!["mx1.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec![]), + policy_type: TlsPolicyType::Tlsa, + total_failed_sessions: 20, + total_successful_sessions: 180, + }, + TlsReportPolicy { + failure_details: List::from(vec![TlsFailureDetails { + additional_information: Some( + "DNSSEC validation failed for _25._tcp.mx1.myserver.com" + .to_string(), + ), + failed_session_count: 5, + failure_reason_code: Some("dnssec-invalid".to_string()), + receiving_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 27, + )))), + receiving_mx_helo: Some("mx1.myserver.com".to_string()), + receiving_mx_hostname: Some("mx1.myserver.com".to_string()), + result_type: TlsResultType::DnssecInvalid, + sending_mta_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 52, + )))), + }]), + mx_hosts: Map::new(vec!["mx1.myserver.com".to_string()]), + policy_domain: "myserver.com".to_string(), + policy_strings: Map::new(vec![]), + policy_type: TlsPolicyType::Tlsa, + total_failed_sessions: 5, + total_successful_sessions: 95, + }, + ]), + report_id: "tls-ext-003-mailhoster".to_string(), + }, + subject: "TLS Report: myserver.com".to_string(), + to: Map::new(vec!["tls-rpt@myserver.com".to_string()]), + }, + ] +} + +fn sample_dmarc_external_reports() -> Vec { + let future = one_year_from_now(); + + vec![ + // Report 1: Google-style aggregate report, all pass + DmarcExternalReport { + expires_at: future, + from: "noreply-dmarc@google.com".to_string(), + member_tenant_id: None, + received_at: now(), + report: DmarcReport { + date_range_begin: days_ago(1), + date_range_end: now(), + email: "noreply-dmarc@google.com".to_string(), + errors: Map::new(vec![]), + extensions: List::from(vec![]), + extra_contact_info: Some( + "https://support.google.com/a/answer/10032169".to_string(), + ), + org_name: "google.com".to_string(), + policy_adkim: DmarcAlignment::Relaxed, + policy_aspf: DmarcAlignment::Relaxed, + policy_disposition: DmarcDisposition::None, + policy_domain: "myserver.com".to_string(), + policy_failure_reporting_options: Map::new(vec![]), + policy_subdomain_disposition: DmarcDisposition::None, + policy_testing_mode: false, + policy_version: Some("DMARC1".to_string()), + records: List::from(vec![ + DmarcReportRecord { + count: 1500, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: None, + result: DkimAuthResult::Pass, + selector: "google".to_string(), + }]), + envelope_from: "myserver.com".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Pass, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "myserver.com".to_string(), + human_result: None, + result: SpfAuthResult::Pass, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 320, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: None, + result: DkimAuthResult::Pass, + selector: "selector2".to_string(), + }]), + envelope_from: "myserver.com".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Pass, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 2)))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "myserver.com".to_string(), + human_result: None, + result: SpfAuthResult::Pass, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 7, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "sub.myserver.com".to_string(), + human_result: Some("DKIM signature uses subdomain".to_string()), + result: DkimAuthResult::Pass, + selector: "google".to_string(), + }]), + envelope_from: "sub.myserver.com".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 3)))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "sub.myserver.com".to_string(), + human_result: Some( + "SPF alignment failed: subdomain mismatch".to_string(), + ), + result: SpfAuthResult::SoftFail, + scope: SpfDomainScope::MailFrom, + }]), + }, + ]), + report_id: "dmarc-ext-001-google".to_string(), + version: Float::from(1.0), + }, + subject: "Report domain: myserver.com Submitter: google.com".to_string(), + to: Map::new(vec!["dmarc-rua@myserver.com".to_string()]), + }, + // Report 2: Report showing spoofing attempts from an unknown source + DmarcExternalReport { + expires_at: future, + from: "dmarc@yahoo.com".to_string(), + member_tenant_id: None, + received_at: now(), + report: DmarcReport { + date_range_begin: days_ago(1), + date_range_end: now(), + email: "dmarc@yahoo.com".to_string(), + errors: Map::new(vec![]), + extensions: List::from(vec![]), + extra_contact_info: None, + org_name: "Yahoo! Inc.".to_string(), + policy_adkim: DmarcAlignment::Strict, + policy_aspf: DmarcAlignment::Strict, + policy_disposition: DmarcDisposition::Reject, + policy_domain: "myserver.com".to_string(), + policy_failure_reporting_options: Map::new(vec![]), + policy_subdomain_disposition: DmarcDisposition::Reject, + policy_testing_mode: false, + policy_version: Some("DMARC1".to_string()), + records: List::from(vec![ + DmarcReportRecord { + count: 8, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: Some("bad signature".to_string()), + result: DkimAuthResult::Fail, + selector: "default".to_string(), + }]), + envelope_from: "attacker.example".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Fail, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 222, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "attacker.example".to_string(), + human_result: Some("domain not found".to_string()), + result: SpfAuthResult::Fail, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 15, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "phisher.net".to_string(), + human_result: Some( + "no valid DKIM signature for myserver.com".to_string(), + ), + result: DkimAuthResult::None, + selector: "".to_string(), + }]), + envelope_from: "phisher.net".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Fail, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 100, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "phisher.net".to_string(), + human_result: Some("SPF domain mismatch".to_string()), + result: SpfAuthResult::Fail, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 2, + dkim_results: List::from(vec![ + DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: Some("signature expired".to_string()), + result: DkimAuthResult::Fail, + selector: "selector1".to_string(), + }, + DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: Some("body hash mismatch".to_string()), + result: DkimAuthResult::Fail, + selector: "selector2".to_string(), + }, + ]), + envelope_from: "compromised-relay.example".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Fail, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 33, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "compromised-relay.example".to_string(), + human_result: None, + result: SpfAuthResult::PermError, + scope: SpfDomainScope::MailFrom, + }]), + }, + ]), + report_id: "dmarc-ext-002-yahoo".to_string(), + version: Float::from(1.0), + }, + subject: "Report domain: myserver.com Submitter: yahoo.com".to_string(), + to: Map::new(vec!["dmarc-rua@myserver.com".to_string()]), + }, + // Report 3: Microsoft report with forwarded mail override + DmarcExternalReport { + expires_at: future, + from: "dmarc-noreply@microsoft.com".to_string(), + member_tenant_id: None, + received_at: now(), + report: DmarcReport { + date_range_begin: days_ago(1), + date_range_end: now(), + email: "dmarc-noreply@microsoft.com".to_string(), + errors: Map::new(vec![]), + extensions: List::from(vec![]), + extra_contact_info: None, + org_name: "Microsoft Corporation".to_string(), + policy_adkim: DmarcAlignment::Relaxed, + policy_aspf: DmarcAlignment::Relaxed, + policy_disposition: DmarcDisposition::Quarantine, + policy_domain: "myserver.com".to_string(), + policy_failure_reporting_options: Map::new(vec![]), + policy_subdomain_disposition: DmarcDisposition::Quarantine, + policy_testing_mode: false, + policy_version: Some("DMARC1".to_string()), + records: List::from(vec![ + DmarcReportRecord { + count: 25, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: None, + result: DkimAuthResult::Pass, + selector: "selector1".to_string(), + }]), + envelope_from: "forwarder.example.com".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 100, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "forwarder.example.com".to_string(), + human_result: None, + result: SpfAuthResult::SoftFail, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 800, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: None, + result: DkimAuthResult::Pass, + selector: "selector1".to_string(), + }]), + envelope_from: "myserver.com".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Pass, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "myserver.com".to_string(), + human_result: None, + result: SpfAuthResult::Pass, + scope: SpfDomainScope::MailFrom, + }]), + }, + DmarcReportRecord { + count: 40, + dkim_results: List::from(vec![DmarcDkimResult { + domain: "myserver.com".to_string(), + human_result: Some( + "DKIM signature OK but SPF failed via mailing list".to_string(), + ), + result: DkimAuthResult::Pass, + selector: "selector1".to_string(), + }]), + envelope_from: "mailinglist.example.org".to_string(), + envelope_to: None, + evaluated_disposition: Default::default(), + evaluated_dkim: DmarcResult::Pass, + evaluated_spf: DmarcResult::Fail, + extensions: List::from(vec![]), + header_from: "myserver.com".to_string(), + policy_override_reasons: List::from(vec![]), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 203, 0, 113, 200, + )))), + spf_results: List::from(vec![DmarcSpfResult { + domain: "mailinglist.example.org".to_string(), + human_result: Some("Mailing list rewrite".to_string()), + result: SpfAuthResult::Fail, + scope: SpfDomainScope::MailFrom, + }]), + }, + ]), + report_id: "dmarc-ext-003-msft".to_string(), + version: Float::from(1.0), + }, + subject: "Report domain: myserver.com Submitter: microsoft.com".to_string(), + to: Map::new(vec!["dmarc-rua@myserver.com".to_string()]), + }, + ] +} + +fn sample_arf_external_reports() -> Vec { + let future = one_year_from_now(); + + vec![ + // Report 1: Abuse complaint from a user + ArfExternalReport { + expires_at: future, + from: "fbl@isp-provider.com".to_string(), + member_tenant_id: None, + received_at: now(), + report: ArfFeedbackReport { + arrival_date: Some(days_ago(1)), + auth_failure: ArfAuthFailureType::Unspecified, + authentication_results: Map::new(vec![ + "dkim=pass header.d=myserver.com".to_string(), + "spf=pass smtp.mailfrom=myserver.com".to_string(), + ]), + delivery_result: ArfDeliveryResult::Delivered, + dkim_adsp_dns: None, + dkim_canonicalized_body: None, + dkim_canonicalized_header: None, + dkim_domain: Some("myserver.com".to_string()), + dkim_identity: None, + dkim_selector: Some("selector1".to_string()), + dkim_selector_dns: None, + feedback_type: ArfFeedbackType::Abuse, + headers: Some( + "From: newsletter@myserver.com\r\nTo: user@isp-provider.com\r\nSubject: Weekly Newsletter\r\nDate: Mon, 10 Apr 2027 10:00:00 +0000" + .to_string(), + ), + identity_alignment: ArfIdentityAlignment::DkimSpf, + incidents: 1, + message: Some("User marked this message as spam".to_string()), + original_envelope_id: Some("env-newsletter-001".to_string()), + original_mail_from: Some("newsletter@myserver.com".to_string()), + original_rcpt_to: Some("user@isp-provider.com".to_string()), + reported_domains: Map::new(vec!["myserver.com".to_string()]), + reported_uris: Map::new(vec![]), + reporting_mta: Some("fbl.isp-provider.com".to_string()), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)))), + source_port: Some(25), + spf_dns: None, + user_agent: Some("ISP-FBL/2.0".to_string()), + version: 1, + }, + subject: "FBL report from isp-provider.com".to_string(), + to: Map::new(vec!["abuse@myserver.com".to_string()]), + }, + // Report 2: Auth failure report (DKIM) + ArfExternalReport { + expires_at: future, + from: "authfail@receiver.org".to_string(), + member_tenant_id: None, + received_at: now(), + report: ArfFeedbackReport { + arrival_date: Some(days_ago(2)), + auth_failure: ArfAuthFailureType::Signature, + authentication_results: Map::new(vec![ + "dkim=fail header.d=myserver.com".to_string(), + ]), + delivery_result: ArfDeliveryResult::Reject, + dkim_adsp_dns: None, + dkim_canonicalized_body: Some("base64bodyhash==".to_string()), + dkim_canonicalized_header: Some("base64headerhash==".to_string()), + dkim_domain: Some("myserver.com".to_string()), + dkim_identity: Some("@myserver.com".to_string()), + dkim_selector: Some("selector1".to_string()), + dkim_selector_dns: Some( + "v=DKIM1; k=rsa; p=MIGfMA0GCSqGSIb3DQEBAQUA".to_string(), + ), + feedback_type: ArfFeedbackType::AuthFailure, + headers: Some( + "From: info@myserver.com\r\nTo: contact@receiver.org\r\nSubject: Important Update" + .to_string(), + ), + identity_alignment: ArfIdentityAlignment::None, + incidents: 3, + message: None, + original_envelope_id: None, + original_mail_from: Some("info@myserver.com".to_string()), + original_rcpt_to: Some("contact@receiver.org".to_string()), + reported_domains: Map::new(vec!["myserver.com".to_string()]), + reported_uris: Map::new(vec![]), + reporting_mta: Some("mx.receiver.org".to_string()), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)))), + source_port: Some(587), + spf_dns: None, + user_agent: Some("ReceiverMTA/1.0".to_string()), + version: 1, + }, + subject: "Authentication failure report for myserver.com".to_string(), + to: Map::new(vec!["abuse@myserver.com".to_string()]), + }, + // Report 3: Fraud/phishing report + ArfExternalReport { + expires_at: future, + from: "reports@antiphish.net".to_string(), + member_tenant_id: None, + received_at: now(), + report: ArfFeedbackReport { + arrival_date: Some(days_ago(3)), + auth_failure: ArfAuthFailureType::Dmarc, + authentication_results: Map::new(vec![ + "dmarc=fail header.from=myserver.com".to_string(), + "spf=fail smtp.mailfrom=spoofed.example".to_string(), + ]), + delivery_result: ArfDeliveryResult::Policy, + dkim_adsp_dns: None, + dkim_canonicalized_body: None, + dkim_canonicalized_header: None, + dkim_domain: None, + dkim_identity: None, + dkim_selector: None, + dkim_selector_dns: None, + feedback_type: ArfFeedbackType::Fraud, + headers: Some( + "From: security@myserver.com\r\nTo: victim@antiphish.net\r\nSubject: Urgent: Verify your account" + .to_string(), + ), + identity_alignment: ArfIdentityAlignment::None, + incidents: 50, + message: Some("Phishing attempt impersonating myserver.com".to_string()), + original_envelope_id: None, + original_mail_from: Some("spoofer@spoofed.example".to_string()), + original_rcpt_to: Some("victim@antiphish.net".to_string()), + reported_domains: Map::new(vec![ + "myserver.com".to_string(), + "spoofed.example".to_string(), + ]), + reported_uris: Map::new(vec![ + "https://evil-site.example/phish".to_string(), + ]), + reporting_mta: Some("gateway.antiphish.net".to_string()), + source_ip: Some(IpAddr(std::net::IpAddr::V4(Ipv4Addr::new( + 198, 51, 100, 77, + )))), + source_port: Some(25), + spf_dns: Some("v=spf1 -all".to_string()), + user_agent: Some("AntiPhish/3.0".to_string()), + version: 1, + }, + subject: "Fraud report: phishing attempt using myserver.com".to_string(), + to: Map::new(vec![ + "abuse@myserver.com".to_string(), + "security@myserver.com".to_string(), + ]), + }, + ] +} + +fn one_year_from_now_u64() -> u64 { + store::write::now() + 365 * 24 * 3600 +} + +fn one_year_from_now() -> UTCDateTime { + UTCDateTime::from_timestamp(store::write::now().cast_signed() + 365 * 24 * 3600) +} + +fn now() -> UTCDateTime { + UTCDateTime::from_timestamp(store::write::now().cast_signed()) +} + +fn days_ago(days: i64) -> UTCDateTime { + UTCDateTime::from_timestamp(store::write::now().cast_signed() - days * 24 * 3600) +} diff --git a/crates/trc/Cargo.toml b/crates/trc/Cargo.toml index 5bb8b2d1..ed45b2a5 100644 --- a/crates/trc/Cargo.toml +++ b/crates/trc/Cargo.toml @@ -21,6 +21,7 @@ hashify = "0.2.7" [features] test_mode = [] +dev_mode = [] enterprise = [] [dev-dependencies] diff --git a/crates/trc/src/ipc/collector.rs b/crates/trc/src/ipc/collector.rs index b1eb8a73..a031d002 100644 --- a/crates/trc/src/ipc/collector.rs +++ b/crates/trc/src/ipc/collector.rs @@ -165,7 +165,7 @@ impl Collector { { event.inner.span = Some(span.clone()); } else { - #[cfg(debug_assertions)] + #[cfg(any(feature = "dev_mode", feature = "test_mode"))] { if event.span_id().unwrap() != 0 { eprintln!("Unregistered span ID: {event:?}"); @@ -179,7 +179,10 @@ impl Collector { if let Some(span) = self.active_spans.get(&span_id) { event.inner.span = Some(span.clone()); } else { - #[cfg(debug_assertions)] + #[cfg(any( + feature = "dev_mode", + feature = "test_mode" + ))] { if span_id != 0 { eprintln!("Unregistered span ID: {event:?}"); diff --git a/crates/trc/src/ipc/metrics.rs b/crates/trc/src/ipc/metrics.rs index c16c5090..3620a58a 100644 --- a/crates/trc/src/ipc/metrics.rs +++ b/crates/trc/src/ipc/metrics.rs @@ -342,6 +342,10 @@ impl Collector { CONNECTION_METRICS[CONN_SMTP_OUT].elapsed.observe(value) } MetricType::DnsLookupTime => DNS_LOOKUP_TIME.observe(value), + MetricType::StoreDataReadTime => STORE_DATA_READ_TIME.observe(value), + MetricType::StoreDataWriteTime => STORE_DATA_WRITE_TIME.observe(value), + MetricType::StoreBlobReadTime => STORE_BLOB_READ_TIME.observe(value), + MetricType::StoreBlobWriteTime => STORE_BLOB_WRITE_TIME.observe(value), _ => {} } } diff --git a/resources/html-templates/login.html b/resources/html-templates/login.html index ae8b8477..234b3dde 100644 --- a/resources/html-templates/login.html +++ b/resources/html-templates/login.html @@ -401,8 +401,8 @@ if (isDevice) { var req = { type: 'authDevice', - account_name: creds.account_name, - account_secret: creds.account_secret, + accountName: creds.account_name, + accountSecret: creds.account_secret, code: ($('device-code').value || '').trim() }; if (otpValue) req.mfaToken = otpValue; @@ -410,9 +410,9 @@ } var r = { type: 'authCode', - account_name: creds.account_name, - account_secret: creds.account_secret, - client_id: oauth.client_id || '' + accountName: creds.account_name, + accountSecret: creds.account_secret, + clientId: oauth.client_id || '' }; if (oauth.redirect_uri) r.redirectUri = oauth.redirect_uri; if (oauth.scope) r.scope = oauth.scope; diff --git a/resources/html-templates/login.html.min b/resources/html-templates/login.html.min index 29df511e..d594e423 100644 --- a/resources/html-templates/login.html.min +++ b/resources/html-templates/login.html.min @@ -1 +1 @@ - Sign in

Sign in

Enter your credentials to continue

\ No newline at end of file + Sign in

Sign in

Enter your credentials to continue

\ No newline at end of file diff --git a/tests/src/telemetry/metrics.rs b/tests/src/telemetry/metrics.rs index 32222965..800aee4c 100644 --- a/tests/src/telemetry/metrics.rs +++ b/tests/src/telemetry/metrics.rs @@ -5,14 +5,10 @@ */ use crate::utils::server::TestServer; -use common::telemetry::metrics::store::{MetricsStore, SharedMetricHistory}; +use common::telemetry::metrics::store::MetricsStore; use registry::{schema::prelude::ObjectType, types::datetime::UTCDateTime}; use std::time::Duration; -use store::{ - rand::{self, Rng}, - write::now, -}; -use trc::*; +use store::write::now; use types::id::Id; pub async fn test(test: &TestServer) { @@ -34,7 +30,7 @@ pub async fn test(test: &TestServer) { ); // Insert test metrics - insert_test_metrics(test).await; + test.server.insert_test_metrics().await; // Fetch all metrics let metric_ids = admin @@ -95,70 +91,3 @@ pub async fn test(test: &TestServer) { Vec::::new() ); } - -async fn insert_test_metrics(test: &TestServer) { - test.server - .metrics_store() - .purge_metrics(Duration::from_secs(0)) - .await - .unwrap(); - let mut start_time = now() - (90 * 24 * 60 * 60); - let timestamp = now(); - let history = SharedMetricHistory::default(); - - while start_time < timestamp { - for event_type in [ - EventType::Smtp(SmtpEvent::ConnectionStart), - EventType::Imap(ImapEvent::ConnectionStart), - EventType::Pop3(Pop3Event::ConnectionStart), - EventType::ManageSieve(ManageSieveEvent::ConnectionStart), - EventType::Http(HttpEvent::ConnectionStart), - EventType::Delivery(DeliveryEvent::AttemptStart), - EventType::Queue(QueueEvent::MessageQueued), - EventType::Queue(QueueEvent::AuthenticatedMessageQueued), - EventType::Queue(QueueEvent::DsnQueued), - EventType::Queue(QueueEvent::ReportQueued), - EventType::MessageIngest(MessageIngestEvent::Ham), - EventType::MessageIngest(MessageIngestEvent::Spam), - EventType::Auth(AuthEvent::Failed), - EventType::Security(SecurityEvent::AuthenticationBan), - EventType::Security(SecurityEvent::ScanBan), - EventType::Security(SecurityEvent::AbuseBan), - EventType::Security(SecurityEvent::LoiterBan), - EventType::Security(SecurityEvent::IpBlocked), - EventType::IncomingReport(IncomingReportEvent::DmarcReport), - EventType::IncomingReport(IncomingReportEvent::DmarcReportWithWarnings), - EventType::IncomingReport(IncomingReportEvent::TlsReport), - EventType::IncomingReport(IncomingReportEvent::TlsReportWithWarnings), - ] { - // Generate a random value between 0 and 100 - Collector::update_event_counter(event_type, rand::rng().random_range(0..=100)) - } - - Collector::update_gauge(MetricType::QueueCount, rand::rng().random_range(0..=1000)); - Collector::update_gauge( - MetricType::ServerMemory, - rand::rng().random_range(100 * 1024 * 1024..=300 * 1024 * 1024), - ); - - for metric_type in [ - MetricType::MessageIngestTime, - MetricType::MessageIngestIndexTime, - MetricType::DeliveryTotalTime, - MetricType::DnsLookupTime, - ] { - Collector::update_histogram(metric_type, rand::rng().random_range(2..=1000)) - } - Collector::update_histogram( - MetricType::DeliveryTotalTime, - rand::rng().random_range(1000..=5000), - ); - - test.server - .metrics_store() - .write_metrics(start_time.into(), history.clone()) - .await - .unwrap(); - start_time += 60 * 60 * 24; - } -}