1use std::collections::HashMap;
2use std::num::NonZeroUsize;
3use std::sync::LazyLock;
4use std::time::Duration;
5
6use async_trait::async_trait;
7use regex::Regex;
8use saluki_core::cache::{Cache, CacheBuilder};
9use saluki_core::{
10 accounting::{MemoryBounds, MemoryBoundsBuilder},
11 components::{
12 transforms::{SynchronousTransform, SynchronousTransformBuilder},
13 BuildContext,
14 },
15 data_model::{
16 event::metric::context::{Context, ContextResolver, ContextResolverBuilder},
17 tags::{SharedTagSet, TagSet},
18 },
19 topology::EventsBuffer,
20};
21use saluki_error::{generic_error, ErrorContext, GenericError};
22use stringtheory::MetaString;
23
24const MATCH_TYPE_WILDCARD: &str = "wildcard";
25const MATCH_TYPE_REGEX: &str = "regex";
26
27static ALLOWED_WILDCARD_MATCH_PATTERN: LazyLock<Regex> =
28 LazyLock::new(|| Regex::new(r"^[a-zA-Z0-9\-_*.]+$").expect("Invalid regex in ALLOWED_WILDCARD_MATCH_PATTERN"));
29
30pub struct DogStatsDMapperConfiguration {
32 context_string_interner_bytes: NonZeroUsize,
34
35 cache_size: usize,
39
40 profiles: Vec<DogStatsDMapperProfile>,
42}
43
44pub struct DogStatsDMapperProfile {
46 pub name: String,
48
49 pub prefix: String,
51
52 pub mappings: Vec<DogStatsDMetricMapping>,
54}
55
56pub struct DogStatsDMetricMapping {
58 pub metric_match: String,
60
61 pub match_type: String,
63
64 pub name: String,
66
67 pub tags: HashMap<String, String>,
69}
70
71impl DogStatsDMapperConfiguration {
72 pub fn new(
74 context_string_interner_bytes: NonZeroUsize, cache_size: usize, profiles: Vec<DogStatsDMapperProfile>,
75 ) -> Self {
76 Self {
77 context_string_interner_bytes,
78 cache_size,
79 profiles,
80 }
81 }
82
83 fn build_mapper(&self, context: BuildContext) -> Result<MetricMapper, GenericError> {
84 let mut profiles = Vec::with_capacity(self.profiles.len());
85 for config_profile in &self.profiles {
86 if config_profile.name.is_empty() {
87 return Err(generic_error!("missing profile name"));
88 }
89 if config_profile.prefix.is_empty() {
90 return Err(generic_error!("missing prefix for profile: {}", config_profile.name));
91 }
92
93 let mut profile = MappingProfile {
94 prefix: config_profile.prefix.clone(),
95 mappings: Vec::with_capacity(config_profile.mappings.len()),
96 };
97
98 for (mapping_index, mapping) in config_profile.mappings.iter().enumerate() {
99 let match_type = match mapping.match_type.as_str() {
100 "" => MATCH_TYPE_WILDCARD,
102 MATCH_TYPE_WILDCARD => MATCH_TYPE_WILDCARD,
103 MATCH_TYPE_REGEX => MATCH_TYPE_REGEX,
104 unknown => {
105 return Err(generic_error!(
106 "profile: {}, mapping num {}: invalid match type `{}`, expected `wildcard` or `regex`",
107 config_profile.name,
108 mapping_index,
109 unknown,
110 ))
111 }
112 };
113 if mapping.name.is_empty() {
114 return Err(generic_error!(
115 "profile: {}, mapping num {}: name is required",
116 config_profile.name,
117 mapping_index
118 ));
119 }
120 if mapping.metric_match.is_empty() {
121 return Err(generic_error!(
122 "profile: {}, mapping num {}: match is required",
123 config_profile.name,
124 mapping_index
125 ));
126 }
127 let regex = build_regex(&mapping.metric_match, match_type)?;
128 profile.mappings.push(MetricMapping {
129 name: mapping.name.clone(),
130 tags: mapping.tags.clone(),
131 regex,
132 });
133 }
134 profiles.push(profile);
135 }
136
137 let context_resolver =
138 ContextResolverBuilder::from_name(format!("{}/dsd_mapper/primary", context.component_id()))
139 .expect("resolver name is not empty")
140 .with_interner_capacity_bytes(self.context_string_interner_bytes)
141 .with_idle_context_expiration(Duration::from_secs(30))
142 .build();
143
144 let cache = match NonZeroUsize::new(self.cache_size) {
145 Some(capacity) => Some(
146 CacheBuilder::from_identifier(format!("{}/dsd_mapper/result_cache", context.component_id()))?
147 .with_capacity(capacity)
148 .build(),
149 ),
150 None => None,
151 };
152
153 Ok(MetricMapper {
154 context_resolver,
155 profiles,
156 cache,
157 })
158 }
159}
160
161fn build_regex(match_re: &str, match_type: &str) -> Result<Regex, GenericError> {
162 let mut pattern = match_re.to_owned();
163 if match_type == MATCH_TYPE_WILDCARD {
164 if !ALLOWED_WILDCARD_MATCH_PATTERN.is_match(&pattern) {
166 return Err(generic_error!(
167 "invalid wildcard match pattern `{}`, it does not match allowed match regex `{}`",
168 pattern,
169 ALLOWED_WILDCARD_MATCH_PATTERN.as_str()
170 ));
171 }
172 if pattern.contains("**") {
173 return Err(generic_error!(
174 "invalid wildcard match pattern `{}`, it should not contain consecutive `*`",
175 pattern
176 ));
177 }
178 pattern = pattern.replace(".", "\\.");
179 pattern = pattern.replace("*", "([^.]*)");
180 }
181
182 let final_pattern = format!("^{}$", pattern);
183
184 Regex::new(&final_pattern).with_error_context(|| {
185 format!(
186 "Failed to compile regular expression `{}` for `{}` match type",
187 final_pattern, match_type
188 )
189 })
190}
191
192struct MappingProfile {
193 prefix: String,
194 mappings: Vec<MetricMapping>,
195}
196
197struct MetricMapping {
198 name: String,
199 tags: HashMap<String, String>,
200 regex: Regex,
201}
202
203#[derive(Clone)]
204struct CachedMapResult {
205 name: MetaString,
206 extra_tags: SharedTagSet,
207}
208
209struct MetricMapper {
210 profiles: Vec<MappingProfile>,
211 context_resolver: ContextResolver,
212 cache: Option<Cache<MetaString, Option<CachedMapResult>>>,
213}
214
215impl MetricMapper {
216 fn try_map(&mut self, context: &Context) -> Option<Context> {
217 let metric_name = context.name();
222 let tags = context.tags();
223 let origin_tags = context.origin_tags();
224 let host = context.host();
227
228 if let Some(cache) = &self.cache {
230 if let Some(cached) = cache.get(metric_name) {
231 return match cached {
232 None => None,
233 Some(result) => {
234 let mut merged_tags = tags.clone();
235 merged_tags.merge_shared(&result.extra_tags);
236
237 self.context_resolver.resolve_with_optional_host_and_origin_tags(
238 result.name.clone(),
239 host,
240 merged_tags,
241 origin_tags.clone(),
242 )
243 }
244 };
245 }
246 }
247
248 let mut new_name = String::new();
250 let mut expanded_tag_value = String::new();
251
252 for profile in &self.profiles {
253 if !metric_name.starts_with(&profile.prefix) && profile.prefix != "*" {
254 continue;
255 }
256
257 for mapping in &profile.mappings {
258 if let Some(captures) = mapping.regex.captures(metric_name) {
259 new_name.clear();
260 captures.expand(&mapping.name, &mut new_name);
261
262 let mut extra_tags = TagSet::with_capacity(mapping.tags.len());
263 for (tag_key, tag_value_expr) in &mapping.tags {
264 expanded_tag_value.clear();
265 expanded_tag_value.push_str(tag_key);
266 expanded_tag_value.push(':');
267 captures.expand(tag_value_expr, &mut expanded_tag_value);
268
269 extra_tags.insert_tag(expanded_tag_value.as_str());
270 }
271
272 let extra_tags = extra_tags.into_shared();
274
275 let mut merged_tags = tags.clone();
276 merged_tags.merge_shared(&extra_tags);
277
278 let resolved = self.context_resolver.resolve_with_optional_host_and_origin_tags(
279 new_name.as_str(),
280 host,
281 merged_tags,
282 origin_tags.clone(),
283 )?;
284
285 if let Some(cache) = &self.cache {
286 cache.insert(
287 metric_name.clone(),
288 Some(CachedMapResult {
289 name: resolved.name().clone(),
290 extra_tags,
291 }),
292 );
293 }
294 return Some(resolved);
295 }
296 }
297 }
298
299 if let Some(cache) = &self.cache {
301 cache.insert(metric_name.clone(), None);
302 }
303 None
304 }
305
306 #[cfg(test)]
307 fn cache_len(&self) -> Option<usize> {
308 self.cache.as_ref().map(|c| c.len())
309 }
310
311 #[cfg(test)]
312 fn cache_contains(&self, metric_name: &str) -> bool {
313 self.cache
314 .as_ref()
315 .is_some_and(|c| c.get(&MetaString::from(metric_name)).is_some())
316 }
317}
318
319#[async_trait]
320impl SynchronousTransformBuilder for DogStatsDMapperConfiguration {
321 async fn build(&self, context: BuildContext) -> Result<Box<dyn SynchronousTransform + Send>, GenericError> {
322 let metric_mapper = self.build_mapper(context)?;
323 Ok(Box::new(DogStatsDMapper { metric_mapper }))
324 }
325}
326
327impl MemoryBounds for DogStatsDMapperConfiguration {
328 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
329 let mut min = builder.minimum();
330 min
331 .with_single_value::<DogStatsDMapper>("component struct")
333 .with_fixed_amount("string interner", self.context_string_interner_bytes.get());
336
337 if self.cache_size > 0 {
339 min.with_array::<(MetaString, Option<CachedMapResult>)>("mapper result cache", self.cache_size);
340 }
341 }
342}
343
344pub struct DogStatsDMapper {
345 metric_mapper: MetricMapper,
346}
347
348impl SynchronousTransform for DogStatsDMapper {
349 fn transform_buffer(&mut self, event_buffer: &mut EventsBuffer) {
350 for event in event_buffer {
351 if let Some(metric) = event.try_as_metric_mut() {
352 if let Some(new_context) = self.metric_mapper.try_map(metric.context()) {
353 *metric.context_mut() = new_context;
354 }
355 }
356 }
357 }
358}
359
360#[cfg(test)]
361mod tests {
362
363 use std::collections::HashMap;
364 use std::num::NonZeroUsize;
365
366 use saluki_core::{
367 components::{transforms::SynchronousTransform, BuildContext},
368 data_model::event::{
369 metric::{
370 context::{Context, ContextResolverBuilder},
371 Metric,
372 },
373 Event,
374 },
375 topology::EventsBuffer,
376 };
377 use saluki_error::GenericError;
378
379 use super::{
380 DogStatsDMapper, DogStatsDMapperConfiguration, DogStatsDMapperProfile, DogStatsDMetricMapping, MetricMapper,
381 };
382
383 macro_rules! metric_mapping {
384 (@match_type) => {
385 String::new()
386 };
387 (@match_type $match_type:expr) => {
388 $match_type.to_string()
389 };
390 ({
391 "match": $metric_match:expr,
392 $("match_type": $match_type:expr,)?
393 "name": $name:expr
394 $(, "tags": { $($tag_name:literal: $tag_value:expr),* $(,)? })?
395 }) => {
396 DogStatsDMetricMapping {
397 metric_match: $metric_match.to_string(),
398 match_type: metric_mapping!(@match_type $($match_type)?),
399 name: $name.to_string(),
400 tags: HashMap::from([$($(($tag_name.to_string(), $tag_value.to_string())),*)?]),
401 }
402 };
403 }
404
405 macro_rules! mapper_profiles {
406 ([$({
407 "name": $name:expr,
408 "prefix": $prefix:expr,
409 "mappings": [$($mapping:tt),* $(,)?]
410 }),* $(,)?]) => {
411 vec![$(DogStatsDMapperProfile {
412 name: $name.to_string(),
413 prefix: $prefix.to_string(),
414 mappings: vec![$(metric_mapping!($mapping)),*],
415 }),*]
416 };
417 }
418
419 fn counter_metric(name: &'static str, tags: &[&'static str]) -> Metric {
420 let context = Context::from_static_parts(name, tags);
421 Metric::counter(context, 1.0)
422 }
423
424 fn mapper(profiles: Vec<DogStatsDMapperProfile>) -> Result<MetricMapper, GenericError> {
425 mapper_with_cache(profiles, 1000)
426 }
427
428 fn mapper_with_cache(
429 profiles: Vec<DogStatsDMapperProfile>, cache_size: usize,
430 ) -> Result<MetricMapper, GenericError> {
431 let config =
432 DogStatsDMapperConfiguration::new(NonZeroUsize::new(64 * 1024).expect("not zero"), cache_size, profiles);
433 config.build_mapper(BuildContext::test_transform("test_mapper"))
434 }
435
436 fn assert_tags(context: &Context, expected_tags: &[&str]) {
437 for tag in expected_tags {
438 assert!(context.tags().has_tag(tag), "missing tag: {}", tag);
439 }
440 assert_eq!(context.tags().len(), expected_tags.len(), "unexpected number of tags");
441 }
442
443 #[track_caller]
444 fn assert_tags_for_case(context: &Context, expected_tags: &[&str], case: &str, input: &str) {
445 for tag in expected_tags {
446 assert!(
447 context.tags().has_tag(tag),
448 "[{case}] input {input:?}: missing tag {tag:?}"
449 );
450 }
451 assert_eq!(
452 context.tags().len(),
453 expected_tags.len(),
454 "[{case}] input {input:?}: unexpected number of tags"
455 );
456 }
457
458 fn simple_mapping_profile() -> Vec<DogStatsDMapperProfile> {
459 mapper_profiles!([{
460 "name": "test",
461 "prefix": "test.",
462 "mappings": [
463 {
464 "match": "test.job.duration.*.*",
465 "name": "test.job.duration",
466 "tags": {
467 "job_type": "$1",
468 "job_name": "$2"
469 }
470 }
471 ]
472 }])
473 }
474
475 #[tokio::test]
476 async fn config_driven_mappings_produce_expected_output() {
477 struct MapperCase {
481 description: &'static str,
482 config: Vec<DogStatsDMapperProfile>,
483 #[allow(clippy::type_complexity)]
484 checks: Vec<(
485 &'static str,
486 &'static [&'static str],
487 Option<(&'static str, &'static [&'static str])>,
488 )>,
489 }
490
491 let cases = vec![
492 MapperCase {
493 description: "wildcard mappings with capture-group tags",
494 config: mapper_profiles!([{
495 "name": "test",
496 "prefix": "test.",
497 "mappings": [
498 { "match": "test.job.duration.*.*", "name": "test.job.duration", "tags": { "job_type": "$1", "job_name": "$2" } },
499 { "match": "test.job.size.*.*", "name": "test.job.size", "tags": { "foo": "$1", "bar": "$2" } }
500 ]
501 }]),
502 checks: vec![
503 (
504 "test.job.duration.my_job_type.my_job_name",
505 &[],
506 Some(("test.job.duration", &["job_type:my_job_type", "job_name:my_job_name"])),
507 ),
508 (
509 "test.job.size.my_job_type.my_job_name",
510 &[],
511 Some(("test.job.size", &["foo:my_job_type", "bar:my_job_name"])),
512 ),
513 ("test.job.size.not_match", &[], None),
514 ],
515 },
516 MapperCase {
517 description: "partial mapping, second mapping has no tags",
518 config: mapper_profiles!([{
519 "name": "test",
520 "prefix": "test.",
521 "mappings": [
522 { "match": "test.job.duration.*.*", "name": "test.job.duration", "tags": { "job_type": "$1" } },
523 { "match": "test.task.duration.*.*", "name": "test.task.duration" }
524 ]
525 }]),
526 checks: vec![
527 (
528 "test.job.duration.my_job_type.my_job_name",
529 &[],
530 Some(("test.job.duration", &["job_type:my_job_type"])),
531 ),
532 (
533 "test.task.duration.my_job_type.my_job_name",
534 &[],
535 Some(("test.task.duration", &[])),
536 ),
537 ],
538 },
539 MapperCase {
540 description: "regex expansion with ${n} syntax",
541 config: mapper_profiles!([{
542 "name": "test",
543 "prefix": "test.",
544 "mappings": [
545 { "match": "test.job.duration.*.*", "name": "test.job.duration", "tags": { "job_type": "${1}_x", "job_name": "${2}_y" } }
546 ]
547 }]),
548 checks: vec![(
549 "test.job.duration.my_job_type.my_job_name",
550 &[],
551 Some((
552 "test.job.duration",
553 &["job_type:my_job_type_x", "job_name:my_job_name_y"],
554 )),
555 )],
556 },
557 MapperCase {
558 description: "capture groups expanded into the metric name",
559 config: mapper_profiles!([{
560 "name": "test",
561 "prefix": "test.",
562 "mappings": [
563 { "match": "test.job.duration.*.*", "name": "test.hello.$2.$1", "tags": { "job_type": "$1", "job_name": "$2" } }
564 ]
565 }]),
566 checks: vec![(
567 "test.job.duration.my_job_type.my_job_name",
568 &[],
569 Some((
570 "test.hello.my_job_name.my_job_type",
571 &["job_type:my_job_type", "job_name:my_job_name"],
572 )),
573 )],
574 },
575 MapperCase {
576 description: "wildcard matches a segment before an underscore",
577 config: mapper_profiles!([{
578 "name": "test",
579 "prefix": "test.",
580 "mappings": [
581 { "match": "test.*_start", "name": "test.start", "tags": { "job": "$1" } }
582 ]
583 }]),
584 checks: vec![("test.my_job_start", &[], Some(("test.start", &["job:my_job"])))],
585 },
586 MapperCase {
587 description: "mappings without any tags",
588 config: mapper_profiles!([{
589 "name": "test",
590 "prefix": "test.",
591 "mappings": [
592 { "match": "test.my-worker.start", "name": "test.worker.start" },
593 { "match": "test.my-worker.stop.*", "name": "test.worker.stop" }
594 ]
595 }]),
596 checks: vec![
597 ("test.my-worker.start", &[], Some(("test.worker.start", &[]))),
598 ("test.my-worker.stop.worker-name", &[], Some(("test.worker.stop", &[]))),
599 ],
600 },
601 MapperCase {
602 description: "all allowed wildcard characters",
603 config: mapper_profiles!([{
604 "name": "test",
605 "prefix": "test.",
606 "mappings": [
607 { "match": "test.abcdefghijklmnopqrstuvwxyz_ABCDEFGHIJKLMNOPQRSTUVWXYZ-01234567.*", "name": "test.alphabet" }
608 ]
609 }]),
610 checks: vec![(
611 "test.abcdefghijklmnopqrstuvwxyz_ABCDEFGHIJKLMNOPQRSTUVWXYZ-01234567.123",
612 &[],
613 Some(("test.alphabet", &[])),
614 )],
615 },
616 MapperCase {
617 description: "regex match type",
618 config: mapper_profiles!([{
619 "name": "test",
620 "prefix": "test.",
621 "mappings": [
622 { "match": "test\\.job\\.duration\\.(.*)", "match_type": "regex", "name": "test.job.duration", "tags": { "job_name": "$1" } },
623 { "match": "test\\.task\\.duration\\.(.*)", "match_type": "regex", "name": "test.task.duration", "tags": { "task_name": "$1" } }
624 ]
625 }]),
626 checks: vec![
627 (
628 "test.job.duration.my.funky.job$name-abc/123",
629 &[],
630 Some(("test.job.duration", &["job_name:my.funky.job$name-abc/123"])),
631 ),
632 (
633 "test.task.duration.MY_task_name",
634 &[],
635 Some(("test.task.duration", &["task_name:MY_task_name"])),
636 ),
637 ],
638 },
639 MapperCase {
640 description: "complex regex match type",
641 config: mapper_profiles!([{
642 "name": "test",
643 "prefix": "test.",
644 "mappings": [
645 { "match": "test\\.job\\.([a-z][0-9]-\\w+)\\.(.*)", "match_type": "regex", "name": "test.job", "tags": { "job_type": "$1", "job_name": "$2" } }
646 ]
647 }]),
648 checks: vec![
649 (
650 "test.job.a5-foo.bar",
651 &[],
652 Some(("test.job", &["job_type:a5-foo", "job_name:bar"])),
653 ),
654 ("test.job.foo.bar-not-match", &[], None),
655 ],
656 },
657 MapperCase {
658 description: "multiple profiles matched by prefix",
659 config: mapper_profiles!([
660 {
661 "name": "test",
662 "prefix": "foo.",
663 "mappings": [ { "match": "foo.duration.*", "name": "foo.duration", "tags": { "name": "$1" } } ]
664 },
665 {
666 "name": "test",
667 "prefix": "bar.",
668 "mappings": [
669 { "match": "bar.count.*", "name": "bar.count", "tags": { "name": "$1" } },
670 { "match": "foo.duration2.*", "name": "foo.duration2", "tags": { "name": "$1" } }
671 ]
672 }
673 ]),
674 checks: vec![
675 (
676 "foo.duration.foo_name1",
677 &[],
678 Some(("foo.duration", &["name:foo_name1"])),
679 ),
680 ("foo.duration2.foo_name1", &[], None),
682 ("bar.count.bar_name1", &[], Some(("bar.count", &["name:bar_name1"]))),
683 ("z.not.mapped", &[], None),
684 ],
685 },
686 MapperCase {
687 description: "wildcard prefix matches any metric",
688 config: mapper_profiles!([{
689 "name": "test",
690 "prefix": "*",
691 "mappings": [ { "match": "foo.duration.*", "name": "foo.duration", "tags": { "name": "$1" } } ]
692 }]),
693 checks: vec![(
694 "foo.duration.foo_name1",
695 &[],
696 Some(("foo.duration", &["name:foo_name1"])),
697 )],
698 },
699 MapperCase {
700 description: "only the first matching wildcard-prefixed profile applies",
701 config: mapper_profiles!([
702 {
703 "name": "test",
704 "prefix": "*",
705 "mappings": [ { "match": "foo.duration.*", "name": "foo.duration", "tags": { "name1": "$1" } } ]
706 },
707 {
708 "name": "test",
709 "prefix": "*",
710 "mappings": [ { "match": "foo.duration.*", "name": "foo.duration", "tags": { "name2": "$1" } } ]
711 }
712 ]),
713 checks: vec![(
715 "foo.duration.foo_name",
716 &[],
717 Some(("foo.duration", &["name1:foo_name"])),
718 )],
719 },
720 MapperCase {
721 description: "only the first matching profile applies across differing prefixes",
722 config: mapper_profiles!([
723 {
724 "name": "test",
725 "prefix": "foo.",
726 "mappings": [ { "match": "foo.*.duration.*", "name": "foo.bar1.duration", "tags": { "bar": "$1", "foo": "$2" } } ]
727 },
728 {
729 "name": "test",
730 "prefix": "foo.bar.",
731 "mappings": [ { "match": "foo.bar.duration.*", "name": "foo.bar2.duration", "tags": { "foo_bar": "$1" } } ]
732 }
733 ]),
734 checks: vec![(
736 "foo.bar.duration.foo_name",
737 &[],
738 Some(("foo.bar1.duration", &["bar:bar", "foo:foo_name"])),
739 )],
740 },
741 MapperCase {
742 description: "regex expansion with (\\w+) groups",
743 config: mapper_profiles!([{
744 "name": "test",
745 "prefix": "test.",
746 "mappings": [
747 { "match": "test.user.(\\w+).action.(\\w+)", "match_type": "regex", "name": "test.user.action", "tags": { "user": "$1", "action": "$2" } }
748 ]
749 }]),
750 checks: vec![(
751 "test.user.john_doe.action.login",
752 &[],
753 Some(("test.user.action", &["user:john_doe", "action:login"])),
754 )],
755 },
756 MapperCase {
757 description: "existing metric tags are retained alongside mapped tags",
758 config: mapper_profiles!([{
759 "name": "test",
760 "prefix": "test.",
761 "mappings": [
762 { "match": "test.job.duration.*.*", "name": "test.job.duration.$2", "tags": { "job_type": "$1", "job_name": "$2" } }
763 ]
764 }]),
765 checks: vec![(
766 "test.job.duration.abc.def",
767 &["foo:bar", "baz"],
768 Some((
769 "test.job.duration.def",
770 &["foo:bar", "baz", "job_type:abc", "job_name:def"],
771 )),
772 )],
773 },
774 ];
775
776 for case in cases {
777 let mut mapper = mapper(case.config)
778 .unwrap_or_else(|e| panic!("[{}] config should parse and build: {e}", case.description));
779
780 for (input_name, input_tags, expected) in case.checks {
781 let metric = counter_metric(input_name, input_tags);
782 match (mapper.try_map(metric.context()), expected) {
783 (Some(context), Some((expected_name, expected_tags))) => {
784 assert_eq!(
785 context.name(),
786 expected_name,
787 "[{}] wrong mapped name for input {input_name:?}",
788 case.description
789 );
790 assert_tags_for_case(&context, expected_tags, case.description, input_name);
791 }
792 (None, None) => {}
793 (mapped, expected) => panic!(
794 "[{}] input {input_name:?}: expected remap={}, got remap={}",
795 case.description,
796 expected.is_some(),
797 mapped.is_some()
798 ),
799 }
800 }
801 }
802 }
803
804 #[test]
805 fn invalid_mapper_configurations_are_rejected() {
806 let cases: Vec<(&str, Vec<DogStatsDMapperProfile>, &str)> = vec![
809 (
811 "profile with an empty name",
812 mapper_profiles!([{ "name": "", "prefix": "test.", "mappings": [] }]),
813 "missing profile name",
814 ),
815 (
816 "profile with an empty prefix",
817 mapper_profiles!([{ "name": "test", "prefix": "", "mappings": [] }]),
818 "missing prefix for profile: test",
819 ),
820 (
821 "mapping with an empty match",
822 mapper_profiles!([{ "name": "test", "prefix": "test.", "mappings": [{ "match": "", "name": "test.mapped" }] }]),
823 "match is required",
824 ),
825 (
826 "mapping with an empty name",
827 mapper_profiles!([{ "name": "test", "prefix": "test.", "mappings": [{ "match": "test.job.duration.*.*", "name": "", "tags": { "job_type": "$1" } }] }]),
828 "name is required",
829 ),
830 (
831 "second mapping with an empty name",
832 mapper_profiles!([{
833 "name": "test",
834 "prefix": "test.",
835 "mappings": [
836 { "match": "test.valid", "name": "mapped" },
837 { "match": "test.invalid", "name": "" }
838 ]
839 }]),
840 "mapping num 1: name is required",
841 ),
842 (
844 "wildcard match with disallowed characters",
845 mapper_profiles!([{ "name": "test", "prefix": "test.", "mappings": [{ "match": "test.[]duration.*.*", "name": "test.job.duration" }] }]),
846 "does not match allowed match regex",
847 ),
848 (
849 "wildcard match anchored with a caret",
850 mapper_profiles!([{ "name": "test", "prefix": "test.", "mappings": [{ "match": "^test.invalid.duration.*.*", "name": "test.job.duration" }] }]),
851 "does not match allowed match regex",
852 ),
853 (
854 "wildcard match with consecutive wildcards",
855 mapper_profiles!([{ "name": "test", "prefix": "test.", "mappings": [{ "match": "test.invalid.duration.**", "name": "test.job.duration" }] }]),
856 "consecutive",
857 ),
858 (
859 "unknown match type",
860 mapper_profiles!([{ "name": "test", "prefix": "test.", "mappings": [{ "match": "test.invalid.duration", "match_type": "invalid", "name": "test.job.duration" }] }]),
861 "invalid match type",
862 ),
863 ];
864
865 for (description, config, expected_substring) in cases {
866 let err = mapper(config)
867 .err()
868 .unwrap_or_else(|| panic!("[{description}] configuration should be rejected"));
869 let message = err.to_string();
870 assert!(
871 message.contains(expected_substring),
872 "[{description}] error {message:?} should contain {expected_substring:?}"
873 );
874 }
875 }
876
877 #[tokio::test]
878 async fn transform_buffer_remaps_matching_metrics_and_passes_others_through() {
879 let mut transform = DogStatsDMapper {
882 metric_mapper: mapper(simple_mapping_profile()).expect("config should parse and build"),
883 };
884
885 let mut events = EventsBuffer::default();
886 assert!(events
887 .try_push(Event::Metric(counter_metric("test.job.duration.my_type.my_name", &[])))
888 .is_none());
889 assert!(events
890 .try_push(Event::Metric(counter_metric("unrelated.metric", &["keep:me"])))
891 .is_none());
892
893 transform.transform_buffer(&mut events);
894
895 let metrics: Vec<Metric> = events.into_iter().filter_map(Event::try_into_metric).collect();
896 assert_eq!(metrics.len(), 2);
897
898 assert_eq!(metrics[0].context().name(), "test.job.duration");
900 assert_tags(metrics[0].context(), &["job_type:my_type", "job_name:my_name"]);
901
902 assert_eq!(metrics[1].context().name(), "unrelated.metric");
904 assert_tags(metrics[1].context(), &["keep:me"]);
905 }
906
907 #[tokio::test]
908 async fn mapper_preserves_host_context_dimension() {
909 let profiles = mapper_profiles!([{
910 "name": "test",
911 "prefix": "test.",
912 "mappings": [
913 {
914 "match": "test.job.duration.*",
915 "name": "test.job.duration",
916 "tags": {
917 "job_name": "$1"
918 }
919 }
920 ]
921 }]);
922
923 let mut resolver = ContextResolverBuilder::for_tests().build();
924 let context_a = resolver
925 .resolve_with_host("test.job.duration.worker", "host-a", &[] as &[&str], None)
926 .expect("context should resolve");
927 let context_b = resolver
928 .resolve_with_host("test.job.duration.worker", "host-b", &[] as &[&str], None)
929 .expect("context should resolve");
930
931 let mut mapper = mapper(profiles).expect("should have built mapper");
932 let mapped_a = mapper.try_map(&context_a).expect("should have remapped");
933 let mapped_b = mapper.try_map(&context_b).expect("should have remapped");
934
935 assert_ne!(mapped_a, mapped_b);
936 assert_eq!(mapped_a.host(), Some("host-a"));
937 assert_eq!(mapped_b.host(), Some("host-b"));
938 assert_eq!(mapped_a.name(), "test.job.duration");
939 assert_tags(&mapped_a, &["job_name:worker"]);
940 assert_tags(&mapped_b, &["job_name:worker"]);
941 }
942
943 #[tokio::test]
944 async fn cache_hit_returns_same_result_as_miss() {
945 let mut mapper = mapper_with_cache(simple_mapping_profile(), 1000).expect("should have parsed mapping config");
946 assert_eq!(mapper.cache_len(), Some(0));
947
948 let metric = counter_metric("test.job.duration.my_type.my_name", &[]);
949 let first = mapper.try_map(metric.context()).expect("should have remapped");
950 assert_eq!(mapper.cache_len(), Some(1));
951
952 let metric = counter_metric("test.job.duration.my_type.my_name", &[]);
953 let second = mapper.try_map(metric.context()).expect("should have remapped");
954 assert_eq!(mapper.cache_len(), Some(1));
955
956 assert_eq!(first.name(), second.name());
957 assert_eq!(first.name(), "test.job.duration");
958 assert_tags(&first, &["job_type:my_type", "job_name:my_name"]);
959 assert_tags(&second, &["job_type:my_type", "job_name:my_name"]);
960 }
961
962 #[tokio::test]
963 async fn negative_results_are_cached() {
964 let mut mapper = mapper_with_cache(simple_mapping_profile(), 1000).expect("should have parsed mapping config");
965
966 let metric = counter_metric("unrelated.metric.name", &[]);
967 assert!(mapper.try_map(metric.context()).is_none());
968 assert_eq!(mapper.cache_len(), Some(1));
969
970 let metric = counter_metric("unrelated.metric.name", &[]);
971 assert!(mapper.try_map(metric.context()).is_none());
972 assert_eq!(mapper.cache_len(), Some(1));
973 }
974
975 #[tokio::test]
976 async fn cache_disabled_when_size_is_zero() {
977 let mut mapper = mapper_with_cache(simple_mapping_profile(), 0).expect("should have parsed mapping config");
978 assert_eq!(mapper.cache_len(), None);
979
980 let metric = counter_metric("test.job.duration.my_type.my_name", &[]);
981 let context = mapper.try_map(metric.context()).expect("should have remapped");
982 assert_eq!(context.name(), "test.job.duration");
983 assert_tags(&context, &["job_type:my_type", "job_name:my_name"]);
984
985 assert!(mapper
986 .try_map(counter_metric("unrelated.metric", &[]).context())
987 .is_none());
988 assert_eq!(mapper.cache_len(), None);
989 }
990
991 #[tokio::test]
992 async fn cache_evicts_older_entry_and_retains_newest_within_capacity() {
993 let mut mapper = mapper_with_cache(simple_mapping_profile(), 2).expect("should have parsed mapping config");
994
995 for suffix in ["a", "b", "c"] {
997 let name = format!("test.job.duration.t.{}", suffix);
998 let metric = counter_metric(Box::leak(name.into_boxed_str()), &[]);
999 mapper.try_map(metric.context()).expect("should have remapped");
1000 }
1001
1002 let a = mapper.cache_contains("test.job.duration.t.a");
1003 let b = mapper.cache_contains("test.job.duration.t.b");
1004 let c = mapper.cache_contains("test.job.duration.t.c");
1005
1006 assert!(
1008 mapper.cache_len().unwrap() <= 2,
1009 "cache should not exceed configured capacity (got {})",
1010 mapper.cache_len().unwrap()
1011 );
1012 assert!(!(a && b && c), "at least one older entry must have been evicted");
1014 assert!(c, "the most-recently-inserted metric name should survive eviction");
1016 assert!(
1018 !(a && b),
1019 "only one older entry may coexist with the newest entry at capacity 2"
1020 );
1021 }
1022
1023 #[tokio::test]
1024 async fn flood_of_identical_names_populates_single_cache_entry() {
1025 let mut profiles: Vec<DogStatsDMapperProfile> = (0..50)
1028 .map(|i| DogStatsDMapperProfile {
1029 name: format!("noise-{i}"),
1030 prefix: format!("noise{i}."),
1031 mappings: vec![DogStatsDMetricMapping {
1032 metric_match: format!("noise{i}.*"),
1033 match_type: String::new(),
1034 name: "noise.mapped".to_string(),
1035 tags: HashMap::new(),
1036 }],
1037 })
1038 .collect();
1039 profiles.push(DogStatsDMapperProfile {
1040 name: "real".to_string(),
1041 prefix: "real.".to_string(),
1042 mappings: vec![DogStatsDMetricMapping {
1043 metric_match: "real.metric.*".to_string(),
1044 match_type: String::new(),
1045 name: "real.mapped".to_string(),
1046 tags: [("x".to_string(), "$1".to_string())].into(),
1047 }],
1048 });
1049
1050 let mut mapper = mapper_with_cache(profiles, 16).expect("should have built mapper");
1051
1052 for _ in 0..10_000 {
1053 let metric = counter_metric("real.metric.flood", &[]);
1054 let context = mapper.try_map(metric.context()).expect("should have remapped");
1055 assert_eq!(context.name(), "real.mapped");
1056 }
1057
1058 assert_eq!(
1059 mapper.cache_len(),
1060 Some(1),
1061 "flood of identical names should populate exactly one cache entry"
1062 );
1063 }
1064}