diff --git a/crates/common/src/config/mod.rs b/crates/common/src/config/mod.rs index a68cd693..64360521 100644 --- a/crates/common/src/config/mod.rs +++ b/crates/common/src/config/mod.rs @@ -9,6 +9,7 @@ use std::sync::Arc; use arc_swap::ArcSwap; use directory::{Directories, Directory}; use store::{BlobBackend, BlobStore, FtsStore, LookupStore, Store, Stores}; +use telemetry::Metrics; use utils::config::Config; use crate::{expr::*, listener::tls::TlsManager, manager::config::ConfigManager, Core, Network}; @@ -138,6 +139,7 @@ impl Core { jmap: JmapConfig::parse(config), imap: ImapConfig::parse(config), tls: TlsManager::parse(config), + metrics: Metrics::parse(config), storage: Storage { data, blob, diff --git a/crates/common/src/config/telemetry.rs b/crates/common/src/config/telemetry.rs index d202d64c..bdd4e9c3 100644 --- a/crates/common/src/config/telemetry.rs +++ b/crates/common/src/config/telemetry.rs @@ -4,7 +4,7 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -use std::{collections::HashMap, str::FromStr, time::Duration}; +use std::{collections::HashMap, str::FromStr, sync::Arc, time::Duration}; use ahash::{AHashMap, AHashSet}; use base64::{engine::general_purpose::STANDARD, Engine}; @@ -12,11 +12,17 @@ use hyper::{ header::{HeaderName, HeaderValue, AUTHORIZATION, CONTENT_TYPE}, HeaderMap, }; +use opentelemetry::{InstrumentationLibrary, KeyValue}; use opentelemetry_otlp::WithExportConfig; use opentelemetry_sdk::{ export::{logs::LogExporter, trace::SpanExporter}, - metrics::exporter::PushMetricsExporter, + metrics::{ + exporter::PushMetricsExporter, + reader::{DefaultAggregationSelector, DefaultTemporalitySelector}, + }, + Resource, }; +use opentelemetry_semantic_conventions::resource::{SERVICE_NAME, SERVICE_VERSION}; use trc::{subscriber::Interests, EventType, Level, TelemetryEvent}; use utils::config::{utils::ParseValue, Config}; @@ -33,7 +39,6 @@ pub enum TelemetrySubscriberType { ConsoleTracer(ConsoleTracer), LogTracer(LogTracer), OtelTracer(OtelTracer), - OtelMetrics(OtelMetrics), Webhook(WebhookTracer), #[cfg(unix)] JournalTracer(crate::telemetry::tracers::journald::Subscriber), @@ -49,8 +54,10 @@ pub struct OtelTracer { } pub struct OtelMetrics { + pub resource: Resource, + pub instrumentation: InstrumentationLibrary, pub exporter: Box, - pub throttle: Duration, + pub interval: Duration, } #[derive(Debug)] @@ -90,12 +97,55 @@ pub enum RotationStrategy { #[derive(Debug)] pub struct Telemetry { - pub global_interests: Interests, - pub custom_levels: AHashMap, - pub tracers: Vec, + pub tracers: Tracers, + pub metrics: Interests, +} + +#[derive(Debug)] +pub struct Tracers { + pub interests: Interests, + pub levels: AHashMap, + pub subscribers: Vec, +} + +#[derive(Debug, Clone, Default)] +pub struct Metrics { + pub prometheus: bool, + pub otel: Option>, } impl Telemetry { + pub fn parse(config: &mut Config) -> Self { + let mut telemetry = Telemetry { + tracers: Tracers::parse(config), + metrics: Interests::default(), + }; + + // Parse metrics + if config + .property_or_default("metrics.prometheus.enable", "true") + .unwrap_or(true) + || config + .property_or_default("metrics.open-telemetry.enable", "false") + .unwrap_or(false) + { + apply_events( + config + .properties::("metrics.disabled-events") + .into_iter() + .map(|(_, e)| e), + false, + |event_type| { + telemetry.metrics.set(event_type); + }, + ); + } + + telemetry + } +} + +impl Tracers { pub fn parse(config: &mut Config) -> Self { // Parse custom logging levels let mut custom_levels = AHashMap::new(); @@ -115,17 +165,6 @@ impl Telemetry { } } - let event_names = EventType::variants() - .into_iter() - .filter_map(|e| { - if e != EventType::Telemetry(TelemetryEvent::WebhookError) { - Some((e, e.name())) - } else { - None - } - }) - .collect::>(); - // Parse tracers let mut tracers: Vec = Vec::new(); let mut global_interests = Interests::default(); @@ -422,67 +461,44 @@ impl Telemetry { .unwrap_or(Level::Info); // Parse disabled events - let mut disabled_events = AHashSet::new(); - match &tracer.typ { - TelemetrySubscriberType::ConsoleTracer(_) => (), + let exclude_event = match &tracer.typ { + TelemetrySubscriberType::ConsoleTracer(_) => None, TelemetrySubscriberType::LogTracer(_) => { - disabled_events.insert(EventType::Telemetry(TelemetryEvent::LogError)); + EventType::Telemetry(TelemetryEvent::LogError).into() } TelemetrySubscriberType::OtelTracer(_) => { - disabled_events.insert(EventType::Telemetry(TelemetryEvent::OtelError)); + EventType::Telemetry(TelemetryEvent::OtelExpoterError).into() } TelemetrySubscriberType::Webhook(_) => { - disabled_events.insert(EventType::Telemetry(TelemetryEvent::WebhookError)); + EventType::Telemetry(TelemetryEvent::WebhookError).into() } #[cfg(unix)] TelemetrySubscriberType::JournalTracer(_) => { - disabled_events.insert(EventType::Telemetry(TelemetryEvent::JournalError)); + EventType::Telemetry(TelemetryEvent::JournalError).into() } - TelemetrySubscriberType::OtelMetrics(_) => todo!(), - } - for (_, event_type) in - config.properties::(("tracer", id, "disabled-events")) - { - match event_type { - EventOrMany::Event(event_type) => { - disabled_events.insert(event_type); - } - EventOrMany::StartsWith(value) => { - for (event_type, name) in event_names.iter() { - if name.starts_with(&value) { - disabled_events.insert(*event_type); - } - } - } - EventOrMany::EndsWith(value) => { - for (event_type, name) in event_names.iter() { - if name.ends_with(&value) { - disabled_events.insert(*event_type); - } - } - } - EventOrMany::All => { - for (event_type, _) in event_names.iter() { - disabled_events.insert(*event_type); - } - break; - } - } - } + }; - // Build interests lists - for event_type in EventType::variants() { - if !disabled_events.contains(&event_type) { - let event_level = custom_levels - .get(&event_type) - .copied() - .unwrap_or(event_type.level()); - if level.is_contained(event_level) { - tracer.interests.set(event_type); - global_interests.set(event_type); + // Parse disabled events + apply_events( + config + .properties::(("tracer", id, "disabled-events")) + .into_iter() + .map(|(_, e)| e), + false, + |event_type| { + if exclude_event != Some(event_type) { + let event_level = custom_levels + .get(&event_type) + .copied() + .unwrap_or(event_type.level()); + if level.is_contained(event_level) { + tracer.interests.set(event_type); + global_interests.set(event_type); + } } - } - } + }, + ); + if !tracer.interests.is_empty() { tracers.push(tracer); } else { @@ -526,14 +542,137 @@ impl Telemetry { }); } - Telemetry { - tracers, - global_interests, - custom_levels, + Tracers { + subscribers: tracers, + interests: global_interests, + levels: custom_levels, } } } +impl Metrics { + pub fn parse(config: &mut Config) -> Self { + let mut metrics = Metrics { + prometheus: config + .property_or_default("metrics.prometheus.enable", "true") + .unwrap_or(true), + otel: None, + }; + + if config + .property_or_default("metrics.open-telemetry.enable", "false") + .unwrap_or(false) + { + let timeout = config + .property::("metrics.open-telemetry.timeout") + .unwrap_or(Duration::from_secs( + opentelemetry_otlp::OTEL_EXPORTER_OTLP_TIMEOUT_DEFAULT, + )); + let interval = config + .property_or_default("metrics.open-telemetry.interval", "1m") + .unwrap_or_else(|| Duration::from_secs(60)); + let resource = Resource::new([ + KeyValue::new(SERVICE_NAME, "stalwart-mail"), + KeyValue::new(SERVICE_VERSION, env!("CARGO_PKG_VERSION")), + ]); + let instrumentation = InstrumentationLibrary::builder("stalwart-mail") + .with_version(env!("CARGO_PKG_VERSION")) + .build(); + + match config + .value_require("metrics.open-telemetry.transport") + .unwrap_or_default() + { + "grpc" => { + let mut exporter = opentelemetry_otlp::new_exporter() + .tonic() + .with_protocol(opentelemetry_otlp::Protocol::Grpc) + .with_timeout(timeout); + if let Some(endpoint) = config.value("metrics.open-telemetry.endpoint") { + exporter = exporter.with_endpoint(endpoint); + } + + match exporter.build_metrics_exporter( + Box::new(DefaultAggregationSelector::new()), + Box::new(DefaultTemporalitySelector::new()), + ) { + Ok(exporter) => { + metrics.otel = Some(Arc::new(OtelMetrics { + exporter: Box::new(exporter), + interval, + resource, + instrumentation, + })); + } + Err(err) => { + config.new_build_error( + "metrics.open-telemetry", + format!("Failed to build OpenTelemetry metrics exporter: {err}"), + ); + } + } + } + "http" => { + if let Some(endpoint) = config + .value_require("metrics.open-telemetry.endpoint") + .map(|s| s.to_string()) + { + let mut headers = HashMap::new(); + let mut err = None; + for (_, value) in config.values("metrics.open-telemetry.headers") { + if let Some((key, value)) = value.split_once(':') { + headers.insert(key.trim().to_string(), value.trim().to_string()); + } else { + err = format!("Invalid open-telemetry header {value:?}").into(); + break; + } + } + if let Some(err) = err { + config.new_parse_error("metrics.open-telemetry.headers", err); + } + + let mut exporter = opentelemetry_otlp::new_exporter() + .http() + .with_endpoint(&endpoint) + .with_timeout(timeout); + if !headers.is_empty() { + exporter = exporter.with_headers(headers); + } + + match exporter.build_metrics_exporter( + Box::new(DefaultAggregationSelector::new()), + Box::new(DefaultTemporalitySelector::new()), + ) { + Ok(exporter) => { + metrics.otel = Some(Arc::new(OtelMetrics { + exporter: Box::new(exporter), + interval, + resource, + instrumentation, + })); + } + Err(err) => { + config.new_build_error( + "metrics.open-telemetry", + format!( + "Failed to build OpenTelemetry metrics exporter: {err}" + ), + ); + } + } + } + } + transport => { + let err = format!("Invalid transport: {transport}"); + config.new_parse_error("metrics.open-telemetry.transport", err); + } + } + } + + metrics + } +} + fn parse_webhook( config: &mut Config, id: &str, @@ -609,49 +748,19 @@ fn parse_webhook( }; // Parse webhook events - let event_names = EventType::variants() - .into_iter() - .filter_map(|e| { - if e != EventType::Telemetry(TelemetryEvent::WebhookError) { - Some((e, e.name())) - } else { - None + apply_events( + config + .properties::(("webhook", id, "events")) + .into_iter() + .map(|(_, e)| e), + true, + |event_type| { + if event_type != EventType::Telemetry(TelemetryEvent::WebhookError) { + tracer.interests.set(event_type); + global_interests.set(event_type); } - }) - .collect::>(); - for (_, event_type) in config.properties::(("webhook", id, "events")) { - match event_type { - EventOrMany::Event(event_type) => { - if event_type != EventType::Telemetry(TelemetryEvent::WebhookError) { - tracer.interests.set(event_type); - global_interests.set(event_type); - } - } - EventOrMany::StartsWith(value) => { - for (event_type, name) in event_names.iter() { - if name.starts_with(&value) { - tracer.interests.set(*event_type); - global_interests.set(*event_type); - } - } - } - EventOrMany::EndsWith(value) => { - for (event_type, name) in event_names.iter() { - if name.ends_with(&value) { - tracer.interests.set(*event_type); - global_interests.set(*event_type); - } - } - } - EventOrMany::All => { - for (event_type, _) in event_names.iter() { - tracer.interests.set(*event_type); - global_interests.set(*event_type); - } - break; - } - } - } + }, + ); if !tracer.interests.is_empty() { Some(tracer) @@ -668,6 +777,66 @@ enum EventOrMany { All, } +fn apply_events( + event_types: impl IntoIterator, + inclusive: bool, + mut apply_fn: impl FnMut(EventType), +) { + let event_names = EventType::variants() + .into_iter() + .map(|e| (e, e.name())) + .collect::>(); + let mut exclude_events = AHashSet::new(); + + for event_or_many in event_types { + match event_or_many { + EventOrMany::Event(event_type) => { + apply_fn(event_type); + } + EventOrMany::StartsWith(value) => { + for (event_type, name) in event_names.iter() { + if name.starts_with(&value) { + if inclusive { + apply_fn(*event_type); + } else { + exclude_events.insert(*event_type); + } + } + } + } + EventOrMany::EndsWith(value) => { + for (event_type, name) in event_names.iter() { + if name.ends_with(&value) { + if inclusive { + apply_fn(*event_type); + } else { + exclude_events.insert(*event_type); + } + } + } + } + EventOrMany::All => { + for (event_type, _) in event_names.iter() { + if inclusive { + apply_fn(*event_type); + } else { + exclude_events.insert(*event_type); + } + } + break; + } + } + } + + if !inclusive { + for (event_type, _) in event_names.iter() { + if !exclude_events.contains(event_type) { + apply_fn(*event_type); + } + } + } +} + impl ParseValue for EventOrMany { fn parse_value(value: &str) -> Result { let value = value.trim(); @@ -686,7 +855,7 @@ impl ParseValue for EventOrMany { impl std::fmt::Debug for OtelMetrics { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("OtelMetrics") - .field("throttle", &self.throttle) + .field("interval", &self.interval) .finish() } } diff --git a/crates/common/src/lib.rs b/crates/common/src/lib.rs index 808eaceb..7eff8bcd 100644 --- a/crates/common/src/lib.rs +++ b/crates/common/src/lib.rs @@ -17,6 +17,7 @@ use config::{ SmtpConfig, }, storage::Storage, + telemetry::Metrics, }; use directory::{core::secret::verify_secret_hash, Directory, Principal, QueryBy, Type}; use expr::if_block::IfBlock; @@ -57,6 +58,7 @@ pub struct Core { pub smtp: SmtpConfig, pub jmap: JmapConfig, pub imap: ImapConfig, + pub metrics: Metrics, #[cfg(feature = "enterprise")] pub enterprise: Option, } diff --git a/crates/common/src/manager/boot.rs b/crates/common/src/manager/boot.rs index d46ec2d3..bb1ff3c5 100644 --- a/crates/common/src/manager/boot.rs +++ b/crates/common/src/manager/boot.rs @@ -161,7 +161,7 @@ impl BootManager { .failed("Failed to read configuration"); } - // Enable tracing + // Enable telemetry Telemetry::parse(&mut config).enable(); match import_export { diff --git a/crates/common/src/telemetry/metrics/mod.rs b/crates/common/src/telemetry/metrics/mod.rs index c8a832f8..f7d82906 100644 --- a/crates/common/src/telemetry/metrics/mod.rs +++ b/crates/common/src/telemetry/metrics/mod.rs @@ -3,3 +3,5 @@ * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ + +pub mod otel; diff --git a/crates/common/src/telemetry/metrics/otel.rs b/crates/common/src/telemetry/metrics/otel.rs new file mode 100644 index 00000000..5b85b5e7 --- /dev/null +++ b/crates/common/src/telemetry/metrics/otel.rs @@ -0,0 +1,138 @@ +/* + * SPDX-FileCopyrightText: 2020 Stalwart Labs Ltd + * + * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + */ + +use std::{sync::Arc, time::SystemTime}; + +use opentelemetry::global::set_error_handler; +use opentelemetry_sdk::metrics::data::{ + DataPoint, Gauge, Histogram, HistogramDataPoint, Metric, ResourceMetrics, ScopeMetrics, Sum, + Temporality, +}; +use trc::{collector::Collector, TelemetryEvent}; + +use crate::{config::telemetry::OtelMetrics, Core}; + +impl OtelMetrics { + pub async fn push_metrics(&self, core: Arc, start_time: SystemTime) { + let mut metrics = Vec::with_capacity(256); + let now = SystemTime::now(); + + #[cfg(feature = "enterprise")] + let is_enterprise = core.is_enterprise_edition(); + + #[cfg(not(feature = "enterprise"))] + let is_enterprise = false; + + // Add counters + for counter in Collector::collect_counters(is_enterprise) { + metrics.push(Metric { + name: counter.id().into(), + description: counter.description().into(), + unit: counter.unit().into(), + data: Box::new(Sum { + data_points: vec![DataPoint { + attributes: vec![], + start_time: start_time.into(), + time: now.into(), + value: counter.get(), + exemplars: vec![], + }], + temporality: Temporality::Cumulative, + is_monotonic: true, + }), + }); + } + + // Add event counters + for counter in Collector::collect_event_counters(is_enterprise) { + metrics.push(Metric { + name: counter.id().into(), + description: counter.description().into(), + unit: "events".into(), + data: Box::new(Sum { + data_points: vec![DataPoint { + attributes: vec![], + start_time: start_time.into(), + time: now.into(), + value: counter.value(), + exemplars: vec![], + }], + temporality: Temporality::Cumulative, + is_monotonic: true, + }), + }); + } + + // Add gauges + for gauge in Collector::collect_gauges(is_enterprise) { + metrics.push(Metric { + name: gauge.id().into(), + description: gauge.description().into(), + unit: gauge.unit().into(), + data: Box::new(Gauge { + data_points: vec![DataPoint { + attributes: vec![], + start_time: start_time.into(), + time: now.into(), + value: gauge.get(), + exemplars: vec![], + }], + }), + }); + } + + // Add histograms + for histogram in Collector::collect_histograms(is_enterprise) { + metrics.push(Metric { + name: histogram.id().into(), + description: histogram.description().into(), + unit: histogram.unit().into(), + data: Box::new(Histogram { + data_points: vec![HistogramDataPoint { + attributes: vec![], + start_time, + time: now, + count: histogram.count(), + bounds: histogram.upper_bounds_vec(), + bucket_counts: histogram.buckets_vec(), + min: histogram.min(), + max: histogram.max(), + sum: histogram.sum(), + exemplars: vec![], + }], + temporality: Temporality::Cumulative, + }), + }); + } + + // Export metrics + if let Err(err) = self + .exporter + .export(&mut ResourceMetrics { + resource: self.resource.clone(), + scope_metrics: vec![ScopeMetrics { + scope: self.instrumentation.clone(), + metrics, + }], + }) + .await + { + trc::event!( + Telemetry(TelemetryEvent::OtelMetricsExporterError), + Reason = err.to_string(), + ); + } + } + + pub fn enable_errors() { + let _ = set_error_handler(|error| { + trc::event!( + Telemetry(TelemetryEvent::OtelMetricsExporterError), + Reason = error.to_string(), + ); + }); + } +} diff --git a/crates/common/src/telemetry/mod.rs b/crates/common/src/telemetry/mod.rs index aad58142..81d5fe68 100644 --- a/crates/common/src/telemetry/mod.rs +++ b/crates/common/src/telemetry/mod.rs @@ -23,7 +23,7 @@ pub const LONG_SLUMBER: Duration = Duration::from_secs(60 * 60 * 24 * 365); impl Telemetry { pub fn enable(self) { // Spawn tracers - for tracer in self.tracers { + for tracer in self.tracers.subscribers { tracer.typ.spawn( SubscriberBuilder::new(tracer.id) .with_interests(tracer.interests) @@ -32,8 +32,9 @@ impl Telemetry { } // Update global collector - Collector::set_interests(self.global_interests); - Collector::update_custom_levels(self.custom_levels); + Collector::set_interests(self.tracers.interests); + Collector::update_custom_levels(self.tracers.levels); + Collector::set_metrics(self.metrics); Collector::reload(); } @@ -43,6 +44,7 @@ impl Telemetry { for subscribed_id in &active_subscribers { if !self .tracers + .subscribers .iter() .any(|tracer| tracer.id == *subscribed_id) { @@ -51,7 +53,7 @@ impl Telemetry { } // Activate new tracers or update existing ones - for tracer in self.tracers { + for tracer in self.tracers.subscribers { if active_subscribers.contains(&tracer.id) { Collector::update_subscriber(tracer.id, tracer.interests, tracer.lossy); } else { @@ -64,8 +66,9 @@ impl Telemetry { } // Update global collector - Collector::set_interests(self.global_interests); - Collector::update_custom_levels(self.custom_levels); + Collector::set_interests(self.tracers.interests); + Collector::update_custom_levels(self.tracers.levels); + Collector::set_metrics(self.metrics); Collector::reload(); } @@ -107,7 +110,6 @@ impl TelemetrySubscriberType { TelemetrySubscriberType::JournalTracer(subscriber) => { tracers::journald::spawn_journald_tracer(builder, subscriber) } - TelemetrySubscriberType::OtelMetrics(_) => todo!(), } } } diff --git a/crates/common/src/telemetry/tracers/otel.rs b/crates/common/src/telemetry/tracers/otel.rs index c0f61a77..e3dd0e87 100644 --- a/crates/common/src/telemetry/tracers/otel.rs +++ b/crates/common/src/telemetry/tracers/otel.rs @@ -89,7 +89,7 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer .await { trc::event!( - Telemetry(TelemetryEvent::OtelError), + Telemetry(TelemetryEvent::OtelExpoterError), Details = "Failed to export spans", Reason = err.to_string() ); @@ -103,7 +103,7 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer .await { trc::event!( - Telemetry(TelemetryEvent::OtelError), + Telemetry(TelemetryEvent::OtelExpoterError), Details = "Failed to export logs", Reason = err.to_string() ); diff --git a/crates/jmap/src/api/management/reload.rs b/crates/jmap/src/api/management/reload.rs index b0d10ac5..6880ff48 100644 --- a/crates/jmap/src/api/management/reload.rs +++ b/crates/jmap/src/api/management/reload.rs @@ -63,15 +63,15 @@ impl JMAP { tracers.update(); } - // Reload ACME + // Reload settings self.inner .housekeeper_tx - .send(Event::AcmeReload) + .send(Event::ReloadSettings) .await .map_err(|err| { trc::EventType::Server(trc::ServerEvent::ThreadError) .reason(err) - .details("Failed to send ACME reload event to housekeeper") + .details("Failed to send settings reload event to housekeeper") .caused_by(trc::location!()) })?; } diff --git a/crates/jmap/src/email/copy.rs b/crates/jmap/src/email/copy.rs index f3d1f7aa..7be5e780 100644 --- a/crates/jmap/src/email/copy.rs +++ b/crates/jmap/src/email/copy.rs @@ -41,10 +41,7 @@ use store::{ use trc::AddContext; use utils::map::vec_map::VecMap; -use crate::{ - api::http::HttpSessionData, auth::AccessToken, mailbox::UidMailbox, - services::housekeeper::Event, JMAP, -}; +use crate::{api::http::HttpSessionData, auth::AccessToken, mailbox::UidMailbox, JMAP}; use super::{ index::{EmailIndexBuilder, TrimTextValue, VisitValues, MAX_ID_LENGTH, MAX_SORT_FIELD_LENGTH}, @@ -432,7 +429,7 @@ impl JMAP { let document_id = ids.last_document_id().caused_by(trc::location!())?; // Request FTS index - let _ = self.inner.housekeeper_tx.send(Event::IndexStart).await; + self.inner.request_fts_index(); // Update response email.id = Id::from_parts(thread_id, document_id); diff --git a/crates/jmap/src/email/ingest.rs b/crates/jmap/src/email/ingest.rs index 4bf61914..618e5bc4 100644 --- a/crates/jmap/src/email/ingest.rs +++ b/crates/jmap/src/email/ingest.rs @@ -37,7 +37,6 @@ use utils::map::vec_map::VecMap; use crate::{ email::index::{IndexMessage, VisitValues, MAX_ID_LENGTH}, mailbox::{UidMailbox, INBOX_ID, JUNK_ID}, - services::housekeeper::Event, JMAP, }; @@ -335,7 +334,7 @@ impl JMAP { let id = Id::from_parts(thread_id, document_id); // Request FTS index - let _ = self.inner.housekeeper_tx.send(Event::IndexStart).await; + self.inner.request_fts_index(); trc::event!( MessageIngest(match params.source { diff --git a/crates/jmap/src/lib.rs b/crates/jmap/src/lib.rs index 1133eae4..c99ff030 100644 --- a/crates/jmap/src/lib.rs +++ b/crates/jmap/src/lib.rs @@ -26,6 +26,7 @@ use jmap_proto::{ use services::{ delivery::spawn_delivery_manager, housekeeper::{self, init_housekeeper, spawn_housekeeper}, + index::spawn_index_task, state::{self, init_state_manager, spawn_state_manager}, }; @@ -41,7 +42,7 @@ use store::{ }, BitmapKey, Deserialize, IterateParams, ValueKey, U32_LEN, }; -use tokio::sync::mpsc; +use tokio::sync::{mpsc, Notify}; use trc::AddContext; use utils::{ config::Config, @@ -95,6 +96,7 @@ pub struct Inner { pub state_tx: mpsc::Sender, pub housekeeper_tx: mpsc::Sender, + pub index_tx: Arc, pub cache_threads: LruCache>, } @@ -109,6 +111,7 @@ impl JMAP { // Init state manager and housekeeper let (state_tx, state_rx) = init_state_manager(); let (housekeeper_tx, housekeeper_rx) = init_housekeeper(); + let index_tx = Arc::new(Notify::new()); let shard_amount = config .property::("cache.shard") .unwrap_or(32) @@ -130,6 +133,7 @@ impl JMAP { ), state_tx, housekeeper_tx, + index_tx: index_tx.clone(), cache_threads: LruCache::with_capacity( config.property("cache.thread.size").unwrap_or(2048), ), @@ -160,6 +164,9 @@ impl JMAP { // Spawn housekeeper spawn_housekeeper(jmap_instance.clone(), housekeeper_rx); + // Spawn index task + spawn_index_task(jmap_instance.clone(), index_tx); + jmap_instance } diff --git a/crates/jmap/src/services/gossip/ping.rs b/crates/jmap/src/services/gossip/ping.rs index 467654d0..1dad82a4 100644 --- a/crates/jmap/src/services/gossip/ping.rs +++ b/crates/jmap/src/services/gossip/ping.rs @@ -71,11 +71,7 @@ impl Gossiper { tokio::spawn(async move { trc::event!(Cluster(ClusterEvent::OneOrMorePeersOffline)); - let _ = core - .jmap_inner - .housekeeper_tx - .send(housekeeper::Event::IndexStart) - .await; + core.jmap_inner.request_fts_index(); let _ = core.smtp_inner.queue_tx.send(queue::Event::Reload).await; }); } @@ -181,13 +177,13 @@ impl Gossiper { // Reload ACME if inner .housekeeper_tx - .send(housekeeper::Event::AcmeReload) + .send(housekeeper::Event::ReloadSettings) .await .is_err() { trc::event!( Server(trc::ServerEvent::ThreadError), - Details = "Failed to send ACME reload event to housekeeper", + Details = "Failed to send setting reload event to housekeeper", CausedBy = trc::location!(), ); } diff --git a/crates/jmap/src/services/housekeeper.rs b/crates/jmap/src/services/housekeeper.rs index 82ea572d..1eee42ab 100644 --- a/crates/jmap/src/services/housekeeper.rs +++ b/crates/jmap/src/services/housekeeper.rs @@ -6,10 +6,10 @@ use std::{ collections::BinaryHeap, - time::{Duration, Instant}, + time::{Duration, Instant, SystemTime}, }; -use common::IPC_CHANNEL_BUFFER; +use common::{config::telemetry::OtelMetrics, IPC_CHANNEL_BUFFER}; use store::{ write::{now, purge::PurgeStore}, BlobStore, LookupStore, Store, @@ -21,16 +21,12 @@ use utils::map::ttl_dashmap::TtlMap; use crate::{Inner, JmapInstance, JMAP, LONG_SLUMBER}; pub enum Event { - IndexStart, - IndexDone, - AcmeReload, AcmeReschedule { provider_id: String, renew_at: Instant, }, Purge(PurgeType), - #[cfg(feature = "test_mode")] - IndexIsActive(tokio::sync::oneshot::Sender), + ReloadSettings, Exit, } @@ -53,8 +49,9 @@ enum ActionClass { Account, Store(usize), Acme(String), + OtelMetrics, #[cfg(feature = "enterprise")] - ReloadLicense, + ReloadSettings, } #[derive(Default)] @@ -65,28 +62,26 @@ struct Queue { pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { tokio::spawn(async move { trc::event!(Housekeeper(HousekeeperEvent::Start)); - - let mut index_busy = true; - let mut index_pending = false; - - // Index any queued messages - let jmap = JMAP::from(core.clone()); - tokio::spawn(async move { - jmap.fts_index_queued().await; - }); + let start_time = SystemTime::now(); // Add all events to queue let mut queue = Queue::default(); { let core_ = core.core.load_full(); + + // Session purge queue.schedule( Instant::now() + core_.jmap.session_purge_frequency.time_to_next(), ActionClass::Session, ); + + // Account purge queue.schedule( Instant::now() + core_.jmap.account_purge_frequency.time_to_next(), ActionClass::Account, ); + + // Store purges for (idx, schedule) in core_.storage.purge_schedules.iter().enumerate() { queue.schedule( Instant::now() + schedule.cron.time_to_next(), @@ -94,6 +89,12 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { ); } + // OTEL Push Metrics + if let Some(otel) = &core_.metrics.otel { + OtelMetrics::enable_errors(); + queue.schedule(Instant::now() + otel.interval, ActionClass::OtelMetrics); + } + // Add all ACME renewals to heap for provider in core_.tls.acme_providers.values() { match core_.init_acme(provider).await { @@ -118,7 +119,7 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { if let Some(enterprise) = &core_.enterprise { queue.schedule( Instant::now() + enterprise.license.expires_in(), - ActionClass::ReloadLicense, + ActionClass::ReloadSettings, ); } // SPDX-SnippetEnd @@ -127,10 +128,24 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { loop { match tokio::time::timeout(queue.wake_up_time(), rx.recv()).await { Ok(Some(event)) => match event { - Event::AcmeReload => { + Event::ReloadSettings => { let core_ = core.core.load_full(); let inner = core.jmap_inner.clone(); + // Reload OTEL push metrics + match &core_.metrics.otel { + Some(otel) if !queue.has_action(&ActionClass::OtelMetrics) => { + OtelMetrics::enable_errors(); + + queue.schedule( + Instant::now() + otel.interval, + ActionClass::OtelMetrics, + ); + } + _ => {} + } + + // Reload ACME certificates tokio::spawn(async move { for provider in core_.tls.acme_providers.values() { match core_.init_acme(provider).await { @@ -160,28 +175,6 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { queue.remove_action(&action); queue.schedule(renew_at, action); } - Event::IndexStart => { - if !index_busy { - index_busy = true; - let jmap = JMAP::from(core.clone()); - tokio::spawn(async move { - jmap.fts_index_queued().await; - }); - } else { - index_pending = true; - } - } - Event::IndexDone => { - if index_pending { - index_pending = false; - let jmap = JMAP::from(core.clone()); - tokio::spawn(async move { - jmap.fts_index_queued().await; - }); - } else { - index_busy = false; - } - } Event::Purge(purge) => match purge { PurgeType::Data(store) => { tokio::spawn(async move { @@ -225,10 +218,6 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { }); } }, - #[cfg(feature = "test_mode")] - Event::IndexIsActive(tx) => { - tx.send(index_busy).ok(); - } Event::Exit => { trc::event!(Housekeeper(HousekeeperEvent::Stop)); @@ -352,12 +341,26 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { }); } } + ActionClass::OtelMetrics => { + if let Some(otel) = &core_.metrics.otel { + queue.schedule( + Instant::now() + otel.interval, + ActionClass::OtelMetrics, + ); + + let otel = otel.clone(); + let core = core_.clone(); + tokio::spawn(async move { + otel.push_metrics(core, start_time).await; + }); + } + } // SPDX-SnippetBegin // SPDX-FileCopyrightText: 2020 Stalwart Labs Ltd // SPDX-License-Identifier: LicenseRef-SEL #[cfg(feature = "enterprise")] - ActionClass::ReloadLicense => { + ActionClass::ReloadSettings => { match core_.reload().await { Ok(result) => { if let Some(new_core) = result.new_core { @@ -365,7 +368,7 @@ pub fn spawn_housekeeper(core: JmapInstance, mut rx: mpsc::Receiver) { queue.schedule( Instant::now() + enterprise.license.expires_in(), - ActionClass::ReloadLicense, + ActionClass::ReloadSettings, ); } @@ -420,6 +423,10 @@ impl Queue { None } } + + pub fn has_action(&self, event: &ActionClass) -> bool { + self.heap.iter().any(|e| &e.event == event) + } } impl Ord for Action { diff --git a/crates/jmap/src/services/index.rs b/crates/jmap/src/services/index.rs index dfaaa9e8..a8501771 100644 --- a/crates/jmap/src/services/index.rs +++ b/crates/jmap/src/services/index.rs @@ -4,7 +4,7 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -use std::time::Instant; +use std::{sync::Arc, time::Instant}; use jmap_proto::types::{collection::Collection, property::Property}; use store::{ @@ -16,16 +16,15 @@ use store::{ Deserialize, IterateParams, Serialize, ValueKey, U32_LEN, U64_LEN, }; +use tokio::sync::Notify; use trc::FtsIndexEvent; use utils::{BlobHash, BLOB_HASH_LEN}; use crate::{ email::{index::IndexMessageText, metadata::MessageMetadata}, - JMAP, + Inner, JmapInstance, JMAP, }; -use super::housekeeper::Event; - #[derive(Debug)] struct IndexEmail { account_id: u32, @@ -37,6 +36,18 @@ struct IndexEmail { const INDEX_LOCK_EXPIRY: u64 = 60 * 5; +pub fn spawn_index_task(core: JmapInstance, rx: Arc) { + tokio::spawn(async move { + loop { + // Index any queued messages + JMAP::from(core.clone()).fts_index_queued().await; + + // Wait for a signal to index more messages + rx.notified().await; + } + }); +} + impl JMAP { pub async fn fts_index_queued(&self) { let from_key = ValueKey::> { @@ -193,20 +204,6 @@ impl JMAP { break; } } - - if self - .inner - .housekeeper_tx - .send(Event::IndexDone) - .await - .is_err() - { - trc::event!( - Server(trc::ServerEvent::ThreadError), - Details = "Failed to send event to Housekeeper", - CausedBy = trc::location!() - ); - } } async fn try_lock_index(&self, event: &IndexEmail) -> bool { @@ -239,6 +236,16 @@ impl JMAP { } } } + + pub fn request_fts_index(&self) { + self.inner.request_fts_index(); + } +} + +impl Inner { + pub fn request_fts_index(&self) { + self.index_tx.notify_one(); + } } impl IndexEmail { diff --git a/crates/trc/src/atomic.rs b/crates/trc/src/atomic.rs index 15d6fa81..ec874bb1 100644 --- a/crates/trc/src/atomic.rs +++ b/crates/trc/src/atomic.rs @@ -16,6 +16,8 @@ pub struct AtomicHistogram { upper_bounds: [u64; N], sum: AtomicU64, count: AtomicU64, + min: AtomicU64, + max: AtomicU64, } pub struct AtomicGauge { id: &'static str, @@ -115,6 +117,8 @@ impl AtomicHistogram { upper_bounds, sum: AtomicU64::new(0), count: AtomicU64::new(0), + min: AtomicU64::new(u64::MAX), + max: AtomicU64::new(0), id, description, unit, @@ -124,6 +128,8 @@ impl AtomicHistogram { pub fn observe(&self, value: u64) { self.sum.fetch_add(value, Ordering::Relaxed); self.count.fetch_add(1, Ordering::Relaxed); + self.min.fetch_min(value, Ordering::Relaxed); + self.max.fetch_max(value, Ordering::Relaxed); for (idx, upper_bound) in self.upper_bounds.iter().enumerate() { if value < *upper_bound { @@ -147,6 +153,63 @@ impl AtomicHistogram { self.unit } + pub fn sum(&self) -> u64 { + self.sum.load(Ordering::Relaxed) + } + + pub fn count(&self) -> u64 { + self.count.load(Ordering::Relaxed) + } + + pub fn min(&self) -> Option { + let min = self.min.load(Ordering::Relaxed); + if min != u64::MAX { + Some(min) + } else { + None + } + } + + pub fn max(&self) -> Option { + let max = self.max.load(Ordering::Relaxed); + if max != 0 { + Some(max) + } else { + None + } + } + + pub fn buckets_iter(&self) -> impl IntoIterator + '_ { + self.buckets + .inner() + .iter() + .map(|bucket| bucket.load(Ordering::Relaxed)) + } + + pub fn buckets_vec(&self) -> Vec { + let mut vec = Vec::with_capacity(N); + for bucket in self.buckets.inner().iter() { + vec.push(bucket.load(Ordering::Relaxed)); + } + vec + } + + pub fn upper_bounds_iter(&self) -> impl IntoIterator + '_ { + self.upper_bounds.iter().copied() + } + + pub fn upper_bounds_vec(&self) -> Vec { + let mut vec = Vec::with_capacity(N - 1); + for upper_bound in self.upper_bounds.iter().take(N - 1) { + vec.push(*upper_bound as f64); + } + vec + } + + pub fn is_active(&self) -> bool { + self.count.load(Ordering::Relaxed) > 0 + } + pub const fn new_message_sizes( id: &'static str, description: &'static str, @@ -294,6 +357,10 @@ impl AtomicCounter { pub fn unit(&self) -> &'static str { self.unit } + + pub fn is_active(&self) -> bool { + self.value.load(Ordering::Relaxed) > 0 + } } impl AtomicGauge { diff --git a/crates/trc/src/imple.rs b/crates/trc/src/imple.rs index fa49fde0..54b9c451 100644 --- a/crates/trc/src/imple.rs +++ b/crates/trc/src/imple.rs @@ -1214,8 +1214,8 @@ impl EventType { | HousekeeperEvent::PurgeAccounts | HousekeeperEvent::PurgeSessions | HousekeeperEvent::PurgeStore - | HousekeeperEvent::Schedule | HousekeeperEvent::Stop => Level::Info, + HousekeeperEvent::Schedule => Level::Debug, }, EventType::FtsIndex(event) => match event { FtsIndexEvent::Index => Level::Info, @@ -1909,9 +1909,9 @@ impl TelemetryEvent { TelemetryEvent::Update => "Tracing update", TelemetryEvent::LogError => "Log collector error", TelemetryEvent::WebhookError => "Webhook collector error", - TelemetryEvent::OtelError => "OpenTelemetry collector error", TelemetryEvent::JournalError => "Journal collector error", - TelemetryEvent::MetricsError => "Metrics collector error", + TelemetryEvent::OtelExpoterError => "OpenTelemetry exporter error", + TelemetryEvent::OtelMetricsExporterError => "OpenTelemetry metrics exporter error", } } } diff --git a/crates/trc/src/lib.rs b/crates/trc/src/lib.rs index 5dae4777..4a95eab1 100644 --- a/crates/trc/src/lib.rs +++ b/crates/trc/src/lib.rs @@ -652,9 +652,9 @@ pub enum TelemetryEvent { Update, LogError, WebhookError, - OtelError, + OtelExpoterError, + OtelMetricsExporterError, JournalError, - MetricsError, } #[event_type] diff --git a/crates/trc/src/metrics.rs b/crates/trc/src/metrics.rs index 886766f0..fd477784 100644 --- a/crates/trc/src/metrics.rs +++ b/crates/trc/src/metrics.rs @@ -226,7 +226,7 @@ impl Collector { METRIC_INTERESTS.update(interests); } - pub fn collect_event_counters() -> impl Iterator { + pub fn collect_event_counters(_is_enterprise: bool) -> impl Iterator { EVENT_COUNTERS .inner() .iter() @@ -247,18 +247,21 @@ impl Collector { }) } - pub fn collect_counters() -> impl Iterator { + pub fn collect_counters(_is_enterprise: bool) -> impl Iterator { CONNECTION_METRICS .iter() .flat_map(|m| [&m.total_connections, &m.bytes_sent, &m.bytes_received]) + .filter(|c| c.is_active()) } - pub fn collect_gauges() -> impl Iterator { + pub fn collect_gauges(_is_enterprise: bool) -> impl Iterator { CONNECTION_METRICS.iter().map(|m| &m.active_connections) } - pub fn collect_histograms() -> impl Iterator> { - [ + pub fn collect_histograms( + is_enterprise: bool, + ) -> impl Iterator> { + static E_HISTOGRAMS: &[&AtomicHistogram<12>] = &[ &MESSAGE_INGESTION_TIME, &MESSAGE_INDEX_TIME, &MESSAGE_DELIVERY_TIME, @@ -270,9 +273,22 @@ impl Collector { &STORE_BLOB_READ_TIME, &STORE_BLOB_WRITE_TIME, &DNS_LOOKUP_TIME, - ] - .into_iter() + ]; + static C_HISTOGRAMS: &[&AtomicHistogram<12>] = &[ + &MESSAGE_DELIVERY_TIME, + &MESSAGE_INCOMING_SIZE, + &MESSAGE_SUBMISSION_SIZE, + ]; + + if is_enterprise { + E_HISTOGRAMS + } else { + C_HISTOGRAMS + } + .iter() + .copied() .chain(CONNECTION_METRICS.iter().map(|m| &m.elapsed)) + .filter(|h| h.is_active()) } } @@ -285,8 +301,8 @@ impl EventCounter { self.description } - pub fn value(&self) -> u32 { - self.value + pub fn value(&self) -> u64 { + self.value as u64 } } diff --git a/tests/resources/jmap/email_set/headers.eml b/tests/resources/jmap/email_set/headers.eml index 53ea559d..e6c688ba 100644 --- a/tests/resources/jmap/email_set/headers.eml +++ b/tests/resources/jmap/email_set/headers.eml @@ -30,6 +30,7 @@ X-References: <1234@local.machine.example> <3456@example.net> X-References: <789@local.machine.example> X-Text: a b X-Text: this is some text +MIME-Version: 1.0 Content-Type: multipart/alternative; boundary="boundary_0" diff --git a/tests/resources/jmap/email_set/headers.jmap b/tests/resources/jmap/email_set/headers.jmap index ebc2f3ea..58fd5ebd 100644 --- a/tests/resources/jmap/email_set/headers.jmap +++ b/tests/resources/jmap/email_set/headers.jmap @@ -162,6 +162,10 @@ "name": "X-Text", "value": " this is some text" }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " multipart/alternative; \r\n\tboundary=\"boundary_0\"" diff --git a/tests/resources/jmap/email_set/minimal.eml b/tests/resources/jmap/email_set/minimal.eml index a9b4cb9d..ab86676b 100644 --- a/tests/resources/jmap/email_set/minimal.eml +++ b/tests/resources/jmap/email_set/minimal.eml @@ -1,5 +1,6 @@ Date: Tue, 10 Jul 2018 01:03:11 +0000 Message-ID: +MIME-Version: 1.0 Content-Type: text/plain; charset="utf-8" Content-Transfer-Encoding: 7bit diff --git a/tests/resources/jmap/email_set/minimal.jmap b/tests/resources/jmap/email_set/minimal.jmap index 55ec7351..1fbd9f24 100644 --- a/tests/resources/jmap/email_set/minimal.jmap +++ b/tests/resources/jmap/email_set/minimal.jmap @@ -21,6 +21,10 @@ "name": "Message-ID", "value": " " }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " text/plain; charset=\"utf-8\"" @@ -54,6 +58,10 @@ "name": "Message-ID", "value": " " }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " text/plain; charset=\"utf-8\"" @@ -81,6 +89,10 @@ "name": "Message-ID", "value": " " }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " text/plain; charset=\"utf-8\"" diff --git a/tests/resources/jmap/email_set/mixed.eml b/tests/resources/jmap/email_set/mixed.eml index a4864272..115cd2ad 100644 --- a/tests/resources/jmap/email_set/mixed.eml +++ b/tests/resources/jmap/email_set/mixed.eml @@ -5,6 +5,7 @@ Subject: =?utf-8?Q?Why_not_both_importing_AND_exporting=3F_=E2=98=BA?= To: "Colleagues": "James Smythe" ; "Friends": , "=?utf-8?Q?John_Sm=C3=AEth?=" +MIME-Version: 1.0 Content-Type: multipart/mixed; boundary="boundary_0" diff --git a/tests/resources/jmap/email_set/mixed.jmap b/tests/resources/jmap/email_set/mixed.jmap index b6063726..723cc1ec 100644 --- a/tests/resources/jmap/email_set/mixed.jmap +++ b/tests/resources/jmap/email_set/mixed.jmap @@ -55,6 +55,10 @@ "name": "To", "value": " \"Colleagues\": \"James Smythe\" ; \r\n\t\"Friends\": , \r\n\t\"=?utf-8?Q?John_Sm=C3=AEth?=\" " }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " multipart/mixed; \r\n\tboundary=\"boundary_0\"" diff --git a/tests/resources/jmap/email_set/nested_body.eml b/tests/resources/jmap/email_set/nested_body.eml index 5d5b2b68..5944a62c 100644 --- a/tests/resources/jmap/email_set/nested_body.eml +++ b/tests/resources/jmap/email_set/nested_body.eml @@ -2,6 +2,7 @@ Date: Tue, 10 Jul 2018 01:03:11 +0000 From: "Joe Bloggs" Message-ID: Subject: RFC 8621 Section 4.1.4 test +MIME-Version: 1.0 Content-Type: multipart/mixed; boundary="boundary_0" diff --git a/tests/resources/jmap/email_set/nested_body.jmap b/tests/resources/jmap/email_set/nested_body.jmap index 501b0964..1d558270 100644 --- a/tests/resources/jmap/email_set/nested_body.jmap +++ b/tests/resources/jmap/email_set/nested_body.jmap @@ -36,6 +36,10 @@ "name": "Subject", "value": " RFC 8621 Section 4.1.4 test" }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " multipart/mixed; \r\n\tboundary=\"boundary_0\"" diff --git a/tests/resources/jmap/email_set/rfc8621_1.eml b/tests/resources/jmap/email_set/rfc8621_1.eml index bc05cf0a..2061a2ec 100644 --- a/tests/resources/jmap/email_set/rfc8621_1.eml +++ b/tests/resources/jmap/email_set/rfc8621_1.eml @@ -2,6 +2,7 @@ Date: Tue, 10 Jul 2018 01:03:11 +0000 From: "Joe Bloggs" Message-ID: Subject: World domination +MIME-Version: 1.0 Content-Language: en Content-Type: text/plain; charset="utf-8" Content-Transfer-Encoding: 7bit diff --git a/tests/resources/jmap/email_set/rfc8621_1.jmap b/tests/resources/jmap/email_set/rfc8621_1.jmap index 76d586f0..adda9a19 100644 --- a/tests/resources/jmap/email_set/rfc8621_1.jmap +++ b/tests/resources/jmap/email_set/rfc8621_1.jmap @@ -39,6 +39,10 @@ "name": "Subject", "value": " World domination" }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Language", "value": " en" @@ -87,6 +91,10 @@ "name": "Subject", "value": " World domination" }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Language", "value": " en" @@ -129,6 +137,10 @@ "name": "Subject", "value": " World domination" }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Language", "value": " en" diff --git a/tests/resources/jmap/email_set/rfc8621_2.eml b/tests/resources/jmap/email_set/rfc8621_2.eml index 3df3066d..2482b4c9 100644 --- a/tests/resources/jmap/email_set/rfc8621_2.eml +++ b/tests/resources/jmap/email_set/rfc8621_2.eml @@ -3,6 +3,7 @@ From: "Joe Bloggs" Message-ID: Subject: World domination To: "John" +MIME-Version: 1.0 Content-Type: multipart/alternative; boundary="boundary_0" diff --git a/tests/resources/jmap/email_set/rfc8621_2.jmap b/tests/resources/jmap/email_set/rfc8621_2.jmap index 49e7bc77..9851e418 100644 --- a/tests/resources/jmap/email_set/rfc8621_2.jmap +++ b/tests/resources/jmap/email_set/rfc8621_2.jmap @@ -46,6 +46,10 @@ "name": "To", "value": " \"John\" " }, + { + "name": "MIME-Version", + "value": " 1.0" + }, { "name": "Content-Type", "value": " multipart/alternative; \r\n\tboundary=\"boundary_0\"" diff --git a/tests/resources/otel/otel-collector-config.yaml b/tests/resources/otel/otel-collector-config.yaml index 5c6d4927..4c320c21 100644 --- a/tests/resources/otel/otel-collector-config.yaml +++ b/tests/resources/otel/otel-collector-config.yaml @@ -1,4 +1,4 @@ -# docker run -p 4317:4317 --network host --rm -v $(pwd)/otel-collector-config.yaml:/etc/otelcol/config.yaml otel/opentelemetry-collector +# docker run -p 4317:4317 --network host --rm -v $(pwd)/tests/resources/otel/otel-collector-config.yaml:/etc/otelcol/config.yaml otel/opentelemetry-collector receivers: otlp: @@ -38,3 +38,7 @@ service: receivers: [otlp] processors: [batch] exporters: [debug] + metrics: + receivers: [otlp] + processors: [batch] + exporters: [debug] diff --git a/tests/resources/otel/stalwart-config.toml b/tests/resources/otel/stalwart-config.toml index 4824833a..99f8df2c 100644 --- a/tests/resources/otel/stalwart-config.toml +++ b/tests/resources/otel/stalwart-config.toml @@ -1,6 +1,8 @@ -[tracer.otel] -type = "otel" -transport = "grpc" -endpoint = "http://127.0.0.1:4317" -level = "trace" - +tracer.otel.type = "otel" +tracer.otel.transport = "grpc" +tracer.otel.endpoint = "http://127.0.0.1:4317" +tracer.otel.level = "trace" +metrics.open-telemetry.interval = "10s" +metrics.open-telemetry.enable = true +metrics.open-telemetry.endpoint = "http://127.0.0.1:4317" +metrics.open-telemetry.transport = "grpc" diff --git a/tests/src/jmap/delivery.rs b/tests/src/jmap/delivery.rs index 4124f758..8cf3bad0 100644 --- a/tests/src/jmap/delivery.rs +++ b/tests/src/jmap/delivery.rs @@ -307,10 +307,10 @@ pub async fn test(params: &mut JMAPTest) { // Check webhook events params.webhook.assert_contains(&[ - "store.ingest", + "message-ingest.", "delivery.dsn", "\"from\": \"bill@example.com\"", - "\"to\": \"john.doe@example.com\"", + "\"john.doe@example.com\"", ]); } diff --git a/tests/src/jmap/mod.rs b/tests/src/jmap/mod.rs index 10e86e33..6ebb3863 100644 --- a/tests/src/jmap/mod.rs +++ b/tests/src/jmap/mod.rs @@ -20,7 +20,7 @@ use common::{ }; use hyper::{header::AUTHORIZATION, Method}; use imap::core::{ImapSessionManager, IMAP}; -use jmap::{api::JmapSessionManager, services::housekeeper::Event, JMAP}; +use jmap::{api::JmapSessionManager, JMAP}; use jmap_client::client::{Client, Credentials}; use jmap_proto::{error::request::RequestError, types::id::Id}; use managesieve::core::ManageSieveSessionManager; @@ -31,11 +31,11 @@ use smtp::core::{SmtpSessionManager, SMTP}; use store::{ roaring::RoaringBitmap, - write::{key::DeserializeBigEndian, AnyKey}, - IterateParams, Stores, SUBSPACE_PROPERTY, + write::{key::DeserializeBigEndian, AnyKey, FtsQueueClass, ValueClass}, + IterateParams, Stores, ValueKey, SUBSPACE_PROPERTY, }; use tokio::sync::{mpsc, watch}; -use utils::config::Config; +use utils::{config::Config, BlobHash}; use webhooks::{spawn_mock_webhook_endpoint, MockWebhookEndpoint}; use crate::{add_test_certs, directory::DirectoryStore, store::TempDir, AssertConfig}; @@ -283,7 +283,7 @@ disabled-events = ["network.*"] [webhook."test"] url = "http://127.0.0.1:8821/hook" -events = ["auth.*", "delivery.dsn*", "store.ingest"] +events = ["auth.*", "delivery.dsn*", "message-ingest.*"] signature-key = "ovos-moles" throttle = "100ms" @@ -356,15 +356,44 @@ pub struct JMAPTest { pub async fn wait_for_index(server: &JMAP) { loop { - let (tx, rx) = tokio::sync::oneshot::channel(); + let mut has_index_tasks = false; server - .inner - .housekeeper_tx - .send(Event::IndexIsActive(tx)) + .core + .storage + .data + .iterate( + IterateParams::new( + ValueKey::> { + account_id: 0, + collection: 0, + document_id: 0, + class: ValueClass::FtsQueue(FtsQueueClass { + seq: 0, + hash: BlobHash::default(), + }), + }, + ValueKey::> { + account_id: u32::MAX, + collection: u8::MAX, + document_id: u32::MAX, + class: ValueClass::FtsQueue(FtsQueueClass { + seq: u64::MAX, + hash: BlobHash::default(), + }), + }, + ) + .ascending(), + |_, _| { + has_index_tasks = true; + + Ok(false) + }, + ) .await .unwrap(); - if rx.await.unwrap() { - tokio::time::sleep(Duration::from_millis(100)).await; + + if has_index_tasks { + tokio::time::sleep(Duration::from_millis(300)).await; } else { break; }