diff --git a/crates/jmap/src/inbuxa/telemetry.rs b/crates/jmap/src/inbuxa/telemetry.rs index bcc8150..ad79ecf 100644 --- a/crates/jmap/src/inbuxa/telemetry.rs +++ b/crates/jmap/src/inbuxa/telemetry.rs @@ -427,9 +427,24 @@ pub(crate) async fn trace_query( } None => false, }, - Property::QueueId => match value.as_str() { + // The queue id column is an integer on every search backend, and + // holds a trace's first queue id; the keywords carry all of them + Property::QueueId => match value + .as_str() + .and_then(|v| v.trim().parse::().ok()) + .or_else(|| value.as_u64()) + { Some(queue_id) => { - search.push(SearchFilter::eq(TracingSearchField::QueueId, queue_id.to_string())); + search.extend([ + SearchFilter::Or, + SearchFilter::eq(TracingSearchField::QueueId, queue_id), + SearchFilter::has_text( + TracingSearchField::Keywords, + queue_id.to_string(), + nlp::language::Language::None, + ), + SearchFilter::End, + ]); true } None => false, diff --git a/crates/services/src/task_manager/index.rs b/crates/services/src/task_manager/index.rs index 8942980..95296d1 100644 --- a/crates/services/src/task_manager/index.rs +++ b/crates/services/src/task_manager/index.rs @@ -567,20 +567,14 @@ async fn build_contact_document( } -// 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 +// inbuxa: MON-16: a trace's search document, when trace search is on 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; + use registry::schema::structs::Search; + use store::write::{TelemetryClass, ValueClass}; let settings = server .registry() @@ -590,7 +584,6 @@ async fn build_tracing_span_document( 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( @@ -601,23 +594,67 @@ async fn build_tracing_span_document( return Ok(None); }; + Ok(Some(trace_search_document( + span_id, + &trace, + &settings + .index_tracing_fields + .iter() + .copied() + .collect::>(), + ))) +} + +/// inbuxa: MON-16: the search document for a stored trace. +/// +/// The event type and queue id columns are integers on every search backend +/// (BIGINT on PostgreSQL and MySQL, long on Elasticsearch), and each holds a +/// single value per trace: the event type is the trace's opening event, the +/// one `x:Trace/query` filters on, and the queue id is the first queue id the +/// trace mentions. Every queue id also goes into the keywords, so a session +/// that queued several messages is found by any of them. +pub fn trace_search_document( + span_id: u64, + trace: ®istry::schema::structs::Trace, + fields: &[registry::schema::enums::SearchTracingField], +) -> IndexDocument { + use registry::schema::{enums::SearchTracingField, structs::TraceValue}; + use store::search::TracingSearchField; + use trc::Key; + + let wants = |field: SearchTracingField| fields.contains(&field); let mut document = IndexDocument::new(SearchIndex::Tracing).with_id(span_id); + if wants(SearchTracingField::EventType) + && let Some(first) = trace.events.iter().next() + { + document.index_unsigned(TracingSearchField::EventType, first.event.to_id() as u64); + } + let mut seen = store::ahash::AHashSet::new(); + let mut queue_id_indexed = false; 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(), + TraceValue::String(v) => v.value.clone(), + TraceValue::UnsignedInt(v) => v.value.to_string(), + 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::QueueId => { + let Ok(queue_id) = text.parse::() else { + continue; + }; + if wants(SearchTracingField::QueueId) && !queue_id_indexed { + document.index_unsigned(TracingSearchField::QueueId, queue_id); + queue_id_indexed = true; + } + if wants(SearchTracingField::Keywords) && seen.insert(format!("k:{text}")) { + document.index_text( + TracingSearchField::Keywords, + &text, + nlp::language::Language::None, + ); } } Key::From @@ -648,7 +685,7 @@ async fn build_tracing_span_document( } } } - Ok(Some(document)) + document } // inbuxa: UD-1, UD-4: archives a deleted file, event or contact noted at diff --git a/docs/spec/features/monitoring.md b/docs/spec/features/monitoring.md index 65ace3e..26913ce 100644 --- a/docs/spec/features/monitoring.md +++ b/docs/spec/features/monitoring.md @@ -212,10 +212,15 @@ unchanged. - **MON-16.** With `indexTelemetry` on, storing a trace schedules an `IndexTrace` task. The task builds one document for `SearchIndex::Tracing` with the fields named in `indexTracingFields`: - - `eventType`: every event type in the trace; - - `queueId`: every `queueId` value; + - `eventType`: the trace's opening event, as its numeric id; + - `queueId`: the first `queueId` value, as an integer; - `keywords`: every address in `from` and `to`, each address's domain, every - `domain`, `hostname`, `remoteIp`, `messageId` and `accountName` value. + `domain`, `hostname`, `remoteIp`, `messageId` and `accountName` value, + and every `queueId` value. + The event type and queue id are single integer columns on every search + backend (BIGINT on PostgreSQL and MySQL), so the `queueId` filter matches + the column or any queue id in the keywords, and a session that queued + several messages is found by each of them. So searching `example.org` finds every trace to or from that domain, as the upstream suite expects. With `indexTelemetry` off nothing is indexed, and the `text` and `queueId` filters are refused (see "Interfaces"). diff --git a/tests/src/store/query.rs b/tests/src/store/query.rs index 1115fb1..9556d43 100644 --- a/tests/src/store/query.rs +++ b/tests/src/store/query.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use crate::{store::deflate_test_resource, utils::server::TestServer}; @@ -122,6 +124,10 @@ pub async fn test(test: &TestServer) { println!("Running global id filtering tests..."); test_global(store.clone()).await; + // inbuxa: trace documents as the index task builds them + println!("Running trace document tests..."); + test_trace_documents(store.clone()).await; + // Large document insert test println!("Running large document insert tests..."); let mut large_text = String::with_capacity(20 * 1024 * 1024); @@ -809,3 +815,160 @@ async fn test_global(store: SearchStore) { AHashSet::from_iter([3, 4, 5]) ); } + +// inbuxa: MON-16: documents built by the index task from stored traces go +// into every search backend (the SQL backends type etyp and qid as BIGINT) +// and are found again by queue id and keyword. +async fn test_trace_documents(store: SearchStore) { + use registry::schema::{ + enums::SearchTracingField, + structs::{ + Trace, TraceEvent, TraceKeyValue, TraceValue, TraceValueString, + TraceValueUnsignedInt, + }, + }; + use services::task_manager::index::trace_search_document; + use trc::{DeliveryEvent, EventType, Key, SmtpEvent}; + + let kv_u = |key: Key, value: u64| TraceKeyValue { + key, + value: TraceValue::UnsignedInt(TraceValueUnsignedInt { value }), + }; + let kv_s = |key: Key, value: &str| TraceKeyValue { + key, + value: TraceValue::String(TraceValueString { + value: value.to_string(), + }), + }; + let event = |event: EventType, key_values: Vec| TraceEvent { + event, + key_values: key_values.into(), + ..Default::default() + }; + let fields = [ + SearchTracingField::EventType, + SearchTracingField::QueueId, + SearchTracingField::Keywords, + ]; + + // An SMTP session that queued two messages, and a delivery attempt + let session = Trace { + events: vec![ + event( + EventType::Smtp(SmtpEvent::ConnectionStart), + vec![kv_s(Key::RemoteIp, "192.0.2.7")], + ), + event( + EventType::Smtp(SmtpEvent::MailFrom), + vec![kv_s(Key::From, "Sender@Example.org")], + ), + event( + EventType::Smtp(SmtpEvent::RcptTo), + vec![kv_u(Key::QueueId, 9_000_000_001), kv_s(Key::To, "rcpt@example.net")], + ), + event( + EventType::Smtp(SmtpEvent::RcptTo), + vec![kv_u(Key::QueueId, 9_000_000_002)], + ), + ] + .into(), + }; + let delivery = Trace { + events: vec![event( + EventType::Delivery(DeliveryEvent::AttemptStart), + vec![kv_u(Key::QueueId, 9_000_000_003), kv_s(Key::Hostname, "relay.example.net")], + )] + .into(), + }; + let documents = vec![ + trace_search_document(100, &session, &fields), + trace_search_document(101, &delivery, &fields), + ]; + assert!( + documents + .iter() + .all(|d| d.has_field(&SearchField::Tracing(TracingSearchField::QueueId)) + && d.has_field(&SearchField::Tracing(TracingSearchField::EventType))), + "trace documents carry a queue id and an event type" + ); + store.index(documents).await.unwrap(); + if let SearchStore::ElasticSearch(store) = &store { + store.refresh_index(SearchIndex::Tracing).await.unwrap(); + } + + let query = |filters: Vec| { + let store = store.clone(); + async move { + store + .query_global( + SearchQuery::new(SearchIndex::Tracing) + .with_filter(SearchFilter::ge(SearchField::Id, 100u64)) + .with_filters(filters), + ) + .await + .unwrap() + .into_iter() + .collect::>() + } + }; + // By queue id, the way x:Trace/query asks: the queue id column, or any + // queue id in the keywords + let by_queue_id = |queue_id: u64| { + vec![ + SearchFilter::Or, + SearchFilter::eq(TracingSearchField::QueueId, queue_id), + SearchFilter::has_text( + TracingSearchField::Keywords, + queue_id.to_string(), + Language::None, + ), + SearchFilter::End, + ] + }; + assert_eq!(query(by_queue_id(9_000_000_001)).await, AHashSet::from_iter([100])); + assert_eq!(query(by_queue_id(9_000_000_002)).await, AHashSet::from_iter([100])); + assert_eq!(query(by_queue_id(9_000_000_003)).await, AHashSet::from_iter([101])); + assert_eq!(query(by_queue_id(9_000_000_004)).await, AHashSet::new()); + assert_eq!( + query(vec![SearchFilter::eq(TracingSearchField::QueueId, 9_000_000_003u64)]).await, + AHashSet::from_iter([101]) + ); + // By opening event type + assert_eq!( + query(vec![SearchFilter::eq( + TracingSearchField::EventType, + EventType::Delivery(DeliveryEvent::AttemptStart).to_id() as u64, + )]) + .await, + AHashSet::from_iter([101]) + ); + // By keyword: an address, lowercased, and its domain + assert_eq!( + query(vec![SearchFilter::has_text( + TracingSearchField::Keywords, + "example.org", + Language::None, + )]) + .await, + AHashSet::from_iter([100]) + ); + assert_eq!( + query(vec![SearchFilter::has_text( + TracingSearchField::Keywords, + "relay.example.net", + Language::None, + )]) + .await, + AHashSet::from_iter([101]) + ); + + for id in [100u64, 101] { + store + .unindex( + SearchQuery::new(SearchIndex::Tracing) + .with_filter(SearchFilter::eq(SearchField::Id, id)), + ) + .await + .unwrap(); + } +} diff --git a/tests/src/system/monitoring.rs b/tests/src/system/monitoring.rs index e8e2fa4..365c82f 100644 --- a/tests/src/system/monitoring.rs +++ b/tests/src/system/monitoring.rs @@ -148,6 +148,73 @@ pub async fn test(test: &mut TestServer) { "test 9: to" ); + // MON-16: the queueId filter finds the traces that name a queue id (the + // session that queued the message and its delivery attempt) through the + // search index, given as a string or a number (the index column is an + // integer) + fn queue_ids(value: &Value, out: &mut Vec) { + match value { + Value::Object(map) => { + if map.get("key").and_then(|k| k.as_str()) == Some("queueId") + && let Some(id) = map + .get("value") + .and_then(|v| v.get("value").unwrap_or(v).as_u64()) + { + out.push(id); + } + map.values().for_each(|v| queue_ids(v, out)); + } + Value::Array(list) => list.iter().for_each(|v| queue_ids(v, out)), + _ => {} + } + } + let with_ids = traces + .iter() + .map(|t| { + let mut ids = Vec::new(); + queue_ids(t, &mut ids); + (t["id"].as_str().unwrap().to_string(), ids) + }) + .collect::>(); + let queue_id = with_ids + .iter() + .find_map(|(_, ids)| ids.first().copied()) + .expect("MON-16: a trace with a queue id"); + let mut expected = with_ids + .iter() + .filter(|(_, ids)| ids.contains(&queue_id)) + .map(|(id, _)| id.clone()) + .collect::>(); + expected.sort(); + for filter in [json!(queue_id.to_string()), json!(queue_id)] { + let response = admin + .jmap_method_call("x:Trace/query", json!({"filter": {"queueId": filter}})) + .await; + let mut found = response + .0 + .pointer("/methodResponses/0/1/ids") + .and_then(|ids| ids.as_array()) + .map(|ids| { + ids.iter() + .filter_map(|id| id.as_str().map(str::to_string)) + .collect::>() + }) + .unwrap_or_default(); + found.sort(); + assert_eq!(found, expected, "MON-16: queueId {filter}: {response:?}"); + } + let response = admin + .jmap_method_call( + "x:Trace/query", + json!({"filter": {"queueId": (queue_id ^ 0x5a5a_5a5a).to_string()}}), + ) + .await; + assert_eq!( + response.0.pointer("/methodResponses/0/1/ids"), + Some(&json!([])), + "MON-16: an unknown queue id" + ); + // Acceptance test 24: destroy removes a trace; create is refused let trace_id = traces[0]["id"].as_str().unwrap().to_string(); let response = admin