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