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_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
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: 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        // 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        let (events_tx, events_rx) = mpsc::channel(8);
278        let (payloads_tx, mut payloads_rx) = mpsc::channel(8);
279
280        // Run our request builder task on the worker pool.
281        //
282        // The request builder task ignores the shutdown signal on purpose: it drains its incoming event buffer channel
283        // until the channel closes, which is what guarantees every buffered metric is encoded and dispatched.
284        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; // 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        // Draining `payloads_rx` to completion already implies the request builder finished: it owns the only sender,
325        // so the channel only closes once that child's future has run to completion (or been dropped).
326        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                    // Encode the trace. If we get it back, that means the current request is full, and we need to
349                    // flush it before we can try to encode the trace again.
350                    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                                // If the error is recoverable, we'll hold on to the trace to retry it later.
377                                continue;
378                            } else {
379                                return Err(GenericError::from(e).context("Failed to flush request."));
380                            }
381                        }
382                    }
383
384                    // Now try to encode the trace again.
385                    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 we're not already pending a flush, we'll start the countdown.
394                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                // Once we've encoded and written all traces, we flush the request builders to generate a request with
405                // anything left over. Again, we'll enqueue those requests to be sent immediately.
406                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            // Event buffers channel has been closed, and we have no pending flushing, so we're all done.
429            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        // Resolve computed metadata strings (may produce strings not directly present on trace fields).
489        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        // Build 128-bit big-endian trace ID bytes for the chunk.
513        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        // Reset the string table; strings are interned on the fly during encoding below.
518        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            // Tracer payload metadata refs (skip default/empty values).
531            // Strings are interned on the fly; the string table is written at the end of this
532            // closure so it is complete by the time tp.strings() is called.
533            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            // Container tags go in the payload-level attributes map.
557            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            // Single TraceChunk containing all spans.
566            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                // Write 128-bit trace ID.
574                chunk.trace_id(&trace_id_bytes)?;
575
576                // Sampling mechanism (only write when non-zero).
577                if trace.sampling_mechanism != 0 {
578                    chunk.sampling_mechanism(trace.sampling_mechanism)?;
579                }
580
581                // Spans.
582                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                        // Unified attribute map (replaces separate meta/metrics/meta_struct).
602                        {
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                        // Span links.
628                        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                        // Span events.
651                        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                // Chunk attributes: decision maker, ETS tag, OTLP sampling rate.
673                {
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            // Write the string table after all refs so the table is complete.
716            // Protobuf allows fields in any order; decoders that do a full parse before
717            // resolving refs handle strings-after-refs correctly.
718            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
773/// Encodes an [`AttributeValue`] into an ETP-format `AnyValue` builder, interning any
774/// string values into `st` on the fly.
775fn 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    // ---------------------------------------------------------------------------
892    // Decode helpers for assertions on the encoded ETP `AgentPayload`.
893    //
894    // The encoder emits the wire format via `piecemeal` builders; here we decode it back
895    // with the independently code-generated `rust-protobuf` types (`idx::TracerPayload`
896    // and friends). A full-message parse resolves protobuf field ordering for us, so the
897    // fact that the encoder writes string refs before the string table is a non-issue.
898    //
899    // The only ETP-specific work left is resolving `u32` string-table refs and unwrapping
900    // the `AnyValue` oneof, which the two helpers below handle.
901    // ---------------------------------------------------------------------------
902
903    /// Resolves a string-table ref against a tracer payload's string table.
904    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    /// Resolves the string-valued entries of an ETP attribute map into a readable
909    /// `key -> value` map, keeping only entries whose `AnyValue` is a string ref and
910    /// whose key resolves to a non-empty string. Keys and values are resolved through
911    /// the payload string table.
912    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    /// Collects the resolved string-valued attributes of every chunk across all ETP
926    /// tracer payloads in an encoded `AgentPayload`, one map per chunk.
927    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    /// Resolves the `tracerVersionRef` of every ETP tracer payload in an encoded `AgentPayload`.
942    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    // The APM sampler defaults: 10 target and error traces per second.
952    const DEFAULT_TARGET_TPS: f64 = 10.0;
953    const DEFAULT_ERRORS_PER_SECOND: f64 = 10.0;
954    // Keep-everything sampling percentage (the OTLP probabilistic sampler default).
955    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,    // span_id
977            0,    // parent_id
978            0,    // start
979            1000, // duration
980            0,    // error
981        );
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,    // span_id
994            0,    // parent_id
995            0,    // start
996            1000, // duration
997            1,    // error
998        );
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(); // no error
1040        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        // `sampling_percentage` is a 0..100 percentage; only strictly in-range values map to a fractional rate, and
1065        // anything <= 0 or >= 100 collapses to 1.0 (sample everything).
1066        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        // A non-empty payload hostname always wins.
1101        assert_eq!(
1102            Some("payload-host"),
1103            resolve_hostname_from_payload("payload-host", Some(&host_source), Some("default"), false)
1104        );
1105        // An empty payload plus `ignore_missing_fields` short-circuits to an empty hostname.
1106        assert_eq!(
1107            Some(""),
1108            resolve_hostname_from_payload("", Some(&host_source), Some("default"), true)
1109        );
1110        // Honoring fields, a hostname-kind source supplies its identifier.
1111        assert_eq!(
1112            Some("resolved-host"),
1113            resolve_hostname_from_payload("", Some(&host_source), Some("default"), false)
1114        );
1115        // A non-hostname (Fargate) source resolves to an empty hostname.
1116        assert_eq!(
1117            Some(""),
1118            resolve_hostname_from_payload("", Some(&fargate_source), Some("default"), false)
1119        );
1120        // With no source, it falls back to the default hostname (which may itself be absent).
1121        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        // Appending an empty segment is a no-op.
1133        append_tags(&mut target, "");
1134        assert_eq!("", target);
1135
1136        // The first non-empty append does not prepend a separator.
1137        append_tags(&mut target, "a:1");
1138        assert_eq!("a:1", target);
1139
1140        // Subsequent non-empty appends are comma-separated.
1141        append_tags(&mut target, "b:2");
1142        assert_eq!("a:1,b:2", target);
1143
1144        // An empty segment remains a no-op even once the target is non-empty.
1145        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        // An explicit, non-empty `datadog.container_tags` attribute is used verbatim.
1170        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        // `ignore_missing_fields` skips the extraction path entirely.
1185        assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, true));
1186        // Honoring fields but with no container attributes and no Fargate source still yields nothing.
1187        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        // The encoder emits the ETP `idxTracerPayloads` format, so decode through the string table.
1201        let tracer_versions = decode_etp_tracer_versions(&buf);
1202        let tracer_version = tracer_versions.first().expect("a tracer payload should be encoded");
1203        // The tracer version is prefixed with `otlp-` to mark the OTLP ingestion path.
1204        assert_eq!("otlp-1.2.3", tracer_version);
1205
1206        // The OTLP sampling rate is written to each chunk formatted to two decimal places.
1207        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}