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_core::{
11    accounting::{MemoryBounds, MemoryBoundsBuilder},
12    components::{encoders::*, BuildContext},
13    data_model::{
14        event::{
15            trace::{AttributeValue, Trace},
16            EventType,
17        },
18        payload::{HttpPayload, Payload, PayloadMetadata, PayloadType},
19        tags::TagSet,
20    },
21    observability::ComponentMetricsExt as _,
22    runtime,
23    topology::{EventsBuffer, PayloadsBuffer},
24};
25use saluki_env::host::providers::BoxedHostProvider;
26use saluki_env::{EnvironmentProvider, HostProvider};
27use saluki_error::generic_error;
28use saluki_error::{ErrorContext as _, GenericError};
29use saluki_io::compression::CompressionScheme;
30use saluki_metrics::MetricsBuilder;
31use stringtheory::MetaString;
32use tokio::pin;
33use tokio::{
34    select,
35    sync::mpsc::{self, Receiver, Sender},
36    time::sleep,
37};
38use tracing::{debug, error};
39
40use crate::common::datadog::{
41    io::RB_BUFFER_CHUNK_SIZE,
42    request_builder::{EndpointEncoder, RequestBuilder},
43    telemetry::ComponentTelemetry,
44    DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT, DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT, TAG_DECISION_MAKER,
45};
46use crate::common::otlp::util::{
47    attributes_to_source, extract_container_tags_from_attributes_map, Source as OtlpSource,
48    SourceKind as OtlpSourceKind, KEY_DATADOG_CONTAINER_TAGS,
49};
50
51const CONTAINER_TAGS_META_KEY: &str = "_dd.tags.container";
52const MAX_TRACES_PER_PAYLOAD: usize = 10000;
53static CONTENT_TYPE_PROTOBUF: HeaderValue = HeaderValue::from_static("application/x-protobuf");
54
55// Sampling metadata keys / values.
56const TAG_OTLP_SAMPLING_RATE: &str = "_dd.otlp_sr";
57const DEFAULT_CHUNK_PRIORITY: i32 = 1; // PRIORITY_AUTO_KEEP
58
59// ETS chunk-level attribute keys / values.
60const TAG_ETS_STANDALONE_ERROR_KEY: &str = "_dd.error_tracking_standalone.error";
61const TAG_ETS_STANDALONE_ERROR_VALUE: &str = "true";
62
63/// String interning table for the ETP (Efficient Trace Payload) tracer payload format.
64///
65/// Index 0 is always the empty string. All other strings are assigned indices in insertion order.
66#[derive(Debug)]
67struct StringTable {
68    indices: FastIndexSet<MetaString>,
69}
70
71impl Default for StringTable {
72    fn default() -> Self {
73        Self::new()
74    }
75}
76
77impl StringTable {
78    fn new() -> Self {
79        let mut t = Self {
80            indices: FastIndexSet::default(),
81        };
82        t.intern("");
83        t
84    }
85
86    fn clear(&mut self) {
87        self.indices.clear();
88        self.intern("");
89    }
90
91    /// Interns `s`, returning its index. If `s` was already interned, returns the existing index.
92    fn intern(&mut self, s: &str) -> u32 {
93        // 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.
94        if let Some(idx) = self.indices.get_index_of(s) {
95            idx as u32
96        } else {
97            let (idx, _) = self.indices.insert_full(MetaString::from(s));
98            idx as u32
99        }
100    }
101}
102
103/// Configuration for the Datadog Traces encoder.
104///
105/// This encoder converts trace events into Datadog's TracerPayload protobuf format and sends them
106/// to the Datadog traces intake endpoint (`/api/v0.2/traces`). It handles batching, compression,
107/// and enrichment with metadata such as hostname, environment, and container tags.
108pub struct DatadogTraceConfiguration {
109    /// Compression algorithm applied to outgoing payloads.
110    compressor_kind: String,
111
112    /// Effective zstd compression level, resolved by the configuration layer.
113    zstd_compressor_level: i32,
114
115    /// How long the encoder waits before flushing a partially filled payload.
116    ///
117    /// A zero duration is treated as "flush almost immediately" during `build`.
118    flush_timeout: Duration,
119
120    /// Global environment tag applied to emitted payloads.
121    env: String,
122
123    /// Target sampled traces per second, forwarded to the intake as `target_tps`.
124    target_traces_per_second: f64,
125
126    /// Target sampled error traces per second, forwarded to the intake as `error_tps`.
127    errors_per_second: f64,
128
129    /// Whether Error Tracking standalone mode is enabled.
130    error_tracking_standalone: bool,
131
132    /// Whether spans missing intake-required fields are ingested rather than rejected.
133    ignore_missing_datadog_fields: bool,
134
135    /// Percentage of OTLP traces the probabilistic sampler keeps.
136    sampling_percentage: f64,
137
138    /// Default hostname, resolved from the environment provider.
139    default_hostname: Option<String>,
140
141    /// ADP version string embedded in emitted payloads.
142    version: String,
143}
144
145impl DatadogTraceConfiguration {
146    /// Creates a new `DatadogTraceConfiguration` from the resolved traces and shared configuration.
147    ///
148    /// The OTLP trace settings live in their own domain, so they arrive as a separate slice rather
149    /// than through the traces domain.
150    pub fn from_configuration(
151        traces: &domains::traces::Domain, otlp_traces: &domains::otlp::Traces, shared: &SharedConfiguration,
152    ) -> Self {
153        let app_details = saluki_metadata::get_app_details();
154        let version = format!("agent-data-plane/{}", app_details.version().raw());
155
156        let compression = &shared.endpoints.compression;
157
158        // ADP defaults the global environment tag to `none` rather than the Core Agent's empty
159        // string, so normalize an empty resolved value back to `none`.
160        let env = if traces.env.is_empty() {
161            DEFAULT_TRACE_ENV.to_owned()
162        } else {
163            traces.env.clone()
164        };
165
166        Self {
167            compressor_kind: compression.compressor_kind.clone(),
168            zstd_compressor_level: compression.effective_zstd_level(),
169            flush_timeout: shared.metrics_encoding.flush_timeout,
170            env,
171            target_traces_per_second: traces.target_traces_per_second,
172            errors_per_second: traces.errors_per_second,
173            error_tracking_standalone: traces.error_tracking_standalone_enabled,
174            ignore_missing_datadog_fields: otlp_traces.ignore_missing_datadog_fields,
175            sampling_percentage: otlp_traces.probabilistic_sampler_sampling_percentage,
176            default_hostname: None,
177            version,
178        }
179    }
180
181    /// Sets the `default_hostname` using the environment provider
182    pub async fn with_environment_provider<E>(mut self, environment_provider: E) -> Result<Self, GenericError>
183    where
184        E: EnvironmentProvider<Host = BoxedHostProvider>,
185    {
186        let host_provider = environment_provider.host();
187        let hostname = host_provider.get_hostname().await?;
188        self.default_hostname = Some(hostname);
189        Ok(self)
190    }
191}
192
193#[async_trait]
194impl EncoderBuilder for DatadogTraceConfiguration {
195    fn input_event_type(&self) -> EventType {
196        EventType::Trace
197    }
198
199    fn output_payload_type(&self) -> PayloadType {
200        PayloadType::Http
201    }
202
203    async fn build(&self, context: BuildContext) -> Result<Box<dyn Encoder + Send>, GenericError> {
204        let metrics_builder = MetricsBuilder::from_component_context(context.component_context());
205        let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
206        let compression_scheme = CompressionScheme::new(&self.compressor_kind, self.zstd_compressor_level);
207
208        let default_hostname = self.default_hostname.clone().unwrap_or_default();
209        let default_hostname = MetaString::from(default_hostname);
210
211        // Create request builder for traces which is used to generate HTTP requests.
212
213        let mut trace_rb = RequestBuilder::new(
214            TraceEndpointEncoder::new(
215                default_hostname,
216                self.version.clone(),
217                self.env.clone(),
218                self.target_traces_per_second,
219                self.errors_per_second,
220                self.error_tracking_standalone,
221                self.ignore_missing_datadog_fields,
222                self.sampling_percentage,
223            ),
224            compression_scheme,
225            RB_BUFFER_CHUNK_SIZE,
226        )
227        .await?;
228        trace_rb.with_max_inputs_per_payload(MAX_TRACES_PER_PAYLOAD);
229
230        let flush_timeout = if self.flush_timeout.is_zero() {
231            // We always give ourselves a minimum flush timeout of 10ms to allow for some very minimal amount of
232            // batching, while still practically flushing things almost immediately.
233            Duration::from_millis(10)
234        } else {
235            self.flush_timeout
236        };
237
238        Ok(Box::new(DatadogTrace {
239            trace_rb,
240            telemetry,
241            flush_timeout,
242        }))
243    }
244}
245
246impl MemoryBounds for DatadogTraceConfiguration {
247    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
248        // TODO: How do we properly represent the requests we can generate that may be sitting around in-flight?
249        builder
250            .minimum()
251            .with_single_value::<DatadogTrace>("component struct")
252            .with_array::<EventsBuffer>("request builder events channel", 8)
253            .with_array::<PayloadsBuffer>("request builder payloads channel", 8);
254
255        builder
256            .firm()
257            .with_array::<Trace>("traces split re-encode buffer", MAX_TRACES_PER_PAYLOAD);
258    }
259}
260
261pub struct DatadogTrace {
262    trace_rb: RequestBuilder<TraceEndpointEncoder>,
263    telemetry: ComponentTelemetry,
264    flush_timeout: Duration,
265}
266
267// Encodes Trace events to TracerPayloads.
268#[async_trait]
269impl Encoder for DatadogTrace {
270    async fn run(mut self: Box<Self>, mut context: EncoderContext) -> Result<(), GenericError> {
271        let Self {
272            trace_rb,
273            telemetry,
274            flush_timeout,
275        } = *self;
276
277        let mut health = context.take_health_handle();
278
279        let (events_tx, events_rx) = mpsc::channel(8);
280        let (payloads_tx, mut payloads_rx) = mpsc::channel(8);
281
282        // Run our request builder task on the worker pool.
283        //
284        // The request builder task ignores the shutdown signal on purpose: it drains its incoming event buffer channel
285        // until the channel closes, which is what guarantees every buffered metric is encoded and dispatched.
286        let request_builder_fut = run_request_builder(trace_rb, telemetry, events_rx, payloads_tx, flush_timeout);
287        runtime::worker("request_builder", request_builder_fut)
288            .on_runtime(context.topology_context().global_thread_pool().clone())
289            .spawn();
290
291        health.mark_ready();
292        debug!("Datadog Trace encoder started.");
293
294        loop {
295            select! {
296                biased; // makes the branches of the select statement be evaluated in order.
297
298                _ = health.live() => continue,
299                maybe_payload = payloads_rx.recv() => match maybe_payload {
300                    Some(payload) => {
301                        // Dispatch an HTTP payload to the dispatcher.
302                        if let Err(e) = context.dispatcher().dispatch(payload).await {
303                            error!("Failed to dispatch payload: {}", e);
304                        }
305                    }
306                    None => break,
307                },
308                maybe_event_buffer = context.events().next() => match maybe_event_buffer {
309                    Some(event_buffer) => events_tx.send(event_buffer).await
310                        .error_context("Failed to send event buffer to request builder task.")?,
311                    None => break,
312                },
313            }
314        }
315
316        // Drop the events sender, which signals the request builder task to stop.
317        drop(events_tx);
318
319        // Continue draining the payloads receiver until it is closed.
320        while let Some(payload) = payloads_rx.recv().await {
321            if let Err(e) = context.dispatcher().dispatch(payload).await {
322                error!("Failed to dispatch payload: {}", e);
323            }
324        }
325
326        // Draining `payloads_rx` to completion already implies the request builder finished: it owns the only sender,
327        // so the channel only closes once that child's future has run to completion (or been dropped).
328        debug!("Datadog Trace encoder stopped.");
329
330        Ok(())
331    }
332}
333
334async fn run_request_builder(
335    mut trace_request_builder: RequestBuilder<TraceEndpointEncoder>, telemetry: ComponentTelemetry,
336    mut events_rx: Receiver<EventsBuffer>, payloads_tx: Sender<PayloadsBuffer>, flush_timeout: std::time::Duration,
337) -> Result<(), GenericError> {
338    let mut pending_flush = false;
339    let pending_flush_timeout = sleep(flush_timeout);
340    pin!(pending_flush_timeout);
341
342    loop {
343        select! {
344            Some(event_buffer) = events_rx.recv() => {
345                for event in event_buffer {
346                    let trace = match event.try_into_trace() {
347                        Some(trace) => trace,
348                        None => continue,
349                    };
350                    // Encode the trace. If we get it back, that means the current request is full, and we need to
351                    // flush it before we can try to encode the trace again.
352                    let trace_to_retry = match trace_request_builder.encode(trace).await {
353                        Ok(None) => continue,
354                        Ok(Some(trace)) => trace,
355                        Err(e) => {
356                            error!(error = %e, "Failed to encode trace.");
357                            telemetry.events_dropped_encoder().increment(1);
358                            continue;
359                        }
360                    };
361
362                    let maybe_requests = trace_request_builder.flush().await;
363                    if maybe_requests.is_empty() {
364                        panic!("builder told us to flush, but gave us nothing");
365                    }
366
367                    for maybe_request in maybe_requests {
368                        match maybe_request {
369                            Ok((events, _data_points, request)) => {
370                                let payload_meta = PayloadMetadata::from_event_count(events);
371                                let http_payload = HttpPayload::new(payload_meta, request);
372                                let payload = Payload::Http(http_payload);
373
374                                payloads_tx.send(payload).await
375                                    .map_err(|_| generic_error!("Failed to send payload to encoder."))?;
376                            },
377                            Err(e) => if e.is_recoverable() {
378                                // If the error is recoverable, we'll hold on to the trace to retry it later.
379                                continue;
380                            } else {
381                                return Err(GenericError::from(e).context("Failed to flush request."));
382                            }
383                        }
384                    }
385
386                    // Now try to encode the trace again.
387                    if let Err(e) = trace_request_builder.encode(trace_to_retry).await {
388                        error!(error = %e, "Failed to encode trace.");
389                        telemetry.events_dropped_encoder().increment(1);
390                    }
391                }
392
393                debug!("Processed event buffer.");
394
395                // If we're not already pending a flush, we'll start the countdown.
396                if !pending_flush {
397                    pending_flush_timeout.as_mut().reset(tokio::time::Instant::now() + flush_timeout);
398                    pending_flush = true;
399                }
400            },
401            _ = &mut pending_flush_timeout, if pending_flush => {
402                debug!("Flushing pending request(s).");
403
404                pending_flush = false;
405
406                // Once we've encoded and written all traces, we flush the request builders to generate a request with
407                // anything left over. Again, we'll enqueue those requests to be sent immediately.
408                let maybe_trace_requests = trace_request_builder.flush().await;
409                for maybe_request in maybe_trace_requests {
410                    match maybe_request {
411                        Ok((events, _data_points, request)) => {
412                            let payload_meta = PayloadMetadata::from_event_count(events);
413                            let http_payload = HttpPayload::new(payload_meta, request);
414                            let payload = Payload::Http(http_payload);
415
416                            payloads_tx.send(payload).await
417                                .map_err(|_| generic_error!("Failed to send payload to encoder."))?;
418                        },
419                        Err(e) => if e.is_recoverable() {
420                            continue;
421                        } else {
422                            return Err(GenericError::from(e).context("Failed to flush request."));
423                        }
424                    }
425                }
426
427                debug!("All flushed requests sent to I/O task. Waiting for next event buffer...");
428            },
429
430            // Event buffers channel has been closed, and we have no pending flushing, so we're all done.
431            else => break,
432        }
433    }
434
435    Ok(())
436}
437
438#[derive(Debug)]
439struct TraceEndpointEncoder {
440    scratch: ScratchWriter<Vec<u8>>,
441    default_hostname: MetaString,
442    agent_hostname: String,
443    version: String,
444    env: String,
445    target_traces_per_second: f64,
446    errors_per_second: f64,
447    ignore_missing_datadog_fields: bool,
448    sampling_percentage: f64,
449    string_builder: StringBuilder,
450    string_table: StringTable,
451    error_tracking_standalone: bool,
452    extra_headers: Vec<(HeaderName, HeaderValue)>,
453}
454
455impl TraceEndpointEncoder {
456    fn new(
457        default_hostname: MetaString, version: String, env: String, target_traces_per_second: f64,
458        errors_per_second: f64, error_tracking_standalone: bool, ignore_missing_datadog_fields: bool,
459        sampling_percentage: f64,
460    ) -> Self {
461        let extra_headers = if error_tracking_standalone {
462            vec![(
463                HeaderName::from_static("x-datadog-error-tracking-standalone"),
464                HeaderValue::from_static("true"),
465            )]
466        } else {
467            Vec::new()
468        };
469        Self {
470            scratch: ScratchWriter::new(Vec::with_capacity(8192)),
471            agent_hostname: default_hostname.as_ref().to_string(),
472            default_hostname,
473            version,
474            env,
475            target_traces_per_second,
476            errors_per_second,
477            ignore_missing_datadog_fields,
478            sampling_percentage,
479            string_builder: StringBuilder::new(),
480            string_table: StringTable::new(),
481            error_tracking_standalone,
482            extra_headers,
483        }
484    }
485
486    fn encode_tracer_payload(&mut self, trace: &Trace, output_buffer: &mut Vec<u8>) -> std::io::Result<()> {
487        let sampling_rate = self.sampling_rate();
488        let source = attributes_to_source(&trace.attributes);
489
490        // Resolve computed metadata strings (may produce strings not directly present on trace fields).
491        let tracer_version = format!("otlp-{}", &trace.payload.tracer_version);
492        let container_tags =
493            resolve_container_tags_from_attrs(&trace.attributes, source.as_ref(), self.ignore_missing_datadog_fields);
494        let env_str: Option<&str> = if !trace.payload.env.is_empty() {
495            Some(&trace.payload.env)
496        } else if self.ignore_missing_datadog_fields {
497            Some("")
498        } else {
499            None
500        };
501        let hostname_str: Option<&str> = resolve_hostname_from_payload(
502            &trace.payload.hostname,
503            source.as_ref(),
504            Some(self.default_hostname.as_ref()),
505            self.ignore_missing_datadog_fields,
506        );
507        let decision_maker = trace.decision_maker.as_deref();
508        let priority = trace.priority.unwrap_or(DEFAULT_CHUNK_PRIORITY);
509        let dropped_trace = trace.dropped_trace;
510        let otlp_sr = trace.otlp_sampling_rate.unwrap_or(sampling_rate);
511        self.string_builder.clear();
512        write!(&mut self.string_builder, "{:.2}", otlp_sr).expect("should never fail to format sampling rate");
513
514        // Build 128-bit big-endian trace ID bytes for the chunk.
515        let mut trace_id_bytes = [0u8; 16];
516        trace_id_bytes[..8].copy_from_slice(&trace.trace_id_high.to_be_bytes());
517        trace_id_bytes[8..].copy_from_slice(&trace.trace_id_low.to_be_bytes());
518
519        // Reset the string table; strings are interned on the fly during encoding below.
520        self.string_table.clear();
521
522        let mut ap_builder = AgentPayloadBuilder::new(&mut self.scratch);
523
524        ap_builder
525            .host_name(&self.agent_hostname)?
526            .env(&self.env)?
527            .agent_version(&self.version)?
528            .target_tps(self.target_traces_per_second)?
529            .error_tps(self.errors_per_second)?;
530
531        ap_builder.add_idx_tracer_payloads(|tp| {
532            // Tracer payload metadata refs (skip default/empty values).
533            // Strings are interned on the fly; the string table is written at the end of this
534            // closure so it is complete by the time tp.strings() is called.
535            if !trace.payload.container_id.is_empty() {
536                tp.container_id_ref(self.string_table.intern(&trace.payload.container_id))?;
537            }
538            if !trace.payload.language_name.is_empty() {
539                tp.language_name_ref(self.string_table.intern(&trace.payload.language_name))?;
540            }
541            if !trace.payload.language_version.is_empty() {
542                tp.language_version_ref(self.string_table.intern(&trace.payload.language_version))?;
543            }
544            tp.tracer_version_ref(self.string_table.intern(tracer_version.as_str()))?;
545            if !trace.payload.runtime_id.is_empty() {
546                tp.runtime_id_ref(self.string_table.intern(&trace.payload.runtime_id))?;
547            }
548            if let Some(e) = env_str {
549                tp.env_ref(self.string_table.intern(e))?;
550            }
551            if let Some(h) = hostname_str {
552                tp.hostname_ref(self.string_table.intern(h))?;
553            }
554            if !trace.payload.app_version.is_empty() {
555                tp.app_version_ref(self.string_table.intern(&trace.payload.app_version))?;
556            }
557
558            // Container tags go in the payload-level attributes map.
559            if let Some(ct) = &container_tags {
560                let k_ref = self.string_table.intern(CONTAINER_TAGS_META_KEY);
561                let v_ref = self.string_table.intern(ct);
562                tp.attributes().write_entry(k_ref, |av: &mut _| {
563                    av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
564                })?;
565            }
566
567            // Single TraceChunk containing all spans.
568            tp.add_chunks(|chunk| {
569                chunk.priority(priority)?;
570
571                if !trace.origin.is_empty() {
572                    chunk.origin_ref(self.string_table.intern(&trace.origin))?;
573                }
574
575                // Write 128-bit trace ID.
576                chunk.trace_id(&trace_id_bytes)?;
577
578                // Sampling mechanism (only write when non-zero).
579                if trace.sampling_mechanism != 0 {
580                    chunk.sampling_mechanism(trace.sampling_mechanism)?;
581                }
582
583                // Spans.
584                for span in trace.spans() {
585                    let service_ref = self.string_table.intern(span.service());
586                    let name_ref = self.string_table.intern(span.name());
587                    let resource_ref = self.string_table.intern(span.resource());
588                    let type_ref = self.string_table.intern(span.span_type());
589                    let env_ref = (!span.env.is_empty()).then(|| self.string_table.intern(&span.env));
590                    let version_ref = (!span.version.is_empty()).then(|| self.string_table.intern(&span.version));
591                    let component_ref = (!span.component.is_empty()).then(|| self.string_table.intern(&span.component));
592
593                    chunk.add_spans(|s| {
594                        s.service_ref(service_ref)?
595                            .name_ref(name_ref)?
596                            .resource_ref(resource_ref)?
597                            .span_id(span.span_id())?
598                            .parent_id(span.parent_id())?
599                            .start(span.start())?
600                            .duration(span.duration())?
601                            .error(span.error() != 0)?;
602
603                        // Unified attribute map (replaces separate meta/metrics/meta_struct).
604                        {
605                            let mut attrs = s.attributes();
606                            for (k, v) in &span.attributes {
607                                let k_ref = self.string_table.intern(k);
608                                attrs.write_entry(k_ref, |av: &mut _| {
609                                    encode_etp_attribute_value(av, v, &mut self.string_table)
610                                })?;
611                            }
612                        }
613
614                        s.type_ref(type_ref)?;
615
616                        if let Some(er) = env_ref {
617                            s.env_ref(er)?;
618                        }
619                        if let Some(vr) = version_ref {
620                            s.version_ref(vr)?;
621                        }
622                        if let Some(cr) = component_ref {
623                            s.component_ref(cr)?;
624                        }
625                        if span.kind != 0 {
626                            s.kind(SpanKind::from(span.kind as i32))?;
627                        }
628
629                        // Span links.
630                        for link in span.span_links() {
631                            let mut link_trace_id_bytes = [0u8; 16];
632                            link_trace_id_bytes[..8].copy_from_slice(&link.trace_id_high().to_be_bytes());
633                            link_trace_id_bytes[8..].copy_from_slice(&link.trace_id().to_be_bytes());
634                            let tracestate_ref = self.string_table.intern(link.tracestate());
635
636                            s.add_links(|sl| {
637                                sl.trace_id(&link_trace_id_bytes)?.span_id(link.span_id())?;
638                                {
639                                    let mut lattrs = sl.attributes();
640                                    for (k, v) in link.attributes() {
641                                        let k_ref = self.string_table.intern(k);
642                                        lattrs.write_entry(k_ref, |av: &mut _| {
643                                            encode_etp_attribute_value(av, v, &mut self.string_table)
644                                        })?;
645                                    }
646                                }
647                                sl.tracestate_ref(tracestate_ref)?.flags(link.flags())?;
648                                Ok(())
649                            })?;
650                        }
651
652                        // Span events.
653                        for event in span.span_events() {
654                            let name_ref = self.string_table.intern(event.name());
655                            s.add_events(|se| {
656                                se.time(event.time_unix_nano())?.name_ref(name_ref)?;
657                                {
658                                    let mut eattrs = se.attributes();
659                                    for (k, v) in event.attributes() {
660                                        let k_ref = self.string_table.intern(k);
661                                        eattrs.write_entry(k_ref, |av: &mut _| {
662                                            encode_etp_attribute_value(av, v, &mut self.string_table)
663                                        })?;
664                                    }
665                                }
666                                Ok(())
667                            })?;
668                        }
669
670                        Ok(())
671                    })?;
672                }
673
674                // Chunk attributes: decision maker, ETS tag, OTLP sampling rate.
675                {
676                    let mut cattrs = chunk.attributes();
677                    if let Some(dm) = decision_maker {
678                        let k_ref = self.string_table.intern(TAG_DECISION_MAKER);
679                        let v_ref = self.string_table.intern(dm);
680                        cattrs.write_entry(k_ref, |av: &mut _| {
681                            av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
682                        })?;
683                    }
684                    if self.error_tracking_standalone {
685                        let trace_has_error = trace.spans().iter().any(|span| {
686                            span.error() != 0
687                                || span
688                                    .attributes
689                                    .get("_dd.span_events.has_exception")
690                                    .and_then(AttributeValue::as_string)
691                                    .is_some_and(|v| v == "true")
692                        });
693                        if trace_has_error {
694                            let k_ref = self.string_table.intern(TAG_ETS_STANDALONE_ERROR_KEY);
695                            let v_ref = self.string_table.intern(TAG_ETS_STANDALONE_ERROR_VALUE);
696                            cattrs.write_entry(k_ref, |av: &mut _| {
697                                av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
698                            })?;
699                        }
700                    }
701                    {
702                        let k_ref = self.string_table.intern(TAG_OTLP_SAMPLING_RATE);
703                        let v_ref = self.string_table.intern(self.string_builder.as_str());
704                        cattrs.write_entry(k_ref, |av: &mut _| {
705                            av.value(|vo| vo.string_value_ref(v_ref)).map(|_| ())
706                        })?;
707                    }
708                }
709
710                if dropped_trace {
711                    chunk.dropped_trace(true)?;
712                }
713
714                Ok(())
715            })?;
716
717            // Write the string table after all refs so the table is complete.
718            // Protobuf allows fields in any order; decoders that do a full parse before
719            // resolving refs handle strings-after-refs correctly.
720            tp.strings(|sb| sb.add_many_mapped(&self.string_table.indices, |s| &**s))?;
721
722            Ok(())
723        })?;
724
725        ap_builder.finish(output_buffer)?;
726
727        Ok(())
728    }
729
730    fn sampling_rate(&self) -> f64 {
731        let rate = self.sampling_percentage / 100.0;
732        if rate <= 0.0 || rate >= 1.0 {
733            return 1.0;
734        }
735        rate
736    }
737}
738
739impl EndpointEncoder for TraceEndpointEncoder {
740    type Input = Trace;
741    type EncodeError = std::io::Error;
742    fn encoder_name() -> &'static str {
743        "traces"
744    }
745
746    fn compressed_size_limit(&self) -> usize {
747        DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT
748    }
749
750    fn uncompressed_size_limit(&self) -> usize {
751        DEFAULT_INTAKE_UNCOMPRESSED_SIZE_LIMIT
752    }
753
754    fn encode(&mut self, trace: &Self::Input, buffer: &mut Vec<u8>) -> Result<(), Self::EncodeError> {
755        self.encode_tracer_payload(trace, buffer)
756    }
757
758    fn endpoint_uri(&self) -> Uri {
759        PathAndQuery::from_static("/api/v0.2/traces").into()
760    }
761
762    fn endpoint_method(&self) -> Method {
763        Method::POST
764    }
765
766    fn content_type(&self) -> HeaderValue {
767        CONTENT_TYPE_PROTOBUF.clone()
768    }
769
770    fn additional_headers(&self) -> &[(HeaderName, HeaderValue)] {
771        &self.extra_headers
772    }
773}
774
775/// Encodes an [`AttributeValue`] into an ETP-format `AnyValue` builder, interning any
776/// string values into `st` on the fly.
777fn encode_etp_attribute_value<S: ScratchBuffer>(
778    builder: &mut datadog_protos::traces::builders::idx::AnyValueBuilder<'_, S>, value: &AttributeValue,
779    st: &mut StringTable,
780) -> std::io::Result<()> {
781    builder
782        .value(|vo| match value {
783            AttributeValue::String(s) => vo.string_value_ref(st.intern(s)),
784            AttributeValue::Bool(b) => vo.bool_value(*b),
785            AttributeValue::Int(i) => vo.int_value(*i),
786            AttributeValue::Float(f) => vo.double_value(*f),
787            AttributeValue::Bytes(b) => vo.bytes_value(b),
788            AttributeValue::Array(values) => vo.array_value(|arr| {
789                for v in values {
790                    arr.add_values(|av| encode_etp_attribute_value(av, v, st))?;
791                }
792                Ok(())
793            }),
794            AttributeValue::KeyValueList(kvs) => vo.key_value_list(|kvl| {
795                for (k, v) in kvs {
796                    kvl.add_key_values(|kv| {
797                        kv.key(st.intern(k))?
798                            .value(|av| encode_etp_attribute_value(av, v, st))?;
799                        Ok(())
800                    })?;
801                }
802                Ok(())
803            }),
804        })
805        .map(|_| ())
806}
807
808fn resolve_hostname_from_payload<'a>(
809    payload_hostname: &'a str, source: Option<&'a OtlpSource>, default_hostname: Option<&'a str>,
810    ignore_missing_fields: bool,
811) -> Option<&'a str> {
812    if !payload_hostname.is_empty() {
813        return Some(payload_hostname);
814    }
815    if ignore_missing_fields {
816        return Some("");
817    }
818    match source {
819        Some(src) => match src.kind {
820            OtlpSourceKind::HostnameKind => Some(src.identifier.as_str()),
821            _ => Some(""),
822        },
823        None => default_hostname,
824    }
825}
826
827fn resolve_container_tags_from_attrs(
828    attributes: &FastHashMap<MetaString, AttributeValue>, source: Option<&OtlpSource>, ignore_missing_fields: bool,
829) -> Option<MetaString> {
830    if let Some(AttributeValue::String(tags)) = attributes.get(KEY_DATADOG_CONTAINER_TAGS) {
831        if !tags.is_empty() {
832            return Some(tags.clone());
833        }
834    }
835
836    if ignore_missing_fields {
837        return None;
838    }
839    let mut container_tags = TagSet::default();
840    extract_container_tags_from_attributes_map(attributes, &mut container_tags);
841    let is_fargate_source = source.is_some_and(|src| src.kind == OtlpSourceKind::AwsEcsFargateKind);
842    if container_tags.is_empty() && !is_fargate_source {
843        return None;
844    }
845
846    let mut flattened = flatten_container_tag(container_tags);
847    if is_fargate_source {
848        if let Some(src) = source {
849            append_tags(&mut flattened, &src.tag());
850        }
851    }
852
853    if flattened.is_empty() {
854        None
855    } else {
856        Some(MetaString::from(flattened))
857    }
858}
859
860fn flatten_container_tag(tags: TagSet) -> String {
861    let mut flattened = String::new();
862    for tag in tags {
863        if !flattened.is_empty() {
864            flattened.push(',');
865        }
866        flattened.push_str(tag.as_str());
867    }
868    flattened
869}
870
871fn append_tags(target: &mut String, tags: &str) {
872    if tags.is_empty() {
873        return;
874    }
875    if !target.is_empty() {
876        target.push(',');
877    }
878    target.push_str(tags);
879}
880
881#[cfg(test)]
882mod tests {
883    use std::collections::{BTreeSet, HashMap};
884
885    use datadog_protos::traces::{idx, AgentPayload};
886    use protobuf::Message as _;
887    use saluki_core::data_model::{
888        event::trace::{Span as DdSpan, Trace},
889        tags::Tag,
890    };
891    use stringtheory::MetaString;
892
893    use super::*;
894
895    // ---------------------------------------------------------------------------
896    // Decode helpers for assertions on the encoded ETP `AgentPayload`.
897    //
898    // The encoder emits the wire format via `piecemeal` builders; here we decode it back
899    // with the independently code-generated `rust-protobuf` types (`idx::TracerPayload`
900    // and friends). A full-message parse resolves protobuf field ordering for us, so the
901    // fact that the encoder writes string refs before the string table is a non-issue.
902    //
903    // The only ETP-specific work left is resolving `u32` string-table refs and unwrapping
904    // the `AnyValue` oneof, which the two helpers below handle.
905    // ---------------------------------------------------------------------------
906
907    /// Resolves a string-table ref against a tracer payload's string table.
908    fn resolve_ref(strings: &[String], string_ref: u32) -> &str {
909        strings.get(string_ref as usize).map(String::as_str).unwrap_or_default()
910    }
911
912    /// Resolves the string-valued entries of an ETP attribute map into a readable
913    /// `key -> value` map, keeping only entries whose `AnyValue` is a string ref and
914    /// whose key resolves to a non-empty string. Keys and values are resolved through
915    /// the payload string table.
916    fn string_attrs(attrs: &HashMap<u32, idx::AnyValue>, strings: &[String]) -> HashMap<String, String> {
917        let mut out = HashMap::new();
918        for (k_ref, value) in attrs {
919            if let Some(idx::any_value::Value::StringValueRef(v_ref)) = &value.value {
920                let key = resolve_ref(strings, *k_ref);
921                if !key.is_empty() {
922                    out.insert(key.to_string(), resolve_ref(strings, *v_ref).to_string());
923                }
924            }
925        }
926        out
927    }
928
929    /// Collects the resolved string-valued attributes of every chunk across all ETP
930    /// tracer payloads in an encoded `AgentPayload`, one map per chunk.
931    fn decode_etp_chunk_attributes(buf: &[u8]) -> Vec<HashMap<String, String>> {
932        let payload = AgentPayload::parse_from_bytes(buf).expect("AgentPayload should decode");
933        payload
934            .idxTracerPayloads
935            .iter()
936            .flat_map(|tp| {
937                let strings = tp.strings();
938                tp.chunks()
939                    .iter()
940                    .map(move |chunk| string_attrs(chunk.attributes(), strings))
941            })
942            .collect()
943    }
944
945    /// Resolves the `tracerVersionRef` of every ETP tracer payload in an encoded `AgentPayload`.
946    fn decode_etp_tracer_versions(buf: &[u8]) -> Vec<String> {
947        let payload = AgentPayload::parse_from_bytes(buf).expect("AgentPayload should decode");
948        payload
949            .idxTracerPayloads
950            .iter()
951            .map(|tp| resolve_ref(tp.strings(), tp.tracerVersionRef()).to_string())
952            .collect()
953    }
954
955    // The APM sampler defaults: 10 target and error traces per second.
956    const DEFAULT_TARGET_TPS: f64 = 10.0;
957    const DEFAULT_ERRORS_PER_SECOND: f64 = 10.0;
958    // Keep-everything sampling percentage (the OTLP probabilistic sampler default).
959    const DEFAULT_SAMPLING_PERCENTAGE: f64 = 100.0;
960
961    fn make_encoder(ets_enabled: bool) -> TraceEndpointEncoder {
962        TraceEndpointEncoder::new(
963            MetaString::from("test-host"),
964            "0.0.0".to_string(),
965            "none".to_string(),
966            DEFAULT_TARGET_TPS,
967            DEFAULT_ERRORS_PER_SECOND,
968            ets_enabled,
969            false,
970            DEFAULT_SAMPLING_PERCENTAGE,
971        )
972    }
973
974    fn make_trace() -> Trace {
975        let span = DdSpan::new(
976            MetaString::from("svc"),
977            MetaString::from("op"),
978            MetaString::from("res"),
979            MetaString::from("web"),
980            1,    // span_id
981            0,    // parent_id
982            0,    // start
983            1000, // duration
984            0,    // error
985        );
986        let mut trace = Trace::new(vec![span]);
987        trace.priority = Some(1);
988        trace
989    }
990
991    fn make_error_trace() -> Trace {
992        let span = DdSpan::new(
993            MetaString::from("svc"),
994            MetaString::from("op"),
995            MetaString::from("res"),
996            MetaString::from("web"),
997            1,    // span_id
998            0,    // parent_id
999            0,    // start
1000            1000, // duration
1001            1,    // error
1002        );
1003        let mut trace = Trace::new(vec![span]);
1004        trace.priority = Some(1);
1005        trace
1006    }
1007
1008    #[test]
1009    fn ets_header_present_when_enabled() {
1010        let encoder = make_encoder(true);
1011        let headers = encoder.additional_headers();
1012        assert_eq!(headers.len(), 1);
1013        assert_eq!(headers[0].0.as_str(), "x-datadog-error-tracking-standalone");
1014        assert_eq!(headers[0].1, "true");
1015    }
1016
1017    #[test]
1018    fn ets_header_absent_when_disabled() {
1019        let encoder = make_encoder(false);
1020        assert!(encoder.additional_headers().is_empty());
1021    }
1022
1023    #[test]
1024    fn ets_chunk_tag_present_for_error_trace() {
1025        let mut encoder = make_encoder(true);
1026        let trace = make_error_trace();
1027        let mut buf = Vec::new();
1028        encoder.encode(&trace, &mut buf).expect("encode should succeed");
1029        let chunk_attrs = decode_etp_chunk_attributes(&buf);
1030        let tag_value = chunk_attrs
1031            .iter()
1032            .find_map(|attrs| attrs.get("_dd.error_tracking_standalone.error").map(|v| v.as_str()));
1033        assert_eq!(
1034            tag_value,
1035            Some("true"),
1036            "ETS chunk tag should be present for error traces when ETS is enabled"
1037        );
1038    }
1039
1040    #[test]
1041    fn ets_chunk_tag_absent_for_non_error_trace() {
1042        let mut encoder = make_encoder(true);
1043        let trace = make_trace(); // no error
1044        let mut buf = Vec::new();
1045        encoder.encode(&trace, &mut buf).expect("encode should succeed");
1046        let chunk_attrs = decode_etp_chunk_attributes(&buf);
1047        let has_tag = chunk_attrs
1048            .iter()
1049            .any(|attrs| attrs.contains_key("_dd.error_tracking_standalone.error"));
1050        assert!(!has_tag, "ETS chunk tag should be absent for non-error traces");
1051    }
1052
1053    #[test]
1054    fn ets_chunk_tag_absent_when_disabled() {
1055        let mut encoder = make_encoder(false);
1056        let trace = make_trace();
1057        let mut buf = Vec::new();
1058        encoder.encode(&trace, &mut buf).expect("encode should succeed");
1059        let chunk_attrs = decode_etp_chunk_attributes(&buf);
1060        let has_tag = chunk_attrs
1061            .iter()
1062            .any(|attrs| attrs.contains_key("_dd.error_tracking_standalone.error"));
1063        assert!(!has_tag, "ETS chunk tag should be absent when ETS is disabled");
1064    }
1065
1066    #[test]
1067    fn sampling_rate_clamps_percentage_to_unit_interval() {
1068        // `sampling_percentage` is a 0..100 percentage; only strictly in-range values map to a fractional rate, and
1069        // anything <= 0 or >= 100 collapses to 1.0 (sample everything).
1070        let cases = [
1071            (25.0, 0.25),
1072            (50.0, 0.5),
1073            (0.0, 1.0),
1074            (-10.0, 1.0),
1075            (100.0, 1.0),
1076            (150.0, 1.0),
1077        ];
1078        for (percentage, expected) in cases {
1079            let encoder = TraceEndpointEncoder::new(
1080                MetaString::from("test-host"),
1081                "0.0.0".to_string(),
1082                "none".to_string(),
1083                DEFAULT_TARGET_TPS,
1084                DEFAULT_ERRORS_PER_SECOND,
1085                false,
1086                false,
1087                percentage,
1088            );
1089            assert_eq!(expected, encoder.sampling_rate(), "sampling_rate for {percentage}%");
1090        }
1091    }
1092
1093    #[test]
1094    fn resolve_hostname_from_payload_prefers_payload_then_source_then_default() {
1095        let host_source = OtlpSource {
1096            kind: OtlpSourceKind::HostnameKind,
1097            identifier: "resolved-host".to_string(),
1098        };
1099        let fargate_source = OtlpSource {
1100            kind: OtlpSourceKind::AwsEcsFargateKind,
1101            identifier: "task-arn".to_string(),
1102        };
1103
1104        // A non-empty payload hostname always wins.
1105        assert_eq!(
1106            Some("payload-host"),
1107            resolve_hostname_from_payload("payload-host", Some(&host_source), Some("default"), false)
1108        );
1109        // An empty payload plus `ignore_missing_fields` short-circuits to an empty hostname.
1110        assert_eq!(
1111            Some(""),
1112            resolve_hostname_from_payload("", Some(&host_source), Some("default"), true)
1113        );
1114        // Honoring fields, a hostname-kind source supplies its identifier.
1115        assert_eq!(
1116            Some("resolved-host"),
1117            resolve_hostname_from_payload("", Some(&host_source), Some("default"), false)
1118        );
1119        // A non-hostname (Fargate) source resolves to an empty hostname.
1120        assert_eq!(
1121            Some(""),
1122            resolve_hostname_from_payload("", Some(&fargate_source), Some("default"), false)
1123        );
1124        // With no source, it falls back to the default hostname (which may itself be absent).
1125        assert_eq!(
1126            Some("default"),
1127            resolve_hostname_from_payload("", None, Some("default"), false)
1128        );
1129        assert_eq!(None, resolve_hostname_from_payload("", None, None, false));
1130    }
1131
1132    #[test]
1133    fn append_tags_joins_non_empty_segments_with_commas() {
1134        let mut target = String::new();
1135
1136        // Appending an empty segment is a no-op.
1137        append_tags(&mut target, "");
1138        assert_eq!("", target);
1139
1140        // The first non-empty append does not prepend a separator.
1141        append_tags(&mut target, "a:1");
1142        assert_eq!("a:1", target);
1143
1144        // Subsequent non-empty appends are comma-separated.
1145        append_tags(&mut target, "b:2");
1146        assert_eq!("a:1,b:2", target);
1147
1148        // An empty segment remains a no-op even once the target is non-empty.
1149        append_tags(&mut target, "");
1150        assert_eq!("a:1,b:2", target);
1151    }
1152
1153    #[test]
1154    fn flatten_container_tag_comma_joins_the_tag_set() {
1155        assert_eq!("", flatten_container_tag(TagSet::default()));
1156
1157        let single: TagSet = std::iter::once(Tag::from_static("image_name:web")).collect();
1158        assert_eq!("image_name:web", flatten_container_tag(single));
1159
1160        let multiple: TagSet = ["image_name:web", "runtime:docker"]
1161            .into_iter()
1162            .map(Tag::from_static)
1163            .collect();
1164        let flattened = flatten_container_tag(multiple);
1165        assert_eq!(
1166            BTreeSet::from(["image_name:web", "runtime:docker"]),
1167            flattened.split(',').collect::<BTreeSet<_>>()
1168        );
1169    }
1170
1171    #[test]
1172    fn resolve_container_tags_prefers_explicit_container_tags_attribute() {
1173        // An explicit, non-empty `datadog.container_tags` attribute is used verbatim.
1174        let mut attributes = FastHashMap::default();
1175        attributes.insert(
1176            MetaString::from(KEY_DATADOG_CONTAINER_TAGS),
1177            AttributeValue::String(MetaString::from("region:us,team:core")),
1178        );
1179        assert_eq!(
1180            Some(MetaString::from("region:us,team:core")),
1181            resolve_container_tags_from_attrs(&attributes, None, false)
1182        );
1183    }
1184
1185    #[test]
1186    fn resolve_container_tags_returns_none_without_container_attributes() {
1187        let attributes = FastHashMap::default();
1188        // `ignore_missing_fields` skips the extraction path entirely.
1189        assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, true));
1190        // Honoring fields but with no container attributes and no Fargate source still yields nothing.
1191        assert_eq!(None, resolve_container_tags_from_attrs(&attributes, None, false));
1192    }
1193
1194    #[test]
1195    fn encode_prefixes_tracer_version_and_writes_otlp_sampling_rate() {
1196        let mut encoder = make_encoder(false);
1197        let mut trace = make_trace();
1198        trace.payload.tracer_version = MetaString::from("1.2.3");
1199        trace.otlp_sampling_rate = Some(0.5);
1200
1201        let mut buf = Vec::new();
1202        encoder.encode(&trace, &mut buf).expect("encode should succeed");
1203
1204        // The encoder emits the ETP `idxTracerPayloads` format, so decode through the string table.
1205        let tracer_versions = decode_etp_tracer_versions(&buf);
1206        let tracer_version = tracer_versions.first().expect("a tracer payload should be encoded");
1207        // The tracer version is prefixed with `otlp-` to mark the OTLP ingestion path.
1208        assert_eq!("otlp-1.2.3", tracer_version);
1209
1210        // The OTLP sampling rate is written to each chunk formatted to two decimal places.
1211        let chunk_attrs = decode_etp_chunk_attributes(&buf);
1212        let otlp_sr = chunk_attrs
1213            .iter()
1214            .find_map(|attrs| attrs.get("_dd.otlp_sr"))
1215            .expect("chunk should carry the _dd.otlp_sr tag");
1216        assert_eq!("0.50", otlp_sr.as_str());
1217    }
1218}