mod attribute_processing;
#[cfg(feature = "xray-subsegment-nesting")]
mod document_builder_tree;
#[cfg(feature = "xray-subsegment-nesting")]
use document_builder_tree::DocumentBuilderHeaderTree;
pub(crate) mod error;
mod utils;
use utils::{sanitize_annotation_key, translate_timestamp};
use opentelemetry::{trace::SpanKind, KeyValue, SpanId};
use opentelemetry_sdk::{trace::SpanData, Resource};
use crate::xray_exporter::{
translator::attribute_processing::{
get_annotation, get_any_value,
value_builder::{
AnyValueBuilder, AwsOperationBuilder, AwsXraySdkBuilder, BeanstalkDeploymentIdBuilder,
CauseBuilder, CloudwatchLogGroupBuilder, HttpRequestUrlBuilder,
HttpResponseContentLengthBuilder, SegmentNameBuilder, SegmentOriginBuilder,
SqlUrlBuilder, SubsegmentNamespaceBuilder, ValueBuilder,
},
DispatchTable, SpanAttributeProcessor,
},
types::{
DocumentBuilder, DocumentBuilderType, SegmentDocument, SegmentDocumentBuilder,
SubsegmentDocumentBuilder,
},
};
use error::{Result, TranslationError};
#[derive(Debug)]
pub struct SegmentTranslator {
indexed_attrs: Vec<String>,
index_all_attrs: bool,
metadata_attrs: Vec<String>,
metadata_all_attrs: bool,
log_group_names: Vec<String>,
skip_timestamp_validation: bool,
#[cfg(feature = "xray-subsegment-nesting")]
always_nest_subsegments: bool,
resource: Option<Resource>,
dispatch_table: DispatchTable,
}
impl Default for SegmentTranslator {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug)]
enum AnyDocumentBuilder<'span> {
Segment(SegmentDocumentBuilder<'span>),
Subsegment(SubsegmentDocumentBuilder<'span>),
}
impl<'span> AnyDocumentBuilder<'span> {
fn build(self) -> Result<SegmentDocument<'span>> {
match self {
AnyDocumentBuilder::Segment(builder) => Ok(builder.build()?),
AnyDocumentBuilder::Subsegment(builder) => Ok(builder.build()?),
}
}
}
const ANY_DOCUMENT_BUILDER_PROCESSOR_ID: usize = AnyValueBuilder::count();
impl SegmentTranslator {
pub fn new() -> Self {
let mut st = Self {
indexed_attrs: Default::default(),
index_all_attrs: Default::default(),
metadata_attrs: Default::default(),
metadata_all_attrs: Default::default(),
log_group_names: Default::default(),
skip_timestamp_validation: Default::default(),
#[cfg(feature = "xray-subsegment-nesting")]
always_nest_subsegments: Default::default(),
resource: Default::default(),
dispatch_table: Default::default(),
};
AnyValueBuilder::register_builders(&mut st.dispatch_table);
st.dispatch_table.register::<{
const fn __len<T, const N: usize>(_: &[T; N]) -> usize {
N
}
__len(&AnyDocumentBuilder::HANDLERS)
}, AnyDocumentBuilder>(ANY_DOCUMENT_BUILDER_PROCESSOR_ID);
st
}
pub fn index_all_attrs(mut self) -> Self {
self.index_all_attrs = true;
self
}
pub fn metadata_all_attrs(mut self) -> Self {
self.metadata_all_attrs = true;
self
}
pub fn skip_timestamp_validation(mut self) -> Self {
self.skip_timestamp_validation = true;
self
}
#[cfg(feature = "xray-subsegment-nesting")]
pub fn always_nest_subsegments(mut self) -> Self {
self.always_nest_subsegments = true;
self
}
pub fn with_indexed_attr(mut self, attr: String) -> Self {
if let Err(i) = self.indexed_attrs.binary_search(&attr) {
self.indexed_attrs.insert(i, attr);
}
self
}
pub fn with_indexed_attrs(mut self, attrs: impl IntoIterator<Item = String>) -> Self {
for attr in attrs {
if let Err(i) = self.indexed_attrs.binary_search(&attr) {
self.indexed_attrs.insert(i, attr);
}
}
self
}
pub fn with_metadata_attr(mut self, attr: String) -> Self {
if let Err(i) = self.metadata_attrs.binary_search(&attr) {
self.metadata_attrs.insert(i, attr);
}
self
}
pub fn with_metadata_attrs(mut self, attrs: impl IntoIterator<Item = String>) -> Self {
for attr in attrs {
if let Err(i) = self.metadata_attrs.binary_search(&attr) {
self.metadata_attrs.insert(i, attr);
}
}
self
}
pub fn set_indexed_attrs(mut self, mut indexed_attrs: Vec<String>) -> Self {
indexed_attrs.sort();
self.indexed_attrs = indexed_attrs;
self
}
pub fn with_log_group_name(mut self, log_group_name: String) -> Self {
self.log_group_names.push(log_group_name);
self
}
pub fn with_log_group_names(
mut self,
log_group_names: impl IntoIterator<Item = String>,
) -> Self {
self.log_group_names.extend(log_group_names);
self
}
pub fn set_log_group_names(mut self, log_group_names: Vec<String>) -> Self {
self.log_group_names = log_group_names;
self
}
pub fn set_resource(&mut self, resource: &Resource) {
self.resource.replace(resource.clone());
}
}
impl SegmentTranslator {
#[cfg_attr(feature = "internal-logs", tracing::instrument(skip(self, batch)))]
pub fn translate_spans<'span, 'translator: 'span>(
&'translator self,
batch: &'span [SpanData],
) -> Vec<SegmentDocument<'span>> {
#[cfg(feature = "internal-logs")]
tracing::debug!("Received {} spans", batch.len());
#[cfg(feature = "xray-subsegment-nesting")]
if self.always_nest_subsegments {
self._translate_spans_nested(batch)
} else {
self._translate_spans_simple(batch)
}
#[cfg(not(feature = "xray-subsegment-nesting"))]
self._translate_spans_simple(batch)
}
fn _translate_spans_simple<'span, 'translator: 'span>(
&'translator self,
batch: &'span [SpanData],
) -> Vec<SegmentDocument<'span>> {
batch
.iter()
.filter_map(|span_data| {
match self
.translate_span(span_data)
.and_then(|builder| builder.build())
{
Ok(segment) => Some(segment),
Err(e) => {
#[cfg(feature = "internal-logs")]
tracing::error!(message="A segment or subsegment was lost", error=?e);
#[cfg(feature = "internal-logs")]
tracing::debug!(error=?e, ?span_data);
None
}
}
})
.collect()
}
#[cfg(feature = "xray-subsegment-nesting")]
fn _translate_spans_nested<'span, 'translator: 'span>(
&'translator self,
batch: &'span [SpanData],
) -> Vec<SegmentDocument<'span>> {
use crate::xray_exporter::types::{
error::ConstraintError, DocumentBuilderHeader, Id, TraceId,
};
use std::collections::HashMap;
let mut document_builders: HashMap<(TraceId, Id), AnyDocumentBuilder> =
HashMap::with_capacity(batch.len());
let mut document_builder_headers_tree = DocumentBuilderHeaderTree::new(batch.len());
for span_data in batch.iter() {
match
self.translate_span(span_data).and_then(|builder| {
let header = match &builder {
AnyDocumentBuilder::Segment(builder) => builder.header(),
AnyDocumentBuilder::Subsegment(builder) => builder.header(),
};
let id = header.id.ok_or(ConstraintError::MissingId)?;
let trace_id = header.trace_id.ok_or(ConstraintError::MissingTraceId)?;
if document_builders.insert((trace_id, id), builder).is_some() {
#[cfg(feature = "internal-logs")]
tracing::error!("Duplicated builder (id: {id}; trace-id: {trace_id:?}), a segment or subsegment was lost");
} else {
document_builder_headers_tree
.add(header)
.expect("id and trace_id always present at this point");
}
Ok(())
}) {
Ok(_) => {},
Err(e) => {
#[cfg(feature = "internal-logs")]
tracing::error!(message="A segment or subsegment was lost", error=?e);
#[cfg(feature = "internal-logs")]
tracing::debug!(error=?e, ?span_data);
}
}
}
let mut last_seen_parent_id = None;
let mut last_seen_endtime = None;
let mut precursors = Vec::new();
for header in document_builder_headers_tree.iter() {
let DocumentBuilderHeader {
id,
parent_id,
trace_id,
start_time,
end_time,
} = header;
let id = id.expect("id always set at this point");
let trace_id = trace_id.expect("trace_id always set at this point");
let start_time = start_time.expect("start_time always set at this point");
let Some(parent_id) = parent_id else {
continue;
};
let document_builders_index = &(trace_id, id);
let AnyDocumentBuilder::Subsegment(subsegment_builder) = document_builders
.get_mut(document_builders_index)
.expect("builders in subsegments are always also in document_builders")
else {
continue;
};
if last_seen_parent_id != Some(parent_id) {
last_seen_parent_id = Some(parent_id);
precursors.clear();
last_seen_endtime = None;
}
if last_seen_endtime.is_none()
|| last_seen_endtime
.is_some_and(|last_seens_endtime| last_seens_endtime <= start_time)
{
if !precursors.is_empty() {
subsegment_builder.precursor_ids(precursors.clone());
}
if let Some(end_time) = end_time {
last_seen_endtime = Some(end_time);
precursors.push(id);
}
}
let document_builders_parent_index = &(trace_id, parent_id);
if document_builders.contains_key(document_builders_parent_index) {
if document_builders_index == document_builders_parent_index {
#[cfg(feature = "internal-logs")]
tracing::warn!(
message =
"Subsegment cannot be nested because it references itself as its parent",
?document_builders_index
);
continue;
}
if let AnyDocumentBuilder::Subsegment(subsegment) = document_builders
.remove(document_builders_index)
.expect("builders in subsegments are always also in document_builders")
{
match document_builders
.get_mut(document_builders_parent_index)
.expect("verified parent_id is present")
{
AnyDocumentBuilder::Segment(parent) => {
parent.subsegment(subsegment);
}
AnyDocumentBuilder::Subsegment(parent) => {
parent.subsegment(subsegment);
}
}
};
}
}
document_builders
.into_values()
.filter_map(|document_builder| match document_builder.build() {
Ok(segment) => Some(segment),
Err(e) => {
#[cfg(feature = "internal-logs")]
tracing::error!(message="A segment or subsegment was lost", error=?e);
None
}
})
.collect()
}
#[cfg_attr(feature = "internal-logs", tracing::instrument(skip(self, span_data)))]
fn translate_span<'span, 'translator: 'span>(
&'translator self,
span_data: &'span SpanData,
) -> Result<AnyDocumentBuilder<'span>> {
#[cfg(feature = "internal-logs")]
tracing::trace!(?span_data);
let SpanData {
span_kind,
name,
attributes,
events,
status,
..
} = span_data;
let span_kind_is_server = matches!(span_kind, SpanKind::Server);
let span_is_remote = matches!(span_kind, SpanKind::Client | SpanKind::Producer);
let mut additionnal_builders = AnyValueBuilder::array(
AwsOperationBuilder::default(),
AwsXraySdkBuilder::default(),
BeanstalkDeploymentIdBuilder::default(),
CauseBuilder::new(events, status, span_is_remote),
CloudwatchLogGroupBuilder::new(&self.log_group_names),
HttpRequestUrlBuilder::new(span_kind_is_server),
HttpResponseContentLengthBuilder::default(),
SegmentNameBuilder::new(name.as_ref(), span_kind_is_server),
SegmentOriginBuilder::default(),
SqlUrlBuilder::default(),
SubsegmentNamespaceBuilder::new(matches!(span_kind, SpanKind::Client)),
);
let mut any_segment_builder = {
match span_data.span_kind {
SpanKind::Server => {
AnyDocumentBuilder::Segment(self.init_document_builder(span_data)?)
}
_ => AnyDocumentBuilder::Subsegment(self.init_document_builder(span_data)?),
}
};
let attribute_iterator = self
.resource
.iter()
.flat_map(|r: &Resource| r.iter())
.chain(attributes.iter().map(|kv: &KeyValue| (&kv.key, &kv.value)));
for (key, value) in attribute_iterator {
let key = key.as_str();
let (should_index, should_metadata, key) =
if let Some(key) = key.strip_prefix("annotation.") {
(true, false, key)
} else if let Some(key) = key.strip_prefix("metadata.") {
(false, true, key)
} else {
(false, false, key)
};
let mut attribute_included = false;
for &(builder_id, handler_index) in self.dispatch_table.dispatch(key) {
let handler_result = match builder_id {
ANY_DOCUMENT_BUILDER_PROCESSOR_ID => AnyDocumentBuilder::HANDLERS
[handler_index]
.1(
&mut any_segment_builder, value
),
builder_id => {
additionnal_builders[builder_id].process(handler_index, value)
}
};
attribute_included = attribute_included || handler_result;
}
let should_index = should_index
|| (!attribute_included && self.index_all_attrs)
|| self
.indexed_attrs
.binary_search_by(|s| s.as_str().cmp(key))
.is_ok();
if should_index {
if let Some(annotation) = get_annotation(value) {
let sanitized_key = sanitize_annotation_key(key);
attribute_included = match &mut any_segment_builder {
AnyDocumentBuilder::Segment(builder) => {
builder.annotation(sanitized_key, annotation).is_ok()
}
AnyDocumentBuilder::Subsegment(builder) => {
builder.annotation(sanitized_key, annotation).is_ok()
}
};
} else {
attribute_included = false;
}
}
let should_metadata = should_metadata && !should_index
|| should_index && !attribute_included
|| (!attribute_included && self.metadata_all_attrs)
|| self
.metadata_attrs
.binary_search_by(|s| s.as_str().cmp(key))
.is_ok();
if should_metadata {
if let Some(value) = get_any_value(value) {
match &mut any_segment_builder {
AnyDocumentBuilder::Segment(builder) => {
builder.metadata(key, value);
}
AnyDocumentBuilder::Subsegment(builder) => {
builder.metadata(key, value);
}
}
}
}
}
for additionnal_builder in additionnal_builders {
additionnal_builder.resolve(&mut any_segment_builder)?;
}
Ok(any_segment_builder)
}
fn init_document_builder<'span, DBT: DocumentBuilderType>(
&self,
span_data: &'span SpanData,
) -> Result<DocumentBuilder<'span, DBT>> {
let SpanData {
span_context,
parent_span_id,
start_time,
end_time,
attributes,
..
} = span_data;
let start_time = *start_time;
let end_time = *end_time;
let parent_span_id = *parent_span_id;
let attribute_count = attributes.len();
let mut builder = DocumentBuilder::default();
let trace_id = span_context.trace_id();
if trace_id != opentelemetry::TraceId::INVALID {
builder.trace_id(trace_id.into(), self.skip_timestamp_validation)?;
}
let span_id = span_context.span_id();
if span_id == SpanId::INVALID {
return Err(TranslationError::MissingSpanId);
}
builder
.id(span_id.into())
.start_time(translate_timestamp(start_time))
.end_time(translate_timestamp(end_time))?;
if parent_span_id != SpanId::INVALID {
builder.parent_id(parent_span_id.into());
}
let (max_annotations, max_metadata) = if self.index_all_attrs {
(attribute_count, 0)
} else {
let max_annotations = attribute_count.min(self.indexed_attrs.len());
(max_annotations, attribute_count - max_annotations)
};
Ok(builder
.with_annotation_capacity(max_annotations)
.with_metadata_capacity(max_metadata))
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use opentelemetry::{
trace::{SpanContext, SpanId, TraceFlags, TraceId, TraceState},
InstrumentationScope, KeyValue,
};
use opentelemetry_sdk::trace::{SpanData, SpanEvents, SpanLinks};
fn create_valid_trace_id() -> TraceId {
let timestamp = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_secs() as u128;
let random_part: u128 = 0xabcdef0123456789abcdef01;
let trace_id = (timestamp << 96) | random_part;
TraceId::from_bytes(trace_id.to_be_bytes())
}
fn create_span(
trace_id: TraceId,
span_id: SpanId,
kind: SpanKind,
attributes: Vec<KeyValue>,
) -> SpanData {
let span_context = SpanContext::new(
trace_id,
span_id,
TraceFlags::SAMPLED,
false,
TraceState::default(),
);
let start_time = UNIX_EPOCH + Duration::from_secs(1_700_000_000);
let end_time = start_time + Duration::from_millis(100);
SpanData {
span_context,
parent_span_id: SpanId::INVALID,
parent_span_is_remote: false,
span_kind: kind,
name: "test-span".into(),
start_time,
end_time,
attributes,
dropped_attributes_count: 0,
events: SpanEvents::default(),
links: SpanLinks::default(),
status: opentelemetry::trace::Status::Unset,
instrumentation_scope: InstrumentationScope::builder("test").build(),
}
}
#[test]
fn test_init_document_builder_invalid_trace_id() {
let translator = SegmentTranslator::new().skip_timestamp_validation();
let span = create_span(
TraceId::INVALID,
SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 1]),
SpanKind::Server,
vec![],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert!(
documents.is_empty(),
"Span with INVALID trace_id should be silently dropped"
);
}
#[test]
fn test_init_document_builder_invalid_span_id() {
let translator = SegmentTranslator::new().skip_timestamp_validation();
let span = create_span(
create_valid_trace_id(),
SpanId::INVALID,
SpanKind::Server,
vec![],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert!(
documents.is_empty(),
"Span with INVALID span_id should be silently dropped (MissingSpanId)"
);
}
#[test]
fn test_annotation_prefix_stripping() {
let translator = SegmentTranslator::new().skip_timestamp_validation();
let span = create_span(
create_valid_trace_id(),
SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 1]),
SpanKind::Server,
vec![KeyValue::new("annotation.my_key", "my_value")],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert_eq!(documents.len(), 1);
let json: serde_json::Value = serde_json::to_value(&documents[0]).unwrap();
let annotations = json.get("annotations").expect("annotations should exist");
assert_eq!(
annotations.get("my_key"),
Some(&serde_json::json!("my_value")),
"annotation.my_key should be stripped to my_key in annotations"
);
assert!(
json.get("metadata").is_none() || json["metadata"].get("my_key").is_none(),
"annotation-prefixed key should not appear in metadata"
);
}
#[test]
fn test_metadata_prefix_stripping() {
let translator = SegmentTranslator::new().skip_timestamp_validation();
let span = create_span(
create_valid_trace_id(),
SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 2]),
SpanKind::Server,
vec![KeyValue::new("metadata.my_key", "my_value")],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert_eq!(documents.len(), 1);
let json: serde_json::Value = serde_json::to_value(&documents[0]).unwrap();
let metadata = json.get("metadata").expect("metadata should exist");
assert_eq!(
metadata.get("my_key"),
Some(&serde_json::json!("my_value")),
"metadata.my_key should be stripped to my_key in metadata"
);
assert!(
json.get("annotations").is_none() || json["annotations"].get("my_key").is_none(),
"metadata-prefixed key should not appear in annotations"
);
}
#[test]
fn test_with_indexed_attr_deduplication() {
let translator = SegmentTranslator::new()
.skip_timestamp_validation()
.with_indexed_attr("same_key".to_string())
.with_indexed_attr("same_key".to_string());
let span = create_span(
create_valid_trace_id(),
SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 3]),
SpanKind::Server,
vec![KeyValue::new("same_key", "the_value")],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert_eq!(documents.len(), 1);
let json: serde_json::Value = serde_json::to_value(&documents[0]).unwrap();
let annotations = json.get("annotations").expect("annotations should exist");
assert_eq!(
annotations.get("same_key"),
Some(&serde_json::json!("the_value")),
"same_key should be indexed as annotation despite duplicate with_indexed_attr calls"
);
}
#[test]
fn test_with_metadata_attr_deduplication() {
let translator = SegmentTranslator::new()
.skip_timestamp_validation()
.with_metadata_attr("same_key".to_string())
.with_metadata_attr("same_key".to_string());
let span = create_span(
create_valid_trace_id(),
SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 4]),
SpanKind::Server,
vec![KeyValue::new("same_key", "the_value")],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert_eq!(documents.len(), 1);
let json: serde_json::Value = serde_json::to_value(&documents[0]).unwrap();
let metadata = json.get("metadata").expect("metadata should exist");
assert_eq!(
metadata.get("same_key"),
Some(&serde_json::json!("the_value")),
"same_key should be in metadata despite duplicate with_metadata_attr calls"
);
let metadata_obj = metadata.as_object().expect("metadata should be an object");
let same_key_count = metadata_obj.keys().filter(|k| *k == "same_key").count();
assert_eq!(
same_key_count, 1,
"same_key should appear exactly once in metadata"
);
}
#[test]
fn test_metadata_prefix_overridden_by_indexed_attr() {
let translator = SegmentTranslator::new()
.skip_timestamp_validation()
.with_indexed_attr("my_indexed_attr".to_string());
let span = create_span(
create_valid_trace_id(),
SpanId::from_bytes([0, 0, 0, 0, 0, 0, 0, 5]),
SpanKind::Server,
vec![KeyValue::new("metadata.my_indexed_attr", "hello")],
);
let batch = [span];
let documents = translator.translate_spans(&batch);
assert_eq!(documents.len(), 1);
let json: serde_json::Value = serde_json::to_value(&documents[0]).unwrap();
let annotations = json.get("annotations").expect("annotations should exist");
assert_eq!(
annotations.get("my_indexed_attr"),
Some(&serde_json::json!("hello")),
"metadata.my_indexed_attr should be routed to annotations when the key is indexed"
);
assert!(
json.get("metadata").is_none() || json["metadata"].get("my_indexed_attr").is_none(),
"indexed attr should not appear in metadata even with metadata. prefix"
);
}
}