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_common::task::HandleExt as _;
11use saluki_context::tags::TagSet;
12use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
13use saluki_core::data_model::event::trace::AttributeValue;
14use saluki_core::topology::{EventsBuffer, PayloadsBuffer};
15use saluki_core::{
16 components::{encoders::*, ComponentContext},
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: ComponentContext) -> Result<Box<dyn Encoder + Send>, GenericError> {
202 let metrics_builder = MetricsBuilder::from_component_context(&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);
280 let (payloads_tx, mut payloads_rx) = mpsc::channel(8);
282 let request_builder_fut = run_request_builder(trace_rb, telemetry, events_rx, payloads_tx, flush_timeout);
283 let request_builder_handle = context
285 .topology_context()
286 .global_thread_pool() .spawn_traced_named("dd-traces-request-builder", request_builder_fut);
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 match request_builder_handle.await {
326 Ok(Ok(())) => debug!("Request builder task stopped."),
327 Ok(Err(e)) => error!(error = %e, "Request builder task failed."),
328 Err(e) => error!(error = %e, "Request builder task panicked."),
329 }
330
331 debug!("Datadog Trace encoder stopped.");
332
333 Ok(())
334 }
335}
336
337async fn run_request_builder(
338 mut trace_request_builder: RequestBuilder<TraceEndpointEncoder>, telemetry: ComponentTelemetry,
339 mut events_rx: Receiver<EventsBuffer>, payloads_tx: Sender<PayloadsBuffer>, flush_timeout: std::time::Duration,
340) -> Result<(), GenericError> {
341 let mut pending_flush = false;
342 let pending_flush_timeout = sleep(flush_timeout);
343 pin!(pending_flush_timeout);
344
345 loop {
346 select! {
347 Some(event_buffer) = events_rx.recv() => {
348 for event in event_buffer {
349 let trace = match event.try_into_trace() {
350 Some(trace) => trace,
351 None => continue,
352 };
353 let trace_to_retry = match trace_request_builder.encode(trace).await {
356 Ok(None) => continue,
357 Ok(Some(trace)) => trace,
358 Err(e) => {
359 error!(error = %e, "Failed to encode trace.");
360 telemetry.events_dropped_encoder().increment(1);
361 continue;
362 }
363 };
364
365 let maybe_requests = trace_request_builder.flush().await;
366 if maybe_requests.is_empty() {
367 panic!("builder told us to flush, but gave us nothing");
368 }
369
370 for maybe_request in maybe_requests {
371 match maybe_request {
372 Ok((events, _data_points, request)) => {
373 let payload_meta = PayloadMetadata::from_event_count(events);
374 let http_payload = HttpPayload::new(payload_meta, request);
375 let payload = Payload::Http(http_payload);
376
377 payloads_tx.send(payload).await
378 .map_err(|_| generic_error!("Failed to send payload to encoder."))?;
379 },
380 Err(e) => if e.is_recoverable() {
381 continue;
383 } else {
384 return Err(GenericError::from(e).context("Failed to flush request."));
385 }
386 }
387 }
388
389 if let Err(e) = trace_request_builder.encode(trace_to_retry).await {
391 error!(error = %e, "Failed to encode trace.");
392 telemetry.events_dropped_encoder().increment(1);
393 }
394 }
395
396 debug!("Processed event buffer.");
397
398 if !pending_flush {
400 pending_flush_timeout.as_mut().reset(tokio::time::Instant::now() + flush_timeout);
401 pending_flush = true;
402 }
403 },
404 _ = &mut pending_flush_timeout, if pending_flush => {
405 debug!("Flushing pending request(s).");
406
407 pending_flush = false;
408
409 let maybe_trace_requests = trace_request_builder.flush().await;
412 for maybe_request in maybe_trace_requests {
413 match maybe_request {
414 Ok((events, _data_points, request)) => {
415 let payload_meta = PayloadMetadata::from_event_count(events);
416 let http_payload = HttpPayload::new(payload_meta, request);
417 let payload = Payload::Http(http_payload);
418
419 payloads_tx.send(payload).await
420 .map_err(|_| generic_error!("Failed to send payload to encoder."))?;
421 },
422 Err(e) => if e.is_recoverable() {
423 continue;
424 } else {
425 return Err(GenericError::from(e).context("Failed to flush request."));
426 }
427 }
428 }
429
430 debug!("All flushed requests sent to I/O task. Waiting for next event buffer...");
431 },
432
433 else => break,
435 }
436 }
437
438 Ok(())
439}
440
441#[derive(Debug)]
442struct TraceEndpointEncoder {
443 scratch: ScratchWriter<Vec<u8>>,
444 default_hostname: MetaString,
445 agent_hostname: String,
446 version: String,
447 env: String,
448 target_traces_per_second: f64,
449 errors_per_second: f64,
450 ignore_missing_datadog_fields: bool,
451 sampling_percentage: f64,
452 string_builder: StringBuilder,
453 string_table: StringTable,
454 error_tracking_standalone: bool,
455 extra_headers: Vec<(HeaderName, HeaderValue)>,
456}
457
458impl TraceEndpointEncoder {
459 fn new(
460 default_hostname: MetaString, version: String, env: String, target_traces_per_second: f64,
461 errors_per_second: f64, error_tracking_standalone: bool, ignore_missing_datadog_fields: bool,
462 sampling_percentage: f64,
463 ) -> Self {
464 let extra_headers = if error_tracking_standalone {
465 vec![(
466 HeaderName::from_static("x-datadog-error-tracking-standalone"),
467 HeaderValue::from_static("true"),
468 )]
469 } else {
470 Vec::new()
471 };
472 Self {
473 scratch: ScratchWriter::new(Vec::with_capacity(8192)),
474 agent_hostname: default_hostname.as_ref().to_string(),
475 default_hostname,
476 version,
477 env,
478 target_traces_per_second,
479 errors_per_second,
480 ignore_missing_datadog_fields,
481 sampling_percentage,
482 string_builder: StringBuilder::new(),
483 string_table: StringTable::new(),
484 error_tracking_standalone,
485 extra_headers,
486 }
487 }
488
489 fn encode_tracer_payload(&mut self, trace: &Trace, output_buffer: &mut Vec<u8>) -> std::io::Result<()> {
490 let sampling_rate = self.sampling_rate();
491 let source = attributes_to_source(&trace.attributes);
492
493 let tracer_version = format!("otlp-{}", &trace.payload.tracer_version);
495 let container_tags =
496 resolve_container_tags_from_attrs(&trace.attributes, source.as_ref(), self.ignore_missing_datadog_fields);
497 let env_str: Option<&str> = if !trace.payload.env.is_empty() {
498 Some(&trace.payload.env)
499 } else if self.ignore_missing_datadog_fields {
500 Some("")
501 } else {
502 None
503 };
504 let hostname_str: Option<&str> = resolve_hostname_from_payload(
505 &trace.payload.hostname,
506 source.as_ref(),
507 Some(self.default_hostname.as_ref()),
508 self.ignore_missing_datadog_fields,
509 );
510 let decision_maker = trace.decision_maker.as_deref();
511 let priority = trace.priority.unwrap_or(DEFAULT_CHUNK_PRIORITY);
512 let dropped_trace = trace.dropped_trace;
513 let otlp_sr = trace.otlp_sampling_rate.unwrap_or(sampling_rate);
514 self.string_builder.clear();
515 write!(&mut self.string_builder, "{:.2}", otlp_sr).expect("should never fail to format sampling rate");
516
517 let mut trace_id_bytes = [0u8; 16];
519 trace_id_bytes[..8].copy_from_slice(&trace.trace_id_high.to_be_bytes());
520 trace_id_bytes[8..].copy_from_slice(&trace.trace_id_low.to_be_bytes());
521
522 self.string_table.clear();
524
525 let mut ap_builder = AgentPayloadBuilder::new(&mut self.scratch);
526
527 ap_builder
528 .host_name(&self.agent_hostname)?
529 .env(&self.env)?
530 .agent_version(&self.version)?
531 .target_tps(self.target_traces_per_second)?
532 .error_tps(self.errors_per_second)?;
533
534 ap_builder.add_idx_tracer_payloads(|tp| {
535 if !trace.payload.container_id.is_empty() {
539 tp.container_id_ref(self.string_table.intern(&trace.payload.container_id))?;
540 }
541 if !trace.payload.language_name.is_empty() {
542 tp.language_name_ref(self.string_table.intern(&trace.payload.language_name))?;
543 }
544 if !trace.payload.language_version.is_empty() {
545 tp.language_version_ref(self.string_table.intern(&trace.payload.language_version))?;
546 }
547 tp.tracer_version_ref(self.string_table.intern(tracer_version.as_str()))?;
548 if !trace.payload.runtime_id.is_empty() {
549 tp.runtime_id_ref(self.string_table.intern(&trace.payload.runtime_id))?;
550 }
551 if let Some(e) = env_str {
552 tp.env_ref(self.string_table.intern(e))?;
553 }
554 if let Some(h) = hostname_str {
555 tp.hostname_ref(self.string_table.intern(h))?;
556 }
557 if !trace.payload.app_version.is_empty() {
558 tp.app_version_ref(self.string_table.intern(&trace.payload.app_version))?;
559 }
560
561 if let Some(ct) = &container_tags {
563 let k_ref = self.string_table.intern(CONTAINER_TAGS_META_KEY);
564 let v_ref = self.string_table.intern(ct);
565 tp.attributes().write_entry(k_ref, |av: &mut _| {
566 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
567 })?;
568 }
569
570 tp.add_chunks(|chunk| {
572 chunk.priority(priority)?;
573
574 if !trace.origin.is_empty() {
575 chunk.origin_ref(self.string_table.intern(&trace.origin))?;
576 }
577
578 chunk.trace_id(&trace_id_bytes)?;
580
581 if trace.sampling_mechanism != 0 {
583 chunk.sampling_mechanism(trace.sampling_mechanism)?;
584 }
585
586 for span in trace.spans() {
588 let service_ref = self.string_table.intern(span.service());
589 let name_ref = self.string_table.intern(span.name());
590 let resource_ref = self.string_table.intern(span.resource());
591 let type_ref = self.string_table.intern(span.span_type());
592 let env_ref = (!span.env.is_empty()).then(|| self.string_table.intern(&span.env));
593 let version_ref = (!span.version.is_empty()).then(|| self.string_table.intern(&span.version));
594 let component_ref = (!span.component.is_empty()).then(|| self.string_table.intern(&span.component));
595
596 chunk.add_spans(|s| {
597 s.service_ref(service_ref)?
598 .name_ref(name_ref)?
599 .resource_ref(resource_ref)?
600 .span_id(span.span_id())?
601 .parent_id(span.parent_id())?
602 .start(span.start())?
603 .duration(span.duration())?
604 .error(span.error() != 0)?;
605
606 {
608 let mut attrs = s.attributes();
609 for (k, v) in &span.attributes {
610 let k_ref = self.string_table.intern(k);
611 attrs.write_entry(k_ref, |av: &mut _| {
612 encode_etp_attribute_value(av, v, &mut self.string_table)
613 })?;
614 }
615 }
616
617 s.type_ref(type_ref)?;
618
619 if let Some(er) = env_ref {
620 s.env_ref(er)?;
621 }
622 if let Some(vr) = version_ref {
623 s.version_ref(vr)?;
624 }
625 if let Some(cr) = component_ref {
626 s.component_ref(cr)?;
627 }
628 if span.kind != 0 {
629 s.kind(SpanKind::from(span.kind as i32))?;
630 }
631
632 for link in span.span_links() {
634 let mut link_trace_id_bytes = [0u8; 16];
635 link_trace_id_bytes[..8].copy_from_slice(&link.trace_id_high().to_be_bytes());
636 link_trace_id_bytes[8..].copy_from_slice(&link.trace_id().to_be_bytes());
637 let tracestate_ref = self.string_table.intern(link.tracestate());
638
639 s.add_links(|sl| {
640 sl.trace_id(&link_trace_id_bytes)?.span_id(link.span_id())?;
641 {
642 let mut lattrs = sl.attributes();
643 for (k, v) in link.attributes() {
644 let k_ref = self.string_table.intern(k);
645 lattrs.write_entry(k_ref, |av: &mut _| {
646 encode_etp_attribute_value(av, v, &mut self.string_table)
647 })?;
648 }
649 }
650 sl.tracestate_ref(tracestate_ref)?.flags(link.flags())?;
651 Ok(())
652 })?;
653 }
654
655 for event in span.span_events() {
657 let name_ref = self.string_table.intern(event.name());
658 s.add_events(|se| {
659 se.time(event.time_unix_nano())?.name_ref(name_ref)?;
660 {
661 let mut eattrs = se.attributes();
662 for (k, v) in event.attributes() {
663 let k_ref = self.string_table.intern(k);
664 eattrs.write_entry(k_ref, |av: &mut _| {
665 encode_etp_attribute_value(av, v, &mut self.string_table)
666 })?;
667 }
668 }
669 Ok(())
670 })?;
671 }
672
673 Ok(())
674 })?;
675 }
676
677 {
679 let mut cattrs = chunk.attributes();
680 if let Some(dm) = decision_maker {
681 let k_ref = self.string_table.intern(TAG_DECISION_MAKER);
682 let v_ref = self.string_table.intern(dm);
683 cattrs.write_entry(k_ref, |av: &mut _| {
684 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
685 })?;
686 }
687 if self.error_tracking_standalone {
688 let trace_has_error = trace.spans().iter().any(|span| {
689 span.error() != 0
690 || span
691 .attributes
692 .get("_dd.span_events.has_exception")
693 .and_then(AttributeValue::as_string)
694 .is_some_and(|v| v == "true")
695 });
696 if trace_has_error {
697 let k_ref = self.string_table.intern(TAG_ETS_STANDALONE_ERROR_KEY);
698 let v_ref = self.string_table.intern(TAG_ETS_STANDALONE_ERROR_VALUE);
699 cattrs.write_entry(k_ref, |av: &mut _| {
700 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
701 })?;
702 }
703 }
704 {
705 let k_ref = self.string_table.intern(TAG_OTLP_SAMPLING_RATE);
706 let v_ref = self.string_table.intern(self.string_builder.as_str());
707 cattrs.write_entry(k_ref, |av: &mut _| {
708 av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
709 })?;
710 }
711 }
712
713 if dropped_trace {
714 chunk.dropped_trace(true)?;
715 }
716
717 Ok(())
718 })?;
719
720 tp.strings(|sb| sb.add_many_mapped(&self.string_table.indices, |s| &**s))?;
724
725 Ok(())
726 })?;
727
728 ap_builder.finish(output_buffer)?;
729
730 Ok(())
731 }
732
733 fn sampling_rate(&self) -> f64 {
734 let rate = self.sampling_percentage / 100.0;
735 if rate <= 0.0 || rate >= 1.0 {
736 return 1.0;
737 }
738 rate
739 }
740}
741
742impl EndpointEncoder for TraceEndpointEncoder {
743 type Input = Trace;
744 type EncodeError = std::io::Error;
745 fn encoder_name() -> &'static str {
746 "traces"
747 }
748
749 fn compressed_size_limit(&self) -> usize {
750 DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT
751 }
752
753 fn uncompressed_size_limit(&self) -> usize {
754 DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT
755 }
756
757 fn encode(&mut self, trace: &Self::Input, buffer: &mut Vec<u8>) -> Result<(), Self::EncodeError> {
758 self.encode_tracer_payload(trace, buffer)
759 }
760
761 fn endpoint_uri(&self) -> Uri {
762 PathAndQuery::from_static("/api/v0.2/traces").into()
763 }
764
765 fn endpoint_method(&self) -> Method {
766 Method::POST
767 }
768
769 fn content_type(&self) -> HeaderValue {
770 CONTENT_TYPE_PROTOBUF.clone()
771 }
772
773 fn additional_headers(&self) -> &[(HeaderName, HeaderValue)] {
774 &self.extra_headers
775 }
776}
777
778fn encode_etp_attribute_value<S: ScratchBuffer>(
781 builder: &mut datadog_protos::traces::builders::idx::AnyValueBuilder<'_, S>, value: &AttributeValue,
782 st: &mut StringTable,
783) -> std::io::Result<()> {
784 builder
785 .value(|vo| match value {
786 AttributeValue::String(s) => vo.string_value_ref(st.intern(s)),
787 AttributeValue::Bool(b) => vo.bool_value(*b),
788 AttributeValue::Int(i) => vo.int_value(*i),
789 AttributeValue::Float(f) => vo.double_value(*f),
790 AttributeValue::Bytes(b) => vo.bytes_value(b),
791 AttributeValue::Array(values) => vo.array_value(|arr| {
792 for v in values {
793 arr.add_values(|av| encode_etp_attribute_value(av, v, st))?;
794 }
795 Ok(())
796 }),
797 AttributeValue::KeyValueList(kvs) => vo.key_value_list(|kvl| {
798 for (k, v) in kvs {
799 kvl.add_key_values(|kv| {
800 kv.key(st.intern(k))?
801 .value(|av| encode_etp_attribute_value(av, v, st))?;
802 Ok(())
803 })?;
804 }
805 Ok(())
806 }),
807 })
808 .map(|_| ())
809}
810
811fn resolve_hostname_from_payload<'a>(
812 payload_hostname: &'a str, source: Option<&'a OtlpSource>, default_hostname: Option<&'a str>,
813 ignore_missing_fields: bool,
814) -> Option<&'a str> {
815 if !payload_hostname.is_empty() {
816 return Some(payload_hostname);
817 }
818 if ignore_missing_fields {
819 return Some("");
820 }
821 match source {
822 Some(src) => match src.kind {
823 OtlpSourceKind::HostnameKind => Some(src.identifier.as_str()),
824 _ => Some(""),
825 },
826 None => default_hostname,
827 }
828}
829
830fn resolve_container_tags_from_attrs(
831 attributes: &FastHashMap<MetaString, AttributeValue>, source: Option<&OtlpSource>, ignore_missing_fields: bool,
832) -> Option<MetaString> {
833 if let Some(AttributeValue::String(tags)) = attributes.get(KEY_DATADOG_CONTAINER_TAGS) {
834 if !tags.is_empty() {
835 return Some(tags.clone());
836 }
837 }
838
839 if ignore_missing_fields {
840 return None;
841 }
842 let mut container_tags = TagSet::default();
843 extract_container_tags_from_attributes_map(attributes, &mut container_tags);
844 let is_fargate_source = source.is_some_and(|src| src.kind == OtlpSourceKind::AwsEcsFargateKind);
845 if container_tags.is_empty() && !is_fargate_source {
846 return None;
847 }
848
849 let mut flattened = flatten_container_tag(container_tags);
850 if is_fargate_source {
851 if let Some(src) = source {
852 append_tags(&mut flattened, &src.tag());
853 }
854 }
855
856 if flattened.is_empty() {
857 None
858 } else {
859 Some(MetaString::from(flattened))
860 }
861}
862
863fn flatten_container_tag(tags: TagSet) -> String {
864 let mut flattened = String::new();
865 for tag in tags {
866 if !flattened.is_empty() {
867 flattened.push(',');
868 }
869 flattened.push_str(tag.as_str());
870 }
871 flattened
872}
873
874fn append_tags(target: &mut String, tags: &str) {
875 if tags.is_empty() {
876 return;
877 }
878 if !target.is_empty() {
879 target.push(',');
880 }
881 target.push_str(tags);
882}
883
884#[cfg(test)]
885mod tests {
886 use std::collections::{BTreeSet, HashMap};
887
888 use datadog_protos::traces::{idx, AgentPayload};
889 use protobuf::Message as _;
890 use saluki_context::tags::Tag;
891 use saluki_core::data_model::event::trace::{Span as DdSpan, Trace};
892 use stringtheory::MetaString;
893
894 use super::*;
895
896 fn resolve_ref(strings: &[String], string_ref: u32) -> &str {
910 strings.get(string_ref as usize).map(String::as_str).unwrap_or_default()
911 }
912
913 fn string_attrs(attrs: &HashMap<u32, idx::AnyValue>, strings: &[String]) -> HashMap<String, String> {
918 let mut out = HashMap::new();
919 for (k_ref, value) in attrs {
920 if let Some(idx::any_value::Value::StringValueRef(v_ref)) = &value.value {
921 let key = resolve_ref(strings, *k_ref);
922 if !key.is_empty() {
923 out.insert(key.to_string(), resolve_ref(strings, *v_ref).to_string());
924 }
925 }
926 }
927 out
928 }
929
930 fn decode_etp_chunk_attributes(buf: &[u8]) -> Vec<HashMap<String, String>> {
933 let payload = AgentPayload::parse_from_bytes(buf).expect("AgentPayload should decode");
934 payload
935 .idxTracerPayloads
936 .iter()
937 .flat_map(|tp| {
938 let strings = tp.strings();
939 tp.chunks()
940 .iter()
941 .map(move |chunk| string_attrs(chunk.attributes(), strings))
942 })
943 .collect()
944 }
945
946 fn decode_etp_tracer_versions(buf: &[u8]) -> Vec<String> {
948 let payload = AgentPayload::parse_from_bytes(buf).expect("AgentPayload should decode");
949 payload
950 .idxTracerPayloads
951 .iter()
952 .map(|tp| resolve_ref(tp.strings(), tp.tracerVersionRef()).to_string())
953 .collect()
954 }
955
956 const DEFAULT_TARGET_TPS: f64 = 10.0;
958 const DEFAULT_ERRORS_PER_SECOND: f64 = 10.0;
959 const DEFAULT_SAMPLING_PERCENTAGE: f64 = 100.0;
961
962 fn make_encoder(ets_enabled: bool) -> TraceEndpointEncoder {
963 TraceEndpointEncoder::new(
964 MetaString::from("test-host"),
965 "0.0.0".to_string(),
966 "none".to_string(),
967 DEFAULT_TARGET_TPS,
968 DEFAULT_ERRORS_PER_SECOND,
969 ets_enabled,
970 false,
971 DEFAULT_SAMPLING_PERCENTAGE,
972 )
973 }
974
975 fn make_trace() -> Trace {
976 let span = DdSpan::new(
977 MetaString::from("svc"),
978 MetaString::from("op"),
979 MetaString::from("res"),
980 MetaString::from("web"),
981 1, 0, 0, 1000, 0, );
987 let mut trace = Trace::new(vec![span]);
988 trace.priority = Some(1);
989 trace
990 }
991
992 fn make_error_trace() -> Trace {
993 let span = DdSpan::new(
994 MetaString::from("svc"),
995 MetaString::from("op"),
996 MetaString::from("res"),
997 MetaString::from("web"),
998 1, 0, 0, 1000, 1, );
1004 let mut trace = Trace::new(vec![span]);
1005 trace.priority = Some(1);
1006 trace
1007 }
1008
1009 #[test]
1010 fn ets_header_present_when_enabled() {
1011 let encoder = make_encoder(true);
1012 let headers = encoder.additional_headers();
1013 assert_eq!(headers.len(), 1);
1014 assert_eq!(headers[0].0.as_str(), "x-datadog-error-tracking-standalone");
1015 assert_eq!(headers[0].1, "true");
1016 }
1017
1018 #[test]
1019 fn ets_header_absent_when_disabled() {
1020 let encoder = make_encoder(false);
1021 assert!(encoder.additional_headers().is_empty());
1022 }
1023
1024 #[test]
1025 fn ets_chunk_tag_present_for_error_trace() {
1026 let mut encoder = make_encoder(true);
1027 let trace = make_error_trace();
1028 let mut buf = Vec::new();
1029 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1030 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1031 let tag_value = chunk_attrs
1032 .iter()
1033 .find_map(|attrs| attrs.get("_dd.error_tracking_standalone.error").map(|v| v.as_str()));
1034 assert_eq!(
1035 tag_value,
1036 Some("true"),
1037 "ETS chunk tag should be present for error traces when ETS is enabled"
1038 );
1039 }
1040
1041 #[test]
1042 fn ets_chunk_tag_absent_for_non_error_trace() {
1043 let mut encoder = make_encoder(true);
1044 let trace = make_trace(); let mut buf = Vec::new();
1046 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1047 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1048 let has_tag = chunk_attrs
1049 .iter()
1050 .any(|attrs| attrs.contains_key("_dd.error_tracking_standalone.error"));
1051 assert!(!has_tag, "ETS chunk tag should be absent for non-error traces");
1052 }
1053
1054 #[test]
1055 fn ets_chunk_tag_absent_when_disabled() {
1056 let mut encoder = make_encoder(false);
1057 let trace = make_trace();
1058 let mut buf = Vec::new();
1059 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1060 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1061 let has_tag = chunk_attrs
1062 .iter()
1063 .any(|attrs| attrs.contains_key("_dd.error_tracking_standalone.error"));
1064 assert!(!has_tag, "ETS chunk tag should be absent when ETS is disabled");
1065 }
1066
1067 #[test]
1068 fn sampling_rate_clamps_percentage_to_unit_interval() {
1069 let cases = [
1072 (25.0, 0.25),
1073 (50.0, 0.5),
1074 (0.0, 1.0),
1075 (-10.0, 1.0),
1076 (100.0, 1.0),
1077 (150.0, 1.0),
1078 ];
1079 for (percentage, expected) in cases {
1080 let encoder = TraceEndpointEncoder::new(
1081 MetaString::from("test-host"),
1082 "0.0.0".to_string(),
1083 "none".to_string(),
1084 DEFAULT_TARGET_TPS,
1085 DEFAULT_ERRORS_PER_SECOND,
1086 false,
1087 false,
1088 percentage,
1089 );
1090 assert_eq!(expected, encoder.sampling_rate(), "sampling_rate for {percentage}%");
1091 }
1092 }
1093
1094 #[test]
1095 fn resolve_hostname_from_payload_prefers_payload_then_source_then_default() {
1096 let host_source = OtlpSource {
1097 kind: OtlpSourceKind::HostnameKind,
1098 identifier: "resolved-host".to_string(),
1099 };
1100 let fargate_source = OtlpSource {
1101 kind: OtlpSourceKind::AwsEcsFargateKind,
1102 identifier: "task-arn".to_string(),
1103 };
1104
1105 assert_eq!(
1107 Some("payload-host"),
1108 resolve_hostname_from_payload("payload-host", Some(&host_source), Some("default"), false)
1109 );
1110 assert_eq!(
1112 Some(""),
1113 resolve_hostname_from_payload("", Some(&host_source), Some("default"), true)
1114 );
1115 assert_eq!(
1117 Some("resolved-host"),
1118 resolve_hostname_from_payload("", Some(&host_source), Some("default"), false)
1119 );
1120 assert_eq!(
1122 Some(""),
1123 resolve_hostname_from_payload("", Some(&fargate_source), Some("default"), false)
1124 );
1125 assert_eq!(
1127 Some("default"),
1128 resolve_hostname_from_payload("", None, Some("default"), false)
1129 );
1130 assert_eq!(None, resolve_hostname_from_payload("", None, None, false));
1131 }
1132
1133 #[test]
1134 fn append_tags_joins_non_empty_segments_with_commas() {
1135 let mut target = String::new();
1136
1137 append_tags(&mut target, "");
1139 assert_eq!("", target);
1140
1141 append_tags(&mut target, "a:1");
1143 assert_eq!("a:1", target);
1144
1145 append_tags(&mut target, "b:2");
1147 assert_eq!("a:1,b:2", target);
1148
1149 append_tags(&mut target, "");
1151 assert_eq!("a:1,b:2", target);
1152 }
1153
1154 #[test]
1155 fn flatten_container_tag_comma_joins_the_tag_set() {
1156 assert_eq!("", flatten_container_tag(TagSet::default()));
1157
1158 let single: TagSet = std::iter::once(Tag::from_static("image_name:web")).collect();
1159 assert_eq!("image_name:web", flatten_container_tag(single));
1160
1161 let multiple: TagSet = ["image_name:web", "runtime:docker"]
1162 .into_iter()
1163 .map(Tag::from_static)
1164 .collect();
1165 let flattened = flatten_container_tag(multiple);
1166 assert_eq!(
1167 BTreeSet::from(["image_name:web", "runtime:docker"]),
1168 flattened.split(',').collect::<BTreeSet<_>>()
1169 );
1170 }
1171
1172 #[test]
1173 fn resolve_container_tags_prefers_explicit_container_tags_attribute() {
1174 let mut attributes = FastHashMap::default();
1176 attributes.insert(
1177 MetaString::from(KEY_DATADOG_CONTAINER_TAGS),
1178 AttributeValue::String(MetaString::from("region:us,team:core")),
1179 );
1180 assert_eq!(
1181 Some(MetaString::from("region:us,team:core")),
1182 resolve_container_tags_from_attrs(&attributes, None, false)
1183 );
1184 }
1185
1186 #[test]
1187 fn resolve_container_tags_returns_none_without_container_attributes() {
1188 let attributes = FastHashMap::default();
1189 assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, true));
1191 assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, false));
1193 }
1194
1195 #[test]
1196 fn encode_prefixes_tracer_version_and_writes_otlp_sampling_rate() {
1197 let mut encoder = make_encoder(false);
1198 let mut trace = make_trace();
1199 trace.payload.tracer_version = MetaString::from("1.2.3");
1200 trace.otlp_sampling_rate = Some(0.5);
1201
1202 let mut buf = Vec::new();
1203 encoder.encode(&trace, &mut buf).expect("encode should succeed");
1204
1205 let tracer_versions = decode_etp_tracer_versions(&buf);
1207 let tracer_version = tracer_versions.first().expect("a tracer payload should be encoded");
1208 assert_eq!("otlp-1.2.3", tracer_version);
1210
1211 let chunk_attrs = decode_etp_chunk_attributes(&buf);
1213 let otlp_sr = chunk_attrs
1214 .iter()
1215 .find_map(|attrs| attrs.get("_dd.otlp_sr"))
1216 .expect("chunk should carry the _dd.otlp_sr tag");
1217 assert_eq!("0.50", otlp_sr.as_str());
1218 }
1219}