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::trace::v1::ResourceSpans as OtlpResourceSpans;
14use prost::Message;
15use saluki_common::collections::FastHashSet;
16use saluki_common::sync::shutdown::{ShutdownCoordinator, ShutdownHandle};
17use saluki_core::{
18    accounting::{MemoryBounds, MemoryBoundsBuilder},
19    components::{
20        sources::{Source, SourceBuilder, SourceContext},
21        BuildContext,
22    },
23    data_model::{
24        event::{metric::context::ContextResolver, EventType},
25        tags::{SharedTagSet, TagSet},
26    },
27    runtime,
28    topology::{interconnect::BufferedDispatcher, EventsBuffer, OutputDefinition},
29};
30use saluki_env::WorkloadProvider;
31use saluki_error::ErrorContext as _;
32use saluki_error::{generic_error, GenericError};
33use saluki_io::net::{server::http::Http2Config, 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::{
42    build_metrics, resolve_grpc_http2_config, CorsConfiguration, Metrics, OtlpHandler, OtlpServerConfiguration,
43    OtlpTlsConfiguration,
44};
45
46mod logs;
47mod metrics;
48mod resolver;
49use self::logs::translator::OtlpLogsTranslator;
50use self::metrics::translator::OtlpMetricsTranslator;
51use self::resolver::build_context_resolver;
52use crate::common::otlp::origin::OtlpOriginTagResolver;
53use crate::common::otlp::traces::translator::OtlpTracesTranslator;
54
55/// Parses `otlp_config.metrics.tags` into a set of tags added to every emitted metric.
56///
57/// The value is a comma-separated list. An empty configuration yields no tags.
58fn parse_configured_metric_tags(raw: &str) -> SharedTagSet {
59    let mut tags = TagSet::default();
60    for tag in raw.split(',') {
61        let tag = tag.trim();
62        if !tag.is_empty() {
63            tags.insert_tag(tag);
64        }
65    }
66    tags.into_shared()
67}
68
69/// Builds component-owned CORS settings from the resolved configuration model.
70fn cors_configuration(cors: &domains::otlp::Cors) -> CorsConfiguration {
71    CorsConfiguration {
72        allowed_origins: cors.allowed_origins.clone(),
73        allowed_headers: cors.allowed_headers.clone(),
74        exposed_headers: cors.exposed_headers.clone(),
75        max_age: cors.max_age,
76    }
77}
78
79/// Builds an `OtlpTlsConfiguration` from resolved TLS settings, if TLS is enabled.
80///
81/// TLS is enabled when both `cert_file` and `key_file` are non-empty. When `ca_file` is also non-empty, the server
82/// requests client certificates and verifies them against the CA certificates in that file, but does not require a
83/// client certificate (optional verification).
84///
85/// # Errors
86///
87/// Returns an error if any TLS field is set without the others required to form a valid TLS configuration. Both
88/// `cert_file` and `key_file` must be provided together to enable TLS, and `ca_file` must not be set without them.
89/// Setting only a subset is treated as a configuration error rather than silently downgrading to plaintext.
90fn build_tls_config(tls: &domains::otlp::Tls) -> Result<Option<OtlpTlsConfiguration>, GenericError> {
91    match (tls.cert_file.is_empty(), tls.key_file.is_empty()) {
92        (true, true) => {
93            if !tls.ca_file.is_empty() {
94                Err(generic_error!(
95                    "OTLP receiver TLS `ca_file` is set but `cert_file` and `key_file` are empty. All three must \
96                     be provided together, or `ca_file` must be omitted when TLS is disabled."
97                ))
98            } else {
99                Ok(None)
100            }
101        }
102        (false, false) => {
103            let mut config = OtlpTlsConfiguration::new(tls.cert_file.clone().into(), tls.key_file.clone().into());
104            if !tls.ca_file.is_empty() {
105                config = config.with_ca_file(tls.ca_file.clone().into());
106            }
107            Ok(Some(config))
108        }
109        (true, false) => Err(generic_error!(
110            "OTLP receiver TLS `key_file` is set but `cert_file` is empty. Both must be provided to enable TLS."
111        )),
112        (false, true) => Err(generic_error!(
113            "OTLP receiver TLS `cert_file` is set but `key_file` is empty. Both must be provided to enable TLS."
114        )),
115    }
116}
117
118/// Applies resolved static tags using replacement semantics.
119fn apply_static_metric_tags(otlp: &mut domains::otlp::Domain, static_tags: Vec<String>) {
120    if !static_tags.is_empty() {
121        otlp.metrics.tags = static_tags.join(",");
122    }
123}
124
125/// Configuration for the OTLP source.
126pub struct OtlpConfiguration {
127    default_hostname: MetaString,
128
129    /// Resolved OTLP domain slice.
130    otlp: domains::otlp::Domain,
131
132    /// Maximum length of a span's resource name, in bytes.
133    ///
134    /// Defaults to `usize::MAX`, meaning resource names are not truncated.
135    max_resource_len: usize,
136
137    /// Workload provider to utilize for origin detection/enrichment.
138    workload_provider: Arc<dyn WorkloadProvider + Send + Sync>,
139}
140
141impl OtlpConfiguration {
142    /// Creates a new `OtlpConfiguration` from the resolved OTLP configuration and workload provider.
143    ///
144    pub fn from_configuration<W>(otlp: &domains::otlp::Domain, workload_provider: W) -> Self
145    where
146        W: WorkloadProvider + Send + Sync + 'static,
147    {
148        Self {
149            default_hostname: MetaString::default(),
150            otlp: otlp.clone(),
151            max_resource_len: usize::MAX,
152            workload_provider: Arc::new(workload_provider),
153        }
154    }
155
156    /// Replaces the configured metric tags when static tags are required.
157    pub fn with_static_metric_tags(mut self, static_tags: Vec<String>) -> Self {
158        apply_static_metric_tags(&mut self.otlp, static_tags);
159        self
160    }
161
162    fn metrics_translator_config(&self) -> metrics::config::OtlpMetricsTranslatorConfig {
163        let mut config = metrics::config::OtlpMetricsTranslatorConfig::default()
164            .with_summary_mode(self.otlp.metrics.summaries.mode)
165            .with_histogram_mode(self.otlp.metrics.histogram_mode)
166            .with_send_histogram_aggregations(self.otlp.metrics.send_histogram_aggregations)
167            .with_cumulative_monotonic_mode(self.otlp.metrics.sums.cumulative_monotonic_mode)
168            .with_initial_cumulative_monotonic_value(self.otlp.metrics.sums.initial_cumulative_monotonic_value)
169            .with_resource_attributes_as_tags(self.otlp.metrics.resource_attributes_as_tags)
170            .with_instrumentation_scope_metadata_as_tags(self.otlp.metrics.instrumentation_scope_metadata_as_tags)
171            .with_delta_ttl(self.otlp.metrics.delta_ttl);
172        config.tag_cardinality = self.otlp.metrics.tag_cardinality;
173        config
174    }
175
176    /// Sets the default hostname used when OTLP metrics do not carry a resource hostname.
177    pub fn with_default_hostname(mut self, hostname: impl Into<MetaString>) -> Self {
178        self.default_hostname = hostname.into();
179        self
180    }
181
182    /// Sets the maximum length of a span's resource name, in bytes.
183    ///
184    /// Resource names longer than this limit are truncated to the limit, on a UTF-8 boundary. Defaults to `usize::MAX`
185    /// (unbounded), which matches the source's behavior before resource name truncation was introduced.
186    ///
187    /// If set to `0`, every resource name is truncated to an empty string. Pass through the configured value
188    /// unmodified, including `0`, so the configured behavior is always honored.
189    pub fn with_max_resource_len(mut self, max_resource_len: usize) -> Self {
190        self.max_resource_len = max_resource_len;
191        self
192    }
193}
194
195#[async_trait]
196impl SourceBuilder for OtlpConfiguration {
197    fn outputs(&self) -> &[OutputDefinition<EventType>] {
198        static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
199            vec![
200                OutputDefinition::named_output("metrics", EventType::Metric),
201                OutputDefinition::named_output("logs", EventType::Log),
202                OutputDefinition::named_output("traces", EventType::Trace),
203            ]
204        });
205
206        &OUTPUTS
207    }
208
209    async fn build(&self, context: BuildContext) -> Result<Box<dyn Source + Send>, GenericError> {
210        if !self.otlp.receiver.metrics_enabled && !self.otlp.receiver.logs_enabled && !self.otlp.traces.enabled {
211            return Err(generic_error!(
212                "OTLP metrics, logs and traces support is disabled. Please enable at least one of them."
213            ));
214        }
215
216        let grpc_listen_str = format!(
217            "{}://{}",
218            self.otlp.receiver.grpc.transport.as_str(),
219            self.otlp.receiver.grpc.endpoint
220        );
221        let grpc_endpoint = ListenAddress::try_from(grpc_listen_str.as_str())
222            .map_err(|e| generic_error!("Invalid gRPC endpoint address '{}': {}", grpc_listen_str, e))?;
223
224        let http_endpoint_str = &self.otlp.receiver.http.endpoint;
225        let http_socket_addr = http_endpoint_str
226            .to_socket_addrs()
227            .map_err(|e| generic_error!("Invalid HTTP endpoint address '{}': {}", http_endpoint_str, e))?
228            .next()
229            .ok_or_else(|| generic_error!("No addresses resolved for HTTP endpoint '{}'", http_endpoint_str))?;
230
231        let origin_tag_resolver = OtlpOriginTagResolver::new(Arc::clone(&self.workload_provider));
232
233        // Metrics resolve their full OTLP entity list at the resource boundary. Keep the context resolver free of an
234        // origin resolver so it cannot apply the legacy RawOrigin-only lookup a second time. Logs retain that resolver.
235        let context_resolver = build_context_resolver(&self.otlp.contexts, context.component_context(), None)?;
236        let metrics_translator_config = self.metrics_translator_config();
237
238        let metric_tags = parse_configured_metric_tags(&self.otlp.metrics.tags);
239        let traces_translator = OtlpTracesTranslator::new(self.otlp.traces.clone(), self.max_resource_len);
240        let grpc_max_recv_msg_size_bytes = self.otlp.receiver.grpc.max_recv_msg_size_mib as usize * 1024 * 1024;
241        let grpc_http2_config = resolve_grpc_http2_config(
242            &self.otlp.receiver.grpc.keepalive,
243            self.otlp.receiver.grpc.max_concurrent_streams,
244        );
245        let http_max_request_body_size = self.otlp.receiver.http.max_request_body_size;
246        let cors = cors_configuration(&self.otlp.receiver.http.cors);
247        let http_tls_config = build_tls_config(&self.otlp.receiver.http.tls)?;
248        let grpc_tls_config = build_tls_config(&self.otlp.receiver.grpc.tls)?;
249        let metrics = build_metrics(context.component_context());
250        let translator_metrics =
251            metrics::telemetry::OtlpMetricsTranslatorMetrics::from_component_context(context.component_context());
252
253        Ok(Box::new(Otlp {
254            context_resolver,
255            origin_tag_resolver,
256            grpc_endpoint,
257            http_endpoint: ListenAddress::Tcp(http_socket_addr),
258            grpc_max_recv_msg_size_bytes,
259            grpc_http2_config,
260            http_max_request_body_size,
261            metrics_translator_config,
262            metric_tags,
263            default_hostname: self.default_hostname.clone(),
264            traces_translator,
265            cors,
266            http_tls_config,
267            grpc_tls_config,
268            metrics,
269            translator_metrics,
270        }))
271    }
272}
273
274impl MemoryBounds for OtlpConfiguration {
275    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
276        builder
277            .minimum()
278            .with_single_value::<Otlp>("source struct")
279            .with_single_value::<SourceHandler>("source handler");
280    }
281}
282
283pub struct Otlp {
284    context_resolver: ContextResolver,
285    origin_tag_resolver: OtlpOriginTagResolver,
286    grpc_endpoint: ListenAddress,
287    http_endpoint: ListenAddress,
288    grpc_max_recv_msg_size_bytes: usize,
289    grpc_http2_config: Http2Config,
290    http_max_request_body_size: u64,
291    metrics_translator_config: metrics::config::OtlpMetricsTranslatorConfig,
292    metric_tags: SharedTagSet,
293    default_hostname: MetaString,
294    traces_translator: OtlpTracesTranslator,
295    cors: CorsConfiguration,
296    http_tls_config: Option<OtlpTlsConfiguration>,
297    grpc_tls_config: Option<OtlpTlsConfiguration>,
298    metrics: Metrics, // Telemetry metrics, not DD native metrics.
299    translator_metrics: metrics::telemetry::OtlpMetricsTranslatorMetrics,
300}
301
302#[async_trait]
303impl Source for Otlp {
304    async fn run(self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
305        let Self {
306            context_resolver,
307            origin_tag_resolver,
308            grpc_endpoint,
309            http_endpoint,
310            grpc_max_recv_msg_size_bytes,
311            grpc_http2_config,
312            http_max_request_body_size,
313            metrics_translator_config,
314            metric_tags,
315            default_hostname,
316            traces_translator,
317            cors,
318            http_tls_config,
319            grpc_tls_config,
320            metrics,
321            translator_metrics,
322        } = *self;
323
324        let global_shutdown = context.take_shutdown_handle();
325        pin!(global_shutdown);
326
327        let mut health = context.take_health_handle();
328        let memory_limiter = context.topology_context().memory_limiter();
329
330        // Create the internal channel for decoupling the servers from the converter.
331        let (tx, rx) = mpsc::channel::<OtlpSignal>(1024);
332
333        let metrics_translator = OtlpMetricsTranslator::new(
334            metrics_translator_config,
335            default_hostname,
336            context_resolver,
337            origin_tag_resolver.clone(),
338            metric_tags,
339            translator_metrics,
340        )?;
341
342        // Build our gRPC and HTTP servers and spawn them.
343        let handler = SourceHandler::new(tx, metrics.clone());
344        let mut server_config =
345            OtlpServerConfiguration::new(http_endpoint, grpc_endpoint, grpc_max_recv_msg_size_bytes)
346                .with_cors(cors)
347                .with_grpc_http2_config(grpc_http2_config)
348                .with_http_max_request_body_size(http_max_request_body_size);
349
350        if let Some(tls) = http_tls_config {
351            server_config = server_config.with_http_tls(tls);
352        }
353        if let Some(tls) = grpc_tls_config {
354            server_config = server_config.with_grpc_tls(tls);
355        }
356
357        server_config
358            .build(
359                handler,
360                memory_limiter.clone(),
361                metrics.clone(),
362                context.topology_context().global_thread_pool(),
363            )
364            .await?;
365
366        // Run the converter task on the worker pool: translating OTLP resources is highly compute-bound.
367        let converter_context = context.clone();
368
369        let mut converter_shutdown_coordinator = ShutdownCoordinator::default();
370        let converter_shutdown = converter_shutdown_coordinator.register();
371
372        runtime::worker(
373            "resource_converter",
374            run_converter(
375                rx,
376                converter_context,
377                origin_tag_resolver,
378                converter_shutdown,
379                metrics_translator,
380                metrics,
381                traces_translator,
382            ),
383        )
384        .on_runtime(context.topology_context().global_thread_pool().clone())
385        .spawn();
386
387        health.mark_ready();
388        debug!("OTLP source started.");
389
390        // Wait for the global shutdown signal, then notify converter to shutdown.
391        loop {
392            select! {
393                _ = &mut global_shutdown => {
394                    debug!("Received shutdown signal.");
395                    break
396                },
397                _ = health.live() => continue,
398            }
399        }
400
401        debug!("Stopping OTLP source...");
402
403        converter_shutdown_coordinator.shutdown_and_wait().await;
404
405        debug!("OTLP source stopped.");
406
407        Ok(())
408    }
409}
410
411enum OtlpSignal {
412    Metrics(ExportMetricsServiceRequest),
413    Logs(OtlpResourceLogs),
414    Traces(OtlpResourceSpans),
415}
416
417/// Handler that decodes OTLP bytes and sends resources to the converter.
418struct SourceHandler {
419    tx: mpsc::Sender<OtlpSignal>,
420    metrics: Metrics,
421}
422
423impl SourceHandler {
424    fn new(tx: mpsc::Sender<OtlpSignal>, metrics: Metrics) -> Self {
425        Self { tx, metrics }
426    }
427}
428
429#[async_trait]
430impl OtlpHandler for SourceHandler {
431    async fn handle_metrics(&self, body: Bytes) -> Result<(), GenericError> {
432        let request = ExportMetricsServiceRequest::decode(body).map_err(|e| {
433            self.metrics.metrics_errors_decode().increment(1);
434            generic_error!("Failed to decode metrics export request: {}", e)
435        })?;
436
437        // Send the entire request as a single channel message so the converter processes it
438        // atomically. This preserves the request boundary for usage beacon emission without
439        // needing control markers or shared state across concurrent requests.
440        self.tx.send(OtlpSignal::Metrics(request)).await.map_err(|e| {
441            self.metrics.metrics_errors_channel().increment(1);
442            generic_error!("Failed to send metrics request to converter: channel is closed: {}", e)
443        })?;
444        Ok(())
445    }
446
447    async fn handle_logs(&self, body: Bytes) -> Result<(), GenericError> {
448        let request = ExportLogsServiceRequest::decode(body).error_context("Failed to decode logs export request.")?;
449
450        for resource_logs in request.resource_logs {
451            self.tx
452                .send(OtlpSignal::Logs(resource_logs))
453                .await
454                .error_context("Failed to send resource logs to converter: channel is closed.")?;
455        }
456        Ok(())
457    }
458
459    async fn handle_traces(&self, body: Bytes) -> Result<(), GenericError> {
460        let request =
461            ExportTraceServiceRequest::decode(body).error_context("Failed to decode trace export request.")?;
462
463        for resource_spans in request.resource_spans {
464            self.tx
465                .send(OtlpSignal::Traces(resource_spans))
466                .await
467                .error_context("Failed to send resource spans to converter: channel is closed.")?;
468        }
469        Ok(())
470    }
471}
472
473async fn run_converter(
474    mut receiver: mpsc::Receiver<OtlpSignal>, source_context: SourceContext,
475    origin_tag_resolver: OtlpOriginTagResolver, shutdown_handle: ShutdownHandle,
476    mut metrics_translator: OtlpMetricsTranslator, metrics: Metrics, mut traces_translator: OtlpTracesTranslator,
477) {
478    pin!(shutdown_handle);
479
480    debug!("OTLP resource converter task started.");
481
482    // Set a buffer flush interval of 100ms, which will ensure we always flush buffered events at least every 100ms if
483    // we're otherwise idle and not receiving packets from the client.
484    let mut buffer_flush = interval(Duration::from_millis(100));
485    buffer_flush.set_missed_tick_behavior(MissedTickBehavior::Delay);
486
487    let mut metrics_dispatcher: Option<BufferedDispatcher<'_, EventsBuffer>> = None;
488    let mut logs_dispatcher: Option<BufferedDispatcher<'_, EventsBuffer>> = None;
489    let mut traces_dispatcher: Option<BufferedDispatcher<'_, EventsBuffer>> = None;
490
491    loop {
492        select! {
493            Some(otlp_signal) = receiver.recv() => {
494                match otlp_signal {
495                    OtlpSignal::Metrics(request) => {
496                        let mut detected_languages = FastHashSet::default();
497
498                        for resource_metrics in request.resource_metrics {
499                            match metrics_translator.translate_metrics(resource_metrics, &metrics) {
500                                Ok((events, languages)) => {
501                                    detected_languages.extend(languages);
502                                    for event in events {
503                                        let dispatcher = metrics_dispatcher.get_or_insert_with(|| {
504                                            source_context
505                                                .dispatcher()
506                                                .buffered_named("metrics")
507                                                .expect("metrics output should exist")
508                                        });
509                                        if let Err(e) = dispatcher.push(event).await {
510                                            error!(error = %e, "Failed to dispatch metric event.");
511                                            metrics.metrics_errors_dispatch().increment(1);
512                                        }
513                                    }
514                                }
515                                Err(e) => {
516                                    error!(error = %e, "Failed to handle resource metrics.");
517                                }
518                            }
519                        }
520
521                        // Emit usage beacon metrics for the completed request.
522                        for event in metrics_translator.emit_usage_beacons(detected_languages) {
523                            let dispatcher = metrics_dispatcher.get_or_insert_with(|| {
524                                source_context
525                                    .dispatcher()
526                                    .buffered_named("metrics")
527                                    .expect("metrics output should exist")
528                            });
529                            if let Err(e) = dispatcher.push(event).await {
530                                error!(error = %e, "Failed to dispatch usage beacon metric event.");
531                            }
532                        }
533                    }
534                    OtlpSignal::Logs(resource_logs) => {
535                        let translator = OtlpLogsTranslator::from_resource_logs(resource_logs, &origin_tag_resolver);
536                        for log_event in translator {
537                            metrics.logs_received().increment(1);
538
539                            let dispatcher = logs_dispatcher.get_or_insert_with(|| {
540                                source_context
541                                    .dispatcher()
542                                    .buffered_named("logs")
543                                    .expect("logs output should exist")
544                            });
545                            if let Err(e) = dispatcher.push(log_event).await {
546                                error!(error = %e, "Failed to dispatch log event.");
547                            }
548                        }
549                    }
550                    OtlpSignal::Traces(resource_spans) => {
551                        for trace_event in traces_translator.translate_spans(resource_spans, &metrics) {
552                            let dispatcher = traces_dispatcher.get_or_insert_with(|| {
553                                source_context
554                                    .dispatcher()
555                                    .buffered_named("traces")
556                                    .expect("traces output should exist")
557                            });
558                            if let Err(e) = dispatcher.push(trace_event).await {
559                                error!(error = %e, "Failed to dispatch trace event.");
560                            }
561                        }
562                    }
563                }
564            },
565            _ = buffer_flush.tick() => {
566                if let Some(dispatcher) = metrics_dispatcher.take() {
567                    if let Err(e) = dispatcher.flush().await {
568                        error!(error = %e, "Failed to flush metric events.");
569                        metrics.metrics_errors_flush().increment(1);
570                    }
571                }
572                if let Some(dispatcher) = logs_dispatcher.take() {
573                    if let Err(e) = dispatcher.flush().await {
574                        error!(error = %e, "Failed to flush log events.");
575                    }
576                }
577                if let Some(dispatcher) = traces_dispatcher.take() {
578                    if let Err(e) = dispatcher.flush().await {
579                        error!(error = %e, "Failed to flush trace events.");
580                    }
581                }
582            },
583            _ = &mut shutdown_handle => {
584                debug!("Converter task received shutdown signal.");
585                break;
586            }
587        }
588    }
589
590    if let Some(dispatcher) = metrics_dispatcher.take() {
591        if let Err(e) = dispatcher.flush().await {
592            error!(error = %e, "Failed to flush metric events.");
593            metrics.metrics_errors_flush().increment(1);
594        }
595    }
596    if let Some(dispatcher) = logs_dispatcher.take() {
597        if let Err(e) = dispatcher.flush().await {
598            error!(error = %e, "Failed to flush log events.");
599        }
600    }
601    if let Some(dispatcher) = traces_dispatcher.take() {
602        if let Err(e) = dispatcher.flush().await {
603            error!(error = %e, "Failed to flush trace events.");
604        }
605    }
606
607    debug!("OTLP resource converter task stopped.");
608}
609
610#[cfg(test)]
611mod tests {
612    use std::time::Duration;
613
614    use agent_data_plane_config::domains;
615    use agent_data_plane_config::domains::otlp::{
616        CumulativeMonotonicMode, HistogramMode, InitialCumulativeMonotonicValue, SummaryMode,
617    };
618    use prost::Message;
619    use saluki_core::components::ComponentContext;
620    use saluki_metrics::test::TestRecorder;
621
622    use super::{apply_static_metric_tags, parse_configured_metric_tags, OtlpConfiguration};
623    use crate::common::otlp::{build_metrics, OtlpHandler};
624
625    fn tags(raw: &str) -> Vec<String> {
626        parse_configured_metric_tags(raw)
627            .into_iter()
628            .map(|t| t.to_string())
629            .collect()
630    }
631
632    fn config_with_metrics(metrics: domains::otlp::Metrics) -> OtlpConfiguration {
633        let otlp = domains::otlp::Domain {
634            metrics,
635            ..Default::default()
636        };
637        OtlpConfiguration::from_configuration(&otlp, saluki_env::workload::providers::NoopWorkloadProvider)
638    }
639
640    #[test]
641    fn empty_static_tags_preserve_explicit_otlp_metric_tags() {
642        let mut otlp = domains::otlp::Domain::default();
643        otlp.metrics.tags = "configured:true".to_string();
644
645        apply_static_metric_tags(&mut otlp, Vec::new());
646
647        assert_eq!(otlp.metrics.tags, "configured:true");
648    }
649
650    #[test]
651    fn static_metric_tags_replace_explicit_otlp_metric_tags() {
652        let mut otlp = domains::otlp::Domain::default();
653        otlp.metrics.tags = "configured:true".to_string();
654
655        apply_static_metric_tags(&mut otlp, vec!["provider_kind:autopilot".to_string()]);
656
657        assert_eq!(otlp.metrics.tags, "provider_kind:autopilot");
658    }
659
660    #[test]
661    fn histogram_mode_flows_to_metrics_translator() {
662        for mode in [
663            HistogramMode::NoBuckets,
664            HistogramMode::Counters,
665            HistogramMode::Distributions,
666        ] {
667            let config = config_with_metrics(domains::otlp::Metrics {
668                histogram_mode: mode,
669                ..Default::default()
670            });
671
672            assert_eq!(config.metrics_translator_config().hist_mode, mode);
673        }
674    }
675
676    #[test]
677    fn summary_mode_flows_to_metrics_translator() {
678        // `gauges` emits one gauge per quantile (quantiles on); `noquantiles` omits them.
679        for (mode, expected_quantiles) in [(SummaryMode::Gauges, true), (SummaryMode::NoQuantiles, false)] {
680            let config = config_with_metrics(domains::otlp::Metrics {
681                summaries: domains::otlp::Summaries { mode },
682                ..Default::default()
683            });
684
685            assert_eq!(config.metrics_translator_config().quantiles, expected_quantiles);
686        }
687    }
688
689    #[test]
690    fn histogram_aggregation_flows_to_metrics_translator() {
691        for send in [false, true] {
692            let config = config_with_metrics(domains::otlp::Metrics {
693                send_histogram_aggregations: send,
694                ..Default::default()
695            });
696
697            assert_eq!(config.metrics_translator_config().send_histogram_aggregations, send);
698        }
699    }
700
701    #[test]
702    fn nobuckets_with_histogram_aggregations_is_valid() {
703        let config = config_with_metrics(domains::otlp::Metrics {
704            histogram_mode: HistogramMode::NoBuckets,
705            send_histogram_aggregations: true,
706            ..Default::default()
707        });
708
709        assert!(config.metrics_translator_config().validate().is_ok());
710    }
711
712    #[test]
713    fn nobuckets_without_histogram_aggregations_is_invalid() {
714        // Match the Agent: `nobuckets` without aggregation metrics emits nothing and is invalid.
715        let config = config_with_metrics(domains::otlp::Metrics {
716            histogram_mode: HistogramMode::NoBuckets,
717            send_histogram_aggregations: false,
718            ..Default::default()
719        });
720
721        assert!(config.metrics_translator_config().validate().is_err());
722    }
723
724    #[test]
725    fn cumulative_monotonic_sum_mode_defaults_to_delta_conversion() {
726        assert_eq!(
727            config_with_metrics(domains::otlp::Metrics::default())
728                .metrics_translator_config()
729                .cumulative_monotonic_mode,
730            CumulativeMonotonicMode::ToDelta
731        );
732    }
733
734    #[test]
735    fn cumulative_monotonic_mode_flows_to_metrics_translator() {
736        for mode in [CumulativeMonotonicMode::ToDelta, CumulativeMonotonicMode::RawValue] {
737            let config = config_with_metrics(domains::otlp::Metrics {
738                sums: domains::otlp::Sums {
739                    cumulative_monotonic_mode: mode,
740                    ..Default::default()
741                },
742                ..Default::default()
743            });
744
745            assert_eq!(config.metrics_translator_config().cumulative_monotonic_mode, mode);
746        }
747    }
748
749    #[test]
750    fn initial_cumulative_monotonic_value_flows_to_metrics_translator() {
751        for value in [
752            InitialCumulativeMonotonicValue::Auto,
753            InitialCumulativeMonotonicValue::Drop,
754            InitialCumulativeMonotonicValue::Keep,
755        ] {
756            let config = config_with_metrics(domains::otlp::Metrics {
757                sums: domains::otlp::Sums {
758                    initial_cumulative_monotonic_value: value,
759                    ..Default::default()
760                },
761                ..Default::default()
762            });
763
764            assert_eq!(
765                config.metrics_translator_config().initial_cumulative_monotonic_value,
766                value
767            );
768        }
769    }
770
771    #[test]
772    fn delta_ttl_flows_to_metrics_translator() {
773        // An explicit TTL overrides the 3600s default end-to-end into the translator config.
774        let config = config_with_metrics(domains::otlp::Metrics {
775            delta_ttl: Duration::from_secs(7200),
776            ..Default::default()
777        });
778
779        assert_eq!(config.metrics_translator_config().delta_ttl, Duration::from_secs(7200));
780    }
781
782    #[test]
783    fn delta_ttl_defaults_to_3600s() {
784        assert_eq!(
785            config_with_metrics(domains::otlp::Metrics::default())
786                .metrics_translator_config()
787                .delta_ttl,
788            Duration::from_secs(3600)
789        );
790    }
791
792    #[test]
793    fn instrumentation_scope_metadata_as_tags_defaults_to_true() {
794        assert!(
795            config_with_metrics(domains::otlp::Metrics::default())
796                .metrics_translator_config()
797                .instrumentation_scope_metadata_as_tags
798        );
799    }
800
801    #[test]
802    fn instrumentation_scope_metadata_as_tags_flows_to_metrics_translator() {
803        let config = config_with_metrics(domains::otlp::Metrics {
804            instrumentation_scope_metadata_as_tags: false,
805            ..Default::default()
806        });
807
808        assert!(
809            !config
810                .metrics_translator_config()
811                .instrumentation_scope_metadata_as_tags
812        );
813    }
814
815    #[test]
816    fn empty_configuration_yields_no_tags() {
817        assert!(tags("").is_empty());
818    }
819
820    #[test]
821    fn single_tag_is_parsed() {
822        assert_eq!(tags("env:prod"), vec!["env:prod".to_string()]);
823    }
824
825    #[test]
826    fn multiple_tags_are_split_on_comma() {
827        assert_eq!(
828            tags("env:prod,team:core"),
829            vec!["env:prod".to_string(), "team:core".to_string()]
830        );
831    }
832
833    #[test]
834    fn duplicate_tags_are_deduplicated() {
835        assert_eq!(tags("env:prod,env:prod"), vec!["env:prod".to_string()]);
836    }
837
838    #[test]
839    fn whitespace_around_commas_is_stripped() {
840        assert_eq!(
841            tags("env:prod, team:core"),
842            vec!["env:prod".to_string(), "team:core".to_string()]
843        );
844    }
845
846    #[test]
847    fn trailing_and_doubled_commas_produce_no_empty_tags() {
848        assert_eq!(tags("env:prod,"), vec!["env:prod".to_string()]);
849        assert_eq!(
850            tags("env:prod,,team:core"),
851            vec!["env:prod".to_string(), "team:core".to_string()]
852        );
853    }
854
855    // -----------------------------------------------------------------------------------------------
856    // Self-telemetry: server-level decode and channel error counters.
857    // -----------------------------------------------------------------------------------------------
858
859    #[tokio::test]
860    async fn source_handler_increments_decode_error_on_malformed_body() {
861        let recorder = TestRecorder::default();
862        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
863
864        let metrics = build_metrics(&ComponentContext::test_source("otlp_test"));
865        let (tx, _rx) = tokio::sync::mpsc::channel::<super::OtlpSignal>(1);
866        let handler = super::SourceHandler::new(tx, metrics);
867
868        // Invalid protobuf bytes cause decode to fail.
869        let result = handler.handle_metrics(bytes::Bytes::from_static(b"not protobuf")).await;
870        assert!(result.is_err());
871
872        let tags: &[(&str, &str)] = &[
873            ("component_id", "otlp_test"),
874            ("component_type", "source"),
875            ("reason", "decode"),
876        ];
877        assert_eq!(recorder.counter(("component_errors_total", tags)), Some(1));
878    }
879
880    #[tokio::test]
881    async fn source_handler_increments_channel_error_on_closed_channel() {
882        let recorder = TestRecorder::default();
883        let _recorder_guard = metrics::set_default_local_recorder(&recorder);
884
885        let metrics = build_metrics(&ComponentContext::test_source("otlp_test"));
886        // Create a channel with no receiver, then drop the receiver so the send fails.
887        let (tx, rx) = tokio::sync::mpsc::channel::<super::OtlpSignal>(1);
888        drop(rx);
889        let handler = super::SourceHandler::new(tx, metrics);
890
891        // A valid (empty) request that decodes fine but can't be sent because the channel is closed.
892        let request = otlp_protos::opentelemetry::proto::collector::metrics::v1::ExportMetricsServiceRequest::default();
893        let body = bytes::Bytes::from(request.encode_to_vec());
894        let result = handler.handle_metrics(body).await;
895        assert!(result.is_err());
896
897        let tags: &[(&str, &str)] = &[
898            ("component_id", "otlp_test"),
899            ("component_type", "source"),
900            ("reason", "channel"),
901        ];
902        assert_eq!(recorder.counter(("component_errors_total", tags)), Some(1));
903
904        // The decode counter should not have been incremented.
905        let decode_tags: &[(&str, &str)] = &[
906            ("component_id", "otlp_test"),
907            ("component_type", "source"),
908            ("reason", "decode"),
909        ];
910        assert_eq!(recorder.counter(("component_errors_total", decode_tags)), Some(0));
911    }
912}