saluki_components/relays/otlp/
mod.rs1use 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#[derive(Default)]
26pub struct OtlpRelayConfiguration {
27 receiver: domains::otlp::Receiver,
28}
29
30impl OtlpRelayConfiguration {
31 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
80pub 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 GrpcPayload::new(
209 PayloadMetadata::from_event_count(1),
210 MetaString::empty(),
211 service_path,
212 FrozenChunkedBytesBuffer::from(self.data),
213 )
214 }
215}
216
217struct 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}