1use bytes::Bytes;
2use opentelemetry_proto::proto::{
3 DESCRIPTOR_BYTES, LOGS_REQUEST_MESSAGE_TYPE, METRICS_REQUEST_MESSAGE_TYPE,
4 RESOURCE_LOGS_JSON_FIELD, RESOURCE_METRICS_JSON_FIELD, RESOURCE_SPANS_JSON_FIELD,
5 TRACES_REQUEST_MESSAGE_TYPE,
6};
7use smallvec::{SmallVec, smallvec};
8use vector_config::{configurable_component, indexmap::IndexSet};
9use vector_core::{
10 config::{DataType, LogNamespace},
11 event::Event,
12 schema,
13};
14use vrl::{event_path, protobuf::parse::Options, value::Kind};
15
16use super::{Deserializer, ProtobufDeserializer};
17
18#[configurable_component]
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
21#[serde(rename_all = "snake_case")]
22pub enum OtlpSignalType {
23 Logs,
25 Metrics,
27 Traces,
29}
30
31#[configurable_component]
33#[derive(Debug, Clone)]
34pub struct OtlpDeserializerConfig {
35 #[serde(default = "default_signal_types")]
44 pub signal_types: IndexSet<OtlpSignalType>,
45}
46
47fn default_signal_types() -> IndexSet<OtlpSignalType> {
48 IndexSet::from([
49 OtlpSignalType::Logs,
50 OtlpSignalType::Metrics,
51 OtlpSignalType::Traces,
52 ])
53}
54
55impl Default for OtlpDeserializerConfig {
56 fn default() -> Self {
57 Self {
58 signal_types: default_signal_types(),
59 }
60 }
61}
62
63impl OtlpDeserializerConfig {
64 pub fn build(&self) -> OtlpDeserializer {
66 OtlpDeserializer::new_with_signals(self.signal_types.clone())
67 }
68
69 pub fn output_type(&self) -> DataType {
71 DataType::Log | DataType::Trace
72 }
73
74 pub fn schema_definition(&self, log_namespace: LogNamespace) -> schema::Definition {
76 match log_namespace {
77 LogNamespace::Legacy => {
78 schema::Definition::empty_legacy_namespace().unknown_fields(Kind::any())
79 }
80 LogNamespace::Vector => {
81 schema::Definition::new_with_default_metadata(Kind::any(), [log_namespace])
82 }
83 }
84 }
85}
86
87#[derive(Debug, Clone)]
104pub struct OtlpDeserializer {
105 logs_deserializer: ProtobufDeserializer,
106 metrics_deserializer: ProtobufDeserializer,
107 traces_deserializer: ProtobufDeserializer,
108 signals: IndexSet<OtlpSignalType>,
110}
111
112impl Default for OtlpDeserializer {
113 fn default() -> Self {
114 Self::new_with_signals(default_signal_types())
115 }
116}
117
118impl OtlpDeserializer {
119 pub fn new_with_signals(signals: IndexSet<OtlpSignalType>) -> Self {
122 let options = Options {
123 use_json_names: true,
124 };
125
126 let logs_deserializer = ProtobufDeserializer::new_from_bytes(
127 DESCRIPTOR_BYTES,
128 LOGS_REQUEST_MESSAGE_TYPE,
129 options.clone(),
130 )
131 .expect("Failed to create logs deserializer");
132
133 let metrics_deserializer = ProtobufDeserializer::new_from_bytes(
134 DESCRIPTOR_BYTES,
135 METRICS_REQUEST_MESSAGE_TYPE,
136 options.clone(),
137 )
138 .expect("Failed to create metrics deserializer");
139
140 let traces_deserializer = ProtobufDeserializer::new_from_bytes(
141 DESCRIPTOR_BYTES,
142 TRACES_REQUEST_MESSAGE_TYPE,
143 options,
144 )
145 .expect("Failed to create traces deserializer");
146
147 Self {
148 logs_deserializer,
149 metrics_deserializer,
150 traces_deserializer,
151 signals,
152 }
153 }
154}
155
156impl Deserializer for OtlpDeserializer {
157 fn parse(
158 &self,
159 bytes: Bytes,
160 log_namespace: LogNamespace,
161 ) -> vector_common::Result<SmallVec<[Event; 1]>> {
162 for signal_type in &self.signals {
164 match signal_type {
165 OtlpSignalType::Logs => {
166 if let Ok(events) = self.logs_deserializer.parse(bytes.clone(), log_namespace)
167 && let Some(Event::Log(log)) = events.first()
168 && log.get(event_path!(RESOURCE_LOGS_JSON_FIELD)).is_some()
169 {
170 return Ok(events);
171 }
172 }
173 OtlpSignalType::Metrics => {
174 if let Ok(events) = self
175 .metrics_deserializer
176 .parse(bytes.clone(), log_namespace)
177 && let Some(Event::Log(log)) = events.first()
178 && log.get(event_path!(RESOURCE_METRICS_JSON_FIELD)).is_some()
179 {
180 return Ok(events);
181 }
182 }
183 OtlpSignalType::Traces => {
184 if let Ok(mut events) =
186 self.traces_deserializer.parse(bytes.clone(), log_namespace)
187 && let Some(Event::Log(log)) = events.first()
188 && log.get(event_path!(RESOURCE_SPANS_JSON_FIELD)).is_some()
189 {
190 if let Some(Event::Log(log)) = events.pop() {
192 let trace_event = Event::Trace(log.into());
193 return Ok(smallvec![trace_event]);
194 }
195 }
196 }
197 }
198 }
199
200 Err(format!("Invalid OTLP data: expected one of {:?}", self.signals).into())
201 }
202}
203
204#[cfg(test)]
205mod tests {
206 use opentelemetry_proto::proto::{
207 collector::{
208 logs::v1::ExportLogsServiceRequest, metrics::v1::ExportMetricsServiceRequest,
209 trace::v1::ExportTraceServiceRequest,
210 },
211 logs::v1::{LogRecord, ResourceLogs, ScopeLogs},
212 metrics::v1::{Metric, ResourceMetrics, ScopeMetrics},
213 resource::v1::Resource,
214 trace::v1::{ResourceSpans, ScopeSpans, Span},
215 };
216 use prost::Message;
217 use vrl::path;
218
219 use super::*;
220
221 const TEST_TRACE_ID: [u8; 16] = [
223 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d, 0x0e, 0x0f,
224 0x10,
225 ];
226 const TEST_SPAN_ID: [u8; 8] = [0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08];
228
229 fn create_logs_request_bytes() -> Bytes {
230 let request = ExportLogsServiceRequest {
231 resource_logs: vec![ResourceLogs {
232 resource: Some(Resource {
233 attributes: vec![],
234 dropped_attributes_count: 0,
235 }),
236 scope_logs: vec![ScopeLogs {
237 scope: None,
238 log_records: vec![LogRecord {
239 time_unix_nano: 1234567890,
240 severity_number: 9,
241 severity_text: "INFO".to_string(),
242 body: None,
243 attributes: vec![],
244 dropped_attributes_count: 0,
245 flags: 0,
246 trace_id: vec![],
247 span_id: vec![],
248 observed_time_unix_nano: 0,
249 }],
250 schema_url: String::new(),
251 }],
252 schema_url: String::new(),
253 }],
254 };
255
256 Bytes::from(request.encode_to_vec())
257 }
258
259 fn create_metrics_request_bytes() -> Bytes {
260 let request = ExportMetricsServiceRequest {
261 resource_metrics: vec![ResourceMetrics {
262 resource: Some(Resource {
263 attributes: vec![],
264 dropped_attributes_count: 0,
265 }),
266 scope_metrics: vec![ScopeMetrics {
267 scope: None,
268 metrics: vec![Metric {
269 name: "test_metric".to_string(),
270 description: String::new(),
271 unit: String::new(),
272 data: None,
273 }],
274 schema_url: String::new(),
275 }],
276 schema_url: String::new(),
277 }],
278 };
279
280 Bytes::from(request.encode_to_vec())
281 }
282
283 fn create_traces_request_bytes() -> Bytes {
284 let request = ExportTraceServiceRequest {
285 resource_spans: vec![ResourceSpans {
286 resource: Some(Resource {
287 attributes: vec![],
288 dropped_attributes_count: 0,
289 }),
290 scope_spans: vec![ScopeSpans {
291 scope: None,
292 spans: vec![Span {
293 trace_id: TEST_TRACE_ID.to_vec(),
294 span_id: TEST_SPAN_ID.to_vec(),
295 trace_state: String::new(),
296 parent_span_id: vec![],
297 name: "test_span".to_string(),
298 kind: 0,
299 start_time_unix_nano: 1234567890,
300 end_time_unix_nano: 1234567900,
301 attributes: vec![],
302 dropped_attributes_count: 0,
303 events: vec![],
304 dropped_events_count: 0,
305 links: vec![],
306 dropped_links_count: 0,
307 status: None,
308 }],
309 schema_url: String::new(),
310 }],
311 schema_url: String::new(),
312 }],
313 };
314
315 Bytes::from(request.encode_to_vec())
316 }
317
318 fn validate_trace_ids(trace: &vrl::value::Value) {
319 let resource_spans = trace
321 .get(path!("resourceSpans"))
322 .and_then(|v| v.as_array())
323 .expect("resourceSpans should be an array");
324
325 let first_rs = resource_spans
326 .first()
327 .expect("should have at least one resource span");
328
329 let scope_spans = first_rs
330 .get(path!("scopeSpans"))
331 .and_then(|v| v.as_array())
332 .expect("scopeSpans should be an array");
333
334 let first_ss = scope_spans
335 .first()
336 .expect("should have at least one scope span");
337
338 let spans = first_ss
339 .get(path!("spans"))
340 .and_then(|v| v.as_array())
341 .expect("spans should be an array");
342
343 let span = spans.first().expect("should have at least one span");
344
345 let trace_id = span
347 .get(path!("traceId"))
348 .and_then(|v| v.as_bytes())
349 .expect("traceId should exist and be bytes");
350
351 assert_eq!(
352 trace_id.as_ref(),
353 &TEST_TRACE_ID,
354 "traceId should match the expected 16 bytes (0102030405060708090a0b0c0d0e0f10)"
355 );
356
357 let span_id = span
359 .get(path!("spanId"))
360 .and_then(|v| v.as_bytes())
361 .expect("spanId should exist and be bytes");
362
363 assert_eq!(
364 span_id.as_ref(),
365 &TEST_SPAN_ID,
366 "spanId should match the expected 8 bytes (0102030405060708)"
367 );
368 }
369
370 fn assert_otlp_event(bytes: Bytes, field: &str, is_trace: bool) {
371 let deserializer = OtlpDeserializer::default();
372 let events = deserializer.parse(bytes, LogNamespace::Legacy).unwrap();
373
374 assert_eq!(events.len(), 1);
375 if is_trace {
376 assert!(matches!(events[0], Event::Trace(_)));
377 let trace = events[0].as_trace();
378 assert!(trace.get(event_path!(field)).is_some());
379 validate_trace_ids(trace.value());
380 } else {
381 assert!(events[0].as_log().get(event_path!(field)).is_some());
382 }
383 }
384
385 #[test]
386 fn deserialize_otlp_logs() {
387 assert_otlp_event(create_logs_request_bytes(), RESOURCE_LOGS_JSON_FIELD, false);
388 }
389
390 #[test]
391 fn deserialize_otlp_metrics() {
392 assert_otlp_event(
393 create_metrics_request_bytes(),
394 RESOURCE_METRICS_JSON_FIELD,
395 false,
396 );
397 }
398
399 #[test]
400 fn deserialize_otlp_traces() {
401 assert_otlp_event(
402 create_traces_request_bytes(),
403 RESOURCE_SPANS_JSON_FIELD,
404 true,
405 );
406 }
407
408 #[test]
409 fn deserialize_invalid_otlp() {
410 let deserializer = OtlpDeserializer::default();
411 let bytes = Bytes::from("invalid protobuf data");
412 let result = deserializer.parse(bytes, LogNamespace::Legacy);
413
414 assert!(result.is_err());
415 assert!(
416 result
417 .unwrap_err()
418 .to_string()
419 .contains("Invalid OTLP data")
420 );
421 }
422
423 #[test]
424 fn deserialize_with_custom_priority_traces_only() {
425 let deserializer =
427 OtlpDeserializer::new_with_signals(IndexSet::from([OtlpSignalType::Traces]));
428
429 let trace_bytes = create_traces_request_bytes();
431 let result = deserializer.parse(trace_bytes, LogNamespace::Legacy);
432 assert!(result.is_ok());
433 assert!(matches!(result.unwrap()[0], Event::Trace(_)));
434
435 let log_bytes = create_logs_request_bytes();
437 let result = deserializer.parse(log_bytes, LogNamespace::Legacy);
438 assert!(result.is_err());
439 }
440}