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)] 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 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(&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 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
172pub 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 proto = payload.tracer_payloads.pop().expect("just pushed");
225 if payload.tracer_payloads.is_empty() {
226 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 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![], transactions: vec![], tracer_payloads: vec![],
262 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 .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 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 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 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 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}