1use std::{fmt::Write, time::Duration};
2
3use agent_data_plane_config::{defaults::DEFAULT_TRACE_ENV, domains, SharedConfiguration};
4use async_trait::async_trait;
5use datadog_protos::traces::builders::{idx::SpanKind, AgentPayloadBuilder};
6use http::{uri::PathAndQuery, HeaderName, HeaderValue, Method, Uri};
7use piecemeal::{ScratchBuffer, ScratchWriter};
8use saluki_common::collections::{FastHashMap, FastIndexSet};
9use saluki_common::strings::StringBuilder;
10use saluki_core::{
11 accounting::{MemoryBounds, MemoryBoundsBuilder},
12 components::{encoders::*, BuildContext},
13 data_model::{
14 event::{
15 trace::{AttributeValue, Trace},
16 EventType,
17 },
18 payload::{HttpPayload, Payload, PayloadMetadata, PayloadType},
19 tags::TagSet,
20 },
21 observability::ComponentMetricsExt as _,
22 runtime,
23 topology::{EventsBuffer, PayloadsBuffer},
24};
25use saluki_env::host::providers::BoxedHostProvider;
26use saluki_env::{EnvironmentProvider, HostProvider};
27use saluki_error::generic_error;
28use saluki_error::{ErrorContext as _, GenericError};
29use saluki_io::compression::CompressionScheme;
30use saluki_metrics::MetricsBuilder;
31use stringtheory::MetaString;
32use tokio::pin;
33use tokio::{
34 select,
35 sync::mpsc::{self, Receiver, Sender},
36 time::sleep,
37};
38use tracing::{debug, error};
39
40use crate::common::datadog::{
41 io::RB_BUFFER_CHUNK_SIZE,
42 request_builder::{EndpointEncoder, RequestBuilder},
43 telemetry::ComponentTelemetry,
44 DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT, DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT, TAG_DECISION_MAKER,
45};
46use crate::common::otlp::util::{
47 attributes_to_source, extract_container_tags_from_attributes_map, Source as OtlpSource,
48 SourceKind as OtlpSourceKind, KEY_DATADOG_CONTAINER_TAGS,
49};
50
51const CONTAINER_TAGS_META_KEY: &str = "_dd.tags.container";
52const MAX_TRACES_PER_PAYLOAD: usize = 10000;
53static CONTENT_TYPE_PROTOBUF: HeaderValue = HeaderValue::from_static("application/x-protobuf");
54
55const TAG_OTLP_SAMPLING_RATE: &str = "_dd.otlp_sr";
57const DEFAULT_CHUNK_PRIORITY: i32 = 1; const TAG_ETS_STANDALONE_ERROR_KEY: &str = "_dd.error_tracking_standalone.error";
61const TAG_ETS_STANDALONE_ERROR_VALUE: &str = "true";
62
63#[derive(Debug)]
67struct StringTable {
68 indices: FastIndexSet<MetaString>,
69}
70
71impl Default for StringTable {
72 fn default() -> Self {
73 Self::new()
74 }
75}
76
77impl StringTable {
78 fn new() -> Self {
79 let mut t = Self {
80 indices: FastIndexSet::default(),
81 };
82 t.intern("");
83 t
84 }
85
86 fn clear(&mut self) {
87 self.indices.clear();
88 self.intern("");
89 }
90
91 fn intern(&mut self, s: &str) -> u32 {
93 if let Some(idx) = self.indices.get_index_of(s) {
95 idx as u32
96 } else {
97 let (idx, _) = self.indices.insert_full(MetaString::from(s));
98 idx as u32
99 }
100 }
101}
102
103pub struct DatadogTraceConfiguration {
109 compressor_kind: String,
111
112 zstd_compressor_level: i32,
114
115 flush_timeout: Duration,
119
120 env: String,
122
123 target_traces_per_second: f64,
125
126 errors_per_second: f64,
128
129 error_tracking_standalone: bool,
131
132 ignore_missing_datadog_fields: bool,
134
135 sampling_percentage: f64,
137
138 default_hostname: Option<String>,
140
141 version: String,
143}
144
145impl DatadogTraceConfiguration {
146 pub fn from_configuration(
151 traces: &domains::traces::Domain, otlp_traces: &domains::otlp::Traces, shared: &SharedConfiguration,
152 ) -> Self {
153 let app_details = saluki_metadata::get_app_details();
154 let version = format!("agent-data-plane/{}", app_details.version().raw());
155
156 let compression = &shared.endpoints.compression;
157
158 let env = if traces.env.is_empty() {
161 DEFAULT_TRACE_ENV.to_owned()
162 } else {
163 traces.env.clone()
164 };
165
166 Self {
167 compressor_kind: compression.compressor_kind.clone(),
168 zstd_compressor_level: compression.effective_zstd_level(),
169 flush_timeout: shared.metrics_encoding.flush_timeout,
170 env,
171 target_traces_per_second: traces.target_traces_per_second,
172 errors_per_second: traces.errors_per_second,
173 error_tracking_standalone: traces.error_tracking_standalone_enabled,
174 ignore_missing_datadog_fields: otlp_traces.ignore_missing_datadog_fields,
175 sampling_percentage: otlp_traces.probabilistic_sampler_sampling_percentage,
176 default_hostname: None,
177 version,
178 }
179 }
180
181 pub async fn with_environment_provider<E>(mut self, environment_provider: E) -> Result<Self, GenericError>
183 where
184 E: EnvironmentProvider<Host = BoxedHostProvider>,
185 {
186 let host_provider = environment_provider.host();
187 let hostname = host_provider.get_hostname().await?;
188 self.default_hostname = Some(hostname);
189 Ok(self)
190 }
191}
192
193#[async_trait]
194impl EncoderBuilder for DatadogTraceConfiguration {
195 fn input_event_type(&self) -> EventType {
196 EventType::Trace
197 }
198
199 fn output_payload_type(&self) -> PayloadType {
200 PayloadType::Http
201 }
202
203 async fn build(&self, context: BuildContext) -> Result<Box<dyn Encoder + Send>, GenericError> {
204 let metrics_builder = MetricsBuilder::from_component_context(context.component_context());
205 let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
206 let compression_scheme = CompressionScheme::new(&self.compressor_kind, self.zstd_compressor_level);
207
208 let default_hostname = self.default_hostname.clone().unwrap_or_default();
209 let default_hostname = MetaString::from(default_hostname);
210
211 let mut trace_rb = RequestBuilder::new(
214 TraceEndpointEncoder::new(
215 default_hostname,
216 self.version.clone(),
217 self.env.clone(),
218 self.target_traces_per_second,
219 self.errors_per_second,
220 self.error_tracking_standalone,
221 self.ignore_missing_datadog_fields,
222 self.sampling_percentage,
223 ),
224 compression_scheme,
225 RB_BUFFER_CHUNK_SIZE,
226 )
227 .await?;
228 trace_rb.with_max_inputs_per_payload(MAX_TRACES_PER_PAYLOAD);
229
230 let flush_timeout = if self.flush_timeout.is_zero() {
231 Duration::from_millis(10)
234 } else {
235 self.flush_timeout
236 };
237
238 Ok(Box::new(DatadogTrace {
239 trace_rb,
240 telemetry,
241 flush_timeout,
242 }))
243 }
244}
245
246impl MemoryBounds for DatadogTraceConfiguration {
247 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
248 builder
250 .minimum()
251 .with_single_value::<DatadogTrace>("component struct")
252 .with_array::<EventsBuffer>("request builder events channel", 8)
253 .with_array::<PayloadsBuffer>("request builder payloads channel", 8);
254
255 builder
256 .firm()
257 .with_array::<Trace>("traces split re-encode buffer", MAX_TRACES_PER_PAYLOAD);
258 }
259}
260
261pub struct DatadogTrace {
262 trace_rb: RequestBuilder<TraceEndpointEncoder>,
263 telemetry: ComponentTelemetry,
264 flush_timeout: Duration,
265}
266
267#[async_trait]
269impl Encoder for DatadogTrace {
270 async fn run(mut self: Box<Self>, mut context: EncoderContext) -> Result<(), GenericError> {
271 let Self {
272 trace_rb,
273 telemetry,
274 flush_timeout,
275 } = *self;
276
277 let mut health = context.take_health_handle();
278
279 let (events_tx, events_rx) = mpsc::channel(8);
280 let (payloads_tx, mut payloads_rx) = mpsc::channel(8);
281
282 let request_builder_fut = run_request_builder(trace_rb, telemetry, events_rx, payloads_tx, flush_timeout);
287 runtime::worker("request_builder", request_builder_fut)
288 .on_runtime(context.topology_context().global_thread_pool().clone())
289 .spawn();
290
291 health.mark_ready();
292 debug!("Datadog Trace encoder started.");
293
294 loop {
295 select! {
296 biased; _ = health.live() => continue,
299 maybe_payload = payloads_rx.recv() => match maybe_payload {
300 Some(payload) => {
301 if let Err(e) = context.dispatcher().dispatch(payload).await {
303 error!("Failed to dispatch payload: {}", e);
304 }
305 }
306 None => break,
307 },
308 maybe_event_buffer = context.events().next() => match maybe_event_buffer {
309 Some(event_buffer) => events_tx.send(event_buffer).await
310 .error_context("Failed to send event buffer to request builder task.")?,
311 None => break,
312 },
313 }
314 }
315
316 drop(events_tx);
318
319 while let Some(payload) = payloads_rx.recv().await {
321 if let Err(e) = context.dispatcher().dispatch(payload).await {
322 error!("Failed to dispatch payload: {}", e);
323 }
324 }
325
326 debug!("Datadog Trace encoder stopped.");
329
330 Ok(())
331 }
332}
333
334async fn run_request_builder(
335 mut trace_request_builder: RequestBuilder<TraceEndpointEncoder>, telemetry: ComponentTelemetry,
336 mut events_rx: Receiver<EventsBuffer>, payloads_tx: Sender<PayloadsBuffer>, flush_timeout: std::time::Duration,
337) -> Result<(), GenericError> {
338 let mut pending_flush = false;
339 let pending_flush_timeout = sleep(flush_timeout);
340 pin!(pending_flush_timeout);
341
342 loop {
343 select! {
344 Some(event_buffer) = events_rx.recv() => {
345 for event in event_buffer {
346 let trace = match event.try_into_trace() {
347 Some(trace) => trace,
348 None => continue,
349 };
350 let trace_to_retry = match trace_request_builder.encode(trace).await {
353 Ok(None) => continue,
354 Ok(Some(trace)) => trace,
355 Err(e) => {
356 error!(error = %e, "Failed to encode trace.");
357 telemetry.events_dropped_encoder().increment(1);
358 continue;
359 }
360 };
361
362 let maybe_requests = trace_request_builder.flush().await;
363 if maybe_requests.is_empty() {
364 panic!("builder told us to flush, but gave us nothing");
365 }
366
367 for maybe_request in maybe_requests {
368 match maybe_request {
369 Ok((events, _data_points, request)) => {
370 let payload_meta = PayloadMetadata::from_event_count(events);
371 let http_payload = HttpPayload::new(payload_meta, request);
372 let payload = Payload::Http(http_payload);
373
374 payloads_tx.send(payload).await
375 .map_err(|_| generic_error!("Failed to send payload to encoder."))?;
376 },
377 Err(e) => if e.is_recoverable() {
378 continue;
380 } else {
381 return Err(GenericError::from(e).context("Failed to flush request."));
382 }
383 }
384 }
385
386 if let Err(e) = trace_request_builder.encode(trace_to_retry).await {
388 error!(error = %e, "Failed to encode trace.");
389 telemetry.events_dropped_encoder().increment(1);
390 }
391 }
392
393 debug!("Processed event buffer.");
394
395 if !pending_flush {
397 pending_flush_timeout.as_mut().reset(tokio::time::Instant::now() + flush_timeout);
398 pending_flush = true;
399 }
400 },
401 _ = &mut pending_flush_timeout, if pending_flush => {
402 debug!("Flushing pending request(s).");
403
404 pending_flush = false;
405
406 let maybe_trace_requests = trace_request_builder.flush().await;
409 for maybe_request in maybe_trace_requests {
410 match maybe_request {
411 Ok((events, _data_points, request)) => {
412 let payload_meta = PayloadMetadata::from_event_count(events);
413 let http_payload = HttpPayload::new(payload_meta, request);
414 let payload = Payload::Http(http_payload);
415
416 payloads_tx.send(payload).await
417 .map_err(|_| generic_error!("Failed to send payload to encoder."))?;
418 },
419 Err(e) => if e.is_recoverable() {
420 continue;
421 } else {
422 return Err(GenericError::from(e).context("Failed to flush request."));
423 }
424 }
425 }
426
427 debug!("All flushed requests sent to I/O task. Waiting for next event buffer...");
428 },
429
430 else => break,
432 }
433 }
434
435 Ok(())
436}
437
438#[derive(Debug)]
439struct TraceEndpointEncoder {
440 scratch: ScratchWriter<Vec<u8>>,
441 default_hostname: MetaString,
442 agent_hostname: String,
443 version: String,
444 env: String,
445 target_traces_per_second: f64,
446 errors_per_second: f64,
447 ignore_missing_datadog_fields: bool,
448 sampling_percentage: f64,
449 string_builder: StringBuilder,
450 string_table: StringTable,
451 error_tracking_standalone: bool,
452 extra_headers: Vec<(HeaderName, HeaderValue)>,
453}
454
455impl TraceEndpointEncoder {
456 fn new(
457 default_hostname: MetaString, version: String, env: String, target_traces_per_second: f64,
458 errors_per_second: f64, error_tracking_standalone: bool, ignore_missing_datadog_fields: bool,
459 sampling_percentage: f64,
460 ) -> Self {
461 let extra_headers = if error_tracking_standalone {
462 vec![(
463 HeaderName::from_static("x-datadog-error-tracking-standalone"),
464 HeaderValue::from_static("true"),
465 )]
466 } else {
467 Vec::new()
468 };
469 Self {
470 scratch: ScratchWriter::new(Vec::with_capacity(8192)),
471 agent_hostname: default_hostname.as_ref().to_string(),
472 default_hostname,
473 version,
474 env,
475 target_traces_per_second,
476 errors_per_second,
477 ignore_missing_datadog_fields,
478 sampling_percentage,
479 string_builder: StringBuilder::new(),
480 string_table: StringTable::new(),
481 error_tracking_standalone,
482 extra_headers,
483 }
484 }
485
486 fn encode_tracer_payload(&mut self, trace: &Trace, output_buffer: &mut Vec<u8>) -> std::io::Result<()> {
487 let sampling_rate = self.sampling_rate();
488 let source = attributes_to_source(&trace.attributes);
489
490 let tracer_version = format!("otlp-{}", &trace.payload.tracer_version);
492 let container_tags =
493 resolve_container_tags_from_attrs(&trace.attributes, source.as_ref(), self.ignore_missing_datadog_fields);
494 let env_str: Option<&str> = if !trace.payload.env.is_empty() {
495 Some(&trace.payload.env)
496 } else if self.ignore_missing_datadog_fields {
497 Some("")
498 } else {
499 None
500 };
501 let hostname_str: Option<&str> = resolve_hostname_from_payload(
502 &trace.payload.hostname,
503 source.as_ref(),
504 Some(self.default_hostname.as_ref()),
505 self.ignore_missing_datadog_fields,
506 );
507 let decision_maker = trace.decision_maker.as_deref();
508 let priority = trace.priority.unwrap_or(DEFAULT_CHUNK_PRIORITY);
509 let dropped_trace = trace.dropped_trace;
510 let otlp_sr = trace.otlp_sampling_rate.unwrap_or(sampling_rate);
511 self.string_builder.clear();
512 write!(&mut self.string_builder, "{:.2}", otlp_sr).expect("should never fail to format sampling rate");
513
514 let mut trace_id_bytes = [0u8; 16];
516 trace_id_bytes[..8].copy_from_slice(&trace.trace_id_high.to_be_bytes());
517 trace_id_bytes[8..].copy_from_slice(&trace.trace_id_low.to_be_bytes());
518
519 self.string_table.clear();
521
522 let mut ap_builder = AgentPayloadBuilder::new(&mut self.scratch);
523
524 ap_builder
525 .host_name(&self.agent_hostname)?
526 .env(&self.env)?
527 .agent_version(&self.version)?
528 .target_tps(self.target_traces_per_second)?
529 .error_tps(self.errors_per_second)?;
530
531 ap_builder.add_idx_tracer_payloads(|tp| {
532 if !trace.payload.container_id.is_empty() {
536 tp.container_id_ref(self.string_table.intern(&trace.payload.container_id))?;
537 }
538 if !trace.payload.language_name.is_empty() {
539 tp.language_name_ref(self.string_table.intern(&trace.payload.language_name))?;
540 }
541 if !trace.payload.language_version.is_empty() {
542 tp.language_version_ref(self.string_table.intern(&trace.payload.language_version))?;
543 }
544 tp.tracer_version_ref(self.string_table.intern(tracer_version.as_str()))?;
545 if !trace.payload.runtime_id.is_empty() {
546 tp.runtime_id_ref(self.string_table.intern(&trace.payload.runtime_id))?;
547 }
548 if let Some(e) = env_str {
549 tp.env_ref(self.string_table.intern(e))?;
550 }
551 if let Some(h) = hostname_str {
552 tp.hostname_ref(self.string_table.intern(h))?;
553 }
554 if !trace.payload.app_version.is_empty() {
555 tp.app_version_ref(self.string_table.intern(&trace.payload.app_version))?;
556 }
557
558 if let Some(ct) = &container_tags {
560 let k_ref = self.string_table.intern(CONTAINER_TAGS_META_KEY);
561 let v_ref = self.string_table.intern(ct);
562 tp.attributes().write_entry(k_ref, |av: &mut _| {
563 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
564 })?;
565 }
566
567 tp.add_chunks(|chunk| {
569 chunk.priority(priority)?;
570
571 if !trace.origin.is_empty() {
572 chunk.origin_ref(self.string_table.intern(&trace.origin))?;
573 }
574
575 chunk.trace_id(&trace_id_bytes)?;
577
578 if trace.sampling_mechanism != 0 {
580 chunk.sampling_mechanism(trace.sampling_mechanism)?;
581 }
582
583 for span in trace.spans() {
585 let service_ref = self.string_table.intern(span.service());
586 let name_ref = self.string_table.intern(span.name());
587 let resource_ref = self.string_table.intern(span.resource());
588 let type_ref = self.string_table.intern(span.span_type());
589 let env_ref = (!span.env.is_empty()).then(|| self.string_table.intern(&span.env));
590 let version_ref = (!span.version.is_empty()).then(|| self.string_table.intern(&span.version));
591 let component_ref = (!span.component.is_empty()).then(|| self.string_table.intern(&span.component));
592
593 chunk.add_spans(|s| {
594 s.service_ref(service_ref)?
595 .name_ref(name_ref)?
596 .resource_ref(resource_ref)?
597 .span_id(span.span_id())?
598 .parent_id(span.parent_id())?
599 .start(span.start())?
600 .duration(span.duration())?
601 .error(span.error() != 0)?;
602
603 {
605 let mut attrs = s.attributes();
606 for (k, v) in &span.attributes {
607 let k_ref = self.string_table.intern(k);
608 attrs.write_entry(k_ref, |av: &mut _| {
609 encode_etp_attribute_value(av, v, &mut self.string_table)
610 })?;
611 }
612 }
613
614 s.type_ref(type_ref)?;
615
616 if let Some(er) = env_ref {
617 s.env_ref(er)?;
618 }
619 if let Some(vr) = version_ref {
620 s.version_ref(vr)?;
621 }
622 if let Some(cr) = component_ref {
623 s.component_ref(cr)?;
624 }
625 if span.kind != 0 {
626 s.kind(SpanKind::from(span.kind as i32))?;
627 }
628
629 for link in span.span_links() {
631 let mut link_trace_id_bytes = [0u8; 16];
632 link_trace_id_bytes[..8].copy_from_slice(&link.trace_id_high().to_be_bytes());
633 link_trace_id_bytes[8..].copy_from_slice(&link.trace_id().to_be_bytes());
634 let tracestate_ref = self.string_table.intern(link.tracestate());
635
636 s.add_links(|sl| {
637 sl.trace_id(&link_trace_id_bytes)?.span_id(link.span_id())?;
638 {
639 let mut lattrs = sl.attributes();
640 for (k, v) in link.attributes() {
641 let k_ref = self.string_table.intern(k);
642 lattrs.write_entry(k_ref, |av: &mut _| {
643 encode_etp_attribute_value(av, v, &mut self.string_table)
644 })?;
645 }
646 }
647 sl.tracestate_ref(tracestate_ref)?.flags(link.flags())?;
648 Ok(())
649 })?;
650 }
651
652 for event in span.span_events() {
654 let name_ref = self.string_table.intern(event.name());
655 s.add_events(|se| {
656 se.time(event.time_unix_nano())?.name_ref(name_ref)?;
657 {
658 let mut eattrs = se.attributes();
659 for (k, v) in event.attributes() {
660 let k_ref = self.string_table.intern(k);
661 eattrs.write_entry(k_ref, |av: &mut _| {
662 encode_etp_attribute_value(av, v, &mut self.string_table)
663 })?;
664 }
665 }
666 Ok(())
667 })?;
668 }
669
670 Ok(())
671 })?;
672 }
673
674 {
676 let mut cattrs = chunk.attributes();
677 if let Some(dm) = decision_maker {
678 let k_ref = self.string_table.intern(TAG_DECISION_MAKER);
679 let v_ref = self.string_table.intern(dm);
680 cattrs.write_entry(k_ref, |av: &mut _| {
681 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
682 })?;
683 }
684 if self.error_tracking_standalone {
685 let trace_has_error = trace.spans().iter().any(|span| {
686 span.error() != 0
687 || span
688 .attributes
689 .get("_dd.span_events.has_exception")
690 .and_then(AttributeValue::as_string)
691 .is_some_and(|v| v == "true")
692 });
693 if trace_has_error {
694 let k_ref = self.string_table.intern(TAG_ETS_STANDALONE_ERROR_KEY);
695 let v_ref = self.string_table.intern(TAG_ETS_STANDALONE_ERROR_VALUE);
696 cattrs.write_entry(k_ref, |av: &mut _| {
697 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
698 })?;
699 }
700 }
701 {
702 let k_ref = self.string_table.intern(TAG_OTLP_SAMPLING_RATE);
703 let v_ref = self.string_table.intern(self.string_builder.as_str());
704 cattrs.write_entry(k_ref, |av: &mut _| {
705 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
706 })?;
707 }
708 }
709
710 if dropped_trace {
711 chunk.dropped_trace(true)?;
712 }
713
714 Ok(())
715 })?;
716
717 tp.strings(|sb| sb.add_many_mapped(&self.string_table.indices, |s| &**s))?;
721
722 Ok(())
723 })?;
724
725 ap_builder.finish(output_buffer)?;
726
727 Ok(())
728 }
729
730 fn sampling_rate(&self) -> f64 {
731 let rate = self.sampling_percentage / 100.0;
732 if rate <= 0.0 || rate >= 1.0 {
733 return 1.0;
734 }
735 rate
736 }
737}
738
739impl EndpointEncoder for TraceEndpointEncoder {
740 type Input = Trace;
741 type EncodeError = std::io::Error;
742 fn encoder_name() -> &'static str {
743 "traces"
744 }
745
746 fn compressed_size_limit(&self) -> usize {
747 DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT
748 }
749
750 fn uncompressed_size_limit(&self) -> usize {
751 DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT
752 }
753
754 fn encode(&mut self, trace: &Self::Input, buffer: &mut Vec<u8>) -> Result<(), Self::EncodeError> {
755 self.encode_tracer_payload(trace, buffer)
756 }
757
758 fn endpoint_uri(&self) -> Uri {
759 PathAndQuery::from_static("/api/v0.2/traces").into()
760 }
761
762 fn endpoint_method(&self) -> Method {
763 Method::POST
764 }
765
766 fn content_type(&self) -> HeaderValue {
767 CONTENT_TYPE_PROTOBUF.clone()
768 }
769
770 fn additional_headers(&self) -> &[(HeaderName, HeaderValue)] {
771 &self.extra_headers
772 }
773}
774
775fn encode_etp_attribute_value<S: ScratchBuffer>(
778 builder: &mut datadog_protos::traces::builders::idx::AnyValueBuilder<'_, S>, value: &AttributeValue,
779 st: &mut StringTable,
780) -> std::io::Result<()> {
781 builder
782 .value(|vo| match value {
783 AttributeValue::String(s) => vo.string_value_ref(st.intern(s)),
784 AttributeValue::Bool(b) => vo.bool_value(*b),
785 AttributeValue::Int(i) => vo.int_value(*i),
786 AttributeValue::Float(f) => vo.double_value(*f),
787 AttributeValue::Bytes(b) => vo.bytes_value(b),
788 AttributeValue::Array(values) => vo.array_value(|arr| {
789 for v in values {
790 arr.add_values(|av| encode_etp_attribute_value(av, v, st))?;
791 }
792 Ok(())
793 }),
794 AttributeValue::KeyValueList(kvs) => vo.key_value_list(|kvl| {
795 for (k, v) in kvs {
796 kvl.add_key_values(|kv| {
797 kv.key(st.intern(k))?
798 .value(|av| encode_etp_attribute_value(av, v, st))?;
799 Ok(())
800 })?;
801 }
802 Ok(())
803 }),
804 })
805 .map(|_| ())
806}
807
808fn resolve_hostname_from_payload<'a>(
809 payload_hostname: &'a str, source: Option<&'a OtlpSource>, default_hostname: Option<&'a str>,
810 ignore_missing_fields: bool,
811) -> Option<&'a str> {
812 if !payload_hostname.is_empty() {
813 return Some(payload_hostname);
814 }
815 if ignore_missing_fields {
816 return Some("");
817 }
818 match source {
819 Some(src) => match src.kind {
820 OtlpSourceKind::HostnameKind => Some(src.identifier.as_str()),
821 _ => Some(""),
822 },
823 None => default_hostname,
824 }
825}
826
827fn resolve_container_tags_from_attrs(
828 attributes: &FastHashMap<MetaString, AttributeValue>, source: Option<&OtlpSource>, ignore_missing_fields: bool,
829) -> Option<MetaString> {
830 if let Some(AttributeValue::String(tags)) = attributes.get(KEY_DATADOG_CONTAINER_TAGS) {
831 if !tags.is_empty() {
832 return Some(tags.clone());
833 }
834 }
835
836 if ignore_missing_fields {
837 return None;
838 }
839 let mut container_tags = TagSet::default();
840 extract_container_tags_from_attributes_map(attributes, &mut container_tags);
841 let is_fargate_source = source.is_some_and(|src| src.kind == OtlpSourceKind::AwsEcsFargateKind);
842 if container_tags.is_empty() && !is_fargate_source {
843 return None;
844 }
845
846 let mut flattened = flatten_container_tag(container_tags);
847 if is_fargate_source {
848 if let Some(src) = source {
849 append_tags(&mut flattened, &src.tag());
850 }
851 }
852
853 if flattened.is_empty() {
854 None
855 } else {
856 Some(MetaString::from(flattened))
857 }
858}
859
860fn flatten_container_tag(tags: TagSet) -> String {
861 let mut flattened = String::new();
862 for tag in tags {
863 if !flattened.is_empty() {
864 flattened.push(',');
865 }
866 flattened.push_str(tag.as_str());
867 }
868 flattened
869}
870
871fn append_tags(target: &mut String, tags: &str) {
872 if tags.is_empty() {
873 return;
874 }
875 if !target.is_empty() {
876 target.push(',');
877 }
878 target.push_str(tags);
879}
880
881#[cfg(test)]
882mod tests {
883 use std::collections::{BTreeSet, HashMap};
884
885 use datadog_protos::traces::{idx, AgentPayload};
886 use protobuf::Message as _;
887 use saluki_core::data_model::{
888 event::trace::{Span as DdSpan, Trace},
889 tags::Tag,
890 };
891 use stringtheory::MetaString;
892
893 use super::*;
894
895 fn resolve_ref(strings: &[String], string_ref: u32) -> &str {
909 strings.get(string_ref as usize).map(String::as_str).unwrap_or_default()
910 }
911
912 fn string_attrs(attrs: &HashMap<u32, idx::AnyValue>, strings: &[String]) -> HashMap<String, String> {
917 let mut out = HashMap::new();
918 for (k_ref, value) in attrs {
919 if let Some(idx::any_value::Value::StringValueRef(v_ref)) = &value.value {
920 let key = resolve_ref(strings, *k_ref);
921 if !key.is_empty() {
922 out.insert(key.to_string(), resolve_ref(strings, *v_ref).to_string());
923 }
924 }
925 }
926 out
927 }
928
929 fn decode_etp_chunk_attributes(buf: &[u8]) -> Vec<HashMap<String, String>> {
932 let payload = AgentPayload::parse_from_bytes(buf).expect("AgentPayload should decode");
933 payload
934 .idxTracerPayloads
935 .iter()
936 .flat_map(|tp| {
937 let strings = tp.strings();
938 tp.chunks()
939 .iter()
940 .map(move |chunk| string_attrs(chunk.attributes(), strings))
941 })
942 .collect()
943 }
944
945 fn decode_etp_tracer_versions(buf: &[u8]) -> Vec<String> {
947 let payload = AgentPayload::parse_from_bytes(buf).expect("AgentPayload should decode");
948 payload
949 .idxTracerPayloads
950 .iter()
951 .map(|tp| resolve_ref(tp.strings(), tp.tracerVersionRef()).to_string())
952 .collect()
953 }
954
955 const DEFAULT_TARGET_TPS: f64 = 10.0;
957 const DEFAULT_ERRORS_PER_SECOND: f64 = 10.0;
958 const DEFAULT_SAMPLING_PERCENTAGE: f64 = 100.0;
960
961 fn make_encoder(ets_enabled: bool) -> TraceEndpointEncoder {
962 TraceEndpointEncoder::new(
963 MetaString::from("test-host"),
964 "0.0.0".to_string(),
965 "none".to_string(),
966 DEFAULT_TARGET_TPS,
967 DEFAULT_ERRORS_PER_SECOND,
968 ets_enabled,
969 false,
970 DEFAULT_SAMPLING_PERCENTAGE,
971 )
972 }
973
974 fn make_trace() -> Trace {
975 let span = DdSpan::new(
976 MetaString::from("svc"),
977 MetaString::from("op"),
978 MetaString::from("res"),
979 MetaString::from("web"),
980 1, 0, 0, 1000, 0, );
986 let mut trace = Trace::new(vec![span]);
987 trace.priority = Some(1);
988 trace
989 }
990
991 fn make_error_trace() -> Trace {
992 let span = DdSpan::new(
993 MetaString::from("svc"),
994 MetaString::from("op"),
995 MetaString::from("res"),
996 MetaString::from("web"),
997 1, 0, 0, 1000, 1, );
1003 let mut trace = Trace::new(vec![span]);
1004 trace.priority = Some(1);
1005 trace
1006 }
1007
1008 #[test]
1009 fn ets_header_present_when_enabled() {
1010 let encoder = make_encoder(true);
1011 let headers = encoder.additional_headers();
1012 assert_eq!(headers.len(), 1);
1013 assert_eq!(headers[0].0.as_str(), "x-datadog-error-tracking-standalone");
1014 assert_eq!(headers[0].1, "true");
1015 }
1016
1017 #[test]
1018 fn ets_header_absent_when_disabled() {
1019 let encoder = make_encoder(false);
1020 assert!(encoder.additional_headers().is_empty());
1021 }
1022
1023 #[test]
1024 fn ets_chunk_tag_present_for_error_trace() {
1025 let mut encoder = make_encoder(true);
1026 let trace = make_error_trace();
1027 let mut buf = Vec::new();
1028 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1029 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1030 let tag_value = chunk_attrs
1031 .iter()
1032 .find_map(|attrs| attrs.get("_dd.error_tracking_standalone.error").map(|v| v.as_str()));
1033 assert_eq!(
1034 tag_value,
1035 Some("true"),
1036 "ETS chunk tag should be present for error traces when ETS is enabled"
1037 );
1038 }
1039
1040 #[test]
1041 fn ets_chunk_tag_absent_for_non_error_trace() {
1042 let mut encoder = make_encoder(true);
1043 let trace = make_trace(); let mut buf = Vec::new();
1045 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1046 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1047 let has_tag = chunk_attrs
1048 .iter()
1049 .any(|attrs| attrs.contains_key("_dd.error_tracking_standalone.error"));
1050 assert!(!has_tag, "ETS chunk tag should be absent for non-error traces");
1051 }
1052
1053 #[test]
1054 fn ets_chunk_tag_absent_when_disabled() {
1055 let mut encoder = make_encoder(false);
1056 let trace = make_trace();
1057 let mut buf = Vec::new();
1058 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1059 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1060 let has_tag = chunk_attrs
1061 .iter()
1062 .any(|attrs| attrs.contains_key("_dd.error_tracking_standalone.error"));
1063 assert!(!has_tag, "ETS chunk tag should be absent when ETS is disabled");
1064 }
1065
1066 #[test]
1067 fn sampling_rate_clamps_percentage_to_unit_interval() {
1068 let cases = [
1071 (25.0, 0.25),
1072 (50.0, 0.5),
1073 (0.0, 1.0),
1074 (-10.0, 1.0),
1075 (100.0, 1.0),
1076 (150.0, 1.0),
1077 ];
1078 for (percentage, expected) in cases {
1079 let encoder = TraceEndpointEncoder::new(
1080 MetaString::from("test-host"),
1081 "0.0.0".to_string(),
1082 "none".to_string(),
1083 DEFAULT_TARGET_TPS,
1084 DEFAULT_ERRORS_PER_SECOND,
1085 false,
1086 false,
1087 percentage,
1088 );
1089 assert_eq!(expected, encoder.sampling_rate(), "sampling_rate for {percentage}%");
1090 }
1091 }
1092
1093 #[test]
1094 fn resolve_hostname_from_payload_prefers_payload_then_source_then_default() {
1095 let host_source = OtlpSource {
1096 kind: OtlpSourceKind::HostnameKind,
1097 identifier: "resolved-host".to_string(),
1098 };
1099 let fargate_source = OtlpSource {
1100 kind: OtlpSourceKind::AwsEcsFargateKind,
1101 identifier: "task-arn".to_string(),
1102 };
1103
1104 assert_eq!(
1106 Some("payload-host"),
1107 resolve_hostname_from_payload("payload-host", Some(&host_source), Some("default"), false)
1108 );
1109 assert_eq!(
1111 Some(""),
1112 resolve_hostname_from_payload("", Some(&host_source), Some("default"), true)
1113 );
1114 assert_eq!(
1116 Some("resolved-host"),
1117 resolve_hostname_from_payload("", Some(&host_source), Some("default"), false)
1118 );
1119 assert_eq!(
1121 Some(""),
1122 resolve_hostname_from_payload("", Some(&fargate_source), Some("default"), false)
1123 );
1124 assert_eq!(
1126 Some("default"),
1127 resolve_hostname_from_payload("", None, Some("default"), false)
1128 );
1129 assert_eq!(None, resolve_hostname_from_payload("", None, None, false));
1130 }
1131
1132 #[test]
1133 fn append_tags_joins_non_empty_segments_with_commas() {
1134 let mut target = String::new();
1135
1136 append_tags(&mut target, "");
1138 assert_eq!("", target);
1139
1140 append_tags(&mut target, "a:1");
1142 assert_eq!("a:1", target);
1143
1144 append_tags(&mut target, "b:2");
1146 assert_eq!("a:1,b:2", target);
1147
1148 append_tags(&mut target, "");
1150 assert_eq!("a:1,b:2", target);
1151 }
1152
1153 #[test]
1154 fn flatten_container_tag_comma_joins_the_tag_set() {
1155 assert_eq!("", flatten_container_tag(TagSet::default()));
1156
1157 let single: TagSet = std::iter::once(Tag::from_static("image_name:web")).collect();
1158 assert_eq!("image_name:web", flatten_container_tag(single));
1159
1160 let multiple: TagSet = ["image_name:web", "runtime:docker"]
1161 .into_iter()
1162 .map(Tag::from_static)
1163 .collect();
1164 let flattened = flatten_container_tag(multiple);
1165 assert_eq!(
1166 BTreeSet::from(["image_name:web", "runtime:docker"]),
1167 flattened.split(',').collect::<BTreeSet<_>>()
1168 );
1169 }
1170
1171 #[test]
1172 fn resolve_container_tags_prefers_explicit_container_tags_attribute() {
1173 let mut attributes = FastHashMap::default();
1175 attributes.insert(
1176 MetaString::from(KEY_DATADOG_CONTAINER_TAGS),
1177 AttributeValue::String(MetaString::from("region:us,team:core")),
1178 );
1179 assert_eq!(
1180 Some(MetaString::from("region:us,team:core")),
1181 resolve_container_tags_from_attrs(&attributes, None, false)
1182 );
1183 }
1184
1185 #[test]
1186 fn resolve_container_tags_returns_none_without_container_attributes() {
1187 let attributes = FastHashMap::default();
1188 assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, true));
1190 assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, false));
1192 }
1193
1194 #[test]
1195 fn encode_prefixes_tracer_version_and_writes_otlp_sampling_rate() {
1196 let mut encoder = make_encoder(false);
1197 let mut trace = make_trace();
1198 trace.payload.tracer_version = MetaString::from("1.2.3");
1199 trace.otlp_sampling_rate = Some(0.5);
1200
1201 let mut buf = Vec::new();
1202 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1203
1204 let tracer_versions = decode_etp_tracer_versions(&buf);
1206 let tracer_version = tracer_versions.first().expect("a tracer payload should be encoded");
1207 assert_eq!("otlp-1.2.3", tracer_version);
1209
1210 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1212 let otlp_sr = chunk_attrs
1213 .iter()
1214 .find_map(|attrs| attrs.get("_dd.otlp_sr"))
1215 .expect("chunk should carry the _dd.otlp_sr tag");
1216 assert_eq!("0.50", otlp_sr.as_str());
1217 }
1218}