use super::*;
#[cfg(feature = "daemon")]
use crate::report::metrics::MetricsState;
use opentelemetry_proto::tonic::common::v1::AnyValue;
use opentelemetry_proto::tonic::resource::v1::Resource;
use opentelemetry_proto::tonic::trace::v1::{ResourceSpans, ScopeSpans};
#[cfg(feature = "daemon")]
fn fresh_metrics_sink() -> (Arc<MetricsState>, Arc<dyn MetricsSink>) {
let state = Arc::new(MetricsState::new());
let sink: Arc<dyn MetricsSink> = state.clone();
(state, sink)
}
fn make_kv(key: &str, value: &str) -> KeyValue {
KeyValue {
key: key.to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue(value.to_string())),
}),
..Default::default()
}
}
fn make_int_kv(key: &str, value: i64) -> KeyValue {
KeyValue {
key: key.to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::IntValue(value)),
}),
..Default::default()
}
}
fn make_sql_span(
trace_id: &[u8],
span_id: &[u8],
parent_span_id: &[u8],
statement: &str,
start_ns: u64,
end_ns: u64,
) -> Span {
Span {
trace_id: trace_id.to_vec(),
span_id: span_id.to_vec(),
parent_span_id: parent_span_id.to_vec(),
name: "db.query".to_string(),
start_time_unix_nano: start_ns,
end_time_unix_nano: end_ns,
attributes: vec![
make_kv("db.statement", statement),
make_kv("db.system", "postgresql"),
],
..Default::default()
}
}
#[allow(clippy::too_many_arguments)] fn make_http_span(
trace_id: &[u8],
span_id: &[u8],
parent_span_id: &[u8],
url: &str,
method: &str,
status: i64,
start_ns: u64,
end_ns: u64,
) -> Span {
Span {
trace_id: trace_id.to_vec(),
span_id: span_id.to_vec(),
parent_span_id: parent_span_id.to_vec(),
name: "http.request".to_string(),
start_time_unix_nano: start_ns,
end_time_unix_nano: end_ns,
attributes: vec![
make_kv("http.url", url),
make_kv("http.method", method),
make_int_kv("http.status_code", status),
],
..Default::default()
}
}
fn make_parent_span(span_id: &[u8], route: &str) -> Span {
Span {
trace_id: vec![1; 16],
span_id: span_id.to_vec(),
parent_span_id: vec![],
name: "HandleRequest".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("http.route", route),
make_kv("code.function", "OrderService::create_order"),
],
..Default::default()
}
}
fn make_request(service: &str, spans: Vec<Span>) -> ExportTraceServiceRequest {
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![make_kv("service.name", service)],
..Default::default()
}),
scope_spans: vec![ScopeSpans {
spans,
..Default::default()
}],
..Default::default()
}],
}
}
#[test]
fn empty_request_returns_empty() {
let req = ExportTraceServiceRequest {
resource_spans: vec![],
};
assert!(convert_otlp_request(&req).is_empty());
}
#[test]
fn sql_span_maps_correctly() {
let span = make_sql_span(
&[1; 16],
&[2; 8],
&[],
"SELECT * FROM order_item WHERE order_id = 42",
1_720_621_921_000_000_000, 1_720_621_921_001_200_000, );
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "postgresql");
assert_eq!(
events[0].target,
"SELECT * FROM order_item WHERE order_id = 42"
);
assert_eq!(&*events[0].service, "order-svc");
assert_eq!(events[0].duration_us, 1200);
assert!(events[0].status_code.is_none());
}
#[test]
fn http_span_maps_correctly() {
let span = make_http_span(
&[1; 16],
&[3; 8],
&[],
"http://user-svc:5000/api/users/123",
"GET",
200,
1_720_621_921_000_000_000,
1_720_621_921_015_000_000, );
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::HttpOut);
assert_eq!(events[0].operation, "GET");
assert_eq!(events[0].target, "http://user-svc:5000/api/users/123");
assert_eq!(events[0].status_code, Some(200));
assert_eq!(events[0].duration_us, 15000);
}
#[test]
fn non_io_span_skipped() {
let span = Span {
trace_id: vec![1; 16],
span_id: vec![4; 8],
name: "internal.processing".to_string(),
start_time_unix_nano: 1_720_621_921_000_000_000,
end_time_unix_nano: 1_720_619_521_000_500_000,
attributes: vec![make_kv("custom.attr", "value")],
..Default::default()
};
let req = make_request("order-svc", vec![span]);
assert!(convert_otlp_request(&req).is_empty());
}
fn make_bare_span(span_id: &[u8], attributes: Vec<KeyValue>) -> Span {
Span {
trace_id: vec![1; 16],
span_id: span_id.to_vec(),
name: "fixture".to_string(),
start_time_unix_nano: 1_720_621_921_000_000_000,
end_time_unix_nano: 1_720_621_921_000_500_000,
attributes,
..Default::default()
}
}
#[test]
fn counted_conversion_classifies_filtered_spans() {
let internal = make_bare_span(&[4; 8], vec![make_kv("custom.attr", "value")]);
let db_no_statement = make_bare_span(&[5; 8], vec![make_kv("db.system", "postgresql")]);
let http_no_url = make_bare_span(&[6; 8], vec![make_kv("http.method", "GET")]);
let sql = make_sql_span(&[1; 16], &[7; 8], &[], "SELECT 1", 0, 1000);
let req = make_request(
"order-svc",
vec![internal, db_no_statement, http_no_url, sql],
);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(
stats,
SpanConversionStats {
received: 4,
filtered_not_io: 1,
filtered_missing_db_statement: 1,
filtered_missing_http_url: 1,
filtered_non_sql_datastore: 0,
filtered_merged_db_span: 0,
}
);
}
#[test]
fn non_sql_datastore_span_is_dropped() {
let redis = make_bare_span(
&[8; 8],
vec![
make_kv("db.system", "redis"),
make_kv("db.statement", "GET user:123"),
],
);
let sql = make_sql_span(&[1; 16], &[7; 8], &[], "SELECT 1", 0, 1000);
let req = make_request("order-svc", vec![redis, sql]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(stats.received, 2);
assert_eq!(stats.filtered_non_sql_datastore, 1);
assert_eq!(stats.filtered_not_io, 0);
}
#[test]
fn non_sql_datastore_span_with_url_is_dropped_not_http() {
let es = make_bare_span(
&[8; 8],
vec![
make_kv("db.system", "elasticsearch"),
make_kv("db.statement", "GET /index/_search"),
make_kv("url.full", "http://es:9200/index/_search"),
],
);
let req = make_request("search-svc", vec![es]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_non_sql_datastore, 1);
}
#[test]
fn non_sql_datastore_span_without_statement_is_not_an_instrumentation_gap() {
let redis = make_bare_span(&[8; 8], vec![make_kv("db.system", "redis")]);
let req = make_request("cache-svc", vec![redis]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_non_sql_datastore, 1);
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn datadog_resource_with_db_type_classifies_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "SELECT * FROM orders WHERE id = ?"),
make_kv("db.type", "postgres"),
],
);
let req = make_request("order-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].target, "SELECT * FROM orders WHERE id = ?");
assert_eq!(events[0].operation, "postgresql");
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn datadog_resource_with_unknown_db_type_is_not_tokenized_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "GET namespace:set:key"),
make_kv("db.type", "aerospike"),
],
);
let req = make_request("cache-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_missing_db_statement, 1);
}
#[test]
fn datadog_empty_resource_is_not_an_empty_sql_event() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", " "),
make_kv("db.type", "postgres"),
],
);
let req = make_request("order-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_missing_db_statement, 1);
}
#[test]
fn datadog_resource_with_otel_db_system_classifies_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "SELECT 1"),
make_kv("db.system", "postgresql"),
],
);
let req = make_request("order-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].target, "SELECT 1");
assert_eq!(events[0].operation, "postgresql");
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn datadog_receiver_stable_semconv_db_system_name_classifies_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv(
"dd.span.Resource",
"SELECT * FROM order_item WHERE order_id = ?",
),
make_kv("db.system.name", "postgres"),
],
);
let req = make_request("dd-shop", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(
events[0].target,
"SELECT * FROM order_item WHERE order_id = ?"
);
assert_eq!(events[0].operation, "postgresql");
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn datadog_stable_namespaced_non_sql_is_dropped() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("db.system.name", "aws.dynamodb"),
make_kv(
"db.statement",
"SELECT * FROM Orders WHERE Id = 'secret-key'",
),
],
);
let req = make_request("shop", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_non_sql_datastore, 1);
}
#[test]
fn datadog_stable_unknown_db_system_name_without_statement_is_a_gap() {
let span = make_bare_span(&[9; 8], vec![make_kv("db.system.name", "aerospike")]);
let req = make_request("cache", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_missing_db_statement, 1);
}
#[test]
fn datadog_http_resource_with_db_tag_is_not_tokenized_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "GET /api/users?token=SECRET"),
make_kv("db.type", "postgres"),
make_kv("http.url", "https://svc/api/users?token=SECRET"),
],
);
let req = make_request("svc", vec![span]);
let (events, _stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::HttpOut);
}
#[test]
fn datadog_stable_namespaced_sql_server_classifies_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "SELECT * FROM orders WHERE id = ?"),
make_kv("db.system.name", "microsoft.sql_server"),
],
);
let req = make_request("shop", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "mssql");
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn datadog_empty_db_system_name_does_not_shadow_db_type() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("db.system.name", ""),
make_kv("db.type", "postgres"),
make_kv("dd.span.Resource", "SELECT 1"),
],
);
let req = make_request("shop", vec![span]);
let (events, _stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "postgresql");
}
#[test]
fn datadog_stable_http_method_blocks_resource_sql_fallback() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "GET /api/users?token=SECRET"),
make_kv("db.type", "postgres"),
make_kv("http.request.method", "GET"),
],
);
let req = make_request("svc", vec![span]);
let (events, _stats) = convert_otlp_request_counted(&req);
assert!(!events.iter().any(|e| e.target.contains("SECRET")));
}
#[test]
fn datadog_whitespace_db_system_name_does_not_shadow_db_type() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("db.system.name", " "),
make_kv("db.type", "postgres"),
make_kv("dd.span.Resource", "SELECT 1"),
],
);
let req = make_request("shop", vec![span]);
let (events, _stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "postgresql");
}
#[test]
fn datadog_cloud_sql_engine_classifies_as_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "SELECT * FROM orders WHERE id = ?"),
make_kv("db.type", "snowflake"),
],
);
let req = make_request("shop", vec![span]);
let (events, _stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].operation, "snowflake");
}
#[test]
fn datadog_resource_whitespace_is_trimmed_in_target() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", " SELECT 1\n"),
make_kv("db.type", "postgres"),
],
);
let req = make_request("shop", vec![span]);
let (events, _stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT 1");
}
#[test]
fn datadog_resource_without_db_signal_is_not_sql() {
let span = make_bare_span(
&[9; 8],
vec![make_kv("dd.span.Resource", "GET /api/orders")],
);
let req = make_request("order-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_not_io, 1);
}
#[test]
fn datadog_redis_resource_is_dropped_non_sql() {
let span = make_bare_span(
&[9; 8],
vec![
make_kv("dd.span.Resource", "GET user:123"),
make_kv("db.type", "redis"),
],
);
let req = make_request("cache-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_non_sql_datastore, 1);
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn server_span_without_url_counts_not_io_not_missing_url() {
use opentelemetry_proto::tonic::trace::v1::span::SpanKind;
let mut server = make_bare_span(&[5; 8], vec![make_kv("http.request.method", "GET")]);
server.kind = SpanKind::Server as i32;
let mut client = make_bare_span(&[6; 8], vec![make_kv("http.request.method", "GET")]);
client.kind = SpanKind::Client as i32;
let req = make_request("order-svc", vec![server, client]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(
stats,
SpanConversionStats {
received: 2,
filtered_not_io: 1,
filtered_missing_db_statement: 0,
filtered_missing_http_url: 1,
filtered_non_sql_datastore: 0,
filtered_merged_db_span: 0,
}
);
}
#[test]
fn counted_conversion_all_filtered_yields_zero_events() {
let internal = make_bare_span(&[4; 8], vec![make_kv("custom.attr", "value")]);
let req = make_request("order-svc", vec![internal]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.received, 1);
assert_eq!(stats.filtered_not_io, 1);
}
#[cfg(feature = "daemon")]
#[test]
fn record_otlp_spans_moves_received_and_filtered_counters() {
let (state, sink) = fresh_metrics_sink();
sink.record_otlp_spans(SpanConversionStats {
received: 5,
filtered_not_io: 2,
filtered_missing_db_statement: 1,
filtered_missing_http_url: 0,
filtered_non_sql_datastore: 3,
filtered_merged_db_span: 4,
});
assert_eq!(state.otlp_spans_received_total.get(), 5);
let filtered = |reason: &str| {
state
.otlp_spans_filtered_total
.with_label_values(&[reason])
.get()
};
assert_eq!(filtered("not_io"), 2);
assert_eq!(filtered("missing_db_statement"), 1);
assert_eq!(filtered("missing_http_url"), 0);
assert_eq!(filtered("non_sql_datastore"), 3);
assert_eq!(filtered("merged_db_span"), 4);
}
#[test]
fn parent_span_provides_source_endpoint() {
let parent = make_parent_span(&[10; 8], "POST /api/orders/{id}/submit");
let child = make_sql_span(
&[1; 16],
&[20; 8],
&[10; 8], "SELECT * FROM order_item WHERE order_id = 42",
1_720_621_921_000_000_000,
1_720_621_921_001_200_000,
);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "POST /api/orders/{id}/submit");
assert_eq!(events[0].source.method, "OrderService::create_order");
}
#[test]
fn parent_span_http_route_takes_precedence_over_http_url() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "HandleRequest".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("http.route", "POST /api/orders/{id}/submit"),
make_kv("http.url", "http://order-svc/api/orders/42/submit"),
make_kv("code.function", "OrderService::create_order"),
],
..Default::default()
};
let child = make_sql_span(
&[1; 16],
&[20; 8],
&[10; 8],
"SELECT * FROM order_item WHERE order_id = 42",
1_720_621_921_000_000_000,
1_720_621_921_001_200_000,
);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "POST /api/orders/{id}/submit");
}
#[test]
fn parent_span_http_url_used_only_when_route_absent() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "HandleRequest".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("http.url", "http://order-svc/api/orders/42/submit")],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "http://order-svc/api/orders/42/submit");
}
const SPAN_KIND_CLIENT: i32 = opentelemetry_proto::tonic::trace::v1::span::SpanKind::Client as i32;
const SPAN_KIND_SERVER: i32 = opentelemetry_proto::tonic::trace::v1::span::SpanKind::Server as i32;
#[test]
fn grpc_client_rpc_span_is_admitted_as_outbound_call() {
let mut span = make_bare_span(
&[7; 8],
vec![
make_kv("rpc.system", "grpc"),
make_kv("rpc.service", "order.v1.OrderService"),
make_kv("rpc.method", "GetOrder"),
],
);
span.kind = SPAN_KIND_CLIENT;
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::HttpOut);
assert_eq!(events[0].target, "order.v1.OrderService/GetOrder");
assert_eq!(events[0].operation, "grpc");
}
#[test]
fn grpc_server_rpc_span_is_not_admitted() {
let mut span = make_bare_span(
&[9; 8],
vec![
make_kv("rpc.system", "grpc"),
make_kv("rpc.service", "order.v1.OrderService"),
make_kv("rpc.method", "GetOrder"),
],
);
span.kind = SPAN_KIND_SERVER;
let req = make_request("order-svc", vec![span]);
assert!(convert_otlp_request(&req).is_empty());
}
#[test]
fn rpc_span_without_service_falls_back_to_span_name() {
let mut span = make_bare_span(&[8; 8], vec![make_kv("rpc.system", "grpc")]);
span.kind = SPAN_KIND_CLIENT;
span.name = "order.v1.OrderService/ListOrders".to_string();
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::HttpOut);
assert_eq!(events[0].target, "order.v1.OrderService/ListOrders");
}
#[test]
fn rpc_span_with_blank_service_and_method_falls_back_to_span_name() {
let mut span = make_bare_span(
&[11; 8],
vec![
make_kv("rpc.system", "grpc"),
make_kv("rpc.service", ""),
make_kv("rpc.method", ""),
],
);
span.kind = SPAN_KIND_CLIENT;
span.name = "health.Check".to_string();
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "health.Check");
}
const SPAN_KIND_PRODUCER: i32 =
opentelemetry_proto::tonic::trace::v1::span::SpanKind::Producer as i32;
const SPAN_KIND_CONSUMER: i32 =
opentelemetry_proto::tonic::trace::v1::span::SpanKind::Consumer as i32;
#[test]
fn producer_messaging_span_is_admitted_as_outbound_call() {
let mut span = make_bare_span(
&[21; 8],
vec![
make_kv("messaging.system", "kafka"),
make_kv("messaging.destination.name", "orders"),
make_kv("messaging.operation.type", "send"),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Messaging);
assert_eq!(events[0].target, "orders");
assert_eq!(events[0].operation, "kafka");
}
#[test]
fn consumer_messaging_span_is_not_admitted() {
let mut span = make_bare_span(
&[22; 8],
vec![
make_kv("messaging.system", "kafka"),
make_kv("messaging.destination.name", "orders"),
make_kv("messaging.operation.type", "process"),
],
);
span.kind = SPAN_KIND_CONSUMER;
let req = make_request("order-svc", vec![span]);
assert!(convert_otlp_request(&req).is_empty());
}
#[test]
fn messaging_span_reads_the_legacy_destination_key() {
let mut span = make_bare_span(
&[23; 8],
vec![
make_kv("messaging.system", "rabbitmq"),
make_kv("messaging.destination", "signature.jobs"),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("signature-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "signature.jobs");
assert_eq!(events[0].operation, "rabbitmq");
}
#[test]
fn messaging_span_without_destination_falls_back_to_span_name() {
let mut span = make_bare_span(
&[24; 8],
vec![
make_kv("messaging.system", "pulsar"),
make_kv("messaging.destination.name", ""),
],
);
span.kind = SPAN_KIND_PRODUCER;
span.name = "persistent://tenant/ns/topic publish".to_string();
let req = make_request("doc-worker", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "persistent://tenant/ns/topic publish");
}
#[test]
fn messaging_span_keeps_the_message_body_size() {
let mut span = make_bare_span(
&[25; 8],
vec![
make_kv("messaging.system", "kafka"),
make_kv("messaging.destination.name", "orders"),
make_int_kv("messaging.message.body.size", 4096),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].response_size_bytes, Some(4096));
}
#[test]
fn sql_span_under_a_producer_still_reads_as_sql() {
let mut span = make_bare_span(
&[26; 8],
vec![
make_kv("db.system", "postgresql"),
make_kv("db.statement", "SELECT id FROM outbox WHERE sent = false"),
make_kv("messaging.system", "kafka"),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
}
#[test]
fn producer_span_with_rpc_attributes_still_reads_as_messaging() {
let mut span = make_bare_span(
&[27; 8],
vec![
make_kv("rpc.system", "aws-api"),
make_kv("rpc.service", "Sqs"),
make_kv("rpc.method", "SendMessage"),
make_kv("messaging.system", "aws.sqs"),
make_kv("messaging.destination.name", "signature-jobs"),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("signature-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Messaging);
assert_eq!(events[0].target, "signature-jobs");
}
#[test]
fn messaging_span_carrying_an_outbound_url_reads_as_http() {
let mut span = make_bare_span(
&[28; 8],
vec![
make_kv(
"url.full",
"https://sqs.eu-west-3.amazonaws.com/123456789012/jobs",
),
make_kv("messaging.system", "aws.sqs"),
make_kv("messaging.destination.name", "jobs"),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("signature-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::HttpOut);
}
#[test]
fn blank_destination_name_falls_back_to_the_legacy_key() {
let mut span = make_bare_span(
&[29; 8],
vec![
make_kv("messaging.system", "kafka"),
make_kv("messaging.destination.name", " "),
make_kv("messaging.destination", "orders"),
],
);
span.kind = SPAN_KIND_PRODUCER;
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "orders");
}
const PRODUCER_TRACE: [u8; 16] = [0xab; 16];
fn make_linked_consumer_pair(link_trace: Vec<u8>, consumer_kind: i32) -> Vec<Span> {
let mut consumer = make_bare_span(
&[31; 8],
vec![
make_kv("messaging.system", "kafka"),
make_kv("messaging.destination.name", "orders"),
],
);
consumer.kind = consumer_kind;
consumer.links = vec![opentelemetry_proto::tonic::trace::v1::span::Link {
trace_id: link_trace,
span_id: vec![7; 8],
..Default::default()
}];
let mut child = make_bare_span(
&[32; 8],
vec![
make_kv("db.system", "postgresql"),
make_kv("db.statement", "SELECT id FROM orders WHERE id = 1"),
],
);
child.parent_span_id = vec![31; 8];
vec![consumer, child]
}
#[test]
fn consumer_link_propagates_to_child_io_span() {
let req = make_request(
"order-svc",
make_linked_consumer_pair(PRODUCER_TRACE.to_vec(), SPAN_KIND_CONSUMER),
);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(events[0].link_trace_id.as_deref(), Some(&*"ab".repeat(16)));
}
fn make_sibling_consumer_trace(link_trace: Vec<u8>) -> Vec<Span> {
let root = make_bare_span(&[24; 8], vec![]);
let mut receive = make_bare_span(
&[43; 8],
vec![
make_kv("messaging.system", "kafka"),
make_kv("messaging.destination.name", "orders"),
],
);
receive.kind = SPAN_KIND_CONSUMER;
receive.parent_span_id = vec![24; 8];
receive.links = vec![opentelemetry_proto::tonic::trace::v1::span::Link {
trace_id: link_trace,
span_id: vec![7; 8],
..Default::default()
}];
let mut sql = make_bare_span(
&[31; 8],
vec![
make_kv("db.system", "postgresql"),
make_kv("db.statement", "INSERT INTO orders (id) VALUES (1)"),
],
);
sql.parent_span_id = vec![24; 8];
vec![root, receive, sql]
}
#[test]
fn sibling_consumer_link_reaches_the_io_span() {
let req = make_request(
"accounting",
make_sibling_consumer_trace(PRODUCER_TRACE.to_vec()),
);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL sibling is admitted");
assert_eq!(
sql.link_trace_id.as_deref(),
Some(&*"ab".repeat(16)),
"the receive span is a sibling, not an ancestor"
);
}
#[test]
fn a_sibling_that_started_before_the_receive_is_not_triggered_by_it() {
let mut spans = make_sibling_consumer_trace(PRODUCER_TRACE.to_vec());
spans[2].start_time_unix_nano = spans[1].start_time_unix_nano - 1;
let req = make_request("accounting", spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(sql.link_trace_id, None);
}
#[test]
fn a_handler_between_receive_and_the_io_span_keeps_the_link() {
let mut spans = make_sibling_consumer_trace(PRODUCER_TRACE.to_vec());
let mut handler = make_bare_span(&[52; 8], vec![]);
handler.parent_span_id = vec![24; 8];
spans[2].parent_span_id = vec![52; 8];
spans.push(handler);
let req = make_request("accounting", spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(
sql.link_trace_id.as_deref(),
Some(&*"ab".repeat(16)),
"the sibling lookup must retry at each ancestor"
);
}
#[test]
fn a_handler_that_started_before_the_receive_shields_its_children() {
let mut spans = make_sibling_consumer_trace(PRODUCER_TRACE.to_vec());
let mut handler = make_bare_span(&[52; 8], vec![]);
handler.parent_span_id = vec![24; 8];
handler.start_time_unix_nano = spans[1].start_time_unix_nano - 1;
spans[2].parent_span_id = vec![52; 8];
spans[2].start_time_unix_nano = spans[1].start_time_unix_nano + 1;
spans.push(handler);
let req = make_request("accounting", spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(sql.link_trace_id, None);
}
#[test]
fn work_follows_the_nearest_preceding_receive_not_the_first_in_the_payload() {
let mut spans = make_sibling_consumer_trace(PRODUCER_TRACE.to_vec());
let mut second = spans[1].clone();
second.span_id = vec![44; 8];
second.start_time_unix_nano = spans[1].start_time_unix_nano + 100;
second.links[0].trace_id = vec![0xcd; 16];
spans[2].start_time_unix_nano = second.start_time_unix_nano + 1;
spans.insert(1, second);
let events = convert_otlp_request(&make_request("accounting", spans.clone()));
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(
sql.link_trace_id.as_deref(),
Some(&*"cd".repeat(16)),
"the later receive is the one that triggered this work"
);
spans.reverse();
let events = convert_otlp_request(&make_request("accounting", spans));
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(sql.link_trace_id.as_deref(), Some(&*"cd".repeat(16)));
}
#[test]
fn an_all_zero_parent_id_does_not_pair_two_roots() {
let mut spans = make_sibling_consumer_trace(PRODUCER_TRACE.to_vec());
spans[1].parent_span_id = vec![0; 8];
spans[2].parent_span_id = vec![0; 8];
let req = make_request("accounting", spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(sql.link_trace_id, None);
}
#[test]
fn a_sibling_consumer_from_another_trace_is_ignored() {
let mut spans = make_sibling_consumer_trace(PRODUCER_TRACE.to_vec());
spans[1].trace_id = vec![9; 16];
let req = make_request("accounting", spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL span is admitted");
assert_eq!(sql.link_trace_id, None);
}
#[test]
fn link_on_a_non_consumer_ancestor_is_ignored() {
let mut spans = make_linked_consumer_pair(PRODUCER_TRACE.to_vec(), SPAN_KIND_CLIENT);
let mut decoy = make_bare_span(&[41; 8], vec![]);
decoy.trace_id = vec![9; 16];
decoy.kind = SPAN_KIND_CONSUMER;
decoy.links = vec![opentelemetry_proto::tonic::trace::v1::span::Link {
trace_id: PRODUCER_TRACE.to_vec(),
span_id: vec![7; 8],
..Default::default()
}];
spans.push(decoy);
let req = make_request("order-svc", spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("the SQL child is admitted");
assert_eq!(sql.link_trace_id, None);
}
#[test]
fn link_less_consumer_does_not_mask_a_linked_one_above() {
let mut receive = make_bare_span(&[51; 8], vec![]);
receive.kind = SPAN_KIND_CONSUMER;
receive.links = vec![opentelemetry_proto::tonic::trace::v1::span::Link {
trace_id: PRODUCER_TRACE.to_vec(),
span_id: vec![7; 8],
..Default::default()
}];
let mut process = make_bare_span(&[52; 8], vec![]);
process.kind = SPAN_KIND_CONSUMER;
process.parent_span_id = vec![51; 8];
let mut child = make_bare_span(
&[53; 8],
vec![
make_kv("db.system", "postgresql"),
make_kv("db.statement", "SELECT 1"),
],
);
child.parent_span_id = vec![52; 8];
let req = make_request("order-svc", vec![receive, process, child]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(
events[0].link_trace_id.as_deref(),
Some(&*"ab".repeat(16)),
"the walk must pass the link-less process span"
);
}
#[test]
fn same_trace_link_is_dropped() {
let req = make_request(
"order-svc",
make_linked_consumer_pair(vec![1; 16], SPAN_KIND_CONSUMER),
);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].link_trace_id, None);
}
#[test]
fn all_zero_link_trace_id_is_dropped() {
let req = make_request(
"order-svc",
make_linked_consumer_pair(vec![0; 16], SPAN_KIND_CONSUMER),
);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].link_trace_id, None);
}
#[test]
fn parent_span_url_full_used_when_neither_route_nor_url_present() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "HandleRequest".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("url.full", "http://order-svc/api/orders/42")],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "http://order-svc/api/orders/42");
}
#[test]
fn parent_url_query_string_and_userinfo_stripped_from_source_endpoint() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "HandleRequest".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv(
"http.url",
"https://user:pass@order-svc/oauth/callback?code=SECRET",
)],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "https://order-svc/oauth/callback");
}
#[test]
fn missing_parent_falls_back() {
let child = make_sql_span(
&[1; 16],
&[20; 8],
&[99; 8], "SELECT * FROM order_item WHERE order_id = 42",
1_720_621_921_000_000_000,
1_720_621_921_001_200_000,
);
let req = make_request("order-svc", vec![child]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].source.endpoint, "unknown");
assert_eq!(events[0].source.method, "db.query");
}
#[test]
fn trace_id_hex_encoding() {
let trace_bytes: Vec<u8> = (0..16).collect();
assert_eq!(
bytes_to_hex(&trace_bytes),
"000102030405060708090a0b0c0d0e0f"
);
}
#[test]
fn timestamp_nanos_to_iso8601() {
let nanos: u64 = 1_720_621_921_123_000_000;
let iso = nanos_to_iso8601(nanos);
assert_eq!(iso, "2024-07-10T14:32:01.123Z");
}
#[test]
fn timestamp_epoch_zero() {
assert_eq!(nanos_to_iso8601(0), "1970-01-01T00:00:00.000Z");
}
#[test]
fn duration_calculation() {
let span = make_sql_span(
&[1; 16],
&[2; 8],
&[],
"SELECT 1",
1_000_000_000, 1_002_500_000, );
let req = make_request("test", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events[0].duration_us, 2500);
}
#[test]
fn status_code_extraction() {
let span = make_http_span(
&[1; 16],
&[3; 8],
&[],
"http://svc/api/health",
"GET",
404,
1_000_000_000,
1_001_000_000,
);
let req = make_request("test", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events[0].status_code, Some(404));
}
#[test]
fn service_name_from_resource() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request("my-service", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(&*events[0].service, "my-service");
}
#[test]
fn span_with_both_db_and_http_prefers_sql() {
use opentelemetry_proto::tonic::common::v1::{AnyValue, KeyValue, any_value};
let mut span = make_sql_span(
&[1; 16],
&[2; 8],
&[],
"SELECT 1",
1_000_000_000,
1_001_000_000,
);
span.attributes.push(KeyValue {
key: "http.url".to_string(),
value: Some(AnyValue {
value: Some(any_value::Value::StringValue("http://svc/api".to_string())),
}),
..Default::default()
});
let req = make_request("test", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events[0].event_type, EventType::Sql);
}
#[test]
fn clock_skew_duration_is_zero() {
let span = make_sql_span(
&[1; 16],
&[2; 8],
&[],
"SELECT 1",
2_000_000_000, 1_000_000_000, );
let req = make_request("test", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events[0].duration_us, 0);
}
#[test]
fn bytes_to_hex_empty() {
assert_eq!(bytes_to_hex(&[]), "");
}
#[test]
fn bytes_to_hex_all_values() {
assert_eq!(bytes_to_hex(&[0x00, 0xff, 0xab]), "00ffab");
}
#[test]
fn nanos_to_iso8601_leap_year() {
let nanos: u64 = 1_709_164_800_000_000_000;
let iso = nanos_to_iso8601(nanos);
assert_eq!(iso, "2024-02-29T00:00:00.000Z");
}
#[test]
fn empty_trace_id_produces_empty_hex() {
assert_eq!(bytes_to_hex(&[]), "");
}
#[test]
fn short_span_id_produces_short_hex() {
assert_eq!(bytes_to_hex(&[0xab]), "ab");
}
#[test]
fn missing_service_name_defaults_to_unknown() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![], ..Default::default()
}),
scope_spans: vec![ScopeSpans {
spans: vec![span],
..Default::default()
}],
..Default::default()
}],
};
let events = convert_otlp_request(&req);
assert_eq!(&*events[0].service, "unknown");
}
#[test]
fn no_resource_defaults_to_unknown_service() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: None,
scope_spans: vec![ScopeSpans {
spans: vec![span],
..Default::default()
}],
..Default::default()
}],
};
let events = convert_otlp_request(&req);
assert_eq!(&*events[0].service, "unknown");
}
fn make_request_with_resource_attrs(
attrs: Vec<KeyValue>,
spans: Vec<Span>,
) -> ExportTraceServiceRequest {
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: attrs,
..Default::default()
}),
scope_spans: vec![ScopeSpans {
spans,
..Default::default()
}],
..Default::default()
}],
}
}
#[test]
fn cloud_region_extracted_from_resource_attributes() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request_with_resource_attrs(
vec![
make_kv("service.name", "order-svc"),
make_kv("cloud.region", "eu-west-3"),
],
vec![span],
);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].cloud_region.as_deref(), Some("eu-west-3"));
}
#[test]
fn cloud_region_falls_back_to_span_attribute() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
span.attributes.push(make_kv("cloud.region", "us-east-1"));
let req =
make_request_with_resource_attrs(vec![make_kv("service.name", "order-svc")], vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].cloud_region.as_deref(), Some("us-east-1"));
}
#[test]
fn cloud_region_resource_wins_over_span() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
span.attributes.push(make_kv("cloud.region", "us-east-1"));
let req = make_request_with_resource_attrs(
vec![
make_kv("service.name", "order-svc"),
make_kv("cloud.region", "eu-west-3"),
],
vec![span],
);
let events = convert_otlp_request(&req);
assert_eq!(events[0].cloud_region.as_deref(), Some("eu-west-3"));
}
#[test]
fn no_cloud_region_yields_none() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert!(events[0].cloud_region.is_none());
}
#[test]
fn cloud_region_with_space_is_sanitized_to_none() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request_with_resource_attrs(
vec![
make_kv("service.name", "order-svc"),
make_kv("cloud.region", "eu west 3"),
],
vec![span],
);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert!(
events[0].cloud_region.is_none(),
"region with space must be sanitized to None"
);
}
#[test]
fn oversized_cloud_region_is_sanitized_to_none() {
let long_region = "a".repeat(65);
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request_with_resource_attrs(
vec![
make_kv("service.name", "order-svc"),
make_kv("cloud.region", &long_region),
],
vec![span],
);
let events = convert_otlp_request(&req);
assert!(events[0].cloud_region.is_none());
}
#[test]
fn cloud_region_with_control_char_is_sanitized_to_none() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request_with_resource_attrs(
vec![
make_kv("service.name", "order-svc"),
make_kv("cloud.region", "eu-west-3\n2026-04-07 ERROR fake"),
],
vec![span],
);
let events = convert_otlp_request(&req);
assert!(events[0].cloud_region.is_none());
}
#[test]
fn cloud_region_span_level_fallback_also_sanitized() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
span.attributes.push(make_kv("cloud.region", "bad region!"));
let req = make_request("order-svc", vec![span]);
let events = convert_otlp_request(&req);
assert!(events[0].cloud_region.is_none());
}
fn scoped_request(service: &str, scoped: Vec<(&str, Vec<Span>)>) -> ExportTraceServiceRequest {
use opentelemetry_proto::tonic::common::v1::InstrumentationScope;
ExportTraceServiceRequest {
resource_spans: vec![ResourceSpans {
resource: Some(Resource {
attributes: vec![make_kv("service.name", service)],
..Default::default()
}),
scope_spans: scoped
.into_iter()
.map(|(name, spans)| ScopeSpans {
scope: Some(InstrumentationScope {
name: name.to_string(),
..Default::default()
}),
spans,
..Default::default()
})
.collect(),
..Default::default()
}],
}
}
#[test]
fn instrumentation_scope_captured_from_leaf_only() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = scoped_request("svc", vec![("io.opentelemetry.jdbc", vec![span])]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
let scopes: Vec<&str> = events[0]
.instrumentation_scopes
.iter()
.map(AsRef::as_ref)
.collect();
assert_eq!(scopes, vec!["io.opentelemetry.jdbc"]);
}
#[test]
fn instrumentation_scopes_walk_parent_chain_deduped() {
let http = make_span_with_code_attrs(
&[10; 8],
&[],
"GET /api/orders",
vec![make_kv("http.route", "GET /api/orders")],
);
let spring_data =
make_span_with_code_attrs(&[11; 8], &[10; 8], "OrderRepository.findById", vec![]);
let hibernate = make_span_with_code_attrs(&[12; 8], &[11; 8], "Session.find", vec![]);
let jdbc = make_sql_span(&[1; 16], &[13; 8], &[12; 8], "SELECT 1", 0, 1_000_000);
let req = scoped_request(
"svc",
vec![
("io.opentelemetry.spring-webmvc-6.0", vec![http]),
("io.opentelemetry.spring-data-3.0", vec![spring_data]),
("io.opentelemetry.hibernate-6.0", vec![hibernate]),
("io.opentelemetry.jdbc", vec![jdbc]),
],
);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1, "only the JDBC span yields a SpanEvent");
let scopes: Vec<&str> = events[0]
.instrumentation_scopes
.iter()
.map(AsRef::as_ref)
.collect();
assert_eq!(
scopes,
vec![
"io.opentelemetry.jdbc",
"io.opentelemetry.hibernate-6.0",
"io.opentelemetry.spring-data-3.0",
"io.opentelemetry.spring-webmvc-6.0",
],
"leaf-to-root order, deduplicated"
);
}
#[test]
fn instrumentation_scopes_empty_when_scope_absent() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert!(events[0].instrumentation_scopes.is_empty());
}
#[test]
fn cloud_region_empty_string_is_sanitized_to_none() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1000);
let req = make_request_with_resource_attrs(
vec![
make_kv("service.name", "order-svc"),
make_kv("cloud.region", ""),
],
vec![span],
);
let events = convert_otlp_request(&req);
assert!(events[0].cloud_region.is_none());
}
fn make_span_with_code_attrs(
span_id: &[u8],
parent_span_id: &[u8],
name: &str,
code_attrs: Vec<KeyValue>,
) -> Span {
Span {
trace_id: vec![1; 16],
span_id: span_id.to_vec(),
parent_span_id: parent_span_id.to_vec(),
name: name.to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000,
attributes: code_attrs,
..Default::default()
}
}
#[test]
fn code_attrs_inherited_from_immediate_parent() {
let parent = make_span_with_code_attrs(
&[10; 8],
&[],
"GET /api/orders",
vec![
make_kv("http.route", "GET /api/orders"),
make_kv("code.namespace", "com.foo.OrderController"),
make_kv("code.function", "list"),
],
);
let child = make_sql_span(
&[1; 16],
&[20; 8],
&[10; 8],
"SELECT * FROM orders",
0,
1_000_000,
);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(
events[0].code_namespace.as_deref(),
Some("com.foo.OrderController")
);
assert_eq!(events[0].code_function.as_deref(), Some("list"));
}
#[test]
fn code_attrs_inherited_from_grandparent() {
let http = make_span_with_code_attrs(
&[10; 8],
&[],
"GET /api/orders",
vec![make_kv("http.route", "GET /api/orders")],
);
let service = make_span_with_code_attrs(
&[11; 8],
&[10; 8],
"OrderService.list",
vec![
make_kv("code.namespace", "com.foo.OrderService"),
make_kv("code.function", "list"),
],
);
let hibernate = make_span_with_code_attrs(&[12; 8], &[11; 8], "Hibernate.query", vec![]);
let jdbc = make_sql_span(&[1; 16], &[13; 8], &[12; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![http, service, hibernate, jdbc]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(
events[0].code_namespace.as_deref(),
Some("com.foo.OrderService")
);
}
#[test]
fn code_attrs_orphan_span_returns_none() {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert!(events[0].code_namespace.is_none());
assert!(events[0].code_function.is_none());
}
#[test]
fn code_attrs_max_depth_safety() {
let depth = u8::try_from(CODE_ATTRS_MAX_DEPTH * 2 + 4).unwrap();
let mut spans = Vec::new();
for i in 0..depth {
let id = [i + 1; 8];
let parent = if i == 0 { vec![] } else { vec![i; 8] };
spans.push(make_span_with_code_attrs(
&id,
&parent,
&format!("level.{i}"),
vec![],
));
}
let leaf = make_sql_span(&[1; 16], &[100; 8], &[depth; 8], "SELECT 1", 0, 1_000_000);
spans.push(leaf);
let req = make_request("svc", spans);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert!(events[0].code_namespace.is_none());
}
#[test]
fn code_attrs_self_takes_precedence() {
let parent = make_span_with_code_attrs(
&[10; 8],
&[],
"GET /api/x",
vec![make_kv("code.namespace", "com.parent")],
);
let mut child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
child
.attributes
.push(make_kv("code.namespace", "com.child"));
let req = make_request("svc", vec![parent, child]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].code_namespace.as_deref(), Some("com.child"));
}
#[test]
fn code_attrs_stable_conventions() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes.extend(vec![
make_kv("code.function.name", "com.foo.OrderService.findItems"),
make_kv("code.file.path", "src/main/java/com/foo/OrderService.java"),
make_int_kv("code.line.number", 42),
]);
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(
events[0].code_function.as_deref(),
Some("com.foo.OrderService.findItems")
);
assert_eq!(
events[0].code_filepath.as_deref(),
Some("src/main/java/com/foo/OrderService.java")
);
assert_eq!(events[0].code_lineno, Some(42));
assert_eq!(
events[0].code_namespace.as_deref(),
Some("com.foo.OrderService")
);
}
#[test]
fn code_attrs_php_backslash_namespace_derivation() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes.push(make_kv(
"code.function.name",
"Doctrine\\DBAL\\Driver\\Connection::query",
));
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(
events[0].code_namespace.as_deref(),
Some("Doctrine\\DBAL\\Driver")
);
}
#[test]
fn code_attrs_legacy_conventions_still_work() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes.extend(vec![
make_kv("code.function", "findItems"),
make_kv("code.namespace", "com.foo.OrderService"),
make_kv("code.filepath", "src/OrderService.java"),
make_int_kv("code.lineno", 99),
]);
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].code_function.as_deref(), Some("findItems"));
assert_eq!(
events[0].code_namespace.as_deref(),
Some("com.foo.OrderService")
);
assert_eq!(events[0].code_lineno, Some(99));
}
#[test]
fn code_attrs_legacy_namespace_wins_over_derivation() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes.extend(vec![
make_kv("code.function.name", "com.foo.X.y"),
make_kv("code.namespace", "com.bar.X"),
]);
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].code_namespace.as_deref(), Some("com.bar.X"));
}
#[test]
fn code_attrs_legacy_function_does_not_derive_namespace() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes
.push(make_kv("code.function", "OrderService.findItems"));
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(
events[0].code_function.as_deref(),
Some("OrderService.findItems")
);
assert!(events[0].code_namespace.is_none());
}
#[test]
fn code_attrs_no_dot_in_fq_name() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes.push(make_kv("code.function.name", "main"));
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].code_function.as_deref(), Some("main"));
assert!(events[0].code_namespace.is_none());
}
#[test]
fn java_rules_match_via_derived_namespace() {
let mut span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
span.attributes.push(make_kv(
"code.function.name",
"org.springframework.data.jpa.repository.support.SimpleJpaRepository.findAll",
));
let req = make_request("svc", vec![span]);
let events = convert_otlp_request(&req);
assert_eq!(events.len(), 1);
let ns = events[0]
.code_namespace
.as_deref()
.expect("namespace derived from FQ name");
assert_eq!(
ns,
"org.springframework.data.jpa.repository.support.SimpleJpaRepository"
);
assert!(ns.contains("org.springframework.data.jpa"));
}
#[test]
fn endpoint_falls_back_to_code_frame() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "PurgeNotificationJob.execute".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("code.function", "execute"),
make_kv("code.namespace", "com.foo.scheduler.PurgeNotificationJob"),
],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("web-notification", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(
sql.source.endpoint,
"com.foo.scheduler.PurgeNotificationJob.execute"
);
}
#[test]
fn endpoint_code_frame_resolves_through_two_ancestors() {
let job = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv(
"code.function.name",
"com.foo.repository.NotificationRepository.findByHash",
)],
..Default::default()
};
let hibernate = Span {
trace_id: vec![1; 16],
span_id: vec![20; 8],
parent_span_id: vec![10; 8],
name: "SELECT com.foo.Notification".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 500_000_000,
attributes: vec![],
..Default::default()
};
let jdbc = make_sql_span(&[1; 16], &[30; 8], &[20; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("web-notification", vec![job, hibernate, jdbc]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(
sql.source.endpoint,
"com.foo.repository.NotificationRepository.findByHash"
);
}
#[test]
fn endpoint_http_route_wins_over_code_frame() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "POST /api/orders".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("http.route", "POST /api/orders"),
make_kv("code.function", "createOrder"),
make_kv("code.namespace", "com.foo.OrderController"),
],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "POST /api/orders");
}
#[test]
fn endpoint_code_frame_survives_secret_stripping() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("code.function", "execute"),
make_kv("code.namespace", "com.foo.Job"),
],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert!(
sql.source.endpoint.ends_with("execute"),
"method half was stripped: {}",
sql.source.endpoint
);
}
#[test]
fn endpoint_code_frame_separates_two_jobs_sharing_a_statement() {
let job_endpoint = |ns: &str, span_seed: u8| {
let parent = Span {
trace_id: vec![span_seed; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("code.function", "execute"),
make_kv("code.namespace", ns),
],
..Default::default()
};
let child = make_sql_span(
&[span_seed; 16],
&[20; 8],
&[10; 8],
"SELECT * FROM notification WHERE id = 1",
0,
1_000_000,
);
let req = make_request("web-notification", vec![parent, child]);
convert_otlp_request(&req)
.into_iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present")
.source
.endpoint
};
let purge = job_endpoint("com.foo.scheduler.PurgeNotificationJob", 1);
let resend = job_endpoint("com.foo.scheduler.ResendNotificationJob", 2);
assert_ne!(purge, resend);
assert_ne!(purge, "unknown");
}
#[test]
fn endpoint_http_route_resolves_through_ancestors() {
let server = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "POST /api/orders".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("http.route", "POST /api/orders")],
..Default::default()
};
let repository = Span {
trace_id: vec![1; 16],
span_id: vec![20; 8],
parent_span_id: vec![10; 8],
name: "OrderRepository.findItems".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 900_000_000,
attributes: vec![
make_kv("code.function", "findItems"),
make_kv("code.namespace", "com.foo.OrderRepository"),
],
..Default::default()
};
let jdbc = make_sql_span(&[1; 16], &[30; 8], &[20; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![server, repository, jdbc]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(sql.source.endpoint, "POST /api/orders");
}
#[test]
fn endpoint_code_frame_resolves_past_a_nameless_leaf_frame() {
let job = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("code.function", "execute"),
make_kv("code.namespace", "com.foo.PurgeJob"),
],
..Default::default()
};
let mut jdbc = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
jdbc.attributes
.push(make_kv("code.filepath", "Driver.java"));
let req = make_request("svc", vec![job, jdbc]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(sql.source.endpoint, "com.foo.PurgeJob.execute");
}
#[test]
fn endpoint_stays_unknown_without_a_usable_frame() {
let parent = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("code.filepath", "Job.java")],
..Default::default()
};
let child = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("svc", vec![parent, child]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql child event present");
assert_eq!(sql.source.endpoint, "unknown");
}
#[test]
fn endpoint_ignores_a_span_carrying_no_span_id() {
let nameless = Span {
trace_id: vec![1; 16],
span_id: vec![],
parent_span_id: vec![],
name: "malformed".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("http.url", "https://attacker.example/x")],
..Default::default()
};
let root_sql = make_sql_span(&[1; 16], &[20; 8], &[], "SELECT 1", 0, 1_000_000);
let req = make_request("svc", vec![nameless, root_sql]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql event present");
assert_eq!(sql.source.endpoint, "unknown");
}
#[test]
fn endpoint_skips_an_outbound_client_url() {
let server = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "POST /checkout".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("http.route", "POST /checkout")],
..Default::default()
};
let client = Span {
trace_id: vec![1; 16],
span_id: vec![20; 8],
parent_span_id: vec![10; 8],
name: "POST /v1/pay".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 900_000_000,
kind: SPAN_KIND_CLIENT,
attributes: vec![make_kv("url.full", "https://api.partner.com/v1/pay")],
..Default::default()
};
let jdbc = make_sql_span(&[1; 16], &[30; 8], &[20; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("checkout-svc", vec![server, client, jdbc]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(sql.source.endpoint, "POST /checkout");
}
#[test]
fn endpoint_walks_past_a_blank_http_route() {
let server = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "POST /api/orders".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("http.route", "POST /api/orders")],
..Default::default()
};
let filter = Span {
trace_id: vec![1; 16],
span_id: vec![20; 8],
parent_span_id: vec![10; 8],
name: "filter".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 900_000_000,
attributes: vec![make_kv("http.route", " ")],
..Default::default()
};
let jdbc = make_sql_span(&[1; 16], &[30; 8], &[20; 8], "SELECT 1", 0, 1_000_000);
let req = make_request("order-svc", vec![server, filter, jdbc]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(sql.source.endpoint, "POST /api/orders");
}
#[test]
fn endpoint_code_frame_keeps_the_outermost_named_frame() {
let job_endpoint = |job_ns: &str, seed: u8| {
let job = Span {
trace_id: vec![seed; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("code.function", "execute"),
make_kv("code.namespace", job_ns),
],
..Default::default()
};
let repository = Span {
trace_id: vec![seed; 16],
span_id: vec![20; 8],
parent_span_id: vec![10; 8],
name: "repo".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 900_000_000,
attributes: vec![make_kv(
"code.function.name",
"com.foo.repository.NotificationRepository.findByHash",
)],
..Default::default()
};
let jdbc = make_sql_span(&[seed; 16], &[30; 8], &[20; 8], "SELECT 1", 0, 1_000_000);
convert_otlp_request(&make_request(
"web-notification",
vec![job, repository, jdbc],
))
.into_iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present")
.source
.endpoint
};
let purge = job_endpoint("com.foo.scheduler.PurgeJob", 1);
let resend = job_endpoint("com.foo.scheduler.ResendJob", 2);
assert_eq!(purge, "com.foo.scheduler.PurgeJob.execute");
assert_ne!(purge, resend);
}
#[test]
fn endpoint_walk_crosses_resource_block_boundaries() {
let server = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "POST /api/fault/redundant-http".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![make_kv("http.route", "POST /api/fault/redundant-http")],
..Default::default()
};
let client = Span {
trace_id: vec![1; 16],
span_id: vec![20; 8],
parent_span_id: vec![10; 8],
name: "GET /api/payments/history".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 900_000_000,
kind: SPAN_KIND_CLIENT,
attributes: vec![make_kv(
"url.full",
"http://localhost:8090/api/payments/history",
)],
..Default::default()
};
let jdbc = make_sql_span(&[1; 16], &[30; 8], &[20; 8], "SELECT 1", 0, 1_000_000);
let mut req = make_request("nest-svc", vec![client, jdbc]);
req.resource_spans
.extend(make_request("nest-svc", vec![server]).resource_spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(sql.source.endpoint, "POST /api/fault/redundant-http");
}
#[test]
fn ancestor_walk_stops_at_the_service_boundary() {
let caller = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "GET /checkout".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("http.route", "GET /checkout"),
make_kv("code.function", "checkout"),
make_kv("code.namespace", "com.foo.frontend.CheckoutController"),
],
..Default::default()
};
let jdbc = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
let mut req = make_request("frontend", vec![caller]);
req.resource_spans
.extend(make_request("orders", vec![jdbc]).resource_spans);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql event present");
assert_eq!(sql.service.as_ref(), "orders");
assert_eq!(sql.source.endpoint, "unknown");
assert!(
sql.code_function.is_none() && sql.code_namespace.is_none(),
"code frame must not come from another service, got {:?} / {:?}",
sql.code_function,
sql.code_namespace
);
}
#[test]
fn endpoint_walks_past_an_unusable_leaf_frame() {
let job = Span {
trace_id: vec![1; 16],
span_id: vec![10; 8],
parent_span_id: vec![],
name: "job".to_string(),
start_time_unix_nano: 0,
end_time_unix_nano: 1_000_000_000,
attributes: vec![
make_kv("code.function", "execute"),
make_kv("code.namespace", "com.foo.PurgeJob"),
],
..Default::default()
};
let mut jdbc = make_sql_span(&[1; 16], &[20; 8], &[10; 8], "SELECT 1", 0, 1_000_000);
jdbc.attributes
.push(make_kv("code.function", "executeQuery"));
let req = make_request("svc", vec![job, jdbc]);
let events = convert_otlp_request(&req);
let sql = events
.iter()
.find(|e| e.event_type == EventType::Sql)
.expect("sql leaf event present");
assert_eq!(sql.source.endpoint, "com.foo.PurgeJob.execute");
}
#[cfg(feature = "daemon")]
mod http_handler {
use super::*;
use axum::body::Body;
use axum::http::{Request, StatusCode, header};
use flate2::Compression;
use flate2::write::GzEncoder;
use prost::Message;
use std::io::Write;
use tokio::sync::mpsc;
use tower::ServiceExt;
fn build_minimal_request_bytes() -> Vec<u8> {
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = make_request("svc", vec![span]);
req.encode_to_vec()
}
fn gzip(body: &[u8]) -> Vec<u8> {
let mut encoder = GzEncoder::new(Vec::new(), Compression::default());
encoder.write_all(body).expect("gzip encode");
encoder.finish().expect("gzip finish")
}
fn unsupported_media_type_request() -> Request<Body> {
Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/json")
.body(Body::from(Vec::<u8>::new()))
.expect("build request")
}
#[tokio::test]
async fn otlp_http_accepts_gzip_request() {
let (tx, mut rx) = mpsc::channel(8);
let router = otlp_http_router(tx, 1_048_576, None);
let body = build_minimal_request_bytes();
let gzipped = gzip(&body);
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.header(header::CONTENT_ENCODING, "gzip")
.body(Body::from(gzipped))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::OK);
let events = rx.recv().await.expect("event batch sent");
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT 1");
}
#[tokio::test]
async fn otlp_http_accepts_uncompressed_request() {
let (tx, mut rx) = mpsc::channel(8);
let router = otlp_http_router(tx, 1_048_576, None);
let body = build_minimal_request_bytes();
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.body(Body::from(body))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::OK);
let events = rx.recv().await.expect("event batch sent");
assert_eq!(events.len(), 1);
}
#[tokio::test]
async fn otlp_http_rejects_unsupported_encoding() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let router = otlp_http_router(tx, 1_048_576, None);
let body = build_minimal_request_bytes();
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.header(header::CONTENT_ENCODING, "br")
.body(Body::from(body))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
}
#[tokio::test]
async fn otlp_http_rejects_oversize_compressed_body() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let cap = 256_usize;
let router = otlp_http_router(tx, cap, None);
let payload: Vec<u8> = vec![0u8; 4096];
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.header(header::CONTENT_LENGTH, payload.len())
.body(Body::from(payload))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE);
}
#[tokio::test]
async fn otlp_http_content_type_check_with_gzip() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let router = otlp_http_router(tx, 1_048_576, None);
let body = build_minimal_request_bytes();
let gzipped = gzip(&body);
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/json")
.header(header::CONTENT_ENCODING, "gzip")
.body(Body::from(gzipped))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
}
#[tokio::test]
async fn http_handler_records_unsupported_media_type() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let (metrics, sink) = fresh_metrics_sink();
let router = otlp_http_router(tx, 1_048_576, Some(sink));
let response = router
.oneshot(unsupported_media_type_request())
.await
.expect("router runs");
assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["unsupported_media_type"])
.get(),
1
);
}
#[tokio::test]
async fn http_handler_records_parse_error() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let (metrics, sink) = fresh_metrics_sink();
let router = otlp_http_router(tx, 1_048_576, Some(sink));
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.body(Body::from(vec![0xff_u8, 0xff, 0xff, 0xff, 0xff, 0xff]))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["parse_error"])
.get(),
1
);
}
#[tokio::test]
async fn http_handler_records_channel_full() {
let (tx, rx) = mpsc::channel::<Vec<SpanEvent>>(1);
drop(rx);
let (metrics, sink) = fresh_metrics_sink();
let router = otlp_http_router(tx, 1_048_576, Some(sink));
let body = build_minimal_request_bytes();
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.body(Body::from(body))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["channel_full"])
.get(),
1
);
}
#[tokio::test]
async fn http_handler_rejects_under_memory_pressure() {
let (tx, mut rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let (metrics, sink) = fresh_metrics_sink();
metrics.set_memory_high_water(true);
let router = otlp_http_router(tx, 1_048_576, Some(sink));
let body = build_minimal_request_bytes();
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.body(Body::from(body))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["memory_pressure"])
.get(),
1
);
assert!(rx.try_recv().is_err(), "ingest must be halted at the door");
}
#[tokio::test(start_paused = true)]
async fn http_handler_rejects_when_channel_full_but_open() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(1);
tx.try_send(vec![]).expect("fill the only slot");
let (metrics, sink) = fresh_metrics_sink();
let router = otlp_http_router(tx, 1_048_576, Some(sink));
let body = build_minimal_request_bytes();
let req = Request::builder()
.method("POST")
.uri("/v1/traces")
.header(header::CONTENT_TYPE, "application/x-protobuf")
.body(Body::from(body))
.expect("build request");
let response = router.oneshot(req).await.expect("router runs");
assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["channel_full"])
.get(),
1
);
}
#[tokio::test]
async fn http_handler_no_metrics_state_does_not_panic() {
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let router = otlp_http_router(tx, 1_048_576, None);
let response = router
.oneshot(unsupported_media_type_request())
.await
.expect("router runs");
assert_eq!(response.status(), StatusCode::UNSUPPORTED_MEDIA_TYPE);
}
#[tokio::test]
async fn grpc_handler_records_channel_full() {
use opentelemetry_proto::tonic::collector::trace::v1::trace_service_server::TraceService;
let (tx, rx) = mpsc::channel::<Vec<SpanEvent>>(1);
drop(rx);
let metrics = Arc::new(MetricsState::new());
let svc = OtlpGrpcService::new(tx, Some(metrics.clone()));
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = tonic::Request::new(make_request("svc", vec![span]));
let result = svc.export(req).await;
assert_eq!(
result.expect_err("closed channel rejects").code(),
tonic::Code::Internal,
"shutdown (closed channel) is a genuine internal state"
);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["channel_full"])
.get(),
1
);
}
#[tokio::test]
async fn grpc_handler_rejects_under_memory_pressure() {
use opentelemetry_proto::tonic::collector::trace::v1::trace_service_server::TraceService;
let (tx, mut rx) = mpsc::channel::<Vec<SpanEvent>>(8);
let metrics = Arc::new(MetricsState::new());
metrics.set_memory_high_water(true);
let svc = OtlpGrpcService::new(tx, Some(metrics.clone()));
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = tonic::Request::new(make_request("svc", vec![span]));
assert_eq!(
svc.export(req)
.await
.expect_err("memory pressure rejects")
.code(),
tonic::Code::Unavailable,
);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["memory_pressure"])
.get(),
1
);
assert!(rx.try_recv().is_err(), "ingest must be halted at the door");
}
#[tokio::test(start_paused = true)]
async fn grpc_handler_returns_unavailable_when_channel_full_but_open() {
use opentelemetry_proto::tonic::collector::trace::v1::trace_service_server::TraceService;
let (tx, _rx) = mpsc::channel::<Vec<SpanEvent>>(1);
tx.try_send(vec![]).expect("fill the only slot");
let metrics = Arc::new(MetricsState::new());
let svc = OtlpGrpcService::new(tx, Some(metrics.clone()));
let span = make_sql_span(&[1; 16], &[2; 8], &[], "SELECT 1", 0, 1_000_000);
let req = tonic::Request::new(make_request("svc", vec![span]));
let result = svc.export(req).await;
assert_eq!(
result.expect_err("full channel rejects").code(),
tonic::Code::Unavailable
);
assert_eq!(
metrics
.otlp_rejected_total
.with_label_values(&["channel_full"])
.get(),
1
);
}
}
fn stitch_span(
span_id: &[u8],
parent_span_id: &[u8],
attributes: Vec<KeyValue>,
start_ns: u64,
end_ns: u64,
) -> Span {
Span {
trace_id: vec![1; 16],
span_id: span_id.to_vec(),
parent_span_id: parent_span_id.to_vec(),
name: "stitch.execute".to_string(),
start_time_unix_nano: start_ns,
end_time_unix_nano: end_ns,
attributes,
..Default::default()
}
}
fn php_split_query(
idx: u8,
parent_span_id: &[u8],
statement: &str,
start_ns: u64,
duration_ns: u64,
) -> (Vec<Span>, Vec<Span>) {
let sid = |tag: u8| vec![idx, tag, 0, 0, 0, 0, 0, 0];
let doctrine_stmt = || {
vec![
make_kv("db.query.text", statement),
make_kv("db.operation.name", "prepare"),
]
};
let pdo_stmt = || {
vec![
make_kv("db.query.text", statement),
make_kv("db.system.name", "postgresql"),
make_kv("db.namespace", "shop"),
]
};
let pdo_sys = || {
vec![
make_kv("db.system.name", "postgresql"),
make_kv("db.namespace", "shop"),
]
};
let doctrine_prepare = stitch_span(
&sid(1),
parent_span_id,
doctrine_stmt(),
start_ns,
start_ns + 500_000,
);
let pdo_prepare = stitch_span(
&sid(2),
&sid(1),
pdo_stmt(),
start_ns + 100_000,
start_ns + 300_000,
);
let doctrine_execute = stitch_span(
&sid(3),
parent_span_id,
vec![make_kv("db.operation.name", "execute")],
start_ns + 600_000,
start_ns + 600_000 + duration_ns,
);
let pdo_execute = stitch_span(
&sid(4),
&sid(3),
pdo_sys(),
start_ns + 700_000,
start_ns + 700_000 + duration_ns,
);
(
vec![doctrine_prepare, doctrine_execute],
vec![pdo_prepare, pdo_execute],
)
}
#[test]
fn php_split_query_stitches_one_event_with_execute_duration() {
let start = 1_720_621_921_000_000_000u64;
let root = stitch_span(
&[9, 9, 0, 0, 0, 0, 0, 0],
&[],
vec![make_kv("http.route", "POST /api/orders")],
start,
start + 12_700_000_000,
);
let (doctrine, pdo) = php_split_query(
1,
&[9, 9, 0, 0, 0, 0, 0, 0],
"SELECT pg_sleep(0.6), orders.id FROM orders ORDER BY orders.id OFFSET 3 LIMIT 1",
start,
602_000_000,
);
let req = scoped_request(
"order-svc",
vec![
("io.opentelemetry.contrib.php.symfony", vec![root]),
("io.opentelemetry.contrib.php.doctrine", doctrine),
("io.opentelemetry.contrib.php.pdo", pdo),
],
);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1, "one event per logical query");
assert_eq!(events[0].event_type, EventType::Sql);
assert_eq!(
events[0].target,
"SELECT pg_sleep(0.6), orders.id FROM orders ORDER BY orders.id OFFSET 3 LIMIT 1"
);
assert_eq!(events[0].operation, "sql");
assert_eq!(
events[0].duration_us, 602_000,
"duration from the execute span"
);
assert_eq!(
events[0].span_id, "0103000000000000",
"event carried by the doctrine execute span"
);
let scopes: Vec<&str> = events[0]
.instrumentation_scopes
.iter()
.map(AsRef::as_ref)
.collect();
assert_eq!(
scopes,
vec![
"io.opentelemetry.contrib.php.doctrine",
"io.opentelemetry.contrib.php.symfony",
],
"doctrine scope preserved for framework tagging"
);
assert_eq!(
stats,
SpanConversionStats {
received: 5,
filtered_not_io: 1,
filtered_missing_db_statement: 0,
filtered_missing_http_url: 0,
filtered_non_sql_datastore: 0,
filtered_merged_db_span: 3,
}
);
}
#[test]
fn php_prepare_once_execute_many_yields_one_event_per_execute() {
let root = make_parent_span(&[10; 8], "GET /api/orders");
let donor = make_sql_span(
&[1; 16],
&[30; 8],
&[10; 8],
"SELECT * FROM orders WHERE id = ?",
1_000_000_000,
1_000_500_000,
);
let sys = || vec![make_kv("db.system", "postgresql")];
let exec1 = stitch_span(&[31; 8], &[10; 8], sys(), 2_000_000_000, 2_600_000_000);
let exec2 = stitch_span(&[32; 8], &[10; 8], sys(), 3_000_000_000, 3_700_000_000);
let exec3 = stitch_span(&[33; 8], &[10; 8], sys(), 4_000_000_000, 4_800_000_000);
let req = make_request("order-svc", vec![root, donor, exec1, exec2, exec3]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 3);
let durations: Vec<u64> = events.iter().map(|e| e.duration_us).collect();
assert_eq!(durations, vec![600_000, 700_000, 800_000]);
for event in &events {
assert_eq!(event.target, "SELECT * FROM orders WHERE id = ?");
}
assert_eq!(stats.filtered_merged_db_span, 1, "donor suppressed once");
assert_eq!(stats.filtered_missing_db_statement, 0);
}
#[test]
fn orphan_without_donor_stays_missing_db_statement() {
let sys = || vec![make_kv("db.system", "postgresql")];
let outer = stitch_span(&[40; 8], &[], sys(), 1_000_000_000, 1_600_000_000);
let inner = stitch_span(&[41; 8], &[40; 8], sys(), 1_100_000_000, 1_500_000_000);
let req = make_request("order-svc", vec![outer, inner]);
let (events, stats) = convert_otlp_request_counted(&req);
assert!(events.is_empty());
assert_eq!(stats.filtered_missing_db_statement, 2);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn laravel_sibling_prepare_and_execute_both_emit_unchanged() {
let root = make_parent_span(&[10; 8], "GET /api/users");
let prepare = make_sql_span(
&[1; 16],
&[50; 8],
&[10; 8],
"SELECT * FROM users WHERE id = ?",
1_000_000_000,
1_000_500_000,
);
let execute = make_sql_span(
&[1; 16],
&[51; 8],
&[10; 8],
"SELECT * FROM users WHERE id = ?",
2_000_000_000,
2_600_000_000,
);
let req = make_request("laravel-svc", vec![root, prepare, execute]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 2);
assert_eq!(events[0].duration_us, 500);
assert_eq!(events[1].duration_us, 600_000);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn nested_duplicate_statement_collapses_to_outermost() {
let outer = make_sql_span(
&[1; 16],
&[60; 8],
&[],
"SELECT * FROM orders WHERE id = ?",
1_000_000_000,
1_700_000_000,
);
let inner = make_sql_span(
&[1; 16],
&[61; 8],
&[60; 8],
"SELECT * FROM orders WHERE id = ?",
1_005_000_000,
1_695_000_000,
);
let req = make_request("order-svc", vec![outer, inner]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].duration_us, 700_000, "outermost span wins");
assert_eq!(stats.filtered_merged_db_span, 1);
}
#[test]
fn nested_donor_with_different_statement_not_collapsed() {
let outer = make_sql_span(
&[1; 16],
&[62; 8],
&[],
"CALL refresh_orders()",
1_000_000_000,
1_700_000_000,
);
let inner = make_sql_span(
&[1; 16],
&[63; 8],
&[62; 8],
"SELECT * FROM orders WHERE id = ?",
1_005_000_000,
1_695_000_000,
);
let req = make_request("order-svc", vec![outer, inner]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 2);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn stitch_picks_nearest_preceding_donor() {
let root = make_parent_span(&[10; 8], "GET /api/things");
let donor_a = make_sql_span(
&[1; 16],
&[70; 8],
&[10; 8],
"SELECT a FROM t",
1_000_000_000,
1_000_500_000,
);
let donor_b = make_sql_span(
&[1; 16],
&[71; 8],
&[10; 8],
"SELECT b FROM t",
3_000_000_000,
3_000_500_000,
);
let orphan = stitch_span(
&[72; 8],
&[10; 8],
vec![make_kv("db.system", "postgresql")],
4_000_000_000,
4_700_000_000,
);
let req = make_request("order-svc", vec![root, donor_a, donor_b, orphan]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 2, "unconsumed donor A still emits");
let stitched = events
.iter()
.find(|e| e.duration_us == 700_000)
.expect("stitched execute event");
assert_eq!(stitched.target, "SELECT b FROM t");
assert_eq!(stats.filtered_merged_db_span, 1);
}
#[test]
fn stitching_never_crosses_traces() {
let donor = make_sql_span(
&[1; 16],
&[80; 8],
&[],
"SELECT 1",
1_000_000_000,
1_500_000_000,
);
let sys = || vec![make_kv("db.system", "postgresql")];
let mut outer = stitch_span(&[81; 8], &[], sys(), 2_000_000_000, 2_600_000_000);
outer.trace_id = vec![2; 16];
let mut inner = stitch_span(&[82; 8], &[81; 8], sys(), 2_100_000_000, 2_500_000_000);
inner.trace_id = vec![2; 16];
let req = make_request("order-svc", vec![donor, outer, inner]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT 1");
assert_eq!(stats.filtered_missing_db_statement, 2);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn non_sql_orphan_not_stitched() {
let root = make_parent_span(&[10; 8], "GET /api/cache");
let redis = stitch_span(
&[90; 8],
&[10; 8],
vec![make_kv("db.system", "redis")],
1_000_000_000,
1_600_000_000,
);
let donor = make_sql_span(
&[1; 16],
&[91; 8],
&[10; 8],
"SELECT 1",
2_000_000_000,
2_000_500_000,
);
let req = make_request("cache-svc", vec![root, redis, donor]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT 1");
assert_eq!(stats.filtered_non_sql_datastore, 1);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn dd_resource_donor_stitches_orphan() {
let root = make_parent_span(&[10; 8], "GET /api/orders");
let donor = stitch_span(
&[95; 8],
&[10; 8],
vec![
make_kv("dd.span.Resource", "SELECT * FROM orders WHERE id = ?"),
make_kv("db.type", "postgres"),
],
1_000_000_000,
1_000_500_000,
);
let orphan = stitch_span(
&[96; 8],
&[10; 8],
vec![make_kv("db.type", "postgres")],
2_000_000_000,
2_600_000_000,
);
let req = make_request("dd-shop", vec![root, donor, orphan]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT * FROM orders WHERE id = ?");
assert_eq!(events[0].operation, "postgresql");
assert_eq!(events[0].duration_us, 600_000);
assert_eq!(stats.filtered_merged_db_span, 1);
}
#[test]
fn orphan_adopts_statement_from_donor_ancestor() {
let donor = make_sql_span(
&[1; 16],
&[100; 8],
&[],
"SELECT * FROM orders WHERE id = ?",
1_000_000_000,
1_700_000_000,
);
let orphan = stitch_span(
&[101; 8],
&[100; 8],
vec![make_kv("db.system", "postgresql")],
1_050_000_000,
1_650_000_000,
);
let req = make_request("order-svc", vec![donor, orphan]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT * FROM orders WHERE id = ?");
assert_eq!(events[0].duration_us, 600_000, "orphan carries the event");
assert_eq!(stats.filtered_merged_db_span, 1);
}
#[test]
fn root_orphan_adopts_statement_from_donor_descendant() {
let orphan = stitch_span(
&[102; 8],
&[],
vec![make_kv("db.system", "postgresql")],
1_000_000_000,
1_700_000_000,
);
let donor = make_sql_span(
&[1; 16],
&[103; 8],
&[102; 8],
"SELECT * FROM orders WHERE id = ?",
1_050_000_000,
1_050_500_000,
);
let req = make_request("order-svc", vec![orphan, donor]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].target, "SELECT * FROM orders WHERE id = ?");
assert_eq!(events[0].duration_us, 700_000, "orphan carries the event");
assert_eq!(stats.filtered_merged_db_span, 1);
}
#[test]
fn orphan_with_only_following_sibling_donor_stays_missing() {
let root = make_parent_span(&[10; 8], "GET /api/things");
let orphan = stitch_span(
&[104; 8],
&[10; 8],
vec![make_kv("db.system", "postgresql")],
1_000_000_000,
1_600_000_000,
);
let donor = make_sql_span(
&[1; 16],
&[105; 8],
&[10; 8],
"SELECT 1",
2_000_000_000,
2_000_500_000,
);
let req = make_request("order-svc", vec![root, orphan, donor]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1, "the donor still emits its own event");
assert_eq!(events[0].target, "SELECT 1");
assert_eq!(events[0].duration_us, 500);
assert_eq!(stats.filtered_missing_db_statement, 1);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn self_parented_donor_still_emits() {
let mut span = make_sql_span(&[1; 16], &[110; 8], &[110; 8], "SELECT 1", 0, 1_000_000);
span.parent_span_id = vec![110; 8];
let req = make_request("order-svc", vec![span]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn mutual_cycle_donors_keep_both_events() {
let a = make_sql_span(
&[1; 16],
&[111; 8],
&[112; 8],
"SELECT 1",
1_000_000_000,
1_001_000_000,
);
let b = make_sql_span(
&[1; 16],
&[112; 8],
&[111; 8],
"SELECT 1",
1_000_000_000,
1_001_000_000,
);
let req = make_request("order-svc", vec![a, b]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 2, "pre-stitch behavior on malformed cycles");
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn connect_span_is_not_stitched() {
let root = make_parent_span(&[10; 8], "GET /api/orders");
let mut connect = stitch_span(
&[120; 8],
&[10; 8],
vec![make_kv("db.system", "postgresql")],
1_000_000_000,
1_300_000_000,
);
connect.name = "pg.connect".to_string();
let donor = make_sql_span(
&[1; 16],
&[121; 8],
&[10; 8],
"SELECT * FROM orders WHERE id = 1",
500_000_000,
500_400_000,
);
let req = make_request("order-svc", vec![root, connect, donor]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1, "the query keeps its own event");
assert_eq!(events[0].duration_us, 400, "query duration untouched");
assert_eq!(stats.filtered_missing_db_statement, 1);
assert_eq!(stats.filtered_merged_db_span, 0);
}
#[test]
fn transaction_wrapper_does_not_swallow_queries() {
let mut wrapper = stitch_span(
&[130; 8],
&[],
vec![make_kv("db.system", "postgresql")],
1_000_000_000,
5_000_000_000,
);
wrapper.name = "PDO::beginTransaction".to_string();
let sys = || vec![make_kv("db.system", "postgresql")];
let d1 = make_sql_span(
&[1; 16],
&[131; 8],
&[130; 8],
"SELECT a FROM t",
1_100_000_000,
1_100_400_000,
);
let e1 = stitch_span(&[132; 8], &[130; 8], sys(), 1_200_000_000, 1_800_000_000);
let d2 = make_sql_span(
&[1; 16],
&[133; 8],
&[130; 8],
"SELECT b FROM t",
2_000_000_000,
2_000_400_000,
);
let e2 = stitch_span(&[134; 8], &[130; 8], sys(), 2_100_000_000, 2_800_000_000);
let req = make_request("order-svc", vec![wrapper, d1, e1, d2, e2]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 2);
let targets: Vec<&str> = events.iter().map(|e| e.target.as_str()).collect();
assert_eq!(targets, vec!["SELECT a FROM t", "SELECT b FROM t"]);
assert_eq!(stats.filtered_missing_db_statement, 1, "the wrapper");
assert_eq!(stats.filtered_merged_db_span, 2);
}
#[test]
fn batch_prepared_statements_pair_off() {
let root = make_parent_span(&[10; 8], "GET /api/things");
let d_a = make_sql_span(
&[1; 16],
&[140; 8],
&[10; 8],
"SELECT a FROM t",
1_000_000_000,
1_000_400_000,
);
let d_b = make_sql_span(
&[1; 16],
&[141; 8],
&[10; 8],
"SELECT b FROM t",
1_100_000_000,
1_100_400_000,
);
let sys = || vec![make_kv("db.system", "postgresql")];
let e1 = stitch_span(&[142; 8], &[10; 8], sys(), 2_000_000_000, 2_600_000_000);
let e2 = stitch_span(&[143; 8], &[10; 8], sys(), 3_000_000_000, 3_700_000_000);
let req = make_request("order-svc", vec![root, d_a, d_b, e1, e2]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 2);
let mut targets: Vec<&str> = events.iter().map(|e| e.target.as_str()).collect();
targets.sort_unstable();
assert_eq!(
targets,
vec!["SELECT a FROM t", "SELECT b FROM t"],
"each donor consumed exactly once"
);
assert_eq!(stats.filtered_merged_db_span, 2);
}
#[test]
fn statement_less_span_without_db_system_needs_a_sibling_donor() {
let root = make_parent_span(&[10; 8], "GET /api/orders");
let wrapper = stitch_span(
&[150; 8],
&[10; 8],
vec![make_kv("db.operation.name", "execute")],
1_000_000_000,
1_600_000_000,
);
let child = make_sql_span(
&[1; 16],
&[151; 8],
&[150; 8],
"SELECT * FROM orders WHERE id = ?",
1_050_000_000,
1_600_000_000,
);
let req = make_request("rails-svc", vec![root, wrapper, child]);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1, "only the child SQL span emits");
assert_eq!(events[0].target, "SELECT * FROM orders WHERE id = ?");
assert_eq!(
events[0].duration_us, 550_000,
"child duration, not the wrapper"
);
assert_eq!(stats.filtered_merged_db_span, 0, "no stitch");
assert_eq!(stats.filtered_not_io, 2);
}
#[test]
fn doctrine_layer_stitches_without_db_system() {
let start = 1_720_621_921_000_000_000u64;
let root = stitch_span(
&[9, 9, 0, 0, 0, 0, 0, 0],
&[],
vec![make_kv("http.route", "POST app_fault_slowsql")],
start,
start + 5_000_000_000,
);
let (doctrine, pdo) = php_split_query(
1,
&[9, 9, 0, 0, 0, 0, 0, 0],
"SELECT pg_sleep(0.6) FROM orders OFFSET 2 LIMIT 1",
start,
700_000_000,
);
let req = scoped_request(
"order-svc",
vec![
("io.opentelemetry.contrib.php.symfony", vec![root]),
("io.opentelemetry.contrib.php.doctrine", doctrine),
("io.opentelemetry.contrib.php.pdo", pdo),
],
);
let (events, stats) = convert_otlp_request_counted(&req);
assert_eq!(events.len(), 1);
assert_eq!(events[0].duration_us, 700_000, "the 700ms execute duration");
assert_eq!(
events[0].target,
"SELECT pg_sleep(0.6) FROM orders OFFSET 2 LIMIT 1"
);
assert!(
events[0]
.instrumentation_scopes
.iter()
.any(|s| s.as_ref() == "io.opentelemetry.contrib.php.doctrine"),
"doctrine scope preserved for php_doctrine tagging"
);
assert_eq!(stats.filtered_merged_db_span, 3);
assert_eq!(stats.filtered_missing_db_statement, 0);
}