Skip to main content

codecs/decoding/format/
otlp.rs

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/// OTLP signal type for prioritized parsing.
19#[configurable_component]
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
21#[serde(rename_all = "snake_case")]
22pub enum OtlpSignalType {
23    /// OTLP logs signal (ExportLogsServiceRequest)
24    Logs,
25    /// OTLP metrics signal (ExportMetricsServiceRequest)
26    Metrics,
27    /// OTLP traces signal (ExportTraceServiceRequest)
28    Traces,
29}
30
31/// Config used to build an `OtlpDeserializer`.
32#[configurable_component]
33#[derive(Debug, Clone)]
34pub struct OtlpDeserializerConfig {
35    /// Signal types to attempt parsing, in priority order.
36    ///
37    /// The deserializer tries to parse signals in the specified order. This allows you to optimize
38    /// performance when you know the expected signal types. For example, if you only receive
39    /// traces, set this to `["traces"]` to avoid attempting to parse as logs or metrics first.
40    ///
41    /// If not specified, defaults to trying all types in this order: logs, metrics, traces.
42    /// Duplicate signal types are automatically removed while preserving order.
43    #[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    /// Build the `OtlpDeserializer` from this configuration.
65    pub fn build(&self) -> OtlpDeserializer {
66        OtlpDeserializer::new_with_signals(self.signal_types.clone())
67    }
68
69    /// Return the type of event build by this deserializer.
70    pub fn output_type(&self) -> DataType {
71        DataType::Log | DataType::Trace
72    }
73
74    /// The schema produced by the deserializer.
75    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/// Deserializer that builds `Event`s from a byte frame containing [OTLP](https://opentelemetry.io/docs/specs/otlp/) protobuf data.
88///
89/// This deserializer decodes events using the OTLP protobuf specification. It handles the three
90/// OTLP signal types: logs, metrics, and traces.
91///
92/// The implementation supports three OTLP message types:
93/// - `ExportLogsServiceRequest` → Log events with `resourceLogs` field
94/// - `ExportMetricsServiceRequest` → Log events with `resourceMetrics` field
95/// - `ExportTraceServiceRequest` → Trace events with `resourceSpans` field
96///
97/// One major caveat here is that the incoming metrics will be parsed as logs but they will preserve the OTLP format.
98/// This means that components that work on metrics, will not be compatible with this output.
99/// However, these events can be forwarded directly to a downstream OTEL collector.
100///
101/// This is the inverse of what the OTLP encoder does, ensuring round-trip compatibility
102/// with the `opentelemetry` source when `use_otlp_decoding` is enabled.
103#[derive(Debug, Clone)]
104pub struct OtlpDeserializer {
105    logs_deserializer: ProtobufDeserializer,
106    metrics_deserializer: ProtobufDeserializer,
107    traces_deserializer: ProtobufDeserializer,
108    /// Signal types to parse, in priority order
109    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    /// Creates a new OTLP deserializer with custom signal support.
120    /// During parsing, each signal type is tried in order until one succeeds.
121    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        // Try parsing in the priority order specified
163        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                    // TODO: <https://github.com/vectordotdev/vector/issues/25045>
185                    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                        // Convert the log event to a trace event by taking ownership
191                        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    // trace_id: 0102030405060708090a0b0c0d0e0f10 (16 bytes)
222    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    // span_id: 0102030405060708 (8 bytes)
227    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        // Navigate to the span and check traceId and spanId
320        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        // Verify traceId - should be raw bytes (16 bytes for trace_id)
346        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        // Verify spanId - should be raw bytes (8 bytes for span_id)
358        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        // Configure to only try traces - should succeed for traces, fail for others
426        let deserializer =
427            OtlpDeserializer::new_with_signals(IndexSet::from([OtlpSignalType::Traces]));
428
429        // Traces should work
430        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        // Logs should fail since we're not trying to parse logs
436        let log_bytes = create_logs_request_bytes();
437        let result = deserializer.parse(log_bytes, LogNamespace::Legacy);
438        assert!(result.is_err());
439    }
440}