vector/sources/datadog_agent/
traces.rs1use 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 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
108fn 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
207fn 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 decoded_payload
220 .traces
221 .into_iter()
222 .map(|dd_trace| {
223 let mut trace_event = TraceEvent::default();
224
225 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 .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 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}