saluki_components/relays/otlp/
mod.rs

1use std::sync::LazyLock;
2
3use agent_data_plane_config::domains;
4use async_trait::async_trait;
5use axum::body::Bytes;
6use saluki_common::buf::FrozenChunkedBytesBuffer;
7use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
8use saluki_core::components::relays::{Relay, RelayBuilder, RelayContext};
9use saluki_core::components::ComponentContext;
10use saluki_core::data_model::payload::{GrpcPayload, Payload, PayloadMetadata, PayloadType};
11use saluki_core::topology::OutputDefinition;
12use saluki_error::{ErrorContext as _, GenericError};
13use saluki_io::net::ListenAddress;
14use stringtheory::MetaString;
15use tokio::sync::mpsc;
16use tokio::{pin, select};
17use tracing::{debug, error};
18
19use crate::common::otlp::{
20    build_metrics, Metrics, OtlpHandler, OtlpServerBuilder, OTLP_LOGS_GRPC_SERVICE_PATH,
21    OTLP_METRICS_GRPC_SERVICE_PATH, OTLP_TRACES_GRPC_SERVICE_PATH,
22};
23
24/// Configuration for the OTLP relay.
25#[derive(Default)]
26pub struct OtlpRelayConfiguration {
27    receiver: domains::otlp::Receiver,
28}
29
30impl OtlpRelayConfiguration {
31    /// Creates relay configuration from typed OTLP receiver settings.
32    pub fn from_configuration(receiver: &domains::otlp::Receiver) -> Self {
33        Self {
34            receiver: receiver.clone(),
35        }
36    }
37
38    fn http_endpoint(&self) -> ListenAddress {
39        let address = format!("{}://{}", self.receiver.http.transport, self.receiver.http.endpoint);
40        ListenAddress::try_from(address).expect("valid HTTP endpoint")
41    }
42
43    fn grpc_endpoint(&self) -> ListenAddress {
44        let address = format!("{}://{}", self.receiver.grpc.transport, self.receiver.grpc.endpoint);
45        ListenAddress::try_from(address).expect("valid gRPC endpoint")
46    }
47
48    fn grpc_max_recv_msg_size_bytes(&self) -> usize {
49        (self.receiver.grpc.max_recv_msg_size_mib * 1024 * 1024) as usize
50    }
51}
52
53impl MemoryBounds for OtlpRelayConfiguration {
54    fn specify_bounds(&self, _builder: &mut MemoryBoundsBuilder) {}
55}
56
57#[async_trait]
58impl RelayBuilder for OtlpRelayConfiguration {
59    fn outputs(&self) -> &[OutputDefinition<PayloadType>] {
60        static OUTPUTS: LazyLock<Vec<OutputDefinition<PayloadType>>> = LazyLock::new(|| {
61            vec![
62                OutputDefinition::named_output("metrics", PayloadType::Grpc),
63                OutputDefinition::named_output("logs", PayloadType::Grpc),
64                OutputDefinition::named_output("traces", PayloadType::Grpc),
65            ]
66        });
67        &OUTPUTS
68    }
69
70    async fn build(&self, context: ComponentContext) -> Result<Box<dyn Relay + Send>, GenericError> {
71        Ok(Box::new(OtlpRelay {
72            http_endpoint: self.http_endpoint(),
73            grpc_endpoint: self.grpc_endpoint(),
74            grpc_max_recv_msg_size_bytes: self.grpc_max_recv_msg_size_bytes(),
75            metrics: build_metrics(&context),
76        }))
77    }
78}
79
80/// OTLP relay.
81///
82/// Receives OTLP metrics and logs via gRPC and HTTP, outputting payloads for downstream processing.
83pub struct OtlpRelay {
84    http_endpoint: ListenAddress,
85    grpc_endpoint: ListenAddress,
86    grpc_max_recv_msg_size_bytes: usize,
87    metrics: Metrics,
88}
89
90#[async_trait]
91impl Relay for OtlpRelay {
92    async fn run(self: Box<Self>, mut context: RelayContext) -> Result<(), GenericError> {
93        let Self {
94            http_endpoint,
95            grpc_endpoint,
96            grpc_max_recv_msg_size_bytes,
97            metrics,
98        } = *self;
99
100        let global_shutdown = context.take_shutdown_handle();
101        pin!(global_shutdown);
102
103        let mut health = context.take_health_handle();
104        let global_thread_pool = context.topology_context().global_thread_pool().clone();
105        let memory_limiter = context.topology_context().memory_limiter().clone();
106        let dispatcher = context.dispatcher();
107
108        let (payload_tx, mut payload_rx) = mpsc::channel(1024);
109
110        let handler = RelayHandler::new(payload_tx);
111        let server_builder = OtlpServerBuilder::new(
112            http_endpoint.clone(),
113            grpc_endpoint.clone(),
114            grpc_max_recv_msg_size_bytes,
115        );
116
117        let (http_shutdown, mut http_error) = server_builder
118            .build(handler, memory_limiter, global_thread_pool, metrics)
119            .await?;
120
121        health.mark_ready();
122        debug!(%http_endpoint, %grpc_endpoint, "OTLP relay started.");
123
124        loop {
125            select! {
126                _ = &mut global_shutdown => {
127                    debug!("Received shutdown signal.");
128                    break
129                },
130                error = &mut http_error => {
131                    if let Some(error) = error {
132                        debug!(%error, "HTTP server error.");
133                    }
134                    break;
135                },
136                Some(otlp_payload) = payload_rx.recv() => {
137                    let output_name = otlp_payload.signal_type.as_str();
138                    let payload = Payload::Grpc(otlp_payload.into_grpc_payload());
139                    if let Err(e) = dispatcher.dispatch_named(output_name, payload).await {
140                        error!(error = %e, output = output_name, "Failed to dispatch OTLP payload.");
141                    }
142                },
143                _ = health.live() => continue,
144            }
145        }
146
147        debug!("Stopping OTLP relay...");
148
149        http_shutdown.shutdown();
150
151        debug!("OTLP relay stopped.");
152
153        Ok(())
154    }
155}
156
157enum OtlpSignalType {
158    Metrics,
159    Logs,
160    Traces,
161}
162
163impl OtlpSignalType {
164    fn as_str(&self) -> &'static str {
165        match self {
166            OtlpSignalType::Metrics => "metrics",
167            OtlpSignalType::Logs => "logs",
168            OtlpSignalType::Traces => "traces",
169        }
170    }
171}
172
173struct OtlpPayload {
174    signal_type: OtlpSignalType,
175    data: Bytes,
176}
177
178impl OtlpPayload {
179    fn metrics(data: Bytes) -> Self {
180        Self {
181            signal_type: OtlpSignalType::Metrics,
182            data,
183        }
184    }
185
186    fn logs(data: Bytes) -> Self {
187        Self {
188            signal_type: OtlpSignalType::Logs,
189            data,
190        }
191    }
192
193    fn traces(data: Bytes) -> Self {
194        Self {
195            signal_type: OtlpSignalType::Traces,
196            data,
197        }
198    }
199
200    fn into_grpc_payload(self) -> GrpcPayload {
201        let service_path = match self.signal_type {
202            OtlpSignalType::Metrics => OTLP_METRICS_GRPC_SERVICE_PATH,
203            OtlpSignalType::Logs => OTLP_LOGS_GRPC_SERVICE_PATH,
204            OtlpSignalType::Traces => OTLP_TRACES_GRPC_SERVICE_PATH,
205        };
206
207        // We provide an empty endpoint because we want any consuming components to fill that in for themselves.
208        GrpcPayload::new(
209            PayloadMetadata::from_event_count(1),
210            MetaString::empty(),
211            service_path,
212            FrozenChunkedBytesBuffer::from(self.data),
213        )
214    }
215}
216
217/// Handler that forwards OTLP payloads to a channel for downstream processing.
218struct RelayHandler {
219    tx: mpsc::Sender<OtlpPayload>,
220}
221
222impl RelayHandler {
223    fn new(tx: mpsc::Sender<OtlpPayload>) -> Self {
224        Self { tx }
225    }
226}
227
228#[async_trait]
229impl OtlpHandler for RelayHandler {
230    async fn handle_metrics(&self, body: Bytes) -> Result<(), GenericError> {
231        self.tx
232            .send(OtlpPayload::metrics(body))
233            .await
234            .error_context("Failed to send OTLP metrics payload to relay dispatcher: channel closed.")
235    }
236
237    async fn handle_logs(&self, body: Bytes) -> Result<(), GenericError> {
238        self.tx
239            .send(OtlpPayload::logs(body))
240            .await
241            .error_context("Failed to send OTLP logs payload to relay dispatcher: channel closed.")
242    }
243
244    async fn handle_traces(&self, body: Bytes) -> Result<(), GenericError> {
245        self.tx
246            .send(OtlpPayload::traces(body))
247            .await
248            .error_context("Failed to send OTLP traces payload to relay dispatcher: channel closed.")
249    }
250}
251
252#[cfg(test)]
253mod tests {
254    use agent_data_plane_config::domains;
255
256    use super::OtlpRelayConfiguration;
257
258    fn relay(receiver: domains::otlp::Receiver) -> OtlpRelayConfiguration {
259        OtlpRelayConfiguration::from_configuration(&receiver)
260    }
261
262    #[test]
263    fn endpoints_combine_transport_and_address() {
264        let config = relay(domains::otlp::Receiver {
265            grpc: domains::otlp::GrpcReceiver {
266                endpoint: "0.0.0.0:4317".to_string(),
267                transport: "tcp".to_string(),
268                max_recv_msg_size_mib: 4,
269            },
270            http: domains::otlp::HttpReceiver {
271                endpoint: "0.0.0.0:4318".to_string(),
272                transport: "tcp".to_string(),
273            },
274            ..Default::default()
275        });
276
277        assert_eq!(config.grpc_endpoint().to_string(), "tcp://0.0.0.0:4317");
278        assert_eq!(config.http_endpoint().to_string(), "tcp://0.0.0.0:4318");
279    }
280
281    #[test]
282    fn grpc_max_recv_msg_size_converts_mib_to_bytes() {
283        let config = relay(domains::otlp::Receiver {
284            grpc: domains::otlp::GrpcReceiver {
285                max_recv_msg_size_mib: 8,
286                ..Default::default()
287            },
288            ..Default::default()
289        });
290
291        assert_eq!(config.grpc_max_recv_msg_size_bytes(), 8 * 1024 * 1024);
292    }
293}