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
55fn 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
69fn 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
79fn 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
118fn 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
125pub struct OtlpConfiguration {
127 default_hostname: MetaString,
128
129 otlp: domains::otlp::Domain,
131
132 max_resource_len: usize,
136
137 workload_provider: Arc<dyn WorkloadProvider + Send + Sync>,
139}
140
141impl OtlpConfiguration {
142 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 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 pub fn with_default_hostname(mut self, hostname: impl Into<MetaString>) -> Self {
178 self.default_hostname = hostname.into();
179 self
180 }
181
182 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 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, 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 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 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 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 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
417struct 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 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 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 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 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 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 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 #[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 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 let (tx, rx) = tokio::sync::mpsc::channel::<super::OtlpSignal>(1);
888 drop(rx);
889 let handler = super::SourceHandler::new(tx, metrics);
890
891 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 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}