Trace search: index event type and queue id as integers #33

Merged
jcoffey-dev merged 2 commits from fix/pg-index-trace-types into main 2026-09-24 14:57:19 +00:00
5 changed files with 312 additions and 25 deletions
+17 -2
View File
@@ -427,9 +427,24 @@ pub(crate) async fn trace_query(
} }
None => false, 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::<u64>().ok())
.or_else(|| value.as_u64())
{
Some(queue_id) => { 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 true
} }
None => false, None => false,
+57 -20
View File
@@ -567,20 +567,14 @@ async fn build_contact_document(
} }
// inbuxa: MON-16: a trace's search document, when trace search is on: // 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( async fn build_tracing_span_document(
server: &Server, server: &Server,
span_id: u64, span_id: u64,
) -> trc::Result<Option<IndexDocument>> { ) -> trc::Result<Option<IndexDocument>> {
use common::telemetry::tracers::store::MaybeTrace; use common::telemetry::tracers::store::MaybeTrace;
use registry::schema::{enums::SearchTracingField, structs::Search}; use registry::schema::structs::Search;
use store::{ use store::write::{TelemetryClass, ValueClass};
search::TracingSearchField,
write::{TelemetryClass, ValueClass},
};
use trc::Key;
let settings = server let settings = server
.registry() .registry()
@@ -590,7 +584,6 @@ async fn build_tracing_span_document(
if !settings.index_telemetry { if !settings.index_telemetry {
return Ok(None); return Ok(None);
} }
let wants = |field: SearchTracingField| settings.index_tracing_fields.iter().any(|f| *f == field);
let Some(MaybeTrace(Some(trace))) = server let Some(MaybeTrace(Some(trace))) = server
.tracing_store() .tracing_store()
.get_value::<MaybeTrace>(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span( .get_value::<MaybeTrace>(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(
@@ -601,23 +594,67 @@ async fn build_tracing_span_document(
return Ok(None); return Ok(None);
}; };
Ok(Some(trace_search_document(
span_id,
&trace,
&settings
.index_tracing_fields
.iter()
.copied()
.collect::<Vec<_>>(),
)))
}
/// 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: &registry::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); 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 seen = store::ahash::AHashSet::new();
let mut queue_id_indexed = false;
for event in trace.events.iter() { 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() { for kv in event.key_values.iter() {
let text = match &kv.value { let text = match &kv.value {
registry::schema::structs::TraceValue::String(v) => v.value.clone(), TraceValue::String(v) => v.value.clone(),
registry::schema::structs::TraceValue::UnsignedInt(v) => v.value.to_string(), TraceValue::UnsignedInt(v) => v.value.to_string(),
registry::schema::structs::TraceValue::IpAddr(v) => v.value.to_string(), TraceValue::IpAddr(v) => v.value.to_string(),
_ => continue, _ => continue,
}; };
match kv.key { match kv.key {
Key::QueueId if wants(SearchTracingField::QueueId) => { Key::QueueId => {
if seen.insert(format!("q:{text}")) { let Ok(queue_id) = text.parse::<u64>() else {
document.index_keyword(TracingSearchField::QueueId, &text); 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 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 // inbuxa: UD-1, UD-4: archives a deleted file, event or contact noted at
+8 -3
View File
@@ -212,10 +212,15 @@ unchanged.
- **MON-16.** With `indexTelemetry` on, storing a trace schedules an - **MON-16.** With `indexTelemetry` on, storing a trace schedules an
`IndexTrace` task. The task builds one document for `SearchIndex::Tracing` `IndexTrace` task. The task builds one document for `SearchIndex::Tracing`
with the fields named in `indexTracingFields`: with the fields named in `indexTracingFields`:
- `eventType`: every event type in the trace; - `eventType`: the trace's opening event, as its numeric id;
- `queueId`: every `queueId` value; - `queueId`: the first `queueId` value, as an integer;
- `keywords`: every address in `from` and `to`, each address's domain, every - `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 So searching `example.org` finds every trace to or from that domain, as the
upstream suite expects. With `indexTelemetry` off nothing is indexed, and upstream suite expects. With `indexTelemetry` off nothing is indexed, and
the `text` and `queueId` filters are refused (see "Interfaces"). the `text` and `queueId` filters are refused (see "Interfaces").
+163
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]> * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* *
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL * 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}; 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..."); println!("Running global id filtering tests...");
test_global(store.clone()).await; 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 // Large document insert test
println!("Running large document insert tests..."); println!("Running large document insert tests...");
let mut large_text = String::with_capacity(20 * 1024 * 1024); 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]) 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<TraceKeyValue>| 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, "[email protected]")],
),
event(
EventType::Smtp(SmtpEvent::RcptTo),
vec![kv_u(Key::QueueId, 9_000_000_001), kv_s(Key::To, "[email protected]")],
),
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<SearchFilter>| {
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::<AHashSet<_>>()
}
};
// 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();
}
}
+67
View File
@@ -148,6 +148,73 @@ pub async fn test(test: &mut TestServer) {
"test 9: to" "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<u64>) {
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::<Vec<_>>();
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::<Vec<_>>();
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::<Vec<_>>()
})
.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 // Acceptance test 24: destroy removes a trace; create is refused
let trace_id = traces[0]["id"].as_str().unwrap().to_string(); let trace_id = traces[0]["id"].as_str().unwrap().to_string();
let response = admin let response = admin