saluki_components/decoders/otlp/
mod.rs1use 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
28pub struct OtlpDecoderConfiguration {
30 traces: domains::otlp::Traces,
32
33 max_resource_len: usize,
37}
38
39impl OtlpDecoderConfiguration {
40 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 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
88pub 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 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}