Skip to main content

vector/sinks/datadog/traces/
request_builder.rs

1use std::{
2    collections::BTreeMap,
3    io::Write,
4    num::NonZeroUsize,
5    sync::{Arc, Mutex},
6};
7
8use bytes::Bytes;
9use prost::Message;
10use snafu::Snafu;
11use vector_lib::{
12    event::{EventFinalizers, Finalizable},
13    request_metadata::RequestMetadata,
14};
15use vrl::event_path;
16
17use super::{
18    apm_stats::{Aggregator, compute_apm_stats},
19    config::{DatadogTracesEndpoint, DatadogTracesEndpointConfiguration},
20    dd_proto,
21    service::TraceApiRequest,
22    sink::PartitionKey,
23};
24use crate::{
25    event::{Event, ObjectMap, TraceEvent, Value},
26    sinks::util::{
27        Compression, Compressor, IncrementalRequestBuilder, metadata::RequestMetadataBuilder,
28    },
29};
30
31#[derive(Debug, Snafu)]
32pub enum RequestBuilderError {
33    #[snafu(display(
34        "Building an APM stats request payload failed ({}, {})",
35        message,
36        reason
37    ))]
38    FailedToBuild {
39        message: &'static str,
40        reason: String,
41        dropped_events: u64,
42    },
43
44    #[allow(dead_code)]
45    #[snafu(display("Unsupported endpoint ({})", reason))]
46    UnsupportedEndpoint { reason: String, dropped_events: u64 },
47}
48
49impl RequestBuilderError {
50    #[allow(clippy::missing_const_for_fn)] // const cannot run destructor
51    pub fn into_parts(self) -> (&'static str, String, u64) {
52        match self {
53            Self::FailedToBuild {
54                message,
55                reason,
56                dropped_events,
57            } => (message, reason, dropped_events),
58            Self::UnsupportedEndpoint {
59                reason,
60                dropped_events,
61            } => ("unsupported endpoint", reason, dropped_events),
62        }
63    }
64}
65
66pub struct DatadogTracesRequestBuilder {
67    api_key: Arc<str>,
68    endpoint_configuration: DatadogTracesEndpointConfiguration,
69    compression: Compression,
70    max_size: usize,
71    /// Contains the Aggregated stats across a time window.
72    stats_aggregator: Arc<Mutex<Aggregator>>,
73}
74
75impl DatadogTracesRequestBuilder {
76    pub const fn new(
77        api_key: Arc<str>,
78        endpoint_configuration: DatadogTracesEndpointConfiguration,
79        compression: Compression,
80        max_size: usize,
81        stats_aggregator: Arc<Mutex<Aggregator>>,
82    ) -> Result<Self, RequestBuilderError> {
83        Ok(Self {
84            api_key,
85            endpoint_configuration,
86            compression,
87            max_size,
88            stats_aggregator,
89        })
90    }
91}
92
93pub struct DDTracesMetadata {
94    pub api_key: Arc<str>,
95    pub endpoint: DatadogTracesEndpoint,
96    pub finalizers: EventFinalizers,
97    pub uncompressed_size: usize,
98    pub content_type: String,
99}
100
101impl IncrementalRequestBuilder<(PartitionKey, Vec<Event>)> for DatadogTracesRequestBuilder {
102    type Metadata = (DDTracesMetadata, RequestMetadata);
103    type Payload = Bytes;
104    type Request = TraceApiRequest;
105    type Error = RequestBuilderError;
106
107    fn encode_events_incremental(
108        &mut self,
109        input: (PartitionKey, Vec<Event>),
110    ) -> Vec<Result<(Self::Metadata, Self::Payload), Self::Error>> {
111        let (key, events) = input;
112        let trace_events = events
113            .into_iter()
114            .filter_map(|e| e.try_into_trace())
115            .collect::<Vec<TraceEvent>>();
116
117        // Compute APM stats from the incoming events. The stats payloads are sent out
118        // separately from the sink framework, by the thread `flush_apm_stats_thread()`
119        compute_apm_stats(&key, Arc::clone(&self.stats_aggregator), &trace_events);
120
121        encode_traces(&key, trace_events, self.max_size)
122            .into_iter()
123            .map(|result| {
124                result.and_then(|(payload, mut processed)| {
125                    let uncompressed_size = payload.len();
126                    let metadata = DDTracesMetadata {
127                        api_key: key
128                            .api_key
129                            .clone()
130                            .unwrap_or_else(|| Arc::clone(&self.api_key)),
131                        endpoint: DatadogTracesEndpoint::Traces,
132                        finalizers: processed.take_finalizers(),
133                        uncompressed_size,
134                        content_type: "application/x-protobuf".to_string(),
135                    };
136
137                    // build RequestMetadata
138                    let builder = RequestMetadataBuilder::from_events(&processed);
139
140                    let mut compressor = Compressor::from(self.compression);
141                    match compressor.write_all(&payload) {
142                        Ok(()) => {
143                            let bytes = compressor.into_inner().freeze();
144
145                            let bytes_len = NonZeroUsize::new(bytes.len())
146                                .expect("payload should never be zero length");
147                            let request_metadata = builder.with_request_size(bytes_len);
148
149                            Ok(((metadata, request_metadata), bytes))
150                        }
151                        Err(e) => Err(RequestBuilderError::FailedToBuild {
152                            message: "Payload compression failed.",
153                            reason: e.to_string(),
154                            dropped_events: processed.len() as u64,
155                        }),
156                    }
157                })
158            })
159            .collect()
160    }
161
162    fn build_request(&mut self, metadata: Self::Metadata, payload: Self::Payload) -> Self::Request {
163        build_request(
164            metadata,
165            payload,
166            self.compression,
167            &self.endpoint_configuration,
168        )
169    }
170}
171
172/// Builds the `TraceApiRequest` from inputs.
173///
174/// # Arguments
175///
176/// * `metadata`                 - Tuple of Datadog traces specific metadata and the generic `RequestMetadata`.
177/// * `payload`                  - Compressed and encoded bytes to send.
178/// * `compression`              - `Compression` used to reference the Content-Encoding header.
179/// * `endpoint_configuration`   - Endpoint configuration to use when creating the HTTP requests.
180pub fn build_request(
181    metadata: (DDTracesMetadata, RequestMetadata),
182    payload: Bytes,
183    compression: Compression,
184    endpoint_configuration: &DatadogTracesEndpointConfiguration,
185) -> TraceApiRequest {
186    let (ddtraces_metadata, request_metadata) = metadata;
187    let mut headers = BTreeMap::<String, String>::new();
188    headers.insert("Content-Type".to_string(), ddtraces_metadata.content_type);
189    headers.insert(
190        "DD-API-KEY".to_string(),
191        ddtraces_metadata.api_key.to_string(),
192    );
193    if let Some(ce) = compression.content_encoding() {
194        headers.insert("Content-Encoding".to_string(), ce.to_string());
195    }
196    TraceApiRequest {
197        body: payload,
198        headers,
199        finalizers: ddtraces_metadata.finalizers,
200        uri: endpoint_configuration
201            .get_uri_for_endpoint(ddtraces_metadata.endpoint)
202            .into_uri(),
203        uncompressed_size: ddtraces_metadata.uncompressed_size,
204        metadata: request_metadata,
205    }
206}
207
208fn encode_traces(
209    key: &PartitionKey,
210    trace_events: Vec<TraceEvent>,
211    max_size: usize,
212) -> Vec<Result<(Vec<u8>, Vec<TraceEvent>), RequestBuilderError>> {
213    let mut results = Vec::new();
214    let mut processed = Vec::new();
215    let mut payload = build_empty_payload(key);
216
217    for trace in trace_events {
218        let mut proto = encode_trace(&trace);
219
220        loop {
221            payload.tracer_payloads.push(proto);
222            if payload.encoded_len() >= max_size {
223                // take it back out
224                proto = payload.tracer_payloads.pop().expect("just pushed");
225                if payload.tracer_payloads.is_empty() {
226                    // this individual trace is too big
227                    results.push(Err(RequestBuilderError::FailedToBuild {
228                        message: "Dropped trace event",
229                        reason: "Trace is larger than allowed payload size".into(),
230                        dropped_events: 1,
231                    }));
232
233                    break;
234                } else {
235                    // try with a fresh payload
236                    results.push(Ok((
237                        payload.encode_to_vec(),
238                        std::mem::take(&mut processed),
239                    )));
240                    payload = build_empty_payload(key);
241                }
242            } else {
243                processed.push(trace);
244                break;
245            }
246        }
247    }
248    results.push(Ok((
249        payload.encode_to_vec(),
250        std::mem::take(&mut processed),
251    )));
252    results
253}
254
255fn build_empty_payload(key: &PartitionKey) -> dd_proto::TracePayload {
256    dd_proto::TracePayload {
257        host_name: key.hostname.clone().unwrap_or_default(),
258        env: key.env.clone().unwrap_or_default(),
259        traces: vec![],       // Field reserved for the older trace payloads
260        transactions: vec![], // Field reserved for the older trace payloads
261        tracer_payloads: vec![],
262        // We only send tags at the Trace level
263        tags: BTreeMap::new(),
264        agent_version: key.agent_version.clone().unwrap_or_default(),
265        target_tps: key.target_tps.map(|tps| tps as f64).unwrap_or_default(),
266        error_tps: key.error_tps.map(|tps| tps as f64).unwrap_or_default(),
267    }
268}
269
270fn encode_trace(trace: &TraceEvent) -> dd_proto::TracerPayload {
271    let tags = trace
272        .get(event_path!("tags"))
273        .and_then(|m| m.as_object())
274        .map(|m| {
275            m.iter()
276                .map(|(k, v)| (k.to_string(), v.to_string_lossy().into_owned()))
277                .collect::<BTreeMap<String, String>>()
278        })
279        .unwrap_or_default();
280
281    let spans = match trace.get(event_path!("spans")) {
282        Some(Value::Array(v)) => v
283            .iter()
284            .filter_map(|s| s.as_object().map(convert_span))
285            .collect(),
286        _ => vec![],
287    };
288
289    let chunk = dd_proto::TraceChunk {
290        priority: trace
291            .get(event_path!("priority"))
292            .and_then(|v| v.as_integer().map(|v| v as i32))
293            // This should not happen for Datadog originated traces, but in case this field is not populated
294            // we default to 1 (https://github.com/DataDog/datadog-agent/blob/eac2327/pkg/trace/sampler/sampler.go#L54-L55),
295            // which is what the Datadog trace-agent is doing for OTLP originated traces, as per
296            // https://github.com/DataDog/datadog-agent/blob/3ea2eb4/pkg/trace/api/otlp.go#L309.
297            .unwrap_or(1i32),
298        origin: trace
299            .get(event_path!("origin"))
300            .map(|v| v.to_string_lossy().into_owned())
301            .unwrap_or_default(),
302        dropped_trace: trace
303            .get(event_path!("dropped"))
304            .and_then(|v| v.as_boolean())
305            .unwrap_or(false),
306        spans,
307        tags: tags.clone(),
308    };
309
310    dd_proto::TracerPayload {
311        container_id: trace
312            .get(event_path!("container_id"))
313            .map(|v| v.to_string_lossy().into_owned())
314            .unwrap_or_default(),
315        language_name: trace
316            .get(event_path!("language_name"))
317            .map(|v| v.to_string_lossy().into_owned())
318            .unwrap_or_default(),
319        language_version: trace
320            .get(event_path!("language_version"))
321            .map(|v| v.to_string_lossy().into_owned())
322            .unwrap_or_default(),
323        tracer_version: trace
324            .get(event_path!("tracer_version"))
325            .map(|v| v.to_string_lossy().into_owned())
326            .unwrap_or_default(),
327        runtime_id: trace
328            .get(event_path!("runtime_id"))
329            .map(|v| v.to_string_lossy().into_owned())
330            .unwrap_or_default(),
331        chunks: vec![chunk],
332        tags,
333        env: trace
334            .get(event_path!("env"))
335            .map(|v| v.to_string_lossy().into_owned())
336            .unwrap_or_default(),
337        hostname: trace
338            .get(event_path!("hostname"))
339            .map(|v| v.to_string_lossy().into_owned())
340            .unwrap_or_default(),
341        app_version: trace
342            .get(event_path!("app_version"))
343            .map(|v| v.to_string_lossy().into_owned())
344            .unwrap_or_default(),
345    }
346}
347
348fn convert_span(span: &ObjectMap) -> dd_proto::Span {
349    let trace_id = match span.get("trace_id") {
350        Some(Value::Integer(val)) => *val,
351        _ => 0,
352    };
353    let span_id = match span.get("span_id") {
354        Some(Value::Integer(val)) => *val,
355        _ => 0,
356    };
357    let parent_id = match span.get("parent_id") {
358        Some(Value::Integer(val)) => *val,
359        _ => 0,
360    };
361    let duration = match span.get("duration") {
362        Some(Value::Integer(val)) => *val,
363        _ => 0,
364    };
365    let error = match span.get("error") {
366        Some(Value::Integer(val)) => *val,
367        _ => 0,
368    };
369    let start = match span.get("start") {
370        Some(Value::Timestamp(val)) => val.timestamp_nanos_opt().expect("Timestamp out of range"),
371        _ => 0,
372    };
373
374    let meta = span
375        .get("meta")
376        .and_then(|m| m.as_object())
377        .map(|m| {
378            m.iter()
379                .map(|(k, v)| (k.to_string(), v.to_string_lossy().into_owned()))
380                .collect::<BTreeMap<String, String>>()
381        })
382        .unwrap_or_default();
383
384    let meta_struct = span
385        .get("meta_struct")
386        .and_then(|m| m.as_object())
387        .map(|m| {
388            m.iter()
389                .map(|(k, v)| (k.to_string(), v.coerce_to_bytes().into_iter().collect()))
390                .collect::<BTreeMap<String, Vec<u8>>>()
391        })
392        .unwrap_or_default();
393
394    let metrics = span
395        .get("metrics")
396        .and_then(|m| m.as_object())
397        .map(|m| {
398            m.iter()
399                .filter_map(|(k, v)| {
400                    if let Value::Float(f) = v {
401                        Some((k.to_string(), f.into_inner()))
402                    } else {
403                        None
404                    }
405                })
406                .collect::<BTreeMap<String, f64>>()
407        })
408        .unwrap_or_default();
409
410    dd_proto::Span {
411        service: span
412            .get("service")
413            .map(|v| v.to_string_lossy().into_owned())
414            .unwrap_or_default(),
415        name: span
416            .get("name")
417            .map(|v| v.to_string_lossy().into_owned())
418            .unwrap_or_default(),
419        resource: span
420            .get("resource")
421            .map(|v| v.to_string_lossy().into_owned())
422            .unwrap_or_default(),
423        r#type: span
424            .get("type")
425            .map(|v| v.to_string_lossy().into_owned())
426            .unwrap_or_default(),
427        trace_id: trace_id as u64,
428        span_id: span_id as u64,
429        parent_id: parent_id as u64,
430        error: error as i32,
431        start,
432        duration,
433        meta,
434        metrics,
435        meta_struct,
436    }
437}
438
439#[cfg(test)]
440mod test {
441    use proptest::prelude::*;
442    use vrl::event_path;
443
444    use super::{PartitionKey, encode_traces};
445    use crate::event::{LogEvent, TraceEvent};
446
447    proptest! {
448        #[test]
449        fn successfully_encode_payloads_smaller_than_max_size(
450            // 476 is the experimentally determined size that will fill a payload after encoding and overhead
451            lengths in proptest::collection::vec(16usize..476, 1usize..256),
452        ) {
453            let max_size = 1024;
454
455            let key = PartitionKey {
456                api_key: Some("x".repeat(128).into()),
457                env: Some("production".into()),
458                hostname: Some("foo.bar.baz.local".into()),
459                agent_version: Some("1.2.3.4.5".into()),
460                target_tps: None,
461                error_tps: None,
462            };
463
464            // We only care about the size of the incoming traces, so just populate a single tag field
465            // that will be copied into the protobuf representation.
466            let traces = lengths
467                .into_iter()
468                .map(|n| {
469                    let mut log = LogEvent::default();
470                    log.insert(event_path!("tags", "foo"), "x".repeat(n));
471                    TraceEvent::from(log)
472                })
473                .collect();
474
475            for result in encode_traces(&key, traces, max_size) {
476                prop_assert!(result.is_ok());
477                let (encoded, _processed) = result.unwrap();
478
479                prop_assert!(
480                    encoded.len() <= max_size,
481                    "encoded len {} longer than max size {}",
482                    encoded.len(),
483                    max_size
484                );
485            }
486        }
487    }
488
489    #[test]
490    fn handles_too_large_events() {
491        let max_size = 1024;
492        // 476 is experimentally determined to be too big to fit into a <1024 byte proto
493        let lengths = [128, 476, 128];
494
495        let key = PartitionKey {
496            api_key: Some("x".repeat(128).into()),
497            env: Some("production".into()),
498            hostname: Some("foo.bar.baz.local".into()),
499            agent_version: Some("1.2.3.4.5".into()),
500            target_tps: None,
501            error_tps: None,
502        };
503
504        // We only care about the size of the incoming traces, so just populate a single tag field
505        // that will be copied into the protobuf representation.
506        let traces = lengths
507            .into_iter()
508            .map(|n| {
509                let mut log = LogEvent::default();
510                log.insert(event_path!("tags", "foo"), "x".repeat(n));
511                TraceEvent::from(log)
512            })
513            .collect();
514
515        let mut results = encode_traces(&key, traces, max_size);
516        assert_eq!(3, results.len());
517
518        match &mut results[..] {
519            [Ok(one), Err(_two), Ok(three)] => {
520                for (encoded, processed) in [one, three] {
521                    assert_eq!(1, processed.len());
522                    assert!(
523                        encoded.len() <= max_size,
524                        "encoded len {} longer than max size {}",
525                        encoded.len(),
526                        max_size
527                    );
528                }
529            }
530            _ => panic!(
531                "unexpected output {:?}",
532                results
533                    .iter()
534                    .map(|r| r.as_ref().map(|(_, p)| p.len()))
535                    .collect::<Vec<_>>()
536            ),
537        }
538    }
539}