diff --git a/Cargo.lock b/Cargo.lock index 6a9f273..20cc7c3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7780,6 +7780,7 @@ dependencies = [ "mail-builder 1.0.0", "mail-parser", "memory-stats", + "nlp", "p256", "psl", "registry", diff --git a/crates/common/src/config/telemetry.rs b/crates/common/src/config/telemetry.rs index 8cedbcd..7f39ae5 100644 --- a/crates/common/src/config/telemetry.rs +++ b/crates/common/src/config/telemetry.rs @@ -40,6 +40,15 @@ pub enum TelemetrySubscriberType { Webhook(WebhookTracer), #[cfg(unix)] JournalTracer(crate::telemetry::tracers::journald::Subscriber), + // inbuxa: MON-10: trace history + StoreTracer(StoreTracer), +} + +/// Where trace history goes: traces to `tracing`, index tasks to `data`. +#[derive(Debug)] +pub struct StoreTracer { + pub tracing: store::Store, + pub data: store::Store, } #[derive(Debug)] @@ -141,8 +150,7 @@ impl Telemetry { } impl Tracers { - // inbuxa: `_storage` is unused until monitoring history (stored traces and metrics) is rebuilt - pub async fn parse(bp: &mut Bootstrap, _storage: &Storage) -> Self { + pub async fn parse(bp: &mut Bootstrap, storage: &Storage) -> Self { let mut custom_levels = AHashMap::new(); let mut tracers: Vec = Vec::new(); let mut global_interests = Interests::default(); @@ -387,6 +395,7 @@ impl Tracers { TelemetrySubscriberType::JournalTracer(_) => { EventType::Telemetry(TelemetryEvent::JournalError).into() } + TelemetrySubscriberType::StoreTracer(_) => None, }; // Parse disabled events @@ -478,6 +487,36 @@ impl Tracers { } } + // inbuxa: MON-10 to MON-12: trace history, when a tracing store is set: + // info and above, the span edges and MAIL FROM, never raw I/O + if !storage.tracing.is_none() { + let mut interests = Interests::default(); + for event_type in EventType::variants() { + let event_level = custom_levels + .get(event_type) + .copied() + .unwrap_or(event_type.level()); + if !event_type.is_raw_io() + && (Level::Info.is_contained(event_level) + || event_type.is_span_start() + || event_type.is_span_end() + || event_type.as_str().starts_with("smtp.mail-from")) + { + interests.set(event_type.to_id() as usize); + global_interests.set(event_type.to_id() as usize); + } + } + tracers.push(TelemetrySubscriber { + id: "trace-history".to_string(), + interests, + typ: TelemetrySubscriberType::StoreTracer(StoreTracer { + tracing: storage.tracing.clone(), + data: storage.data.clone(), + }), + lossy: true, + }); + } + #[cfg(feature = "dev_mode")] if let Ok(level) = std::env::var("LOG") { let level = Level::from_str(&level).expect("Invalid LOG level"); diff --git a/crates/common/src/telemetry/mod.rs b/crates/common/src/telemetry/mod.rs index 8052163..1f96589 100644 --- a/crates/common/src/telemetry/mod.rs +++ b/crates/common/src/telemetry/mod.rs @@ -102,6 +102,10 @@ impl TelemetrySubscriberType { TelemetrySubscriberType::LogTracer(settings) => spawn_log_tracer(builder, settings), TelemetrySubscriberType::Webhook(settings) => spawn_webhook_tracer(builder, settings), TelemetrySubscriberType::OtelTracer(settings) => spawn_otel_tracer(builder, settings), + // inbuxa: MON-10: trace history + TelemetrySubscriberType::StoreTracer(settings) => { + tracers::store::spawn_store_tracer(builder, settings.tracing, settings.data) + } #[cfg(unix)] TelemetrySubscriberType::JournalTracer(subscriber) => { tracers::journald::spawn_journald_tracer(builder, subscriber) diff --git a/crates/common/src/telemetry/tracers/mod.rs b/crates/common/src/telemetry/tracers/mod.rs index 3461756..570df64 100644 --- a/crates/common/src/telemetry/tracers/mod.rs +++ b/crates/common/src/telemetry/tracers/mod.rs @@ -9,6 +9,7 @@ pub mod journald; pub mod log; pub mod otel; pub mod stdout; +pub mod store; // inbuxa: monitoring history (MON-10 to MON-17) use registry::{ diff --git a/crates/common/src/telemetry/tracers/store.rs b/crates/common/src/telemetry/tracers/store.rs new file mode 100644 index 0000000..a9c81e5 --- /dev/null +++ b/crates/common/src/telemetry/tracers/store.rs @@ -0,0 +1,228 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Trace history (monitoring spec MON-10 to MON-17, MON-34). A lossy +//! collector subscriber gathers each inbound SMTP session and delivery +//! attempt, and writes it once, when the span closes, as an `x:Trace` in the +//! registry's own encoding under `TelemetryClass::Span(span_id)`. + +use crate::telemetry::tracers::TraceEvents; +use ahash::AHashMap; +use registry::{ + pickle::PickledStream, + schema::{ + prelude::{ObjectInner, ObjectType}, + structs::{ + Task, TaskIndexTrace, TaskStatus, Trace, TraceKeyValue, TraceValue, + TraceValueString, TraceValueUnsignedInt, + }, + }, +}; +use std::{future::Future, sync::Arc, time::Duration}; +use store::{ + SearchStore, Store, ValueKey, + search::{SearchFilter, SearchQuery}, + write::{BatchBuilder, SearchIndex, TelemetryClass, ValueClass, now}, +}; +use trc::{ + AddContext, DeliveryEvent, Event, EventDetails, EventType, Key, Level, SmtpEvent, + ipc::subscriber::SubscriberBuilder, +}; +use utils::snowflake::SnowflakeIdGenerator; + +/// Events kept per trace (MON-15). +pub const MAX_EVENTS: usize = 1000; +/// The longest string value kept (MON-15). +pub const MAX_STRING: usize = 4096; +/// A span still open after this is dropped (MON-13). +const SPAN_MAX_HOLD: u64 = 86_400; + +pub trait TracingStore: Sync + Send { + /// Deletes traces older than `keep`, and their search documents + /// (MON-17). + fn purge_spans( + &self, + keep: Duration, + search: Option<&SearchStore>, + ) -> impl Future> + Send; +} + +impl TracingStore for Store { + async fn purge_spans(&self, keep: Duration, search: Option<&SearchStore>) -> trc::Result<()> { + let Some(until) = SnowflakeIdGenerator::from_duration(keep) else { + return Ok(()); + }; + self.delete_range( + ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(0))), + ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(until))), + ) + .await + .caused_by(trc::location!())?; + if let Some(search) = search { + search + .unindex( + SearchQuery::new(SearchIndex::Tracing) + .with_filter(SearchFilter::lt(store::search::SearchField::Id, until)), + ) + .await + .caused_by(trc::location!())?; + } + Ok(()) + } +} + +/// Decodes a stored trace; `None` for records in any other encoding. +pub fn decode_trace(bytes: &[u8]) -> Option { + PickledStream::new(bytes) + .and_then(|mut stream| ObjectInner::unpickle(ObjectType::Trace, &mut stream)) + .and_then(|inner| match inner { + ObjectInner::Trace(trace) => Some(trace), + _ => None, + }) +} + +/// A stored trace as read by key: `None` when it can't be decoded. +pub struct MaybeTrace(pub Option); + +impl store::Deserialize for MaybeTrace { + fn deserialize(bytes: &[u8]) -> trc::Result { + Ok(MaybeTrace(decode_trace(bytes))) + } +} + +fn is_stored_span(event: EventType) -> bool { + matches!( + event, + EventType::Smtp(SmtpEvent::ConnectionStart) | EventType::Delivery(DeliveryEvent::AttemptStart) + ) +} + +fn is_mail_from(event: EventType) -> bool { + event.as_str().starts_with("smtp.mail-from") || event == EventType::Smtp(SmtpEvent::MultipleMailFrom) +} + +struct Span { + started: u64, + is_smtp: bool, + has_mail_from: bool, + events: Vec>>, + cut: usize, +} + +fn truncate_values_list(values: &mut registry::types::list::List) { + for kv in values.values_mut() { + match &mut kv.value { + TraceValue::String(TraceValueString { value }) if value.len() > MAX_STRING => { + let mut end = MAX_STRING; + while !value.is_char_boundary(end) { + end -= 1; + } + value.truncate(end); + } + TraceValue::Event(event) => truncate_values_list(&mut event.value), + _ => {} + } + } +} + +/// The trace a closed span leaves (MON-12, MON-15). +fn build_trace(span: &Span) -> Trace { + let mut trace = Trace::from_events(span.events.iter().map(|e| e.as_ref()), span.events.len()); + for event in trace.events.values_mut() { + truncate_values_list(&mut event.key_values); + } + if span.cut > 0 + && let Some(last) = trace.events.values_mut().last() + { + // The count of events cut rides on the closing event + last.key_values.push(TraceKeyValue { + key: Key::Total, + value: TraceValue::UnsignedInt(TraceValueUnsignedInt { + value: span.cut as u64, + }), + }); + } + trace +} + +/// Starts the subscriber that stores traces in `tracing`, scheduling their +/// indexing in `data` (MON-16). Lossy: a slow store loses history, never +/// delays mail (MON-34, MON-35). +pub(crate) fn spawn_store_tracer(builder: SubscriberBuilder, tracing: Store, data: Store) { + let (_, mut rx) = builder.register(); + tokio::spawn(async move { + let mut spans: AHashMap = AHashMap::new(); + while let Some(events) = rx.recv().await { + let mut closed = Vec::new(); + for event in events { + let typ = event.inner.typ; + let Some(span_id) = event.span_id() else { + continue; + }; + if is_stored_span(typ) { + spans.insert( + span_id, + Span { + started: event.inner.timestamp, + is_smtp: matches!(typ, EventType::Smtp(_)), + has_mail_from: false, + events: vec![event], + cut: 0, + }, + ); + continue; + } + let Some(span) = spans.get_mut(&span_id) else { + continue; + }; + if is_mail_from(typ) { + span.has_mail_from = true; + } + let is_end = typ.is_span_end(); + // MON-12: info and above, never raw I/O + if !typ.is_raw_io() && (is_end || event.inner.level as usize >= Level::Info as usize) { + if span.events.len() < MAX_EVENTS - 1 || is_end { + span.events.push(event); + } else { + span.cut += 1; + } + } + if is_end && let Some(span) = spans.remove(&span_id) { + // MON-11: a session that never reached MAIL FROM isn't kept + if !span.is_smtp || span.has_mail_from { + closed.push((span_id, span)); + } + } + } + + if !closed.is_empty() { + let mut batch = BatchBuilder::new(); + let mut tasks = BatchBuilder::new(); + for (span_id, span) in &closed { + batch.set( + ValueClass::Telemetry(TelemetryClass::Span(*span_id)), + ObjectInner::Trace(build_trace(span)).to_pickled_vec(), + ); + tasks.schedule_task(Task::IndexTrace(TaskIndexTrace { + trace_id: (*span_id).into(), + status: TaskStatus::now(), + })); + } + if let Err(err) = tracing.write(batch.build_all()).await { + trc::error!(err.details("Failed to store trace history")); + } else if let Err(err) = data.write(tasks.build_all()).await { + trc::error!(err.details("Failed to schedule trace indexing")); + } + } + + // MON-13: spans open for over a day are dropped + if spans.len() > 1000 { + let now = now(); + spans.retain(|_, span| now.saturating_sub(span.started) < SPAN_MAX_HOLD); + } + } + }); +} diff --git a/crates/jmap/src/inbuxa/telemetry.rs b/crates/jmap/src/inbuxa/telemetry.rs index 8e7097b..bcc8150 100644 --- a/crates/jmap/src/inbuxa/telemetry.rs +++ b/crates/jmap/src/inbuxa/telemetry.rs @@ -212,3 +212,303 @@ pub(crate) async fn metric_query( } Ok(response) } + +// ---- Traces (MON-10 to MON-17, MON-31, MON-32) ---- + +use common::telemetry::tracers::store::{MaybeTrace, decode_trace}; +use registry::schema::structs::{Search, Trace, TraceValue}; +use store::{ + IterateParams, + search::{SearchFilter, SearchQuery, TracingSearchField}, + write::{BatchBuilder, SearchIndex, key::DeserializeBigEndian}, +}; +use trc::{AddContext, EventType, Key}; + +/// The trace's server-set fields (MON-14): the first event's time, the +/// first `from`, every distinct `to`, and the message size or 0. +fn trace_to_value(id: u64, trace: Trace) -> JmapValue<'static> { + let mut from = None; + let mut to: Vec = Vec::new(); + let mut size = None; + let mut first_timestamp = None; + for event in trace.events.iter() { + first_timestamp.get_or_insert(event.timestamp.timestamp()); + for kv in event.key_values.iter() { + let mut texts = Vec::new(); + match &kv.value { + TraceValue::String(v) => texts.push(v.value.clone()), + TraceValue::List(list) => { + for item in list.value.iter() { + if let TraceValue::String(v) = item { + texts.push(v.value.clone()); + } + } + } + TraceValue::UnsignedInt(v) if kv.key == Key::Size => { + size.get_or_insert(v.value); + } + _ => {} + } + match kv.key { + Key::From => { + if from.is_none() { + from = texts.into_iter().next(); + } + } + Key::To => { + for text in texts { + if !to.contains(&text) { + to.push(text); + } + } + } + _ => {} + } + } + } + let timestamp = first_timestamp + .unwrap_or_else(|| SnowflakeIdGenerator::to_timestamp(id) as i64); + let mut value = trace.into_value(); + if let JmapValue::Object(obj) = &mut value { + obj.insert_unchecked( + Property::Timestamp, + JmapValue::Str(UTCDateTime::from_timestamp(timestamp).to_string().into()), + ); + obj.insert_unchecked( + Property::From, + match from { + Some(from) => JmapValue::Str(from.into()), + None => JmapValue::Null, + }, + ); + obj.insert_unchecked(Property::To, JmapValue::Str(to.join(", ").into())); + obj.insert_unchecked(Property::Size, JmapValue::Number(size.unwrap_or(0).into())); + } + value +} + +/// The oldest trace id still visible (MON-17). +async fn trace_floor(server: &common::Server) -> u64 { + match common::telemetry::metrics::store::retention(server) + .await + .hold_traces_for + { + Some(keep) => SnowflakeIdGenerator::from_duration(keep.into_inner()).unwrap_or(0), + None => 0, + } +} + +async fn read_trace(server: &common::Server, id: u64) -> trc::Result> { + if id < trace_floor(server).await { + return Ok(None); + } + Ok(server + .tracing_store() + .get_value::(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(id)))) + .await? + .and_then(|MaybeTrace(trace)| trace)) +} + +/// `x:Trace/get`. +pub(crate) async fn trace_get( + mut get: RegistryGetResponse<'_>, +) -> trc::Result> { + assert_server_level(get.access_token)?; + let server = get.server; + if server.tracing_store().is_none() { + if let Some(ids) = get.ids.take() { + for id in ids { + get.not_found(id); + } + } + return Ok(get); + } + let ids = match get.ids.take() { + Some(ids) => ids, + None => trace_ids(server, 0, u64::MAX, false, None) + .await? + .into_iter() + .take(server.core.jmap.get_max_objects) + .map(Id::from) + .collect(), + }; + for id in ids { + match read_trace(server, id.id()).await? { + Some(trace) => get.insert(id, trace_to_value(id.id(), trace)), + None => get.not_found(id), + } + } + Ok(get) +} + +/// Trace ids in a range, newest first unless `ascending`, optionally only +/// those whose opening event is `event`. +async fn trace_ids( + server: &common::Server, + from_id: u64, + to_id: u64, + ascending: bool, + event: Option, +) -> trc::Result> { + let from_id = from_id.max(trace_floor(server).await); + let mut ids = Vec::new(); + if from_id > to_id || server.tracing_store().is_none() { + return Ok(ids); + } + let params = IterateParams::new( + ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(from_id))), + ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(to_id))), + ); + let params = if ascending { + params.ascending() + } else { + params.descending() + }; + server + .tracing_store() + .iterate(params, |key, value| { + let id = key.deserialize_be_u64(0)?; + if let Some(trace) = decode_trace(value) + && event.is_none_or(|event| { + trace.events.iter().next().is_some_and(|first| first.event == event) + }) + { + ids.push(id); + } + Ok(true) + }) + .await + .caused_by(trc::location!())?; + Ok(ids) +} + +/// `x:Trace/query`: `event` (the opening event), `text` and `queueId` +/// (through the search index, refused when trace search is off), and the +/// timestamp comparisons. +pub(crate) async fn trace_query( + mut req: RegistryQueryResponse<'_>, +) -> trc::Result { + assert_server_level(req.access_token)?; + let (mut from_id, mut to_id) = (0u64, u64::MAX); + let mut event = None; + let mut search: Vec = Vec::new(); + req.request.extract_filters(|property, op, value| match property { + Property::Timestamp => { + let Some(ts) = timestamp_of(&value) else { + return false; + }; + let at = SnowflakeIdGenerator::first_id_at; + match op { + RegistryFilterOp::GreaterThan => from_id = from_id.max(at(ts + 1)), + RegistryFilterOp::GreaterEqualThan => from_id = from_id.max(at(ts)), + RegistryFilterOp::LowerThan => to_id = to_id.min(at(ts).saturating_sub(1)), + RegistryFilterOp::LowerEqualThan => { + to_id = to_id.min(at(ts + 1).saturating_sub(1)) + } + _ => return false, + } + true + } + Property::Event => match value.as_str().and_then(EventType::parse) { + Some(parsed) => { + event = Some(parsed); + true + } + None => false, + }, + Property::Text => match value.as_str() { + Some(text) => { + search.push(SearchFilter::has_text( + TracingSearchField::Keywords, + text.to_lowercase(), + nlp::language::Language::None, + )); + true + } + None => false, + }, + Property::QueueId => match value.as_str() { + Some(queue_id) => { + search.push(SearchFilter::eq(TracingSearchField::QueueId, queue_id.to_string())); + true + } + None => false, + }, + _ => false, + })?; + + let params = req + .request + .extract_parameters(req.server.core.jmap.query_max_results, None)?; + if !matches!(params.sort_by, Property::Timestamp | Property::Id) { + return Err(trc::JmapEvent::UnsupportedSort + .into_err() + .details("Traces sort by timestamp only")); + } + + let mut ids = trace_ids(req.server, from_id, to_id, params.sort_ascending, event).await?; + if !search.is_empty() { + let settings = req + .server + .registry() + .object::(Id::singleton()) + .await? + .unwrap_or_default(); + if !settings.index_telemetry { + return Err(trc::JmapEvent::UnsupportedFilter + .into_err() + .details("Trace search is off (indexTelemetry)")); + } + let found = req + .server + .search_store() + .query_global(SearchQuery::new(SearchIndex::Tracing).with_filters(search)) + .await?; + let found = found.into_iter().collect::>(); + ids.retain(|id| found.contains(id)); + } + + let mut response = QueryResponseBuilder::new( + ids.len(), + req.server.core.jmap.query_max_results, + State::Initial, + &req.request, + ); + for id in ids { + if !response.add_id(Id::from(id)) { + break; + } + } + Ok(response) +} + +/// `x:Trace/set`: create and update are refused; destroy removes the trace +/// and its search document (MON-32). +pub(crate) async fn trace_set( + mut set: crate::registry::mapping::RegistrySetResponse<'_>, +) -> trc::Result> { + assert_server_level(set.access_token)?; + set.fail_all_create("Traces cannot be created"); + set.fail_all_update("Traces cannot be modified"); + let server = set.server; + for id in std::mem::take(&mut set.destroy) { + if read_trace(server, id.id()).await?.is_none() { + set.response + .not_destroyed + .append(id, jmap_proto::error::set::SetError::not_found()); + continue; + } + let mut batch = BatchBuilder::new(); + batch.clear(ValueClass::Telemetry(TelemetryClass::Span(id.id()))); + server.tracing_store().write(batch.build_all()).await?; + server + .search_store() + .unindex( + SearchQuery::new(SearchIndex::Tracing) + .with_filter(SearchFilter::eq(store::search::SearchField::Id, id.id())), + ) + .await?; + set.response.destroyed.push(id); + } + Ok(set) +} diff --git a/crates/jmap/src/registry/get.rs b/crates/jmap/src/registry/get.rs index a8f2692..451c38f 100644 --- a/crates/jmap/src/registry/get.rs +++ b/crates/jmap/src/registry/get.rs @@ -381,6 +381,9 @@ impl RegistryGet for Server { } ObjectType::Log => log_get(get).await.map(|get| get.into_response()), // inbuxa: monitoring history (MON-17, MON-31) + ObjectType::Trace => crate::inbuxa::telemetry::trace_get(get) + .await + .map(|get| get.into_response()), ObjectType::Metric => crate::inbuxa::telemetry::metric_get(get) .await .map(|get| get.into_response()), diff --git a/crates/jmap/src/registry/query.rs b/crates/jmap/src/registry/query.rs index c0f9288..c873542 100644 --- a/crates/jmap/src/registry/query.rs +++ b/crates/jmap/src/registry/query.rs @@ -123,6 +123,14 @@ impl RegistryQuery for Server { .and_then(|response| response.build()), // inbuxa: monitoring history (MON-17, MON-31) + ObjectType::Trace => crate::inbuxa::telemetry::trace_query(RegistryQueryResponse { + server: self, + access_token, + object_type, + request, + }) + .await + .and_then(|response| response.build()), ObjectType::Metric => crate::inbuxa::telemetry::metric_query(RegistryQueryResponse { server: self, access_token, diff --git a/crates/jmap/src/registry/set.rs b/crates/jmap/src/registry/set.rs index 9a151cd..fabbb0a 100644 --- a/crates/jmap/src/registry/set.rs +++ b/crates/jmap/src/registry/set.rs @@ -881,7 +881,11 @@ impl RegistrySet for Server { .await .map(|set| set.into_response()), - ObjectType::Log | ObjectType::Metric | ObjectType::Trace | ObjectType::ClusterNode => { + // inbuxa: MON-32: a trace can be destroyed, never created or changed + ObjectType::Trace => crate::inbuxa::telemetry::trace_set(set) + .await + .map(|set| set.into_response()), + ObjectType::Log | ObjectType::Metric | ObjectType::ClusterNode => { set.fail_all_create("Telemetry objects cannot be created"); set.fail_all_update("Telemetry objects cannot be modified"); set.fail_all_destroy("Telemetry objects cannot be deleted"); diff --git a/crates/services/Cargo.toml b/crates/services/Cargo.toml index 9f3e714..b2374e2 100644 --- a/crates/services/Cargo.toml +++ b/crates/services/Cargo.toml @@ -5,6 +5,7 @@ edition = "2024" [dependencies] store = { path = "../store" } +nlp = { path = "../nlp" } common = { path = "../common" } utils = { path = "../utils" } trc = { path = "../trc" } diff --git a/crates/services/src/task_manager/index.rs b/crates/services/src/task_manager/index.rs index 492c822..603058e 100644 --- a/crates/services/src/task_manager/index.rs +++ b/crates/services/src/task_manager/index.rs @@ -582,8 +582,88 @@ async fn build_contact_document( #[cfg(not(feature = "enterprise"))] -async fn build_tracing_span_document(_: &Server, _: u64) -> trc::Result> { - Ok(None) +// inbuxa: MON-16: a trace's search document, when trace search is on: +// its event types, queue ids, and addresses, their domains, hosts, IPs, +// message ids and account names as keywords +async fn build_tracing_span_document( + server: &Server, + span_id: u64, +) -> trc::Result> { + use common::telemetry::tracers::store::MaybeTrace; + use registry::schema::{enums::SearchTracingField, structs::Search}; + use store::{ + search::TracingSearchField, + write::{TelemetryClass, ValueClass}, + }; + use trc::Key; + + let settings = server + .registry() + .object::(types::id::Id::singleton()) + .await? + .unwrap_or_default(); + if !settings.index_telemetry { + return Ok(None); + } + let wants = |field: SearchTracingField| settings.index_tracing_fields.iter().any(|f| *f == field); + let Some(MaybeTrace(Some(trace))) = server + .tracing_store() + .get_value::(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span( + span_id, + )))) + .await? + else { + return Ok(None); + }; + + let mut document = IndexDocument::new(SearchIndex::Tracing).with_id(span_id); + let mut seen = store::ahash::AHashSet::new(); + for event in trace.events.iter() { + if wants(SearchTracingField::EventType) && seen.insert(event.event.as_str().to_string()) { + document.index_keyword(TracingSearchField::EventType, event.event.as_str()); + } + for kv in event.key_values.iter() { + let text = match &kv.value { + registry::schema::structs::TraceValue::String(v) => v.value.clone(), + registry::schema::structs::TraceValue::UnsignedInt(v) => v.value.to_string(), + registry::schema::structs::TraceValue::IpAddr(v) => v.value.to_string(), + _ => continue, + }; + match kv.key { + Key::QueueId if wants(SearchTracingField::QueueId) => { + if seen.insert(format!("q:{text}")) { + document.index_keyword(TracingSearchField::QueueId, &text); + } + } + Key::From + | Key::To + | Key::Domain + | Key::Hostname + | Key::RemoteIp + | Key::MessageId + | Key::AccountName + if wants(SearchTracingField::Keywords) => + { + let text = text.to_lowercase(); + if seen.insert(format!("k:{text}")) { + document.index_text(TracingSearchField::Keywords, &text, nlp::language::Language::None); + // An address's domain, so a domain search finds it + if let Some((_, domain)) = text.rsplit_once('@') + && seen.insert(format!("k:{domain}")) + { + document.index_text( + TracingSearchField::Keywords, + domain, + nlp::language::Language::None, + ); + } + } + } + _ => {} + } + } + } + Ok(Some(document)) } // inbuxa: UD-1, UD-4: archives a deleted file, event or contact noted at diff --git a/crates/services/src/task_manager/maintenance.rs b/crates/services/src/task_manager/maintenance.rs index 743b036..fc82e19 100644 --- a/crates/services/src/task_manager/maintenance.rs +++ b/crates/services/src/task_manager/maintenance.rs @@ -248,6 +248,18 @@ async fn store_maintenance( { trc::error!(err.details("Failed to purge metric history")); } + if let Some(keep) = retention.hold_traces_for + && !server.tracing_store().is_none() + { + use common::telemetry::tracers::store::TracingStore; + if let Err(err) = server + .tracing_store() + .purge_spans(keep.into_inner(), Some(server.search_store())) + .await + { + trc::error!(err.details("Failed to purge trace history")); + } + } trc::event!( Store(StoreEvent::DataStorePurged), diff --git a/tests/src/telemetry/mod.rs b/tests/src/telemetry/mod.rs index b403067..bf5b97c 100644 --- a/tests/src/telemetry/mod.rs +++ b/tests/src/telemetry/mod.rs @@ -7,9 +7,7 @@ #[cfg(feature = "pending-rebuild")] // inbuxa: pending-rebuild, see docs/spec/features/ pub mod alerts; pub mod metrics; -#[cfg(feature = "pending-rebuild")] // inbuxa: pending-rebuild, see docs/spec/features/ pub mod tracing; -#[cfg(feature = "pending-rebuild")] // inbuxa: pending-rebuild, see docs/spec/features/ pub mod webhooks; use crate::utils::server::TestServerBuilder; @@ -61,9 +59,7 @@ pub async fn telemetry_tests() { #[cfg(feature = "pending-rebuild")] alerts::test(&test).await; metrics::test(&test).await; - #[cfg(feature = "pending-rebuild")] tracing::test(&test).await; - #[cfg(feature = "pending-rebuild")] webhooks::test(&test).await; if test.is_reset() {