saluki_components/decoders/otlp/
mod.rs

1use std::time::Duration;
2
3use agent_data_plane_config::domains;
4use async_trait::async_trait;
5use otlp_protos::opentelemetry::proto::collector::trace::v1::ExportTraceServiceRequest;
6use prost::Message;
7use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
8use saluki_core::{
9    components::{
10        decoders::{Decoder, DecoderBuilder, DecoderContext},
11        BuildContext,
12    },
13    data_model::{event::EventType, payload::PayloadType},
14    topology::interconnect::EventBufferManager,
15};
16use saluki_error::GenericError;
17use tokio::{
18    select,
19    time::{interval, MissedTickBehavior},
20};
21use tracing::{debug, error, warn};
22
23use crate::common::otlp::traces::translator::OtlpTracesTranslator;
24use crate::common::otlp::{
25    build_metrics, Metrics, OTLP_LOGS_GRPC_SERVICE_PATH, OTLP_METRICS_GRPC_SERVICE_PATH, OTLP_TRACES_GRPC_SERVICE_PATH,
26};
27
28/// Configuration for the OTLP decoder.
29pub struct OtlpDecoderConfiguration {
30    /// Resolved OTLP trace ingestion settings.
31    traces: domains::otlp::Traces,
32
33    /// Maximum length of a span's resource name, in bytes.
34    ///
35    /// Defaults to `usize::MAX`, meaning resource names are not truncated.
36    max_resource_len: usize,
37}
38
39impl OtlpDecoderConfiguration {
40    /// Creates a new `OtlpDecoderConfiguration` from the resolved OTLP trace configuration.
41    pub fn from_configuration(traces: &domains::otlp::Traces) -> Self {
42        Self {
43            traces: traces.clone(),
44            max_resource_len: usize::MAX,
45        }
46    }
47
48    /// Sets the maximum length of a span's resource name, in bytes.
49    ///
50    /// Resource names longer than this limit are truncated to the limit, on a UTF-8 boundary. Defaults to `usize::MAX`
51    /// (unbounded), which matches the decoder's behavior before resource name truncation was introduced.
52    ///
53    /// If set to `0`, every resource name is truncated to an empty string. Pass through the configured value
54    /// unmodified, including `0`, so the configured behavior is always honored.
55    pub fn with_max_resource_len(mut self, max_resource_len: usize) -> Self {
56        self.max_resource_len = max_resource_len;
57        self
58    }
59}
60
61#[async_trait]
62impl DecoderBuilder for OtlpDecoderConfiguration {
63    fn input_payload_type(&self) -> PayloadType {
64        PayloadType::Grpc
65    }
66
67    fn output_event_type(&self) -> EventType {
68        EventType::Trace
69    }
70
71    async fn build(&self, context: BuildContext) -> Result<Box<dyn Decoder + Send>, GenericError> {
72        let metrics = build_metrics(context.component_context());
73        let traces_translator = OtlpTracesTranslator::new(self.traces.clone(), self.max_resource_len);
74
75        Ok(Box::new(OtlpDecoder {
76            traces_translator,
77            metrics,
78        }))
79    }
80}
81
82impl MemoryBounds for OtlpDecoderConfiguration {
83    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
84        builder.minimum().with_single_value::<OtlpDecoder>("decoder struct");
85    }
86}
87
88/// OTLP decoder.
89pub struct OtlpDecoder {
90    traces_translator: OtlpTracesTranslator,
91    metrics: Metrics,
92}
93
94#[async_trait]
95impl Decoder for OtlpDecoder {
96    async fn run(self: Box<Self>, mut context: DecoderContext) -> Result<(), GenericError> {
97        let Self {
98            mut traces_translator,
99            metrics,
100        } = *self;
101        let mut health = context.take_health_handle();
102        health.mark_ready();
103
104        debug!("OTLP decoder started.");
105
106        // Set a buffer flush interval of 100ms to ensure we flush buffered events periodically.
107        let mut buffer_flush = interval(Duration::from_millis(100));
108        buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
109
110        let mut event_buffer_manager = EventBufferManager::default();
111
112        loop {
113            select! {
114                maybe_payload = context.payloads().next() => {
115                    let payload = match maybe_payload {
116                        Some(payload) => payload,
117                        None => {
118                            debug!("Payloads stream closed, shutting down decoder.");
119                            break;
120                        }
121                    };
122
123                    let grpc_payload = match payload.try_into_grpc_payload() {
124                        Some(grpc) => grpc,
125                        None => {
126                            warn!("Received non-gRPC payload in OTLP decoder. Dropping payload.");
127                            continue;
128                        }
129                    };
130
131                    match grpc_payload.service_path() {
132                        path if path == &*OTLP_TRACES_GRPC_SERVICE_PATH => {
133                            let (_, _, _, body) = grpc_payload.into_parts();
134                            let request = match ExportTraceServiceRequest::decode(body.into_bytes()) {
135                                Ok(req) => req,
136                                Err(e) => {
137                                    error!(error = %e, "Failed to decode OTLP trace request.");
138                                    continue;
139                                }
140                            };
141
142                            for resource_spans in request.resource_spans {
143                                for trace_event in traces_translator.translate_spans(resource_spans, &metrics) {
144                                    if let Some(event_buffer) = event_buffer_manager.try_push(trace_event) {
145                                        if let Err(e) = context.dispatcher().dispatch(event_buffer).await {
146                                            error!(error = %e, "Failed to dispatch trace events.");
147                                        }
148                                    }
149                                }
150                            }
151                        }
152                        path if path == &*OTLP_METRICS_GRPC_SERVICE_PATH => {
153                            warn!("OTLP metrics decoding not yet implemented. Dropping metrics payload.");
154                        }
155                        path if path == &*OTLP_LOGS_GRPC_SERVICE_PATH => {
156                            warn!("OTLP logs decoding not yet implemented. Dropping logs payload.");
157                        }
158                        path => {
159                            warn!(service_path = path, "Received gRPC payload with unknown service path. Dropping payload.");
160                        }
161                    }
162                },
163                _ = buffer_flush.tick() => {
164                    if let Some(event_buffer) = event_buffer_manager.consume() {
165                        if let Err(e) = context.dispatcher().dispatch(event_buffer).await {
166                            error!(error = %e, "Failed to dispatch buffered trace events.");
167                        }
168                    }
169                },
170                _ = health.live() => continue,
171            }
172        }
173
174        if let Some(event_buffer) = event_buffer_manager.consume() {
175            if let Err(e) = context.dispatcher().dispatch(event_buffer).await {
176                error!(error = %e, "Failed to dispatch final trace events.");
177            }
178        }
179
180        debug!("OTLP decoder stopped.");
181
182        Ok(())
183    }
184}