use std::collections::HashMap;
use std::sync::Arc;
use serde::Deserialize;
use crate::event::{EventSource, SpanEvent};
use crate::ingest::IngestSource;
use crate::time::micros_to_iso8601;
pub struct JaegerIngest {
max_size: usize,
grouping_attributes: Option<Vec<Arc<str>>>,
}
impl JaegerIngest {
#[must_use]
pub const fn new(max_size: usize) -> Self {
Self {
max_size,
grouping_attributes: None,
}
}
#[must_use]
pub fn with_grouping_attributes(mut self, keys: Vec<Arc<str>>) -> Self {
self.grouping_attributes = Some(keys);
self
}
}
impl IngestSource for JaegerIngest {
type Error = JaegerIngestError;
fn ingest(&self, raw: &[u8]) -> Result<Vec<SpanEvent>, Self::Error> {
if raw.len() > self.max_size {
return Err(JaegerIngestError::PayloadTooLarge {
size: raw.len(),
max: self.max_size,
});
}
let export: JaegerExport = serde_json::from_slice(raw).map_err(JaegerIngestError::Parse)?;
Ok(convert_jaeger_export(
&export,
self.grouping_attributes.as_deref(),
))
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum JaegerIngestError {
#[error("payload too large: {size} bytes exceeds maximum of {max} bytes")]
PayloadTooLarge { size: usize, max: usize },
#[error("JSON parse error: {0}")]
Parse(#[from] serde_json::Error),
}
#[derive(Deserialize)]
pub(super) struct JaegerExport {
pub(super) data: Vec<JaegerTrace>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub(super) struct JaegerTrace {
#[serde(rename = "traceID")]
trace_id: String,
spans: Vec<JaegerSpan>,
processes: HashMap<String, JaegerProcess>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JaegerSpan {
#[serde(rename = "spanID")]
span_id: String,
operation_name: String,
#[serde(default)]
references: Vec<JaegerReference>,
start_time: u64,
duration: u64,
#[serde(rename = "processID")]
process_id: String,
#[serde(default)]
tags: Vec<JaegerTag>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JaegerReference {
ref_type: String,
#[serde(rename = "spanID")]
span_id: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct JaegerProcess {
service_name: String,
#[serde(default)]
tags: Vec<JaegerTag>,
}
type ProcessMetadata<'a> = (Arc<str>, &'a [JaegerTag]);
#[derive(Deserialize)]
struct JaegerTag {
key: String,
value: serde_json::Value,
}
pub(super) fn convert_jaeger_export(
export: &JaegerExport,
grouping_attributes: Option<&[Arc<str>]>,
) -> Vec<SpanEvent> {
let cap: usize = export.data.iter().map(|t| t.spans.len()).sum();
let mut events = Vec::with_capacity(cap);
for trace in &export.data {
let process_metadata: HashMap<&str, ProcessMetadata> = trace
.processes
.iter()
.map(|(pid, process)| {
(
pid.as_str(),
(
Arc::from(process.service_name.as_str()),
process.tags.as_slice(),
),
)
})
.collect();
let span_index: HashMap<&str, &JaegerSpan> = trace
.spans
.iter()
.filter(|s| !s.span_id.is_empty())
.map(|s| (s.span_id.as_str(), s))
.collect();
for span in &trace.spans {
if let Some(event) = convert_jaeger_span(
span,
&trace.trace_id,
&process_metadata,
&span_index,
grouping_attributes,
) {
events.push(event);
}
}
}
events
}
fn child_of(span: &JaegerSpan) -> Option<&str> {
span.references
.iter()
.find(|r| r.ref_type == "CHILD_OF")
.map(|r| r.span_id.as_str())
}
fn inbound_http_endpoint(span: &JaegerSpan) -> Option<String> {
http_endpoint(
span,
find_tag(&span.tags, "span.kind").as_deref() != Some("client"),
)
}
fn own_inbound_http_endpoint(span: &JaegerSpan) -> Option<String> {
http_endpoint(
span,
find_tag(&span.tags, "span.kind").as_deref() == Some("server"),
)
}
fn http_endpoint(span: &JaegerSpan, allow_url_fallback: bool) -> Option<String> {
let usable = |s: &String| !s.trim().is_empty();
let route = find_tag(&span.tags, "http.route");
let url_path = find_tag(&span.tags, "url.path");
crate::ingest::http_route_endpoint(route.as_deref(), url_path.as_deref(), allow_url_fallback)
.or_else(|| {
if !allow_url_fallback {
return None;
}
find_tag(&span.tags, "http.target")
.filter(usable)
.or_else(|| find_tag(&span.tags, "http.url").filter(usable))
.or_else(|| find_tag(&span.tags, "url.full").filter(usable))
.or_else(|| find_tag(&span.tags, "url.path").filter(usable))
})
}
fn own_sql_http_endpoint(span: &JaegerSpan) -> Option<String> {
crate::ingest::http_route_endpoint(
find_tag(&span.tags, "http.route").as_deref(),
find_tag(&span.tags, "url.path").as_deref(),
find_tag(&span.tags, "span.kind").as_deref() == Some("server"),
)
.or_else(|| find_tag(&span.tags, "http.target"))
.filter(|endpoint| !endpoint.trim().is_empty())
}
fn tag_code_frame(tags: &[JaegerTag]) -> Option<String> {
let function_name = find_tag(tags, "code.function.name");
let function = function_name
.clone()
.or_else(|| find_tag(tags, "code.function"));
let namespace = find_tag(tags, "code.namespace").or_else(|| {
function_name
.as_deref()
.and_then(crate::ingest::namespace_from_qualified_name)
.map(ToString::to_string)
});
crate::ingest::code_frame_endpoint(namespace.as_deref(), function.as_deref())
}
fn same_jaeger_service(
leaf: &JaegerSpan,
ancestor: &JaegerSpan,
process_metadata: &HashMap<&str, ProcessMetadata>,
) -> bool {
if leaf.process_id == ancestor.process_id {
return true;
}
match (
process_metadata
.get(leaf.process_id.as_str())
.map(|metadata| metadata.0.as_ref())
.filter(|service| !service.is_empty()),
process_metadata
.get(ancestor.process_id.as_str())
.map(|metadata| metadata.0.as_ref())
.filter(|service| !service.is_empty()),
) {
(Some(leaf_service), Some(ancestor_service)) => leaf_service == ancestor_service,
_ => false,
}
}
fn resolve_source_endpoint(
own_endpoint: Option<String>,
leaf_frame: Option<String>,
leaf: &JaegerSpan,
process_metadata: &HashMap<&str, ProcessMetadata>,
span_index: &HashMap<&str, &JaegerSpan>,
) -> String {
let mut outermost_endpoint = own_endpoint;
let mut outermost_frame = leaf_frame;
let mut current = child_of(leaf);
for _ in 0..crate::ingest::ANCESTOR_WALK_MAX_DEPTH {
let Some(pid) = current else {
break;
};
let Some(parent) = span_index.get(pid) else {
break;
};
if !same_jaeger_service(leaf, parent, process_metadata) {
break;
}
if let Some(route) = inbound_http_endpoint(parent) {
outermost_endpoint = Some(route);
}
if let Some(frame) = tag_code_frame(&parent.tags) {
outermost_frame = Some(frame);
}
current = child_of(parent);
}
outermost_endpoint
.or(outermost_frame)
.unwrap_or_else(|| "unknown".to_string())
}
fn convert_jaeger_span(
span: &JaegerSpan,
trace_id: &str,
process_metadata: &HashMap<&str, ProcessMetadata>,
span_index: &HashMap<&str, &JaegerSpan>,
grouping_attributes: Option<&[Arc<str>]>,
) -> Option<SpanEvent> {
let tags = &span.tags;
let db_system_raw = find_tag(tags, "db.system.name").or_else(|| find_tag(tags, "db.system"));
let db_system = db_system_raw
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty())
.map(super::canonical_db_system);
if db_system.is_some_and(super::is_non_sql_db_system) {
return None;
}
let (io_kind, target) = if let Some(stmt) =
find_tag(tags, "db.statement").or_else(|| find_tag(tags, "db.query.text"))
{
(super::TagIoKind::Sql, stmt)
} else {
if find_tag(tags, "span.kind").as_deref() == Some("server") {
return None;
}
(
super::TagIoKind::HttpOut,
find_tag(tags, "http.url").or_else(|| find_tag(tags, "url.full"))?,
)
};
let operation = match io_kind {
super::TagIoKind::Sql => db_system.unwrap_or("sql").to_string(),
super::TagIoKind::HttpOut => find_tag(tags, "http.method")
.or_else(|| find_tag(tags, "http.request.method"))
.unwrap_or_else(|| "GET".to_string()),
};
let process = process_metadata.get(span.process_id.as_str());
let service: Arc<str> = process.map_or_else(|| Arc::from(""), |m| Arc::clone(&m.0));
let grouping = crate::ingest::collect_grouping(grouping_attributes, |key| {
process
.and_then(|m| find_tag(m.1, key).map(Arc::from))
.or_else(|| find_tag(tags, key).map(Arc::from))
});
let parent_span_id = child_of(span).map(ToString::to_string);
let status_code = match io_kind {
super::TagIoKind::HttpOut => find_tag(tags, "http.status_code")
.or_else(|| find_tag(tags, "http.response.status_code"))
.and_then(|s| s.parse().ok()),
super::TagIoKind::Sql => None,
};
let code_function_name = find_tag(tags, "code.function.name");
let code_function: Option<Arc<str>> = code_function_name
.clone()
.or_else(|| find_tag(tags, "code.function"))
.map(Arc::from);
let code_filepath: Option<Arc<str>> = find_tag(tags, "code.file.path")
.or_else(|| find_tag(tags, "code.filepath"))
.map(Arc::from);
let code_lineno = find_tag(tags, "code.line.number")
.or_else(|| find_tag(tags, "code.lineno"))
.and_then(|s| s.parse::<u32>().ok());
let code_namespace: Option<Arc<str>> = find_tag(tags, "code.namespace")
.or_else(|| {
code_function_name
.as_deref()
.and_then(crate::ingest::namespace_from_qualified_name)
.map(ToString::to_string)
})
.map(Arc::from);
let endpoint = resolve_source_endpoint(
match io_kind {
super::TagIoKind::Sql => own_sql_http_endpoint(span),
super::TagIoKind::HttpOut => own_inbound_http_endpoint(span),
},
tag_code_frame(tags),
span,
process_metadata,
span_index,
);
let method = find_tag(tags, "code.function").unwrap_or_else(|| span.operation_name.clone());
let mut event = SpanEvent {
timestamp: micros_to_iso8601(span.start_time),
trace_id: trace_id.to_string(),
span_id: span.span_id.clone(),
link_trace_id: None,
parent_span_id,
service,
grouping,
cloud_region: None,
event_type: io_kind.event_type(),
operation,
target,
duration_us: span.duration,
source: EventSource { endpoint, method },
status_code,
response_size_bytes: None,
code_function,
code_filepath,
code_lineno,
code_namespace,
instrumentation_scopes: Vec::new(),
};
crate::event::sanitize_span_event(&mut event);
Some(event)
}
fn find_tag(tags: &[JaegerTag], key: &str) -> Option<String> {
tags.iter().find(|t| t.key == key).map(|t| match &t.value {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::EventType;
fn sample_jaeger_json() -> &'static str {
r#"{
"data": [{
"traceID": "abc123",
"spans": [
{
"spanID": "span-1",
"operationName": "OrderService::create_order",
"references": [],
"startTime": 1720621921123000,
"duration": 1200,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT * FROM order_item WHERE order_id = 42" },
{ "key": "db.system", "value": "postgresql" }
]
},
{
"spanID": "span-2",
"operationName": "http-call",
"references": [{ "refType": "CHILD_OF", "spanID": "span-1" }],
"startTime": 1720621921200000,
"duration": 15000,
"processID": "p1",
"tags": [
{ "key": "http.url", "value": "http://user-svc:5000/api/users/123" },
{ "key": "http.method", "value": "GET" },
{ "key": "http.status_code", "value": "200" }
]
},
{
"spanID": "span-3",
"operationName": "internal-op",
"references": [],
"startTime": 1720621921300000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "internal.type", "value": "processing" }
]
}
],
"processes": {
"p1": { "serviceName": "order-svc" }
}
}]
}"#
}
#[test]
fn namespaces_are_extracted_from_process_tags() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 1200,
"processID": "p1",
"tags": [{ "key": "db.statement", "value": "SELECT 1" }]
}],
"processes": { "p1": {
"serviceName": "svc",
"tags": [
{ "key": "service.namespace", "value": "payments" },
{ "key": "k8s.namespace.name", "value": "prod-eu" }
]
}}
}]
}"#;
let events = JaegerIngest::new(64 * 1024)
.ingest(json.as_bytes())
.unwrap();
let captured = crate::test_helpers::grouping_pairs(&events[0].grouping);
assert_eq!(
captured,
vec![
("k8s.namespace.name", "prod-eu"),
("service.namespace", "payments"),
],
"both values are kept, config order, Kubernetes first"
);
assert_eq!(events[0].grouping_value(), Some("prod-eu"));
}
#[test]
fn parses_jaeger_export() {
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(sample_jaeger_json().as_bytes()).unwrap();
assert_eq!(events.len(), 2, "non-IO span should be skipped");
}
#[test]
fn non_sql_datastore_span_is_dropped() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "s1", "operationName": "redis-get",
"references": [], "startTime": 1, "duration": 10, "processID": "p1",
"tags": [
{ "key": "db.system", "value": "redis" },
{ "key": "db.statement", "value": "GET user:123" }
]
},
{
"spanID": "s2", "operationName": "sql",
"references": [], "startTime": 1, "duration": 10, "processID": "p1",
"tags": [
{ "key": "db.system", "value": "postgresql" },
{ "key": "db.statement", "value": "SELECT 1" }
]
}
],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "postgresql");
}
#[test]
fn db_system_alias_is_canonicalized() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1", "operationName": "q",
"startTime": 0, "duration": 100, "processID": "p1",
"tags": [
{ "key": "db.system", "value": "postgres" },
{ "key": "db.statement", "value": "SELECT 1" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "postgresql");
}
#[test]
fn stable_db_system_name_non_sql_is_dropped() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1", "operationName": "q",
"startTime": 0, "duration": 100, "processID": "p1",
"tags": [
{ "key": "db.system.name", "value": "aws.dynamodb" },
{ "key": "db.statement", "value": "SELECT * FROM Orders WHERE Id = 'secret'" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert!(events.is_empty());
}
#[test]
fn sql_span_maps_correctly() {
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(sample_jaeger_json().as_bytes()).unwrap();
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.unwrap();
assert_eq!(sql.trace_id, "abc123");
assert_eq!(sql.span_id, "span-1");
assert_eq!(&*sql.service, "order-svc");
assert_eq!(sql.operation, "postgresql");
assert_eq!(sql.target, "SELECT * FROM order_item WHERE order_id = 42");
assert_eq!(sql.duration_us, 1200);
assert!(sql.parent_span_id.is_none());
assert!(sql.status_code.is_none());
assert_eq!(sql.timestamp, "2024-07-10T14:32:01.123Z");
}
#[test]
fn http_span_maps_correctly() {
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(sample_jaeger_json().as_bytes()).unwrap();
let http = events
.iter()
.find(|e| e.event_type == EventType::HttpOut)
.unwrap();
assert_eq!(http.trace_id, "abc123");
assert_eq!(http.span_id, "span-2");
assert_eq!(http.operation, "GET");
assert_eq!(http.target, "http://user-svc:5000/api/users/123");
assert_eq!(http.duration_us, 15000);
assert_eq!(http.status_code, Some(200));
assert_eq!(http.parent_span_id.as_deref(), Some("span-1"));
}
#[test]
fn rejects_oversized_payload() {
let ingest = JaegerIngest::new(10);
let result = ingest.ingest(sample_jaeger_json().as_bytes());
assert!(result.is_err());
}
#[test]
fn malformed_json_missing_data_key() {
let json = r#"{"traces": []}"#;
let ingest = JaegerIngest::new(1_048_576);
assert!(ingest.ingest(json.as_bytes()).is_err());
}
#[test]
fn malformed_json_missing_trace_id() {
let json = r#"{"data": [{"spans": [], "processes": {}}]}"#;
let ingest = JaegerIngest::new(1_048_576);
assert!(ingest.ingest(json.as_bytes()).is_err());
}
#[test]
fn malformed_json_missing_spans() {
let json = r#"{"data": [{"traceID": "t1", "processes": {}}]}"#;
let ingest = JaegerIngest::new(1_048_576);
assert!(ingest.ingest(json.as_bytes()).is_err());
}
#[test]
fn malformed_json_missing_span_id() {
let json = r#"{"data": [{"traceID": "t1", "spans": [{"operationName": "op", "startTime": 0, "duration": 0, "processID": "p1", "tags": []}], "processes": {"p1": {"serviceName": "svc"}}}]}"#;
let ingest = JaegerIngest::new(1_048_576);
assert!(ingest.ingest(json.as_bytes()).is_err());
}
#[test]
fn empty_data_array_produces_no_events() {
let json = r#"{"data": []}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert!(events.is_empty());
}
#[test]
fn empty_spans_array_produces_no_events() {
let json = r#"{"data": [{"traceID": "t1", "spans": [], "processes": {"p1": {"serviceName": "svc"}}}]}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert!(events.is_empty());
}
#[test]
fn unknown_process_id_produces_empty_service() {
let json = r#"{"data": [{"traceID": "t1", "spans": [{"spanID": "s1", "operationName": "op", "startTime": 0, "duration": 100, "processID": "unknown", "tags": [{"key": "db.statement", "value": "SELECT 1"}]}], "processes": {"p1": {"serviceName": "svc"}}}]}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(&*events[0].service, "");
}
#[test]
fn numeric_tag_value_converted_to_string() {
let json = r#"{"data": [{"traceID": "t1", "spans": [{"spanID": "s1", "operationName": "op", "startTime": 0, "duration": 100, "processID": "p1", "tags": [{"key": "http.url", "value": "http://svc/api"}, {"key": "http.status_code", "value": 200}]}], "processes": {"p1": {"serviceName": "svc"}}}]}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].status_code, Some(200));
}
#[test]
fn slashless_http_route_is_canonicalized_before_http_target() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "http.route", "value": "api/orders/{id}" },
{ "key": "http.target", "value": "/api/orders/42" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "/api/orders/{id}");
}
#[test]
fn named_route_uses_url_path_for_own_and_ancestor_endpoints() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "root",
"operationName": "request",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "app_fault_nplusonesql" },
{ "key": "url.path", "value": "/api/fault/n-plus-one-sql" }
]
},
{
"spanID": "child",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "root" }],
"startTime": 1720621921123100,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" }
]
},
{
"traceID": "t2",
"spanID": "own",
"operationName": "query",
"references": [],
"startTime": 1720621921123200,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "db.statement", "value": "SELECT 2" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "http.route", "value": "app_fault_nplusonesql" },
{ "key": "url.path", "value": "/api/fault/n-plus-one-sql" }
]
}
],
"processes": { "p1": { "serviceName": "symfony-svc" } }
}]
}"#;
let events = JaegerIngest::new(1_048_576)
.ingest(json.as_bytes())
.unwrap();
assert_eq!(events.len(), 2);
assert!(
events
.iter()
.all(|event| event.source.endpoint == "/api/fault/n-plus-one-sql")
);
}
#[test]
fn server_route_and_legacy_url_is_context_not_http_out() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "POST /api/orders/{id}",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "api/orders/{id}" },
{ "key": "http.url", "value": "http://order-svc/api/orders/42" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let events = JaegerIngest::new(1_048_576)
.ingest(json.as_bytes())
.unwrap();
assert!(events.is_empty());
}
#[test]
fn server_url_full_without_route_is_context_not_http_out() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "POST /api/orders/42",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "url.full", "value": "http://order-svc/api/orders/42" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let events = JaegerIngest::new(1_048_576)
.ingest(json.as_bytes())
.unwrap();
assert!(events.is_empty());
}
#[test]
fn unspecified_outgoing_url_uses_its_parent_server_route() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "s1",
"operationName": "POST /api/orders",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "api/orders" }
]
},
{
"spanID": "s2",
"operationName": "GET",
"references": [{ "refType": "CHILD_OF", "spanID": "s1" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "http.url", "value": "https://partner.example/v1/pay" }
]
}
],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let events = JaegerIngest::new(1_048_576)
.ingest(json.as_bytes())
.unwrap();
let outgoing = events
.iter()
.find(|event| event.event_type == EventType::HttpOut)
.expect("outgoing event present");
assert_eq!(outgoing.source.endpoint, "/api/orders");
}
#[test]
fn unspecified_root_url_does_not_self_source() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "GET",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "http.url", "value": "https://partner.example/v1/pay" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let events = JaegerIngest::new(1_048_576)
.ingest(json.as_bytes())
.unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::HttpOut);
assert_eq!(events[0].source.endpoint, "unknown");
}
#[test]
fn http_target_used_only_when_route_absent() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "http.target", "value": "/api/orders/42" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "/api/orders/42");
}
#[test]
fn code_frame_used_when_no_http_tag() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "code.function", "value": "execute" },
{ "key": "code.namespace", "value": "com.foo.PurgeJob" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "com.foo.PurgeJob.execute");
}
#[test]
fn endpoint_resolves_through_ancestors() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "s1",
"operationName": "POST /api/orders",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "api/orders" }
]
},
{
"spanID": "s2",
"operationName": "GET",
"references": [{ "refType": "CHILD_OF", "spanID": "s1" }],
"startTime": 1720621921123100,
"duration": 3000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "client" },
{ "key": "http.url", "value": "https://partner.example/v1/pay" },
{ "key": "url.path", "value": "/v1/pay" }
]
},
{
"spanID": "s3",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "s2" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" }
]
}
],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
let outbound = events
.iter()
.find(|event| event.event_type == EventType::HttpOut)
.expect("client outbound event present");
assert_eq!(outbound.target, "https://partner.example/v1/pay");
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf present");
assert_eq!(sql.source.endpoint, "/api/orders");
}
#[test]
#[allow(clippy::too_many_lines)] fn outermost_route_stops_at_the_service_boundary() {
let json = r#"{
"data": [
{
"traceID": "same-service",
"spans": [
{
"spanID": "outer",
"operationName": "POST /api/fault/pool-saturation",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "laravel",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "/api/fault/pool-saturation" }
]
},
{
"spanID": "nested",
"operationName": "GET /api/payments/history",
"references": [{ "refType": "CHILD_OF", "spanID": "outer" }],
"startTime": 1720621921123100,
"duration": 3000,
"processID": "laravel",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "/api/payments/history" }
]
},
{
"spanID": "sql",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "nested" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "laravel",
"tags": [
{ "key": "db.statement", "value": "SELECT * FROM payments" },
{ "key": "db.system", "value": "postgresql" }
]
}
],
"processes": { "laravel": { "serviceName": "laravel-svc" } }
},
{
"traceID": "cross-service",
"spans": [
{
"spanID": "caller",
"operationName": "POST /api/orders",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "orders",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "/api/orders" }
]
},
{
"spanID": "callee",
"operationName": "GET /api/payments/history",
"references": [{ "refType": "CHILD_OF", "spanID": "caller" }],
"startTime": 1720621921123100,
"duration": 3000,
"processID": "payments",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.route", "value": "/api/payments/history" }
]
},
{
"spanID": "sql",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "callee" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "payments",
"tags": [
{ "key": "db.statement", "value": "SELECT * FROM payments" },
{ "key": "db.system", "value": "postgresql" }
]
}
],
"processes": {
"orders": { "serviceName": "orders-svc" },
"payments": { "serviceName": "payments-svc" }
}
}
]
}"#;
let events = JaegerIngest::new(1_048_576)
.ingest(json.as_bytes())
.unwrap();
let same_service = events
.iter()
.find(|event| event.trace_id == "same-service")
.expect("same-service SQL event present");
let cross_service = events
.iter()
.find(|event| event.trace_id == "cross-service")
.expect("cross-service SQL event present");
assert_eq!(same_service.source.endpoint, "/api/fault/pool-saturation");
assert_eq!(cross_service.source.endpoint, "/api/payments/history");
assert_eq!(&*cross_service.service, "payments-svc");
}
#[test]
fn route_at_hop_nine_is_outside_the_shared_depth_limit() {
let mut spans = vec![serde_json::json!({
"spanID": "p9", "operationName": "server", "references": [],
"startTime": 1, "duration": 1, "processID": "svc",
"tags": [{ "key": "http.route", "value": "/too-deep" }]
})];
for id in (1_u8..9).rev() {
spans.push(serde_json::json!({
"spanID": format!("p{id}"), "operationName": "internal",
"references": [{ "refType": "CHILD_OF", "spanID": format!("p{}", id + 1) }],
"startTime": 1, "duration": 1, "processID": "svc",
"tags": if id == 8 {
serde_json::json!([{ "key": "http.route", "value": "/at-limit" }])
} else {
serde_json::json!([])
}
}));
}
spans.push(serde_json::json!({
"spanID": "sql", "operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "p1" }],
"startTime": 1, "duration": 1, "processID": "svc",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" }
]
}));
let payload = serde_json::json!({
"data": [{
"traceID": "trace", "spans": spans,
"processes": { "svc": { "serviceName": "svc" } }
}]
})
.to_string();
let events = JaegerIngest::new(1_048_576)
.ingest(payload.as_bytes())
.unwrap();
assert_eq!(events[0].source.endpoint, "/at-limit");
}
#[test]
fn walk_accepts_http_target_on_an_ancestor() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "s1",
"operationName": "POST /api/orders",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.target", "value": "/api/orders/42" }
]
},
{
"spanID": "s2",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "s1" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" }
]
}
],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf present");
assert_eq!(sql.source.endpoint, "/api/orders/42");
}
#[test]
fn parent_stable_url_path_provides_source_endpoint() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "s1",
"operationName": "POST /api/fault/pool-saturation",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "url.path", "value": "/api/fault/pool-saturation" }
]
},
{
"spanID": "s2",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "s1" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "code.function.name", "value": "com.foo.FaultPool.query" }
]
}
],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "/api/fault/pool-saturation");
}
#[test]
fn empty_http_fallback_does_not_block_url_path() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [
{
"spanID": "s1",
"operationName": "POST /api/fault/pool-saturation",
"references": [],
"startTime": 1720621921123000,
"duration": 5000,
"processID": "p1",
"tags": [
{ "key": "span.kind", "value": "server" },
{ "key": "http.target", "value": "" },
{ "key": "http.url", "value": "" },
{ "key": "url.full", "value": "" },
{ "key": "url.path", "value": "/api/fault/pool-saturation" }
]
},
{
"spanID": "s2",
"operationName": "query",
"references": [{ "refType": "CHILD_OF", "spanID": "s1" }],
"startTime": 1720621921123200,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "code.function.name", "value": "com.foo.FaultPool.query" }
]
}
],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
let sql = events
.iter()
.find(|event| event.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "/api/fault/pool-saturation");
}
#[test]
fn endpoint_falls_back_to_unknown_not_empty() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events[0].source.endpoint, "unknown");
}
#[test]
fn code_frame_endpoint_reads_stable_semconv() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.statement", "value": "SELECT 1" },
{ "key": "db.system", "value": "postgresql" },
{ "key": "code.function.name", "value": "com.foo.PurgeJob.execute" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "com.foo.PurgeJob.execute");
assert_eq!(
events[0].code_namespace.as_deref(),
Some("com.foo.PurgeJob"),
"namespace must be derived from the FQ name, as the OTLP path does"
);
}
#[test]
fn stable_semconv_tags() {
let json = r#"{
"data": [{
"traceID": "t1",
"spans": [{
"spanID": "s1",
"operationName": "query",
"references": [],
"startTime": 1720621921123000,
"duration": 500,
"processID": "p1",
"tags": [
{ "key": "db.query.text", "value": "SELECT 1" },
{ "key": "db.system", "value": "mysql" }
]
}, {
"spanID": "s2",
"operationName": "fetch",
"references": [],
"startTime": 1720621921200000,
"duration": 1000,
"processID": "p1",
"tags": [
{ "key": "url.full", "value": "http://api/items" },
{ "key": "http.request.method", "value": "POST" },
{ "key": "http.response.status_code", "value": "201" }
]
}],
"processes": { "p1": { "serviceName": "svc" } }
}]
}"#;
let ingest = JaegerIngest::new(1_048_576);
let events = ingest.ingest(json.as_bytes()).unwrap();
assert_eq!(events.len(), 2);
let sql = &events[0];
assert_eq!(sql.target, "SELECT 1");
assert_eq!(sql.operation, "mysql");
let http = &events[1];
assert_eq!(http.target, "http://api/items");
assert_eq!(http.operation, "POST");
assert_eq!(http.status_code, Some(201));
}
}