use crate::exporter::intern::StringInterner;
use crate::exporter::model::{DD_MEASURED_KEY, SAMPLING_PRIORITY_KEY};
use crate::exporter::{Error, ModelConfig};
use crate::propagator::DatadogTraceState;
use opentelemetry::trace::Status;
use opentelemetry_sdk::trace::SpanData;
use opentelemetry_sdk::Resource;
use std::time::SystemTime;
use super::unified_tags::{UnifiedTagField, UnifiedTags};
const SPAN_NUM_ELEMENTS: u32 = 12;
const METRICS_LEN: u32 = 2;
const GIT_META_TAGS_COUNT: u32 = if matches!(
(
option_env!("DD_GIT_REPOSITORY_URL"),
option_env!("DD_GIT_COMMIT_SHA")
),
(Some(_), Some(_))
) {
2
} else {
0
};
pub(crate) fn encode<S, N, R>(
model_config: &ModelConfig,
traces: Vec<&[SpanData]>,
get_service_name: S,
get_name: N,
get_resource: R,
unified_tags: &UnifiedTags,
resource: Option<&Resource>,
) -> Result<Vec<u8>, Error>
where
for<'a> S: Fn(&'a SpanData, &'a ModelConfig) -> &'a str,
for<'a> N: Fn(&'a SpanData, &'a ModelConfig) -> &'a str,
for<'a> R: Fn(&'a SpanData, &'a ModelConfig) -> &'a str,
{
let mut interner = StringInterner::new();
let mut encoded_traces = encode_traces(
&mut interner,
model_config,
get_service_name,
get_name,
get_resource,
&traces,
unified_tags,
resource,
)?;
let mut payload = Vec::with_capacity(traces.len() * 512);
rmp::encode::write_array_len(&mut payload, 2)?;
interner.write_dictionary(&mut payload)?;
payload.append(&mut encoded_traces);
Ok(payload)
}
fn write_unified_tags<'a>(
encoded: &mut Vec<u8>,
interner: &mut StringInterner<'a>,
unified_tags: &'a UnifiedTags,
) -> Result<(), Error> {
write_unified_tag(encoded, interner, &unified_tags.service)?;
write_unified_tag(encoded, interner, &unified_tags.env)?;
write_unified_tag(encoded, interner, &unified_tags.version)?;
Ok(())
}
fn write_unified_tag<'a>(
encoded: &mut Vec<u8>,
interner: &mut StringInterner<'a>,
tag: &'a UnifiedTagField,
) -> Result<(), Error> {
if let Some(tag_value) = &tag.value {
rmp::encode::write_u32(encoded, interner.intern(tag.get_tag_name()))?;
rmp::encode::write_u32(encoded, interner.intern(tag_value.as_str().as_ref()))?;
}
Ok(())
}
#[cfg(not(feature = "agent-sampling"))]
fn get_sampling_priority(_span: &SpanData) -> f64 {
1.0
}
#[cfg(feature = "agent-sampling")]
fn get_sampling_priority(span: &SpanData) -> f64 {
if span.span_context.trace_state().priority_sampling_enabled() {
1.0
} else {
0.0
}
}
fn get_measuring(span: &SpanData) -> f64 {
if span.span_context.trace_state().measuring_enabled() {
1.0
} else {
0.0
}
}
#[allow(clippy::too_many_arguments)]
fn encode_traces<'interner, S, N, R>(
interner: &mut StringInterner<'interner>,
model_config: &'interner ModelConfig,
get_service_name: S,
get_name: N,
get_resource: R,
traces: &'interner [&[SpanData]],
unified_tags: &'interner UnifiedTags,
resource: Option<&'interner Resource>,
) -> Result<Vec<u8>, Error>
where
for<'a> S: Fn(&'a SpanData, &'a ModelConfig) -> &'a str,
for<'a> N: Fn(&'a SpanData, &'a ModelConfig) -> &'a str,
for<'a> R: Fn(&'a SpanData, &'a ModelConfig) -> &'a str,
{
let mut encoded = Vec::new();
rmp::encode::write_array_len(&mut encoded, traces.len() as u32)?;
for trace in traces.iter() {
rmp::encode::write_array_len(&mut encoded, trace.len() as u32)?;
for span in trace.iter() {
let start = span
.start_time
.duration_since(SystemTime::UNIX_EPOCH)
.unwrap()
.as_nanos() as i64;
let duration = span
.end_time
.duration_since(span.start_time)
.map(|x| x.as_nanos() as i64)
.unwrap_or(0);
let mut span_type = interner.intern("");
for kv in &span.attributes {
if kv.key.as_str() == "span.type" {
span_type = interner.intern_value(&kv.value);
break;
}
}
rmp::encode::write_array_len(&mut encoded, SPAN_NUM_ELEMENTS)?;
rmp::encode::write_u32(
&mut encoded,
interner.intern(get_service_name(span, model_config)),
)?;
rmp::encode::write_u32(&mut encoded, interner.intern(get_name(span, model_config)))?;
rmp::encode::write_u32(
&mut encoded,
interner.intern(get_resource(span, model_config)),
)?;
rmp::encode::write_u64(
&mut encoded,
u128::from_be_bytes(span.span_context.trace_id().to_bytes()) as u64,
)?;
rmp::encode::write_u64(
&mut encoded,
u64::from_be_bytes(span.span_context.span_id().to_bytes()),
)?;
rmp::encode::write_u64(
&mut encoded,
u64::from_be_bytes(span.parent_span_id.to_bytes()),
)?;
rmp::encode::write_i64(&mut encoded, start)?;
rmp::encode::write_i64(&mut encoded, duration)?;
rmp::encode::write_i32(
&mut encoded,
match span.status {
Status::Error { .. } => 1,
_ => 0,
},
)?;
rmp::encode::write_map_len(
&mut encoded,
(span.attributes.len() + resource.map(|r| r.len()).unwrap_or(0)) as u32
+ unified_tags.compute_attribute_size()
+ GIT_META_TAGS_COUNT,
)?;
if let Some(resource) = resource {
for (key, value) in resource.iter() {
rmp::encode::write_u32(&mut encoded, interner.intern(key.as_str()))?;
rmp::encode::write_u32(&mut encoded, interner.intern_value(value))?;
}
}
write_unified_tags(&mut encoded, interner, unified_tags)?;
for kv in span.attributes.iter() {
rmp::encode::write_u32(&mut encoded, interner.intern(kv.key.as_str()))?;
rmp::encode::write_u32(&mut encoded, interner.intern_value(&kv.value))?;
}
if let (Some(repository_url), Some(commit_sha)) = (
option_env!("DD_GIT_REPOSITORY_URL"),
option_env!("DD_GIT_COMMIT_SHA"),
) {
rmp::encode::write_u32(&mut encoded, interner.intern("git.repository_url"))?;
rmp::encode::write_u32(&mut encoded, interner.intern(repository_url))?;
rmp::encode::write_u32(&mut encoded, interner.intern("git.commit.sha"))?;
rmp::encode::write_u32(&mut encoded, interner.intern(commit_sha))?;
}
rmp::encode::write_map_len(&mut encoded, METRICS_LEN)?;
rmp::encode::write_u32(&mut encoded, interner.intern(SAMPLING_PRIORITY_KEY))?;
let sampling_priority = get_sampling_priority(span);
rmp::encode::write_f64(&mut encoded, sampling_priority)?;
rmp::encode::write_u32(&mut encoded, interner.intern(DD_MEASURED_KEY))?;
let measuring = get_measuring(span);
rmp::encode::write_f64(&mut encoded, measuring)?;
rmp::encode::write_u32(&mut encoded, span_type)?;
}
}
Ok(encoded)
}