saluki_components/sources/otlp/
mod.rs

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
53/// Parses `otlp_config.metrics.tags` into a set of tags added to every emitted metric.
54///
55/// The value is a comma-separated list. An empty configuration yields no tags.
56fn 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/// Configuration for the OTLP source.
68#[derive(Default)]
69pub struct OtlpConfiguration {
70    default_hostname: MetaString,
71
72    /// Resolved OTLP domain slice.
73    otlp: domains::otlp::Domain,
74
75    /// Workload provider to utilize for origin detection/enrichment.
76    workload_provider: Option<Arc<dyn WorkloadProvider + Send + Sync>>,
77}
78
79impl OtlpConfiguration {
80    /// Creates a new `OtlpConfiguration` from the resolved OTLP configuration.
81    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    /// Sets the default hostname used when OTLP metrics do not carry a resource hostname.
100    pub fn with_default_hostname(mut self, hostname: impl Into<MetaString>) -> Self {
101        self.default_hostname = hostname.into();
102        self
103    }
104
105    /// Sets the workload provider to use for configuring origin detection/enrichment.
106    ///
107    /// A workload provider must be set otherwise origin detection/enrichment won't be enabled.
108    ///
109    /// Defaults to unset.
110    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        // Enforce the current limitation that we only support TCP for gRPC.
148        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, // Telemetry metrics, not DD native metrics.
210}
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        // Create the internal channel for decoupling the servers from the converter.
235        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        // Spawn the converter task. This task is shared by both servers.
249        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        // Wait for the global shutdown signal, then notify converter to shutdown.
273        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
306/// Handler that decodes OTLP bytes and sends resources to the converter.
307struct 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    // Set a buffer flush interval of 100ms, which will ensure we always flush buffered events at least every 100ms if
368    // we're otherwise idle and not receiving packets from the client.
369    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        // `gauges` emits one gauge per quantile (quantiles on); `noquantiles` omits them.
517        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        // Match the Agent: `nobuckets` without aggregation metrics emits nothing and is invalid.
553        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}