1use std::net::ToSocketAddrs as _;
2use std::sync::Arc;
3use std::sync::LazyLock;
4use std::time::Duration;
5
6use agent_data_plane_config::domains;
7use async_trait::async_trait;
8use axum::body::Bytes;
9use otlp_protos::opentelemetry::proto::collector::logs::v1::ExportLogsServiceRequest;
10use otlp_protos::opentelemetry::proto::collector::metrics::v1::ExportMetricsServiceRequest;
11use otlp_protos::opentelemetry::proto::collector::trace::v1::ExportTraceServiceRequest;
12use otlp_protos::opentelemetry::proto::logs::v1::ResourceLogs as OtlpResourceLogs;
13use otlp_protos::opentelemetry::proto::metrics::v1::ResourceMetrics as OtlpResourceMetrics;
14use otlp_protos::opentelemetry::proto::trace::v1::ResourceSpans as OtlpResourceSpans;
15use prost::Message;
16use saluki_common::sync::shutdown::{ShutdownCoordinator, ShutdownHandle};
17use saluki_common::task::HandleExt as _;
18use saluki_context::tags::{SharedTagSet, TagSet};
19use saluki_context::ContextResolver;
20use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder};
21use saluki_core::topology::interconnect::BufferedDispatcher;
22use saluki_core::{
23 components::{
24 sources::{Source, SourceBuilder, SourceContext},
25 ComponentContext,
26 },
27 data_model::event::EventType,
28 topology::{EventsBuffer, OutputDefinition},
29};
30use saluki_env::WorkloadProvider;
31use saluki_error::ErrorContext as _;
32use saluki_error::{generic_error, GenericError};
33use saluki_io::net::ListenAddress;
34use stringtheory::MetaString;
35use tokio::pin;
36use tokio::select;
37use tokio::sync::mpsc;
38use tokio::time::{interval, MissedTickBehavior};
39use tracing::{debug, error};
40
41use crate::common::otlp::config::TracesConfig;
42use crate::common::otlp::{build_metrics, Metrics, OtlpHandler, OtlpServerBuilder};
43
44mod logs;
45mod metrics;
46mod resolver;
47use self::logs::translator::OtlpLogsTranslator;
48use self::metrics::translator::OtlpMetricsTranslator;
49use self::resolver::build_context_resolver;
50use crate::common::otlp::origin::OtlpOriginTagResolver;
51use crate::common::otlp::traces::translator::OtlpTracesTranslator;
52
53fn parse_configured_metric_tags(raw: &str) -> SharedTagSet {
57 let mut tags = TagSet::default();
58 for tag in raw.split(',') {
59 let tag = tag.trim();
60 if !tag.is_empty() {
61 tags.insert_tag(tag);
62 }
63 }
64 tags.into_shared()
65}
66
67#[derive(Default)]
69pub struct OtlpConfiguration {
70 default_hostname: MetaString,
71
72 otlp: domains::otlp::Domain,
74
75 workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
77}
78
79impl OtlpConfiguration {
80 pub fn from_configuration(otlp: &domains::otlp::Domain) -> Self {
82 Self {
83 default_hostname: MetaString::default(),
84 otlp: otlp.clone(),
85 workload_provider: None,
86 }
87 }
88
89 fn metrics_translator_config(&self) -> metrics::config::OtlpMetricsTranslatorConfig {
90 metrics::config::OtlpMetricsTranslatorConfig::default()
91 .with_summary_mode(self.otlp.metrics.summaries.mode)
92 .with_histogram_mode(self.otlp.metrics.histogram_mode)
93 .with_send_histogram_aggregations(self.otlp.metrics.send_histogram_aggregations)
94 .with_cumulative_monotonic_mode(self.otlp.metrics.sums.cumulative_monotonic_mode)
95 .with_initial_cumulative_monotonic_value(self.otlp.metrics.sums.initial_cumulative_monotonic_value)
96 .with_resource_attributes_as_tags(self.otlp.metrics.resource_attributes_as_tags)
97 }
98
99 pub fn with_default_hostname(mut self, hostname: impl Into<MetaString>) -> Self {
101 self.default_hostname = hostname.into();
102 self
103 }
104
105 pub fn with_workload_provider<W>(mut self, workload_provider: W) -> Self
111 where
112 W: WorkloadProvider + Send + Sync + 'static,
113 {
114 self.workload_provider = Some(Arc::new(workload_provider));
115 self
116 }
117}
118
119#[async_trait]
120impl SourceBuilder for OtlpConfiguration {
121 fn outputs(&self) -> &[OutputDefinition<EventType>] {
122 static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
123 vec![
124 OutputDefinition::named_output("metrics", EventType::Metric),
125 OutputDefinition::named_output("logs", EventType::Log),
126 OutputDefinition::named_output("traces", EventType::Trace),
127 ]
128 });
129
130 &OUTPUTS
131 }
132
133 async fn build(&self, context: ComponentContext) -> Result<Box<dyn Source + Send>, GenericError> {
134 if !self.otlp.receiver.metrics_enabled && !self.otlp.receiver.logs_enabled && !self.otlp.traces.enabled {
135 return Err(generic_error!(
136 "OTLP metrics, logs and traces support is disabled. Please enable at least one of them."
137 ));
138 }
139
140 let grpc_listen_str = format!(
141 "{}://{}",
142 self.otlp.receiver.grpc.transport, self.otlp.receiver.grpc.endpoint
143 );
144 let grpc_endpoint = ListenAddress::try_from(grpc_listen_str.as_str())
145 .map_err(|e| generic_error!("Invalid gRPC endpoint address '{}': {}", grpc_listen_str, e))?;
146
147 if !matches!(grpc_endpoint, ListenAddress::Tcp(_)) {
149 return Err(generic_error!("Only 'tcp' transport is supported for OTLP gRPC"));
150 }
151
152 let http_endpoint_str = &self.otlp.receiver.http.endpoint;
153 let http_socket_addr = http_endpoint_str
154 .to_socket_addrs()
155 .map_err(|e| generic_error!("Invalid HTTP endpoint address '{}': {}", http_endpoint_str, e))?
156 .next()
157 .ok_or_else(|| generic_error!("No addresses resolved for HTTP endpoint '{}'", http_endpoint_str))?;
158
159 let maybe_origin_tags_resolver = self.workload_provider.clone().map(OtlpOriginTagResolver::new);
160
161 let context_resolver =
162 build_context_resolver(&self.otlp.contexts, &context, maybe_origin_tags_resolver.clone())?;
163 let metrics_translator_config = self.metrics_translator_config();
164
165 let metric_tags = parse_configured_metric_tags(&self.otlp.metrics.tags);
166 let traces_config = TracesConfig {
167 enable_otlp_compute_top_level_by_span_kind: self.otlp.traces.enable_compute_top_level_by_span_kind,
168 ignore_missing_datadog_fields: self.otlp.traces.ignore_missing_datadog_fields,
169 ..Default::default()
170 };
171 let traces_translator = OtlpTracesTranslator::new(traces_config, self.otlp.traces.string_interner_size);
172 let grpc_max_recv_msg_size_bytes = self.otlp.receiver.grpc.max_recv_msg_size_mib as usize * 1024 * 1024;
173 let metrics = build_metrics(&context);
174
175 Ok(Box::new(Otlp {
176 context_resolver,
177 origin_tag_resolver: maybe_origin_tags_resolver,
178 grpc_endpoint,
179 http_endpoint: ListenAddress::Tcp(http_socket_addr),
180 grpc_max_recv_msg_size_bytes,
181 metrics_translator_config,
182 metric_tags,
183 default_hostname: self.default_hostname.clone(),
184 traces_translator,
185 metrics,
186 }))
187 }
188}
189
190impl MemoryBounds for OtlpConfiguration {
191 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
192 builder
193 .minimum()
194 .with_single_value::<Otlp>("source struct")
195 .with_single_value::<SourceHandler>("source handler");
196 }
197}
198
199pub struct Otlp {
200 context_resolver: ContextResolver,
201 origin_tag_resolver: Option<OtlpOriginTagResolver>,
202 grpc_endpoint: ListenAddress,
203 http_endpoint: ListenAddress,
204 grpc_max_recv_msg_size_bytes: usize,
205 metrics_translator_config: metrics::config::OtlpMetricsTranslatorConfig,
206 metric_tags: SharedTagSet,
207 default_hostname: MetaString,
208 traces_translator: OtlpTracesTranslator,
209 metrics: Metrics, }
211
212#[async_trait]
213impl Source for Otlp {
214 async fn run(self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
215 let Self {
216 context_resolver,
217 origin_tag_resolver,
218 grpc_endpoint,
219 http_endpoint,
220 grpc_max_recv_msg_size_bytes,
221 metrics_translator_config,
222 metric_tags,
223 default_hostname,
224 traces_translator,
225 metrics,
226 } = *self;
227
228 let global_shutdown = context.take_shutdown_handle();
229 pin!(global_shutdown);
230
231 let mut health = context.take_health_handle();
232 let memory_limiter = context.topology_context().memory_limiter();
233
234 let (tx, rx) = mpsc::channel::<OtlpResource>(1024);
236
237 let mut converter_shutdown_coordinator = ShutdownCoordinator::default();
238
239 let metrics_translator = OtlpMetricsTranslator::new(
240 metrics_translator_config,
241 default_hostname,
242 context_resolver,
243 metric_tags,
244 )?;
245
246 let thread_pool_handle = context.topology_context().global_thread_pool().clone();
247
248 thread_pool_handle.spawn_traced_named(
250 "otlp-resource-converter",
251 run_converter(
252 rx,
253 context.clone(),
254 origin_tag_resolver,
255 converter_shutdown_coordinator.register(),
256 metrics_translator,
257 metrics.clone(),
258 traces_translator,
259 ),
260 );
261
262 let handler = SourceHandler::new(tx);
263 let server_builder = OtlpServerBuilder::new(http_endpoint, grpc_endpoint, grpc_max_recv_msg_size_bytes);
264
265 let (http_shutdown, mut http_error) = server_builder
266 .build(handler, memory_limiter.clone(), thread_pool_handle, metrics)
267 .await?;
268
269 health.mark_ready();
270 debug!("OTLP source started.");
271
272 loop {
274 select! {
275 _ = &mut global_shutdown => {
276 debug!("Received shutdown signal.");
277 break
278 },
279 error = &mut http_error => {
280 if let Some(error) = error {
281 debug!(%error, "HTTP server error.");
282 }
283 break;
284 },
285 _ = health.live() => continue,
286 }
287 }
288
289 debug!("Stopping OTLP source...");
290
291 http_shutdown.shutdown();
292 converter_shutdown_coordinator.shutdown_and_wait().await;
293
294 debug!("OTLP source stopped.");
295
296 Ok(())
297 }
298}
299
300enum OtlpResource {
301 Metrics(OtlpResourceMetrics),
302 Logs(OtlpResourceLogs),
303 Traces(OtlpResourceSpans),
304}
305
306struct SourceHandler {
308 tx: mpsc::Sender<OtlpResource>,
309}
310
311impl SourceHandler {
312 fn new(tx: mpsc::Sender<OtlpResource>) -> Self {
313 Self { tx }
314 }
315}
316
317#[async_trait]
318impl OtlpHandler for SourceHandler {
319 async fn handle_metrics(&self, body: Bytes) -> Result<(), GenericError> {
320 let request =
321 ExportMetricsServiceRequest::decode(body).error_context("Failed to decode metrics export request.")?;
322
323 for resource_metrics in request.resource_metrics {
324 self.tx
325 .send(OtlpResource::Metrics(resource_metrics))
326 .await
327 .error_context("Failed to send resource metrics to converter: channel is closed.")?;
328 }
329 Ok(())
330 }
331
332 async fn handle_logs(&self, body: Bytes) -> Result<(), GenericError> {
333 let request = ExportLogsServiceRequest::decode(body).error_context("Failed to decode logs export request.")?;
334
335 for resource_logs in request.resource_logs {
336 self.tx
337 .send(OtlpResource::Logs(resource_logs))
338 .await
339 .error_context("Failed to send resource logs to converter: channel is closed.")?;
340 }
341 Ok(())
342 }
343
344 async fn handle_traces(&self, body: Bytes) -> Result<(), GenericError> {
345 let request =
346 ExportTraceServiceRequest::decode(body).error_context("Failed to decode trace export request.")?;
347
348 for resource_spans in request.resource_spans {
349 self.tx
350 .send(OtlpResource::Traces(resource_spans))
351 .await
352 .error_context("Failed to send resource spans to converter: channel is closed.")?;
353 }
354 Ok(())
355 }
356}
357
358async fn run_converter(
359 mut receiver: mpsc::Receiver<OtlpResource>, source_context: SourceContext,
360 origin_tag_resolver: Option<OtlpOriginTagResolver>, shutdown_handle: ShutdownHandle,
361 mut metrics_translator: OtlpMetricsTranslator, metrics: Metrics, mut traces_translator: OtlpTracesTranslator,
362) {
363 pin!(shutdown_handle);
364
365 debug!("OTLP resource converter task started.");
366
367 let mut buffer_flush = interval(Duration::from_millis(100));
370 buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
371
372 let mut metrics_dispatcher: Option<BufferedDispatcher<'_, EventsBuffer>> = None;
373 let mut logs_dispatcher: Option<BufferedDispatcher<'_, EventsBuffer>> = None;
374 let mut traces_dispatcher: Option<BufferedDispatcher<'_, EventsBuffer>> = None;
375
376 loop {
377 select! {
378 Some(otlp_resource) = receiver.recv() => {
379 match otlp_resource {
380 OtlpResource::Metrics(resource_metrics) => {
381 match metrics_translator.translate_metrics(resource_metrics, &metrics) {
382 Ok(events) => {
383 for event in events {
384 let dispatcher = metrics_dispatcher.get_or_insert_with(|| {
385 source_context
386 .dispatcher()
387 .buffered_named("metrics")
388 .expect("metrics output should exist")
389 });
390 if let Err(e) = dispatcher.push(event).await {
391 error!(error = %e, "Failed to dispatch metric event.");
392 }
393 }
394 }
395 Err(e) => {
396 error!(error = %e, "Failed to handle resource metrics.");
397 }
398 }
399 }
400 OtlpResource::Logs(resource_logs) => {
401 let translator = OtlpLogsTranslator::from_resource_logs(resource_logs, origin_tag_resolver.as_ref());
402 for log_event in translator {
403 metrics.logs_received().increment(1);
404
405 let dispatcher = logs_dispatcher.get_or_insert_with(|| {
406 source_context
407 .dispatcher()
408 .buffered_named("logs")
409 .expect("logs output should exist")
410 });
411 if let Err(e) = dispatcher.push(log_event).await {
412 error!(error = %e, "Failed to dispatch log event.");
413 }
414 }
415 }
416 OtlpResource::Traces(resource_spans) => {
417 for trace_event in traces_translator.translate_spans(resource_spans, &metrics) {
418 let dispatcher = traces_dispatcher.get_or_insert_with(|| {
419 source_context
420 .dispatcher()
421 .buffered_named("traces")
422 .expect("traces output should exist")
423 });
424 if let Err(e) = dispatcher.push(trace_event).await {
425 error!(error = %e, "Failed to dispatch trace event.");
426 }
427 }
428 }
429 }
430 },
431 _ = buffer_flush.tick() => {
432 if let Some(dispatcher) = metrics_dispatcher.take() {
433 if let Err(e) = dispatcher.flush().await {
434 error!(error = %e, "Failed to flush metric events.");
435 }
436 }
437 if let Some(dispatcher) = logs_dispatcher.take() {
438 if let Err(e) = dispatcher.flush().await {
439 error!(error = %e, "Failed to flush log events.");
440 }
441 }
442 if let Some(dispatcher) = traces_dispatcher.take() {
443 if let Err(e) = dispatcher.flush().await {
444 error!(error = %e, "Failed to flush trace events.");
445 }
446 }
447 },
448 _ = &mut shutdown_handle => {
449 debug!("Converter task received shutdown signal.");
450 break;
451 }
452 }
453 }
454
455 if let Some(dispatcher) = metrics_dispatcher.take() {
456 if let Err(e) = dispatcher.flush().await {
457 error!(error = %e, "Failed to flush metric events.");
458 }
459 }
460 if let Some(dispatcher) = logs_dispatcher.take() {
461 if let Err(e) = dispatcher.flush().await {
462 error!(error = %e, "Failed to flush log events.");
463 }
464 }
465 if let Some(dispatcher) = traces_dispatcher.take() {
466 if let Err(e) = dispatcher.flush().await {
467 error!(error = %e, "Failed to flush trace events.");
468 }
469 }
470
471 debug!("OTLP resource converter task stopped.");
472}
473
474#[cfg(test)]
475mod tests {
476 use agent_data_plane_config::domains;
477 use agent_data_plane_config::domains::otlp::{
478 CumulativeMonotonicMode, HistogramMode, InitialCumulativeMonotonicValue, SummaryMode,
479 };
480
481 use super::{parse_configured_metric_tags, OtlpConfiguration};
482
483 fn tags(raw: &str) -> Vec<String> {
484 parse_configured_metric_tags(raw)
485 .into_iter()
486 .map(|t| t.to_string())
487 .collect()
488 }
489
490 fn config_with_metrics(metrics: domains::otlp::Metrics) -> OtlpConfiguration {
491 let otlp = domains::otlp::Domain {
492 metrics,
493 ..Default::default()
494 };
495 OtlpConfiguration::from_configuration(&otlp)
496 }
497
498 #[test]
499 fn histogram_mode_flows_to_metrics_translator() {
500 for mode in [
501 HistogramMode::NoBuckets,
502 HistogramMode::Counters,
503 HistogramMode::Distributions,
504 ] {
505 let config = config_with_metrics(domains::otlp::Metrics {
506 histogram_mode: mode,
507 ..Default::default()
508 });
509
510 assert_eq!(config.metrics_translator_config().hist_mode, mode);
511 }
512 }
513
514 #[test]
515 fn summary_mode_flows_to_metrics_translator() {
516 for (mode, expected_quantiles) in [(SummaryMode::Gauges, true), (SummaryMode::NoQuantiles, false)] {
518 let config = config_with_metrics(domains::otlp::Metrics {
519 summaries: domains::otlp::Summaries { mode },
520 ..Default::default()
521 });
522
523 assert_eq!(config.metrics_translator_config().quantiles, expected_quantiles);
524 }
525 }
526
527 #[test]
528 fn histogram_aggregation_flows_to_metrics_translator() {
529 for send in [false, true] {
530 let config = config_with_metrics(domains::otlp::Metrics {
531 send_histogram_aggregations: send,
532 ..Default::default()
533 });
534
535 assert_eq!(config.metrics_translator_config().send_histogram_aggregations, send);
536 }
537 }
538
539 #[test]
540 fn nobuckets_with_histogram_aggregations_is_valid() {
541 let config = config_with_metrics(domains::otlp::Metrics {
542 histogram_mode: HistogramMode::NoBuckets,
543 send_histogram_aggregations: true,
544 ..Default::default()
545 });
546
547 assert!(config.metrics_translator_config().validate().is_ok());
548 }
549
550 #[test]
551 fn nobuckets_without_histogram_aggregations_is_invalid() {
552 let config = config_with_metrics(domains::otlp::Metrics {
554 histogram_mode: HistogramMode::NoBuckets,
555 send_histogram_aggregations: false,
556 ..Default::default()
557 });
558
559 assert!(config.metrics_translator_config().validate().is_err());
560 }
561
562 #[test]
563 fn cumulative_monotonic_sum_mode_defaults_to_delta_conversion() {
564 assert_eq!(
565 OtlpConfiguration::default()
566 .metrics_translator_config()
567 .cumulative_monotonic_mode,
568 CumulativeMonotonicMode::ToDelta
569 );
570 }
571
572 #[test]
573 fn cumulative_monotonic_mode_flows_to_metrics_translator() {
574 for mode in [CumulativeMonotonicMode::ToDelta, CumulativeMonotonicMode::RawValue] {
575 let config = config_with_metrics(domains::otlp::Metrics {
576 sums: domains::otlp::Sums {
577 cumulative_monotonic_mode: mode,
578 ..Default::default()
579 },
580 ..Default::default()
581 });
582
583 assert_eq!(config.metrics_translator_config().cumulative_monotonic_mode, mode);
584 }
585 }
586
587 #[test]
588 fn initial_cumulative_monotonic_value_flows_to_metrics_translator() {
589 for value in [
590 InitialCumulativeMonotonicValue::Auto,
591 InitialCumulativeMonotonicValue::Drop,
592 InitialCumulativeMonotonicValue::Keep,
593 ] {
594 let config = config_with_metrics(domains::otlp::Metrics {
595 sums: domains::otlp::Sums {
596 initial_cumulative_monotonic_value: value,
597 ..Default::default()
598 },
599 ..Default::default()
600 });
601
602 assert_eq!(
603 config.metrics_translator_config().initial_cumulative_monotonic_value,
604 value
605 );
606 }
607 }
608
609 #[test]
610 fn empty_configuration_yields_no_tags() {
611 assert!(tags("").is_empty());
612 }
613
614 #[test]
615 fn single_tag_is_parsed() {
616 assert_eq!(tags("env:prod"), vec!["env:prod".to_string()]);
617 }
618
619 #[test]
620 fn multiple_tags_are_split_on_comma() {
621 assert_eq!(
622 tags("env:prod,team:core"),
623 vec!["env:prod".to_string(), "team:core".to_string()]
624 );
625 }
626
627 #[test]
628 fn duplicate_tags_are_deduplicated() {
629 assert_eq!(tags("env:prod,env:prod"), vec!["env:prod".to_string()]);
630 }
631
632 #[test]
633 fn whitespace_around_commas_is_stripped() {
634 assert_eq!(
635 tags("env:prod, team:core"),
636 vec!["env:prod".to_string(), "team:core".to_string()]
637 );
638 }
639
640 #[test]
641 fn trailing_and_doubled_commas_produce_no_empty_tags() {
642 assert_eq!(tags("env:prod,"), vec!["env:prod".to_string()]);
643 assert_eq!(
644 tags("env:prod,,team:core"),
645 vec!["env:prod".to_string(), "team:core".to_string()]
646 );
647 }
648}