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#[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 pub fn trace_id(&self) -> u64 {
163 self.trace_id
164 }
165
166 pub fn span_id(&self) -> u64 {
168 self.span_id
169 }
170
171 pub fn get_meta_field(&self, meta_key: &str) -> Option<&str> {
173 self.meta.get(meta_key).map(|s| &**s)
174 }
175
176 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 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
432type SpanAttributeSplit = (
438 FastHashMap<MetaString, MetaString>,
439 FastHashMap<MetaString, WrappedFloat>,
440 FastHashMap<MetaString, Vec<u8>>,
441);
442
443fn 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
451fn 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
464fn trace_id_low_from_bytes(bytes: &[u8]) -> u64 {
466 trace_id_parts_from_bytes(bytes).0
467}
468
469fn 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
485fn 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
518fn 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
541fn 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}