stele/traces/
tracer.rs

1use datadog_protos::traces as proto;
2use datadog_protos::traces::idx as etp_proto;
3use ordered_float::OrderedFloat;
4use saluki_common::collections::FastHashMap;
5use serde::{Deserialize, Serialize};
6use stringtheory::MetaString;
7
8#[derive(Clone, Debug, Eq, PartialEq)]
9struct WrappedFloat(OrderedFloat<f64>);
10
11impl From<f64> for WrappedFloat {
12    fn from(value: f64) -> Self {
13        WrappedFloat(OrderedFloat(value))
14    }
15}
16
17impl<'de> Deserialize<'de> for WrappedFloat {
18    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
19    where
20        D: serde::Deserializer<'de>,
21    {
22        let value = f64::deserialize(deserializer)?;
23        Ok(WrappedFloat(OrderedFloat(value)))
24    }
25}
26
27impl Serialize for WrappedFloat {
28    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
29    where
30        S: serde::Serializer,
31    {
32        self.0 .0.serialize(serializer)
33    }
34}
35
36#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
37struct AgentMetadata {
38    hostname: MetaString,
39    env: MetaString,
40    tags: FastHashMap<MetaString, MetaString>,
41    agent_version: MetaString,
42    target_tps: WrappedFloat,
43    error_tps: WrappedFloat,
44    rare_sampler_enabled: bool,
45}
46
47impl From<&proto::AgentPayload> for AgentMetadata {
48    fn from(payload: &proto::AgentPayload) -> Self {
49        Self {
50            hostname: (*payload.hostName).into(),
51            env: (*payload.env).into(),
52            tags: payload.tags.iter().map(|(k, v)| ((**k).into(), (**v).into())).collect(),
53            agent_version: (*payload.agentVersion).into(),
54            target_tps: payload.targetTPS.into(),
55            error_tps: payload.errorTPS.into(),
56            rare_sampler_enabled: payload.rareSamplerEnabled,
57        }
58    }
59}
60
61#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
62struct TracerMetadata {
63    container_id: MetaString,
64    language_name: MetaString,
65    language_version: MetaString,
66    tracer_version: MetaString,
67    runtime_id: MetaString,
68    tags: FastHashMap<MetaString, MetaString>,
69    env: MetaString,
70    hostname: MetaString,
71    app_version: MetaString,
72}
73
74impl From<&proto::TracerPayload> for TracerMetadata {
75    fn from(payload: &proto::TracerPayload) -> Self {
76        Self {
77            container_id: (*payload.containerID).into(),
78            language_name: (*payload.languageName).into(),
79            language_version: (*payload.languageVersion).into(),
80            tracer_version: (*payload.tracerVersion).into(),
81            runtime_id: (*payload.runtimeID).into(),
82            tags: payload.tags.iter().map(|(k, v)| ((**k).into(), (**v).into())).collect(),
83            env: (*payload.env).into(),
84            hostname: (*payload.hostname).into(),
85            app_version: (*payload.appVersion).into(),
86        }
87    }
88}
89
90impl TracerMetadata {
91    fn from_etp_payload(payload: &etp_proto::TracerPayload) -> Self {
92        let strings = &payload.strings;
93        Self {
94            container_id: resolve_ref(strings, payload.containerIDRef),
95            language_name: resolve_ref(strings, payload.languageNameRef),
96            language_version: resolve_ref(strings, payload.languageVersionRef),
97            tracer_version: resolve_ref(strings, payload.tracerVersionRef),
98            runtime_id: resolve_ref(strings, payload.runtimeIDRef),
99            tags: string_attrs_from_etp(&payload.attributes, strings),
100            env: resolve_ref(strings, payload.envRef),
101            hostname: resolve_ref(strings, payload.hostnameRef),
102            app_version: resolve_ref(strings, payload.appVersionRef),
103        }
104    }
105}
106
107#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
108struct TraceChunkMetadata {
109    priority: i32,
110    origin: MetaString,
111    tags: FastHashMap<MetaString, MetaString>,
112    dropped_trace: bool,
113}
114
115impl From<&proto::TraceChunk> for TraceChunkMetadata {
116    fn from(value: &proto::TraceChunk) -> Self {
117        Self {
118            priority: value.priority,
119            origin: (*value.origin).into(),
120            tags: value.tags.iter().map(|(k, v)| ((**k).into(), (**v).into())).collect(),
121            dropped_trace: value.droppedTrace,
122        }
123    }
124}
125
126impl TraceChunkMetadata {
127    fn from_etp_chunk(chunk: &etp_proto::TraceChunk, strings: &[String]) -> Self {
128        Self {
129            priority: chunk.priority,
130            origin: resolve_ref(strings, chunk.originRef),
131            tags: string_attrs_from_etp(&chunk.attributes, strings),
132            dropped_trace: chunk.droppedTrace,
133        }
134    }
135}
136
137/// A simplified span representation.
138#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
139pub struct Span {
140    agent_metadata: AgentMetadata,
141    tracer_metadata: TracerMetadata,
142    trace_chunk_metadata: TraceChunkMetadata,
143    service: MetaString,
144    name: MetaString,
145    resource: MetaString,
146    trace_id: u64,
147    span_id: u64,
148    parent_id: u64,
149    start: i64,
150    duration: i64,
151    error: i32,
152    meta: FastHashMap<MetaString, MetaString>,
153    metrics: FastHashMap<MetaString, WrappedFloat>,
154    type_: MetaString,
155    meta_struct: FastHashMap<MetaString, Vec<u8>>,
156    span_links: Vec<SpanLink>,
157    span_events: Vec<SpanEvent>,
158}
159
160impl Span {
161    /// Returns the trace ID this span belongs to.
162    pub fn trace_id(&self) -> u64 {
163        self.trace_id
164    }
165
166    /// Returns the ID of this span.
167    pub fn span_id(&self) -> u64 {
168        self.span_id
169    }
170
171    /// Returns the value of the metadata entry of the given key, if it exists.
172    pub fn get_meta_field(&self, meta_key: &str) -> Option<&str> {
173        self.meta.get(meta_key).map(|s| &**s)
174    }
175
176    /// Gets all spans from the classic and indexed tracer payloads in the given `AgentPayload`.
177    pub fn get_spans_from_agent_payload(payload: &proto::AgentPayload) -> Vec<Self> {
178        let agent_metadata = AgentMetadata::from(payload);
179
180        let mut spans = Vec::new();
181        for tracer_payload in payload.tracerPayloads() {
182            let tracer_metadata = TracerMetadata::from(tracer_payload);
183
184            for trace_chunk in tracer_payload.chunks() {
185                let trace_chunk_metadata = TraceChunkMetadata::from(trace_chunk);
186
187                for span in trace_chunk.spans() {
188                    let span = Self::from_proto(
189                        agent_metadata.clone(),
190                        tracer_metadata.clone(),
191                        trace_chunk_metadata.clone(),
192                        span,
193                    );
194                    spans.push(span);
195                }
196            }
197        }
198
199        for tracer_payload in payload.idxTracerPayloads() {
200            let strings = &tracer_payload.strings;
201            let tracer_metadata = TracerMetadata::from_etp_payload(tracer_payload);
202
203            for chunk in &tracer_payload.chunks {
204                let trace_chunk_metadata = TraceChunkMetadata::from_etp_chunk(chunk, strings);
205                let trace_id = trace_id_low_from_bytes(&chunk.traceID);
206
207                for span in &chunk.spans {
208                    spans.push(Self::from_etp_proto(
209                        agent_metadata.clone(),
210                        tracer_metadata.clone(),
211                        trace_chunk_metadata.clone(),
212                        span,
213                        trace_id,
214                        strings,
215                    ));
216                }
217            }
218        }
219
220        spans
221    }
222
223    fn from_proto(
224        agent_metadata: AgentMetadata, tracer_metadata: TracerMetadata, trace_chunk_metadata: TraceChunkMetadata,
225        value: &proto::Span,
226    ) -> Self {
227        let mut span_links = value.spanLinks.iter().map(SpanLink::from).collect::<Vec<_>>();
228        span_links.sort_by_key(|link| (link.trace_id, link.trace_id_high, link.span_id));
229
230        let mut span_events = value.spanEvents.iter().map(SpanEvent::from).collect::<Vec<_>>();
231        span_events.sort_by_key(|event| event.time_unix_nano);
232
233        Self {
234            agent_metadata,
235            tracer_metadata,
236            trace_chunk_metadata,
237            service: (*value.service).into(),
238            name: (*value.name).into(),
239            resource: (*value.resource).into(),
240            trace_id: value.traceID,
241            span_id: value.spanID,
242            parent_id: value.parentID,
243            start: value.start,
244            duration: value.duration,
245            error: value.error,
246            meta: value.meta.iter().map(|(k, v)| ((**k).into(), (**v).into())).collect(),
247            metrics: value.metrics.iter().map(|(k, v)| ((**k).into(), (*v).into())).collect(),
248            type_: (*value.type_).into(),
249            meta_struct: value
250                .meta_struct
251                .iter()
252                .map(|(k, v)| ((**k).into(), v.to_vec()))
253                .collect(),
254            span_links,
255            span_events,
256        }
257    }
258
259    fn from_etp_proto(
260        agent_metadata: AgentMetadata, tracer_metadata: TracerMetadata, trace_chunk_metadata: TraceChunkMetadata,
261        span: &etp_proto::Span, trace_id: u64, strings: &[String],
262    ) -> Self {
263        let (meta, metrics, meta_struct) = split_etp_span_attributes(&span.attributes, strings);
264
265        let mut span_links = span
266            .links
267            .iter()
268            .map(|l| SpanLink::from_etp(l, strings))
269            .collect::<Vec<_>>();
270        span_links.sort_by_key(|link| (link.trace_id, link.trace_id_high, link.span_id));
271
272        let mut span_events = span
273            .events
274            .iter()
275            .map(|e| SpanEvent::from_etp(e, strings))
276            .collect::<Vec<_>>();
277        span_events.sort_by_key(|event| event.time_unix_nano);
278
279        Self {
280            agent_metadata,
281            tracer_metadata,
282            trace_chunk_metadata,
283            service: resolve_ref(strings, span.serviceRef),
284            name: resolve_ref(strings, span.nameRef),
285            resource: resolve_ref(strings, span.resourceRef),
286            trace_id,
287            span_id: span.spanID,
288            parent_id: span.parentID,
289            start: span.start as i64,
290            duration: span.duration as i64,
291            error: i32::from(span.error),
292            meta,
293            metrics,
294            type_: resolve_ref(strings, span.typeRef),
295            meta_struct,
296            span_links,
297            span_events,
298        }
299    }
300}
301
302#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
303struct SpanLink {
304    trace_id: u64,
305    trace_id_high: u64,
306    span_id: u64,
307    attributes: FastHashMap<MetaString, MetaString>,
308    tracestate: MetaString,
309    flags: u32,
310}
311
312impl From<&proto::SpanLink> for SpanLink {
313    fn from(value: &proto::SpanLink) -> Self {
314        Self {
315            trace_id: value.traceID,
316            trace_id_high: value.traceID_high,
317            span_id: value.spanID,
318            attributes: value
319                .attributes
320                .iter()
321                .map(|(k, v)| ((**k).into(), (**v).into()))
322                .collect(),
323            tracestate: (*value.tracestate).into(),
324            flags: value.flags,
325        }
326    }
327}
328
329impl SpanLink {
330    fn from_etp(link: &etp_proto::SpanLink, strings: &[String]) -> Self {
331        let (trace_id, trace_id_high) = trace_id_parts_from_bytes(&link.traceID);
332        Self {
333            trace_id,
334            trace_id_high,
335            span_id: link.spanID,
336            attributes: string_attrs_from_etp(&link.attributes, strings),
337            tracestate: resolve_ref(strings, link.tracestateRef),
338            flags: link.flags,
339        }
340    }
341}
342
343#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
344struct SpanEvent {
345    time_unix_nano: u64,
346    name: MetaString,
347    attributes: FastHashMap<MetaString, AttributeAnyValue>,
348}
349
350impl From<&proto::SpanEvent> for SpanEvent {
351    fn from(value: &proto::SpanEvent) -> Self {
352        Self {
353            time_unix_nano: value.time_unix_nano,
354            name: (*value.name).into(),
355            attributes: value
356                .attributes
357                .iter()
358                .map(|(k, v)| ((**k).into(), AttributeAnyValue::from(v)))
359                .collect(),
360        }
361    }
362}
363
364impl SpanEvent {
365    fn from_etp(event: &etp_proto::SpanEvent, strings: &[String]) -> Self {
366        Self {
367            // The ETP proto renamed time_unix_nano to `time` (same semantics).
368            time_unix_nano: event.time,
369            name: resolve_ref(strings, event.nameRef),
370            attributes: event
371                .attributes
372                .iter()
373                .filter_map(|(k_ref, v)| {
374                    let attr = etp_anyvalue_to_event_attr(v, strings)?;
375                    Some((resolve_ref(strings, *k_ref), attr))
376                })
377                .collect(),
378        }
379    }
380}
381
382#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
383enum AttributeAnyValue {
384    String(MetaString),
385    Boolean(bool),
386    Integer(i64),
387    Double(WrappedFloat),
388    Array(Vec<AttributeArrayValue>),
389}
390
391impl From<&proto::AttributeAnyValue> for AttributeAnyValue {
392    fn from(value: &proto::AttributeAnyValue) -> Self {
393        let maybe_proto_value_type = value.type_.enum_value().expect("unknown/invalid anyvalue type");
394        match maybe_proto_value_type {
395            proto::AttributeAnyValueType::STRING_VALUE => AttributeAnyValue::String((*value.string_value).into()),
396            proto::AttributeAnyValueType::BOOL_VALUE => AttributeAnyValue::Boolean(value.bool_value),
397            proto::AttributeAnyValueType::INT_VALUE => AttributeAnyValue::Integer(value.int_value),
398            proto::AttributeAnyValueType::DOUBLE_VALUE => AttributeAnyValue::Double(value.double_value.into()),
399            proto::AttributeAnyValueType::ARRAY_VALUE => {
400                let mut values = Vec::new();
401                if let Some(array_value) = value.array_value.as_ref() {
402                    for value in array_value.values.iter() {
403                        values.push(AttributeArrayValue::from(value));
404                    }
405                }
406                AttributeAnyValue::Array(values)
407            }
408        }
409    }
410}
411
412#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
413enum AttributeArrayValue {
414    String(MetaString),
415    Boolean(bool),
416    Integer(i64),
417    Double(WrappedFloat),
418}
419
420impl From<&proto::AttributeArrayValue> for AttributeArrayValue {
421    fn from(value: &proto::AttributeArrayValue) -> Self {
422        let maybe_proto_value_type = value.type_.enum_value().expect("unknown/invalid arrayvalue type");
423        match maybe_proto_value_type {
424            proto::AttributeArrayValueType::STRING_VALUE => AttributeArrayValue::String((*value.string_value).into()),
425            proto::AttributeArrayValueType::BOOL_VALUE => AttributeArrayValue::Boolean(value.bool_value),
426            proto::AttributeArrayValueType::INT_VALUE => AttributeArrayValue::Integer(value.int_value),
427            proto::AttributeArrayValueType::DOUBLE_VALUE => AttributeArrayValue::Double(value.double_value.into()),
428        }
429    }
430}
431
432// ---------------------------------------------------------------------------
433// ETP (Efficient Trace Payload) helpers
434// ---------------------------------------------------------------------------
435
436/// Split return type for [`split_etp_span_attributes`]: (meta, metrics, meta struct).
437type SpanAttributeSplit = (
438    FastHashMap<MetaString, MetaString>,
439    FastHashMap<MetaString, WrappedFloat>,
440    FastHashMap<MetaString, Vec<u8>>,
441);
442
443/// Resolves a string table reference to a `MetaString`.
444fn resolve_ref(strings: &[String], r: u32) -> MetaString {
445    strings
446        .get(r as usize)
447        .map(|s| MetaString::from(s.as_str()))
448        .unwrap_or_default()
449}
450
451/// Extracts the low and high 64-bit halves from a 16-byte big-endian trace ID.
452///
453/// Returns `(trace_id_low, trace_id_high)`. Both are zero when the byte slice is not 16 bytes.
454fn trace_id_parts_from_bytes(bytes: &[u8]) -> (u64, u64) {
455    if bytes.len() == 16 {
456        let high = u64::from_be_bytes(bytes[..8].try_into().unwrap_or([0u8; 8]));
457        let low = u64::from_be_bytes(bytes[8..].try_into().unwrap_or([0u8; 8]));
458        (low, high)
459    } else {
460        (0, 0)
461    }
462}
463
464/// Returns only the low 64 bits of a 16-byte big-endian trace ID.
465fn trace_id_low_from_bytes(bytes: &[u8]) -> u64 {
466    trace_id_parts_from_bytes(bytes).0
467}
468
469/// Collects only the string-valued entries from a string interned attribute map.
470fn string_attrs_from_etp(
471    attrs: &std::collections::HashMap<u32, etp_proto::AnyValue>, strings: &[String],
472) -> FastHashMap<MetaString, MetaString> {
473    attrs
474        .iter()
475        .filter_map(|(k_ref, v)| {
476            let val = match &v.value {
477                Some(etp_proto::any_value::Value::StringValueRef(r)) => resolve_ref(strings, *r),
478                _ => return None,
479            };
480            Some((resolve_ref(strings, *k_ref), val))
481        })
482        .collect()
483}
484
485/// Splits a string interned span attribute map into the three stele maps: meta (string), metrics (float),
486/// and `meta_struct` (bytes).
487fn split_etp_span_attributes(
488    attrs: &std::collections::HashMap<u32, etp_proto::AnyValue>, strings: &[String],
489) -> SpanAttributeSplit {
490    let mut meta = FastHashMap::default();
491    let mut metrics = FastHashMap::default();
492    let mut meta_struct = FastHashMap::default();
493
494    for (k_ref, v) in attrs {
495        let key = resolve_ref(strings, *k_ref);
496        match &v.value {
497            Some(etp_proto::any_value::Value::StringValueRef(r)) => {
498                meta.insert(key, resolve_ref(strings, *r));
499            }
500            Some(etp_proto::any_value::Value::DoubleValue(f)) => {
501                metrics.insert(key, WrappedFloat(OrderedFloat(*f)));
502            }
503            Some(etp_proto::any_value::Value::IntValue(i)) => {
504                metrics.insert(key, WrappedFloat(OrderedFloat(*i as f64)));
505            }
506            Some(etp_proto::any_value::Value::BoolValue(b)) => {
507                meta.insert(key, MetaString::from(if *b { "true" } else { "false" }));
508            }
509            Some(etp_proto::any_value::Value::BytesValue(b)) => {
510                meta_struct.insert(key, b.clone());
511            }
512            _ => {}
513        }
514    }
515    (meta, metrics, meta_struct)
516}
517
518/// Converts an `AnyValue` to a stele `AttributeAnyValue` for use in span events.
519fn etp_anyvalue_to_event_attr(v: &etp_proto::AnyValue, strings: &[String]) -> Option<AttributeAnyValue> {
520    match &v.value {
521        Some(etp_proto::any_value::Value::StringValueRef(r)) => {
522            Some(AttributeAnyValue::String(resolve_ref(strings, *r)))
523        }
524        Some(etp_proto::any_value::Value::BoolValue(b)) => Some(AttributeAnyValue::Boolean(*b)),
525        Some(etp_proto::any_value::Value::IntValue(i)) => Some(AttributeAnyValue::Integer(*i)),
526        Some(etp_proto::any_value::Value::DoubleValue(f)) => {
527            Some(AttributeAnyValue::Double(WrappedFloat(OrderedFloat(*f))))
528        }
529        Some(etp_proto::any_value::Value::ArrayValue(arr)) => {
530            let values = arr
531                .values
532                .iter()
533                .filter_map(|inner| etp_anyvalue_to_array_attr(inner, strings))
534                .collect();
535            Some(AttributeAnyValue::Array(values))
536        }
537        _ => None,
538    }
539}
540
541/// Converts an `AnyValue` to a stele `AttributeArrayValue` for use inside array-typed
542/// span event attributes.
543fn etp_anyvalue_to_array_attr(v: &etp_proto::AnyValue, strings: &[String]) -> Option<AttributeArrayValue> {
544    match &v.value {
545        Some(etp_proto::any_value::Value::StringValueRef(r)) => {
546            Some(AttributeArrayValue::String(resolve_ref(strings, *r)))
547        }
548        Some(etp_proto::any_value::Value::BoolValue(b)) => Some(AttributeArrayValue::Boolean(*b)),
549        Some(etp_proto::any_value::Value::IntValue(i)) => Some(AttributeArrayValue::Integer(*i)),
550        Some(etp_proto::any_value::Value::DoubleValue(f)) => {
551            Some(AttributeArrayValue::Double(WrappedFloat(OrderedFloat(*f))))
552        }
553        _ => None,
554    }
555}