saluki_io/deser/codec/dogstatsd/
metric.rs

1use nom::{
2    branch::alt,
3    bytes::complete::{tag, take_while1},
4    combinator::{all_consuming, map, map_res},
5    error::{Error, ErrorKind},
6    number::complete::double,
7    sequence::{preceded, separated_pair, terminated},
8    IResult, Parser as _,
9};
10use saluki_core::data_model::{event::metric::*, origin::OriginTagCardinality, tags::RawTags};
11use tracing::{debug, warn};
12
13use super::{helpers::*, DogStatsDCodecConfiguration, NomParserError};
14
15enum MetricType {
16    Count,
17    Gauge,
18    Set,
19    Timer,
20    Histogram,
21    Distribution,
22}
23
24/// A DogStatsD metric packet.
25///
26/// See the [DogStatsD datagram format][datagram] reference for the wire format and the protocol versions that
27/// introduced each field.
28///
29/// [datagram]: https://docs.datadoghq.com/extend/dogstatsd/datagram_shell/?tab=metrics
30pub struct MetricPacket<'a> {
31    /// Name of the metric.
32    pub metric_name: &'a str,
33
34    /// Tags attached to the metric.
35    pub tags: RawTags<'a>,
36
37    /// The metric kind (counter, gauge, rate, etc.) and its sample points.
38    pub values: MetricValues,
39
40    /// Number of sample points represented by `values`.
41    pub num_points: u64,
42
43    /// Optional Unix timestamp for the sample, in seconds (protocol v1.3).
44    pub timestamp: Option<u64>,
45
46    /// Local Data attached to the metric, carried in the `c:` field (protocol v1.2, extended in v1.4).
47    ///
48    /// Identifies the workload that emitted the metric. Carries a container ID (`ci-<id>`) or, when unavailable, a
49    /// cgroup node inode (`in-<inode>`).
50    pub local_data: Option<&'a str>,
51
52    /// External Data attached to the metric, carried in the `e:` field (protocol v1.5).
53    ///
54    /// Used to convey a richer blob of workload identity data resolved by the receiver.
55    pub external_data: Option<&'a str>,
56
57    /// Cardinality hint for origin tag enrichment, carried in the `card:` field (protocol v1.6).
58    ///
59    /// Specifies which origin tags the receiver should attach to the metric.
60    pub cardinality: Option<OriginTagCardinality>,
61
62    /// Unit for this metric, if any.
63    pub unit: Option<&'static str>,
64}
65
66#[inline]
67pub fn parse_dogstatsd_metric<'a>(
68    input: &'a [u8], config: &DogStatsDCodecConfiguration,
69) -> IResult<&'a [u8], MetricPacket<'a>> {
70    // We always parse the metric name and value(s) first, where value is both the kind (counter, gauge, etc) and the
71    // actual value itself.
72    let metric_name_parser = if config.permissive {
73        permissive_metric_name
74    } else {
75        ascii_alphanum_and_seps
76    };
77    let (remaining, (metric_name, (metric_type, raw_metric_values))) =
78        separated_pair(metric_name_parser, tag(":"), raw_metric_values).parse(input)?;
79
80    // At this point, we may have some of this additional data, and if so, we also then would have a pipe separator at
81    // the very front, which we'd want to consume before going further.
82    //
83    // After that, we simply split the remaining bytes by the pipe separator, and then try and parse each chunk to see
84    // if it's any of the protocol extensions we know of.
85    let mut maybe_sample_rate = None;
86    let mut maybe_tags = None;
87    let mut maybe_local_data = None;
88    let mut maybe_timestamp = None;
89    let mut maybe_external_data = None;
90    let mut maybe_cardinality = None;
91
92    let remaining = if !remaining.is_empty() {
93        let (mut remaining, _) = tag("|")(remaining)?;
94
95        while let Some((chunk, tail)) = split_at_delimiter(remaining, b'|') {
96            if chunk.is_empty() {
97                break;
98            }
99
100            match chunk[0] {
101                // Sample rate: indicates client-side sampling of this metric which will need to be "reinflated" at some
102                // point downstream to calculate the true metric value.
103                b'@' => {
104                    let (_, sample_rate) =
105                        all_consuming(preceded(tag("@"), map_res(double, sample_rate(metric_name, config))))
106                            .parse(chunk)?;
107
108                    maybe_sample_rate = Some(sample_rate);
109                }
110                // Tags: additional tags to be added to the metric.
111                b'#' => {
112                    let (_, tags) = all_consuming(preceded(tag("#"), tags(config))).parse(chunk)?;
113                    maybe_tags = Some(tags);
114                }
115                // Local Data: client-provided data used for resolving the entity ID that this metric originated from.
116                b'c' if chunk.len() > 1 && chunk[1] == b':' && config.client_origin_detection => {
117                    let (_, local_data) = all_consuming(preceded(tag("c:"), local_data)).parse(chunk)?;
118                    maybe_local_data = Some(local_data);
119                }
120                // Timestamp: client-provided timestamp for the metric, relative to the Unix epoch, in seconds.
121                b'T' if config.timestamps => {
122                    let (_, timestamp) = all_consuming(preceded(tag("T"), unix_timestamp)).parse(chunk)?;
123                    maybe_timestamp = Some(timestamp);
124                }
125                // External Data: client-provided data used for resolving the entity ID that this metric originated from.
126                b'e' if chunk.len() > 1 && chunk[1] == b':' && config.client_origin_detection => {
127                    let (_, external_data) = all_consuming(preceded(tag("e:"), external_data)).parse(chunk)?;
128                    maybe_external_data = Some(external_data);
129                }
130                // Cardinality: client-provided cardinality for the metric.
131                b'c' if chunk.starts_with(CARDINALITY_PREFIX) && config.client_origin_detection => {
132                    let (_, cardinality) = cardinality(chunk)?;
133                    maybe_cardinality = cardinality;
134                }
135                _ => {
136                    // We don't know what this is, so we just skip it.
137                    //
138                    // TODO: Should we throw an error, warn, or be silently permissive?
139                }
140            }
141
142            remaining = tail;
143        }
144
145        // TODO: Similarly to the above comment, should having any remaining data here cause us to throw an error, warn,
146        // or be silently permissive?
147
148        remaining
149    } else {
150        remaining
151    };
152
153    // Capture the unit from the metric type before it is erased into a MetricValues variant.
154    // Only timing metrics carry an implicit unit; all other types have no unit.
155    let maybe_unit = if matches!(metric_type, MetricType::Timer) {
156        Some("millisecond")
157    } else {
158        None
159    };
160
161    let effective_sample_rate = if matches!(&metric_type, MetricType::Count) && maybe_timestamp.is_some() {
162        // Match the Datadog Agent no-aggregation pipeline: timestamped DogStatsD counts are forwarded as
163        // pre-aggregated points, so their sample rate is not used to reinflate the count value.
164        None
165    } else {
166        maybe_sample_rate
167    };
168
169    let (num_points, mut metric_values) =
170        metric_values_from_raw(raw_metric_values, metric_type, effective_sample_rate)?;
171
172    // If we got a timestamp, apply it to all metric values.
173    if let Some(timestamp) = maybe_timestamp {
174        metric_values.set_timestamp(timestamp);
175    }
176
177    let tags = maybe_tags.unwrap_or_else(RawTags::empty);
178
179    Ok((
180        remaining,
181        MetricPacket {
182            metric_name,
183            tags,
184            values: metric_values,
185            num_points,
186            timestamp: maybe_timestamp,
187            local_data: maybe_local_data,
188            external_data: maybe_external_data,
189            cardinality: maybe_cardinality,
190            unit: maybe_unit,
191        },
192    ))
193}
194
195#[inline]
196fn permissive_metric_name(input: &[u8]) -> IResult<&[u8], &str> {
197    // Essentially, any ASCII character that is printable and isn't `:` is allowed here.
198    let valid_char = |c: u8| c > 31 && c < 128 && c != b':';
199    map(take_while1(valid_char), |b| {
200        // SAFETY: We know the bytes in `b` can only be comprised of ASCII characters, which ensures that it's valid to
201        // interpret the bytes directly as UTF-8.
202        unsafe { std::str::from_utf8_unchecked(b) }
203    })
204    .parse(input)
205}
206
207#[inline]
208fn sample_rate<'a>(
209    metric_name: &'a str, config: &'a DogStatsDCodecConfiguration,
210) -> impl Fn(f64) -> Result<SampleRate, &'static str> + 'a {
211    let minimum_sample_rate = config.minimum_sample_rate;
212
213    move |mut raw_sample_rate| {
214        if raw_sample_rate < minimum_sample_rate {
215            raw_sample_rate = minimum_sample_rate;
216            warn!(
217                "Sample rate for metric '{}' is below minimum of {}. Clamping to minimum.",
218                metric_name, minimum_sample_rate
219            );
220        }
221        SampleRate::try_from(raw_sample_rate)
222    }
223}
224
225#[inline]
226fn raw_metric_values(input: &[u8]) -> IResult<&[u8], (MetricType, &[u8])> {
227    let (remaining, raw_values) = terminated(take_while1(|b| b != b'|'), tag("|")).parse(input)?;
228    let (remaining, raw_kind) = alt((tag("g"), tag("c"), tag("ms"), tag("h"), tag("s"), tag("d"))).parse(remaining)?;
229
230    // Make sure the raw value(s) are valid UTF-8 before we use them later on.
231    if raw_values.is_empty() || simdutf8::basic::from_utf8(raw_values).is_err() {
232        return Err(nom::Err::Error(Error::new(raw_values, ErrorKind::Verify)));
233    }
234
235    let metric_type = match raw_kind {
236        b"c" => MetricType::Count,
237        b"g" => MetricType::Gauge,
238        b"s" => MetricType::Set,
239        b"ms" => MetricType::Timer,
240        b"h" => MetricType::Histogram,
241        b"d" => MetricType::Distribution,
242        _ => unreachable!("should be constrained by alt parser"),
243    };
244
245    Ok((remaining, (metric_type, raw_values)))
246}
247
248#[inline]
249fn metric_values_from_raw(
250    input: &[u8], metric_type: MetricType, sample_rate: Option<SampleRate>,
251) -> Result<(u64, MetricValues), NomParserError<'_>> {
252    let mut num_points = 0;
253    let floats = FloatIter::new(input).inspect(|_| num_points += 1);
254
255    let values = match metric_type {
256        MetricType::Count => MetricValues::counter_sampled_fallible(floats, sample_rate)?,
257        MetricType::Gauge => MetricValues::gauge_fallible(floats)?,
258        MetricType::Set => {
259            num_points = 1;
260
261            // SAFETY: We've already checked above that `input` is valid UTF-8.
262            let value = unsafe { std::str::from_utf8_unchecked(input) };
263            MetricValues::set(value.to_string())
264        }
265        MetricType::Timer | MetricType::Histogram => MetricValues::histogram_sampled_fallible(floats, sample_rate)?,
266        MetricType::Distribution => MetricValues::distribution_sampled_fallible(floats, sample_rate)?,
267    };
268
269    Ok((num_points, values))
270}
271
272struct FloatIter<'a> {
273    raw_values: &'a [u8],
274}
275
276impl<'a> FloatIter<'a> {
277    fn new(raw_values: &'a [u8]) -> Self {
278        Self { raw_values }
279    }
280}
281
282impl<'a> Iterator for FloatIter<'a> {
283    type Item = Result<f64, NomParserError<'a>>;
284
285    fn next(&mut self) -> Option<Self::Item> {
286        loop {
287            if self.raw_values.is_empty() {
288                return None;
289            }
290
291            let (raw_value, tail) = split_at_delimiter(self.raw_values, b':')?;
292            self.raw_values = tail;
293
294            // SAFETY: The caller that creates `ValueIter` is responsible for ensuring that the entire byte slice is valid
295            // UTF-8.
296            let value_s = unsafe { std::str::from_utf8_unchecked(raw_value) };
297            match value_s.parse::<f64>() {
298                Ok(value) if value.is_finite() => return Some(Ok(value)),
299                Ok(_) => {
300                    debug!(value = value_s, "Dropping non-finite DogStatsD metric value.");
301                }
302                Err(_) => return Some(Err(nom::Err::Error(Error::new(raw_value, ErrorKind::Float)))),
303            }
304        }
305    }
306}
307
308#[cfg(test)]
309mod tests {
310    use proptest::{collection::vec as arb_vec, prelude::*};
311    use saluki_core::data_model::{
312        event::metric::{context::Context, *},
313        origin::OriginTagCardinality,
314        tags::{SharedTagSet, Tag},
315    };
316
317    use super::{parse_dogstatsd_metric, DogStatsDCodecConfiguration};
318
319    type OptionalNomResult<'input, T> = Result<Option<T>, nom::Err<nom::error::Error<&'input [u8]>>>;
320
321    fn parse_dsd_metric(input: &[u8]) -> OptionalNomResult<'_, Metric> {
322        let default_config = DogStatsDCodecConfiguration::default();
323        parse_dsd_metric_with_conf(input, &default_config)
324    }
325
326    fn parse_dsd_metric_with_conf<'input>(
327        input: &'input [u8], config: &DogStatsDCodecConfiguration,
328    ) -> OptionalNomResult<'input, Metric> {
329        let (remaining, packet) = parse_dogstatsd_metric(input, config)?;
330        assert!(remaining.is_empty());
331
332        let tags = packet.tags.iter().map(Tag::from).collect::<SharedTagSet>();
333        let context = Context::from_parts(packet.metric_name, tags);
334
335        Ok(Some(Metric::from_parts(
336            context,
337            packet.values,
338            MetricMetadata::default(),
339        )))
340    }
341
342    #[track_caller]
343    fn check_basic_metric_eq(expected: Metric, actual: Option<Metric>) -> Metric {
344        let actual = actual.expect("event should not have been None");
345        assert_eq!(expected.context(), actual.context());
346        assert_eq!(expected.values(), actual.values());
347        assert_eq!(expected.metadata(), actual.metadata());
348        actual
349    }
350
351    #[test]
352    fn basic_metric() {
353        let name = "my.counter";
354        let value = 1.0;
355        let raw = format!("{}:{}|c", name, value);
356        let expected = Metric::counter(name, value);
357        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
358        check_basic_metric_eq(expected, actual);
359
360        let name = "my.gauge";
361        let value = 2.0;
362        let raw = format!("{}:{}|g", name, value);
363        let expected = Metric::gauge(name, value);
364        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
365        check_basic_metric_eq(expected, actual);
366
367        // Special case where we check this for both timers and histograms since we treat them both the same when
368        // parsing.
369        let name = "my.timer_or_histogram";
370        let value = 3.0;
371        for kind in &["ms", "h"] {
372            let raw = format!("{}:{}|{}", name, value, kind);
373            let expected = Metric::histogram(name, value);
374            let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
375            check_basic_metric_eq(expected, actual);
376        }
377
378        let distribution_name = "my.distribution";
379        let distribution_value = 3.0;
380        let distribution_raw = format!("{}:{}|d", distribution_name, distribution_value);
381        let distribution_expected = Metric::distribution(distribution_name, distribution_value);
382        let distribution_actual = parse_dsd_metric(distribution_raw.as_bytes()).expect("should not fail to parse");
383        check_basic_metric_eq(distribution_expected, distribution_actual);
384
385        let set_name = "my.set";
386        let set_value = "value";
387        let set_raw = format!("{}:{}|s", set_name, set_value);
388        let set_expected = Metric::set(set_name, set_value);
389        let set_actual = parse_dsd_metric(set_raw.as_bytes()).expect("should not fail to parse");
390        check_basic_metric_eq(set_expected, set_actual);
391    }
392
393    #[test]
394    fn metric_unit() {
395        let config = DogStatsDCodecConfiguration::default();
396
397        // Timing metrics must carry an implicit millisecond unit.
398        let (_, packet) = parse_dogstatsd_metric(b"my.timer:1.0|ms", &config).expect("should not fail to parse");
399        assert_eq!(packet.unit, Some("millisecond"));
400
401        // All other metric types must have no unit.
402        for kind in &["c", "g", "h", "d", "s"] {
403            let raw = format!("my.metric:1.0|{}", kind);
404            let (_, packet) = parse_dogstatsd_metric(raw.as_bytes(), &config).expect("should not fail to parse");
405            assert_eq!(packet.unit, None, "expected no unit for metric type '{}'", kind);
406        }
407    }
408
409    #[test]
410    fn metric_tags() {
411        let name = "my.counter";
412        let value = 1.0;
413        let tags = ["tag1", "tag2"];
414        let raw = format!("{}:{}|c|#{}", name, value, tags.join(","));
415        let expected = Metric::counter((name, &tags[..]), value);
416
417        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
418        check_basic_metric_eq(expected, actual);
419    }
420
421    #[test]
422    fn metric_tags_with_invalid_utf8_are_normalized() {
423        let mut input = b"my.gauge:1|g|#env:prod,tag:".to_vec();
424        input.extend_from_slice(&[0xff, 0xfe]);
425
426        let expected = Metric::gauge(("my.gauge", &["env:prod", "tag:\u{FFFD}"][..]), 1.0);
427        let actual = parse_dsd_metric(&input).expect("metric with invalid tag bytes should parse");
428
429        check_basic_metric_eq(expected, actual);
430    }
431
432    #[test]
433    fn metric_sample_rate() {
434        let name = "my.counter";
435        let value = 1.0;
436        let sample_rate = 0.5;
437        let raw = format!("{}:{}|c|@{}", name, value, sample_rate);
438
439        let value_sample_rate_adjusted = value * (1.0 / sample_rate);
440        let expected = Metric::counter(name, value_sample_rate_adjusted);
441
442        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
443        let actual = check_basic_metric_eq(expected, actual);
444        let values = match actual.values() {
445            MetricValues::Counter(values) => values
446                .into_iter()
447                .map(|(ts, v)| (ts.map(|v| v.get()).unwrap_or(0), v))
448                .collect::<Vec<_>>(),
449            _ => panic!("expected counter values"),
450        };
451
452        assert_eq!(values.len(), 1);
453        assert_eq!(values[0], (0, value_sample_rate_adjusted));
454    }
455
456    #[test]
457    fn metric_timestamped_count_sample_rate_matches_no_aggregation_pipeline() {
458        let name = "my.counter";
459        let value = 2.0;
460        let sample_rate = 0.25;
461        let timestamp = 1234567890;
462        let raw = format!("{}:{}|c|@{}|T{}", name, value, sample_rate, timestamp);
463
464        let mut expected = Metric::counter(name, value);
465        expected.values_mut().set_timestamp(timestamp);
466
467        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
468        let actual = check_basic_metric_eq(expected, actual);
469        let values = match actual.values() {
470            MetricValues::Counter(values) => values
471                .into_iter()
472                .map(|(ts, v)| (ts.map(|v| v.get()).unwrap_or(0), v))
473                .collect::<Vec<_>>(),
474            _ => panic!("expected counter values"),
475        };
476
477        assert_eq!(values.len(), 1);
478        assert_eq!(values[0], (timestamp, value));
479    }
480
481    #[test]
482    fn metric_local_data() {
483        let name = "my.counter";
484        let value = 1.0;
485        let local_data = "abcdef123456";
486        let raw = format!("{}:{}|c|c:{}", name, value, local_data);
487        let expected = Metric::counter(name, value);
488
489        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
490        check_basic_metric_eq(expected, actual);
491
492        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
493        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
494        let (_, packet) = parse_dogstatsd_metric(raw.as_bytes(), &config).expect("should not fail to parse");
495        assert_eq!(packet.local_data, Some(local_data));
496    }
497
498    #[test]
499    fn metric_unix_timestamp() {
500        let name = "my.counter";
501        let value = 1.0;
502        let timestamp = 1234567890;
503        let raw = format!("{}:{}|c|T{}", name, value, timestamp);
504        let mut expected = Metric::counter(name, value);
505        expected.values_mut().set_timestamp(timestamp);
506
507        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
508        check_basic_metric_eq(expected, actual);
509    }
510
511    #[test]
512    fn metric_external_data() {
513        let name = "my.counter";
514        let value = 1.0;
515        let external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
516        let raw = format!("{}:{}|c|e:{}", name, value, external_data);
517        let expected = Metric::counter(name, value);
518
519        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
520        check_basic_metric_eq(expected, actual);
521
522        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
523        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
524        let (_, packet) = parse_dogstatsd_metric(raw.as_bytes(), &config).expect("should not fail to parse");
525        assert_eq!(packet.external_data, Some(external_data));
526    }
527
528    #[test]
529    fn metric_cardinality() {
530        let name = "my.counter";
531        let value = 1.0;
532        let cardinality = "high";
533        let raw = format!("{}:{}|c|card:{}", name, value, cardinality);
534        let expected = Metric::counter(name, value);
535
536        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
537        check_basic_metric_eq(expected, actual);
538
539        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
540        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
541        let (_, packet) = parse_dogstatsd_metric(raw.as_bytes(), &config).expect("should not fail to parse");
542        assert_eq!(packet.cardinality, Some(OriginTagCardinality::High));
543    }
544
545    #[test]
546    fn metric_multiple_extensions() {
547        let name = "my.counter";
548        let value = 1.0;
549        let sample_rate = 0.5;
550        let tags = ["tag1", "tag2"];
551        let local_data = "abcdef123456";
552        let external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
553        let cardinality = "orchestrator";
554        let timestamp = 1234567890;
555        let raw = format!(
556            "{}:{}|c|#{}|@{}|c:{}|e:{}|card:{}|T{}",
557            name,
558            value,
559            tags.join(","),
560            sample_rate,
561            local_data,
562            external_data,
563            cardinality,
564            timestamp
565        );
566
567        let mut expected = Metric::counter((name, &tags[..]), value);
568        expected.values_mut().set_timestamp(timestamp);
569
570        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
571        let actual = check_basic_metric_eq(expected, actual);
572        let values = match actual.values() {
573            MetricValues::Counter(values) => values
574                .into_iter()
575                .map(|(ts, v)| (ts.map(|v| v.get()).unwrap_or(0), v))
576                .collect::<Vec<_>>(),
577            _ => panic!("expected counter values"),
578        };
579
580        assert_eq!(values.len(), 1);
581        assert_eq!(values[0], (timestamp, value));
582
583        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
584        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
585        let (_, packet) = parse_dogstatsd_metric(raw.as_bytes(), &config).expect("should not fail to parse");
586        assert_eq!(packet.local_data, Some(local_data));
587        assert_eq!(packet.external_data, Some(external_data));
588        assert_eq!(packet.cardinality, Some(OriginTagCardinality::Orchestrator));
589    }
590
591    #[test]
592    fn multivalue_metrics() {
593        let name = "my.counter";
594        let values = [1.0, 2.0, 3.0];
595        let values_stringified = values.iter().map(|v| v.to_string()).collect::<Vec<_>>();
596        let raw = format!("{}:{}|c", name, values_stringified.join(":"));
597        let expected = Metric::counter(name, values);
598        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
599        check_basic_metric_eq(expected, actual);
600
601        let name = "my.gauge";
602        let values = [42.0, 5.0, -18.0];
603        let values_stringified = values.iter().map(|v| v.to_string()).collect::<Vec<_>>();
604        let raw = format!("{}:{}|g", name, values_stringified.join(":"));
605        let expected = Metric::gauge(name, values);
606        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
607        check_basic_metric_eq(expected, actual);
608
609        // Special case where we check this for both timers and histograms since we treat them both the same when
610        // parsing.
611        //
612        // Additionally, we have an optimization to return a single distribution metric from multi-value payloads, so we
613        // also check here that only one metric is generated for multi-value timers/histograms/distributions.
614        let name = "my.timer_or_histogram";
615        let values = [27.5, 4.20, 80.085];
616        let values_stringified = values.iter().map(|v| v.to_string()).collect::<Vec<_>>();
617        for kind in &["ms", "h"] {
618            let raw = format!("{}:{}|{}", name, values_stringified.join(":"), kind);
619            let expected = Metric::histogram(name, values);
620            let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
621            check_basic_metric_eq(expected, actual);
622        }
623
624        let name = "my.distribution";
625        let raw = format!("{}:{}|d", name, values_stringified.join(":"));
626        let expected = Metric::distribution(name, values);
627        let actual = parse_dsd_metric(raw.as_bytes()).expect("should not fail to parse");
628        check_basic_metric_eq(expected, actual);
629    }
630
631    #[test]
632    fn respects_maximum_tag_count() {
633        let input = b"foo:1|c|#tag1:value1,tag2:value2,tag3:value3";
634
635        let cases = [3, 2, 1];
636        for max_tag_count in cases {
637            let config = DogStatsDCodecConfiguration::default().with_maximum_tag_count(max_tag_count);
638
639            let metric = parse_dsd_metric_with_conf(input, &config)
640                .expect("should not fail to parse")
641                .expect("should not fail to intern");
642            assert_eq!(metric.context().tags().len(), max_tag_count);
643        }
644    }
645
646    #[test]
647    fn respects_maximum_tag_length() {
648        let input = b"foo:1|c|#tag1:short,tag2:medium,tag3:longlong";
649
650        let cases = [6, 5, 4];
651        for max_tag_length in cases {
652            let config = DogStatsDCodecConfiguration::default().with_maximum_tag_length(max_tag_length);
653
654            let metric = parse_dsd_metric_with_conf(input, &config)
655                .expect("should not fail to parse")
656                .expect("should not fail to intern");
657            for tag in metric.context().tags().into_iter() {
658                assert!(tag.len() <= max_tag_length);
659            }
660        }
661    }
662
663    #[test]
664    fn respects_read_timestamps() {
665        let input = b"foo:1|c|T1234567890";
666
667        let config = DogStatsDCodecConfiguration::default().with_timestamps(false);
668
669        let metric = parse_dsd_metric_with_conf(input, &config)
670            .expect("should not fail to parse")
671            .expect("should not fail to intern");
672
673        let value_timestamps = match metric.values() {
674            MetricValues::Counter(values) => values
675                .into_iter()
676                .map(|(ts, _)| ts.map(|v| v.get()).unwrap_or(0))
677                .collect::<Vec<_>>(),
678            _ => panic!("expected counter values"),
679        };
680
681        assert_eq!(value_timestamps.len(), 1);
682        assert_eq!(value_timestamps[0], 0);
683    }
684
685    #[test]
686    fn permissive_mode() {
687        let payload = b"codeheap 'non-nmethods'.usage:0.3054|g|#env:dev,service:foobar,datacenter:localhost.dev";
688
689        let config = DogStatsDCodecConfiguration::default().with_permissive_mode(true);
690        match parse_dsd_metric_with_conf(payload, &config) {
691            Ok(result) => assert!(result.is_some(), "should not fail to materialize metric after decoding"),
692            Err(e) => panic!("should not have errored: {:?}", e),
693        }
694    }
695
696    #[test]
697    fn minimum_sample_rate() {
698        // Sample rate of 0.01 should lead to a count of 100 when handling a single value.
699        let minimum_sample_rate = SampleRate::try_from(0.01).unwrap();
700        let config = DogStatsDCodecConfiguration::default().with_minimum_sample_rate(minimum_sample_rate.rate());
701
702        let cases = [
703            // Worst case scenario: sample rate of zero, or "infinitely sampled".
704            "test:1|d|@0".to_string(),
705            // Bunch of values with different sample rates all below the minimum sample rate.
706            "test:1|d|@0.001".to_string(),
707            "test:1|d|@0.0005".to_string(),
708            "test:1|d|@0.00001".to_string(),
709            // Control: use the minimum sample rate.
710            format!("test:1|d|@{}", minimum_sample_rate.rate()),
711            // Bunch of values with _greater_ sampling rates than the minimum.
712            "test:1|d|@0.1".to_string(),
713            "test:1|d|@0.5".to_string(),
714            "test:1|d".to_string(),
715        ];
716
717        for input in cases {
718            let metric = parse_dsd_metric_with_conf(input.as_bytes(), &config)
719                .expect("Should not fail to parse metric.")
720                .expect("Metric should be present.");
721
722            let sketch = match metric.values() {
723                MetricValues::Distribution(points) => {
724                    points
725                        .into_iter()
726                        .next()
727                        .expect("Should have at least one sketch point.")
728                        .1
729                }
730                _ => panic!("Unexpected metric type."),
731            };
732
733            assert!(sketch.count() as u64 <= minimum_sample_rate.weight());
734        }
735    }
736
737    #[test]
738    fn client_origin_fields_ignored_when_disabled() {
739        let local_data = "cn-name-a";
740        let external_data = "it-false,cn-name-b,pu-810fe89d-da47-410b-8979-9154a40f8183";
741        let raw = format!("foo:1|c|c:{local_data}|e:{external_data}|card:high");
742
743        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(false);
744        let (_, packet) = parse_dogstatsd_metric(raw.as_bytes(), &config).expect("should not fail to parse");
745
746        assert_eq!(packet.local_data, None);
747        assert_eq!(packet.external_data, None);
748        assert_eq!(packet.cardinality, None);
749    }
750
751    #[test]
752    fn non_finite_metric_values_are_dropped() {
753        // Non-finite float values (NaN, ±Inf) are silently dropped at parse time with a debug log.
754        // The Datadog Agent's trace agent sends NaN gauges (e.g. encode_ms.avg) when a flush
755        // window has zero operations, producing 0.0/0.0 in Go. The parse succeeds but yields
756        // zero valid points; handle_frame then returns Ok(None) for zero-point packets.
757        let config = DogStatsDCodecConfiguration::default();
758        let cases = ["my.gauge:NaN|g", "my.gauge:inf|g", "my.gauge:-inf|g"];
759
760        for input in &cases {
761            let (_, packet) = parse_dogstatsd_metric(input.as_bytes(), &config)
762                .unwrap_or_else(|_| panic!("should parse without error: {input}"));
763            assert_eq!(
764                packet.num_points, 0,
765                "non-finite value should be dropped, leaving 0 valid points: {input}"
766            );
767        }
768    }
769
770    proptest! {
771        #![proptest_config(ProptestConfig::with_cases(1000))]
772        #[test]
773        fn property_test_malicious_input_non_exhaustive(input in arb_vec(0..255u8, 0..1000)) {
774            // We're testing that the parser is resilient to malicious input, which means that it should not panic or
775            // crash when given input that's not well-formed.
776            //
777            // As this is a property test, it is _not_ exhaustive but generally should catch simple issues that manage
778            // to escape the unit tests. This is left here for the sole reason of incrementally running this every time
779            // all tests are run, in the hopes of potentially catching an issue that might have been missed.
780            //
781            // Beyond "does not panic", we also encode the structural guarantee the parser makes on success: the metric
782            // name is matched with `take_while1` (so it is always non-empty) and the name parser stops at the `:`
783            // value delimiter (so the name can never contain one). Whenever the parser accepts an input, the produced
784            // metric must satisfy both.
785            if let Ok(Some(metric)) = parse_dsd_metric(&input) {
786                let name = metric.context().name();
787                let name = name.as_ref();
788                prop_assert!(!name.is_empty(), "parsed metric name must be non-empty for input {input:?}");
789                prop_assert!(
790                    !name.contains(':'),
791                    "parsed metric name must not contain the ':' value delimiter, got {name:?}"
792                );
793            }
794        }
795    }
796}