Skip to main content

vector/sources/datadog_agent/
traces.rs

1use std::{collections::BTreeMap, sync::Arc};
2
3use bytes::Bytes;
4use chrono::{TimeZone, Utc};
5use futures::future;
6use http::StatusCode;
7use ordered_float::NotNan;
8use prost::Message;
9use vector_lib::{
10    EstimatedJsonEncodedSizeOf,
11    internal_event::{CountByteSize, InternalEventHandle as _},
12};
13use vrl::event_path;
14use warp::{Filter, Rejection, Reply, filters::BoxedFilter, path, path::FullPath, reply::Response};
15
16use super::{ApiKeyQueryParams, DatadogAgentSource, RequestHandler, ddtrace_proto};
17use crate::{
18    common::http::ErrorMessage,
19    event::{Event, ObjectMap, TraceEvent, Value},
20    sources::util::http::capped_body,
21};
22
23pub(super) fn build_warp_filter(
24    handler: RequestHandler,
25    source: DatadogAgentSource,
26) -> BoxedFilter<(Response,)> {
27    build_trace_filter(handler, source)
28        .or(build_stats_filter())
29        .unify()
30        .boxed()
31}
32
33fn build_trace_filter(
34    handler: RequestHandler,
35    source: DatadogAgentSource,
36) -> BoxedFilter<(Response,)> {
37    warp::post()
38        .and(path!("api" / "v0.2" / "traces" / ..))
39        .and(warp::path::full())
40        .and(warp::header::optional::<String>("content-encoding"))
41        .and(warp::header::optional::<String>("dd-api-key"))
42        .and(warp::header::optional::<String>(
43            "X-Datadog-Reported-Languages",
44        ))
45        .and(warp::query::<ApiKeyQueryParams>())
46        .and(capped_body())
47        .and_then({
48            move |path: FullPath,
49                  encoding_header: Option<String>,
50                  api_token: Option<String>,
51                  reported_language: Option<String>,
52                  query_params: ApiKeyQueryParams,
53                  body: Bytes| {
54                let events = source
55                    .decode(&encoding_header, body, path.as_str())
56                    .and_then(|body| {
57                        handle_dd_trace_payload(
58                            body,
59                            source.api_key_extractor.extract(
60                                path.as_str(),
61                                api_token,
62                                query_params.dd_api_key,
63                            ),
64                            reported_language.as_ref(),
65                            &source,
66                        )
67                        .map_err(|error| {
68                            ErrorMessage::new(
69                                StatusCode::UNPROCESSABLE_ENTITY,
70                                format!("Error decoding Datadog traces: {error:?}"),
71                            )
72                        })
73                    });
74                handler.clone().handle_request(events, super::TRACES)
75            }
76        })
77        .boxed()
78}
79
80fn build_stats_filter() -> BoxedFilter<(Response,)> {
81    warp::post()
82        .and(path!("api" / "v0.2" / "stats" / ..))
83        .and_then(|| {
84            // APM stats are discarded on purpose, they will be computed in the `datadog_traces` sink
85            // thus we simply reply with a 200/OK response.
86            let response: Result<Response, Rejection> = Ok(warp::reply().into_response());
87            future::ready(response)
88        })
89        .boxed()
90}
91
92fn handle_dd_trace_payload(
93    frame: Bytes,
94    api_key: Option<Arc<str>>,
95    lang: Option<&String>,
96    source: &DatadogAgentSource,
97) -> crate::Result<Vec<Event>> {
98    let decoded_payload = ddtrace_proto::TracePayload::decode(frame)?;
99    if decoded_payload.tracer_payloads.is_empty() {
100        debug!("Older trace payload decoded.");
101        handle_dd_trace_payload_v0(decoded_payload, api_key, lang, source)
102    } else {
103        debug!("Newer trace payload decoded.");
104        handle_dd_trace_payload_v1(decoded_payload, api_key, source)
105    }
106}
107
108/// Decode Datadog newer protobuf schema
109fn handle_dd_trace_payload_v1(
110    decoded_payload: ddtrace_proto::TracePayload,
111    api_key: Option<Arc<str>>,
112    source: &DatadogAgentSource,
113) -> crate::Result<Vec<Event>> {
114    let env = decoded_payload.env;
115    let hostname = decoded_payload.host_name;
116    let agent_version = decoded_payload.agent_version;
117    let target_tps = decoded_payload.target_tps;
118    let error_tps = decoded_payload.error_tps;
119    let tags = convert_tags(decoded_payload.tags);
120
121    let trace_events: Vec<TraceEvent> = decoded_payload
122        .tracer_payloads
123        .into_iter()
124        .flat_map(convert_dd_tracer_payload)
125        .collect();
126
127    source.events_received.emit(CountByteSize(
128        trace_events.len(),
129        trace_events.estimated_json_encoded_size_of(),
130    ));
131
132    let enriched_events = trace_events
133        .into_iter()
134        .map(|mut trace_event| {
135            if let Some(k) = &api_key {
136                trace_event
137                    .metadata_mut()
138                    .set_datadog_api_key(Arc::clone(k));
139            }
140            trace_event.insert(
141                &source.log_schema_source_type_key,
142                Bytes::from("datadog_agent"),
143            );
144            trace_event.insert(event_path!("payload_version"), "v2".to_string());
145            trace_event.insert(&source.log_schema_host_key, hostname.clone());
146            trace_event.insert(event_path!("env"), env.clone());
147            trace_event.insert(event_path!("agent_version"), agent_version.clone());
148            trace_event.insert(
149                event_path!("target_tps"),
150                Value::Float(NotNan::new(target_tps).expect("target_tps cannot be Nan")),
151            );
152            trace_event.insert(
153                event_path!("error_tps"),
154                Value::Float(NotNan::new(error_tps).expect("error_tps cannot be Nan")),
155            );
156            if let Some(Value::Object(span_tags)) = trace_event.get_mut(event_path!("tags")) {
157                span_tags.extend(tags.clone());
158            } else {
159                trace_event.insert(event_path!("tags"), Value::from(tags.clone()));
160            }
161            Event::Trace(trace_event)
162        })
163        .collect();
164    Ok(enriched_events)
165}
166
167fn convert_dd_tracer_payload(payload: ddtrace_proto::TracerPayload) -> Vec<TraceEvent> {
168    let tags = convert_tags(payload.tags);
169    payload
170        .chunks
171        .into_iter()
172        .map(|trace| {
173            let mut trace_event = TraceEvent::default();
174            trace_event.insert(event_path!("priority"), trace.priority as i64);
175            trace_event.insert(event_path!("origin"), trace.origin);
176            trace_event.insert(event_path!("dropped"), trace.dropped_trace);
177            let mut trace_tags = convert_tags(trace.tags);
178            trace_tags.extend(tags.clone());
179            trace_event.insert(event_path!("tags"), Value::from(trace_tags));
180
181            trace_event.insert(
182                event_path!("spans"),
183                trace
184                    .spans
185                    .into_iter()
186                    .map(|s| Value::from(convert_span(s)))
187                    .collect::<Vec<Value>>(),
188            );
189
190            trace_event.insert(event_path!("container_id"), payload.container_id.clone());
191            trace_event.insert(event_path!("language_name"), payload.language_name.clone());
192            trace_event.insert(
193                event_path!("language_version"),
194                payload.language_version.clone(),
195            );
196            trace_event.insert(
197                event_path!("tracer_version"),
198                payload.tracer_version.clone(),
199            );
200            trace_event.insert(event_path!("runtime_id"), payload.runtime_id.clone());
201            trace_event.insert(event_path!("app_version"), payload.app_version.clone());
202            trace_event
203        })
204        .collect()
205}
206
207// Decode Datadog older protobuf schema
208fn handle_dd_trace_payload_v0(
209    decoded_payload: ddtrace_proto::TracePayload,
210    api_key: Option<Arc<str>>,
211    lang: Option<&String>,
212    source: &DatadogAgentSource,
213) -> crate::Result<Vec<Event>> {
214    let env = decoded_payload.env;
215    let hostname = decoded_payload.host_name;
216
217    let trace_events: Vec<TraceEvent> =
218    // Each traces is mapped to one event...
219    decoded_payload
220        .traces
221        .into_iter()
222        .map(|dd_trace| {
223            let mut trace_event = TraceEvent::default();
224
225            // TODO trace_id is being forced into an i64 but
226            // the incoming payload is u64. This is a bug and needs to be fixed per:
227            // https://github.com/vectordotdev/vector/issues/14687
228            trace_event.insert(event_path!("trace_id"), dd_trace.trace_id as i64);
229            trace_event.insert(event_path!("start_time"), Utc.timestamp_nanos(dd_trace.start_time));
230            trace_event.insert(event_path!("end_time"), Utc.timestamp_nanos(dd_trace.end_time));
231            trace_event.insert(
232                event_path!("spans"),
233                dd_trace
234                    .spans
235                    .into_iter()
236                    .map(|s| Value::from(convert_span(s)))
237                    .collect::<Vec<Value>>(),
238            );
239            trace_event
240        })
241        //... and each APM event is also mapped into its own event
242        .chain(decoded_payload.transactions.into_iter().map(|s| {
243            let mut trace_event = TraceEvent::default();
244            trace_event.insert(event_path!("spans"), vec![Value::from(convert_span(s))]);
245            trace_event.insert(event_path!("dropped"), true);
246            trace_event
247        })).collect();
248
249    source.events_received.emit(CountByteSize(
250        trace_events.len(),
251        trace_events.estimated_json_encoded_size_of(),
252    ));
253
254    let enriched_events = trace_events
255        .into_iter()
256        .map(|mut trace_event| {
257            if let Some(k) = &api_key {
258                trace_event
259                    .metadata_mut()
260                    .set_datadog_api_key(Arc::clone(k));
261            }
262            if let Some(lang) = lang {
263                trace_event.insert(event_path!("language_name"), lang.clone());
264            }
265            trace_event.insert(
266                &source.log_schema_source_type_key,
267                Bytes::from("datadog_agent"),
268            );
269            trace_event.insert(event_path!("payload_version"), "v1".to_string());
270            trace_event.insert(&source.log_schema_host_key, hostname.clone());
271            trace_event.insert(event_path!("env"), env.clone());
272            Event::Trace(trace_event)
273        })
274        .collect();
275
276    Ok(enriched_events)
277}
278
279fn convert_span(dd_span: ddtrace_proto::Span) -> ObjectMap {
280    let mut span = ObjectMap::new();
281    span.insert("service".into(), Value::from(dd_span.service));
282    span.insert("name".into(), Value::from(dd_span.name));
283
284    span.insert("resource".into(), Value::from(dd_span.resource));
285
286    // TODO trace_id, span_id and parent_id are being forced into an i64 but
287    // the incoming payload is u64. This is a bug and needs to be fixed per:
288    // https://github.com/vectordotdev/vector/issues/14687
289    span.insert("trace_id".into(), Value::from(dd_span.trace_id as i64));
290    span.insert("span_id".into(), Value::from(dd_span.span_id as i64));
291    span.insert("parent_id".into(), Value::from(dd_span.parent_id as i64));
292    span.insert(
293        "start".into(),
294        Value::from(Utc.timestamp_nanos(dd_span.start)),
295    );
296    span.insert("duration".into(), Value::from(dd_span.duration));
297    span.insert("error".into(), Value::from(dd_span.error as i64));
298    span.insert("meta".into(), Value::from(convert_tags(dd_span.meta)));
299    span.insert(
300        "metrics".into(),
301        Value::from(
302            dd_span
303                .metrics
304                .into_iter()
305                .map(|(k, v)| {
306                    (
307                        k.into(),
308                        NotNan::new(v).map(Value::Float).unwrap_or(Value::Null),
309                    )
310                })
311                .collect::<ObjectMap>(),
312        ),
313    );
314    span.insert("type".into(), Value::from(dd_span.r#type));
315    span.insert(
316        "meta_struct".into(),
317        Value::from(
318            dd_span
319                .meta_struct
320                .into_iter()
321                .map(|(k, v)| (k.into(), Value::from(bytes::Bytes::from(v))))
322                .collect::<ObjectMap>(),
323        ),
324    );
325
326    span
327}
328
329fn convert_tags(original_map: BTreeMap<String, String>) -> ObjectMap {
330    original_map
331        .into_iter()
332        .map(|(k, v)| (k.into(), Value::from(v)))
333        .collect::<ObjectMap>()
334}