saluki_components/encoders/datadog/traces/
mod.rs

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
53// Sampling metadata keys / values.
54const TAG_OTLP_SAMPLING_RATE: &str = "_dd.otlp_sr";
55const DEFAULT_CHUNK_PRIORITY: i32 = 1; // PRIORITY_AUTO_KEEP
56
57// ETS chunk-level attribute keys / values.
58const TAG_ETS_STANDALONE_ERROR_KEY: &str = "_dd.error_tracking_standalone.error";
59const TAG_ETS_STANDALONE_ERROR_VALUE: &str = "true";
60
61/// String interning table for the ETP (Efficient Trace Payload) tracer payload format.
62///
63/// Index 0 is always the empty string. All other strings are assigned indices in insertion order.
64#[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    /// Interns `s`, returning its index. If `s` was already interned, returns the existing index.
90    fn intern(&mut self, s: &str) -> u32 {
91        // We use get_index_of to check if the string is already interned to avoid reconstructing a new meta string if it's already interned.
92        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
101/// Configuration for the Datadog Traces encoder.
102///
103/// This encoder converts trace events into Datadog's TracerPayload protobuf format and sends them
104/// to the Datadog traces intake endpoint (`/api/v0.2/traces`). It handles batching, compression,
105/// and enrichment with metadata such as hostname, environment, and container tags.
106pub struct DatadogTraceConfiguration {
107    /// Compression algorithm applied to outgoing payloads.
108    compressor_kind: String,
109
110    /// Effective zstd compression level, resolved by the configuration layer.
111    zstd_compressor_level: i32,
112
113    /// How long the encoder waits before flushing a partially filled payload.
114    ///
115    /// A zero duration is treated as "flush almost immediately" during `build`.
116    flush_timeout: Duration,
117
118    /// Global environment tag applied to emitted payloads.
119    env: String,
120
121    /// Target sampled traces per second, forwarded to the intake as `target_tps`.
122    target_traces_per_second: f64,
123
124    /// Target sampled error traces per second, forwarded to the intake as `error_tps`.
125    errors_per_second: f64,
126
127    /// Whether Error Tracking standalone mode is enabled.
128    error_tracking_standalone: bool,
129
130    /// Whether spans missing intake-required fields are ingested rather than rejected.
131    ignore_missing_datadog_fields: bool,
132
133    /// Percentage of OTLP traces the probabilistic sampler keeps.
134    sampling_percentage: f64,
135
136    /// Default hostname, resolved from the environment provider.
137    default_hostname: Option<String>,
138
139    /// ADP version string embedded in emitted payloads.
140    version: String,
141}
142
143impl DatadogTraceConfiguration {
144    /// Creates a new `DatadogTraceConfiguration` from the resolved traces and shared configuration.
145    ///
146    /// The OTLP trace settings live in their own domain, so they arrive as a separate slice rather
147    /// than through the traces domain.
148    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        // ADP defaults the global environment tag to `none` rather than the Core Agent's empty
157        // string, so normalize an empty resolved value back to `none`.
158        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    /// Sets the `default_hostname` using the environment provider
180    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        // Create request builder for traces which is used to generate HTTP requests.
210
211        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            // We always give ourselves a minimum flush timeout of 10ms to allow for some very minimal amount of
230            // batching, while still practically flushing things almost immediately.
231            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        // TODO: How do we properly represent the requests we can generate that may be sitting around in-flight?
247        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// Encodes Trace events to TracerPayloads.
266#[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        // The encoder runs two async loops, the main encoder loop and the request builder loop,
278        // this channel is used to send events from the main encoder loop to the request builder loop safely.
279        let (events_tx, events_rx) = mpsc::channel(8);
280        // adds a channel to send payloads to the dispatcher and a channel to receive them.
281        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        // Spawn the request builder task on the global thread pool, this task is responsible for encoding traces and flushing requests.
284        let request_builder_handle = context
285            .topology_context()
286            .global_thread_pool() // Use the shared Tokio runtime thread pool.
287            .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; // makes the branches of the select statement be evaluated in order.
295
296                _ = health.live() => continue,
297                maybe_payload = payloads_rx.recv() => match maybe_payload {
298                    Some(payload) => {
299                        // Dispatch an HTTP payload to the dispatcher.
300                        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 the events sender, which signals the request builder task to stop.
315        drop(events_tx);
316
317        // Continue draining the payloads receiver until it is closed.
318        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        // Request build task should now be stopped.
325        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                    // Encode the trace. If we get it back, that means the current request is full, and we need to
354                    // flush it before we can try to encode the trace again.
355                    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                                // If the error is recoverable, we'll hold on to the trace to retry it later.
382                                continue;
383                            } else {
384                                return Err(GenericError::from(e).context("Failed to flush request."));
385                            }
386                        }
387                    }
388
389                    // Now try to encode the trace again.
390                    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 we're not already pending a flush, we'll start the countdown.
399                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                // Once we've encoded and written all traces, we flush the request builders to generate a request with
410                // anything left over. Again, we'll enqueue those requests to be sent immediately.
411                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            // Event buffers channel has been closed, and we have no pending flushing, so we're all done.
434            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        // Resolve computed metadata strings (may produce strings not directly present on trace fields).
494        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        // Build 128-bit big-endian trace ID bytes for the chunk.
518        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        // Reset the string table; strings are interned on the fly during encoding below.
523        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            // Tracer payload metadata refs (skip default/empty values).
536            // Strings are interned on the fly; the string table is written at the end of this
537            // closure so it is complete by the time tp.strings() is called.
538            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            // Container tags go in the payload-level attributes map.
562            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            // Single TraceChunk containing all spans.
571            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                // Write 128-bit trace ID.
579                chunk.trace_id(&trace_id_bytes)?;
580
581                // Sampling mechanism (only write when non-zero).
582                if trace.sampling_mechanism != 0 {
583                    chunk.sampling_mechanism(trace.sampling_mechanism)?;
584                }
585
586                // Spans.
587                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                        // Unified attribute map (replaces separate meta/metrics/meta_struct).
607                        {
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                        // Span links.
633                        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                        // Span events.
656                        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                // Chunk attributes: decision maker, ETS tag, OTLP sampling rate.
678                {
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            // Write the string table after all refs so the table is complete.
721            // Protobuf allows fields in any order; decoders that do a full parse before
722            // resolving refs handle strings-after-refs correctly.
723            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
778/// Encodes an [`AttributeValue`] into an ETP-format `AnyValue` builder, interning any
779/// string values into `st` on the fly.
780fn 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    // ---------------------------------------------------------------------------
897    // Decode helpers for assertions on the encoded ETP `AgentPayload`.
898    //
899    // The encoder emits the wire format via `piecemeal` builders; here we decode it back
900    // with the independently code-generated `rust-protobuf` types (`idx::TracerPayload`
901    // and friends). A full-message parse resolves protobuf field ordering for us, so the
902    // fact that the encoder writes string refs before the string table is a non-issue.
903    //
904    // The only ETP-specific work left is resolving `u32` string-table refs and unwrapping
905    // the `AnyValue` oneof, which the two helpers below handle.
906    // ---------------------------------------------------------------------------
907
908    /// Resolves a string-table ref against a tracer payload's string table.
909    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    /// Resolves the string-valued entries of an ETP attribute map into a readable
914    /// `key -> value` map, keeping only entries whose `AnyValue` is a string ref and
915    /// whose key resolves to a non-empty string. Keys and values are resolved through
916    /// the payload string table.
917    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    /// Collects the resolved string-valued attributes of every chunk across all ETP
931    /// tracer payloads in an encoded `AgentPayload`, one map per chunk.
932    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    /// Resolves the `tracerVersionRef` of every ETP tracer payload in an encoded `AgentPayload`.
947    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    // The APM sampler defaults: 10 target and error traces per second.
957    const DEFAULT_TARGET_TPS: f64 = 10.0;
958    const DEFAULT_ERRORS_PER_SECOND: f64 = 10.0;
959    // Keep-everything sampling percentage (the OTLP probabilistic sampler default).
960    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,    // span_id
982            0,    // parent_id
983            0,    // start
984            1000, // duration
985            0,    // error
986        );
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,    // span_id
999            0,    // parent_id
1000            0,    // start
1001            1000, // duration
1002            1,    // error
1003        );
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(); // no error
1045        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        // `sampling_percentage` is a 0..100 percentage; only strictly in-range values map to a fractional rate, and
1070        // anything <= 0 or >= 100 collapses to 1.0 (sample everything).
1071        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        // A non-empty payload hostname always wins.
1106        assert_eq!(
1107            Some("payload-host"),
1108            resolve_hostname_from_payload("payload-host", Some(&host_source), Some("default"), false)
1109        );
1110        // An empty payload plus `ignore_missing_fields` short-circuits to an empty hostname.
1111        assert_eq!(
1112            Some(""),
1113            resolve_hostname_from_payload("", Some(&host_source), Some("default"), true)
1114        );
1115        // Honoring fields, a hostname-kind source supplies its identifier.
1116        assert_eq!(
1117            Some("resolved-host"),
1118            resolve_hostname_from_payload("", Some(&host_source), Some("default"), false)
1119        );
1120        // A non-hostname (Fargate) source resolves to an empty hostname.
1121        assert_eq!(
1122            Some(""),
1123            resolve_hostname_from_payload("", Some(&fargate_source), Some("default"), false)
1124        );
1125        // With no source, it falls back to the default hostname (which may itself be absent).
1126        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        // Appending an empty segment is a no-op.
1138        append_tags(&mut target, "");
1139        assert_eq!("", target);
1140
1141        // The first non-empty append does not prepend a separator.
1142        append_tags(&mut target, "a:1");
1143        assert_eq!("a:1", target);
1144
1145        // Subsequent non-empty appends are comma-separated.
1146        append_tags(&mut target, "b:2");
1147        assert_eq!("a:1,b:2", target);
1148
1149        // An empty segment remains a no-op even once the target is non-empty.
1150        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        // An explicit, non-empty `datadog.container_tags` attribute is used verbatim.
1175        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        // `ignore_missing_fields` skips the extraction path entirely.
1190        assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, true));
1191        // Honoring fields but with no container attributes and no Fargate source still yields nothing.
1192        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        // The encoder emits the ETP `idxTracerPayloads` format, so decode through the string table.
1206        let tracer_versions = decode_etp_tracer_versions(&buf);
1207        let tracer_version = tracer_versions.first().expect("a tracer payload should be encoded");
1208        // The tracer version is prefixed with `otlp-` to mark the OTLP ingestion path.
1209        assert_eq!("otlp-1.2.3", tracer_version);
1210
1211        // The OTLP sampling rate is written to each chunk formatted to two decimal places.
1212        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}