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