saluki_components/transforms/dogstatsd_mapper/
mod.rs

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
30/// DogStatsD mapper transform.
31pub struct DogStatsDMapperConfiguration {
32    /// Total size of the string interner used for contexts, in bytes.
33    context_string_interner_bytes: NonZeroUsize,
34
35    /// Maximum number of mapped results to cache.
36    ///
37    /// When set to `0`, the cache is disabled.
38    cache_size: usize,
39
40    /// Metric mapping profiles.
41    profiles: Vec<DogStatsDMapperProfile>,
42}
43
44/// A DogStatsD metric mapping profile.
45pub struct DogStatsDMapperProfile {
46    /// Profile name used in configuration errors.
47    pub name: String,
48
49    /// Metric-name prefix this profile applies to.
50    pub prefix: String,
51
52    /// Metric mappings evaluated for this profile.
53    pub mappings: Vec<DogStatsDMetricMapping>,
54}
55
56/// A metric-name mapping in a [`DogStatsDMapperProfile`].
57pub struct DogStatsDMetricMapping {
58    /// Pattern matched against the incoming metric name.
59    pub metric_match: String,
60
61    /// Pattern type: `wildcard`, `regex`, or empty for the default wildcard behavior.
62    pub match_type: String,
63
64    /// Metric name emitted after a match.
65    pub name: String,
66
67    /// Tags emitted after expanding capture groups.
68    pub tags: HashMap<String, String>,
69}
70
71impl DogStatsDMapperConfiguration {
72    /// Creates a DogStatsD mapper configuration.
73    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                    // Default to wildcard when not set.
101                    "" => 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        // Check it against the allowed wildcard pattern
165        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        // TODO: We should really be able to immutably borrow both the incoming tag set and the cached extra tags and
218        // chain them together for our call into `resolve_with_origin_tags`, avoiding any allocations... but we need
219        // some supporting work on the `TagSet` side to make it possible.
220
221        let metric_name = context.name();
222        let tags = context.tags();
223        let origin_tags = context.origin_tags();
224        // TODO: If host-bearing remaps show measurable allocation overhead, preserve the context's underlying host
225        // representation through the resolver instead of rematerializing it from `&str`.
226        let host = context.host();
227
228        // See if we have a cached result for this metric name.
229        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        // Slow path: iterate profiles and run regexes.
249        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                    // Freeze the tags here so they can be shared / cached.
273                    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        // We also cache "negative" results -- no match for this metric in the configured profiles -- to save ourselves some work.
300        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            // Capture the size of the heap allocation when the component is built.
332            .with_single_value::<DogStatsDMapper>("component struct")
333            // We also allocate the backing storage for the string interner up front, which is used by our context
334            // resolver.
335            .with_fixed_amount("string interner", self.context_string_interner_bytes.get());
336
337        // Account for the per-name result cache when enabled.
338        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        // Each case builds one mapper from `config`, then checks a series of inputs against it. A check is
478        // `(input_name, input_tags, expected)`, where `expected` is `Some((mapped_name, mapped_tags))` when the metric
479        // should be remapped, or `None` when it must pass through unmapped.
480        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` only exists under the `bar.` prefix, so it can't be reached by a `foo.` metric.
681                    ("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                // The single expected tag (and exact tag count) proves the second profile's `name2` tag was not applied.
714                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                // The exact tag count proves the second profile's `foo_bar` tag was not applied.
735                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        // Each case is `(description, config, expected_error_substring)`. Empty required values reach the
807        // component's validation; missing fields are rejected at the configuration boundary.
808        let cases: Vec<(&str, Vec<DogStatsDMapperProfile>, &str)> = vec![
809            // Custom (present-but-empty) validation branches.
810            (
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            // Match compilation / type validation.
843            (
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        // Drives the public `SynchronousTransform::transform_buffer` entry point (every other test exercises the
880        // internal `try_map`). Matching metrics are remapped in place; non-matching metrics pass through untouched.
881        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        // The matching metric is remapped in place (order is preserved).
899        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        // The non-matching metric is left untouched.
903        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        // Insert three distinct metric names into a capacity-2 result cache, in order a, b, then c.
996        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        // The cache must respect its configured capacity...
1007        assert!(
1008            mapper.cache_len().unwrap() <= 2,
1009            "cache should not exceed configured capacity (got {})",
1010            mapper.cache_len().unwrap()
1011        );
1012        // ...eviction must actually have happened (three distinct names cannot all fit in a capacity-2 cache)...
1013        assert!(!(a && b && c), "at least one older entry must have been evicted");
1014        // ...the most-recently-inserted name ("c") must be the entry that survives eviction...
1015        assert!(c, "the most-recently-inserted metric name should survive eviction");
1016        // ...and with "c" retained at capacity 2, at most one of the two older names may remain.
1017        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        // Many profiles, only the last one matches the test metric. A flood of identical
1026        // names should be served from the cache after the first call.
1027        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}