saluki_io/deser/codec/dogstatsd/
event.rs

1use nom::{
2    bytes::complete::{tag, take},
3    character::complete::u32 as parse_u32,
4    combinator::all_consuming,
5    error::{Error, ErrorKind},
6    sequence::{delimited, preceded, separated_pair},
7    IResult, Parser as _,
8};
9use saluki_core::data_model::{event::eventd::*, origin::OriginTagCardinality, tags::RawTags};
10use stringtheory::MetaString;
11
12use super::{helpers::*, DogStatsDCodecConfiguration};
13
14/// A DogStatsD event packet.
15pub struct EventPacket<'a> {
16    pub title: MetaString,
17    pub text: MetaString,
18    pub timestamp: Option<u64>,
19    pub hostname: Option<&'a str>,
20    pub aggregation_key: Option<&'a str>,
21    pub priority: Option<Priority>,
22    pub alert_type: Option<AlertType>,
23    pub source_type_name: Option<&'a str>,
24    pub tags: RawTags<'a>,
25    pub local_data: Option<&'a str>,
26    pub external_data: Option<&'a str>,
27    pub cardinality: Option<OriginTagCardinality>,
28}
29
30#[inline]
31pub fn parse_dogstatsd_event<'a>(
32    input: &'a [u8], config: &DogStatsDCodecConfiguration,
33) -> IResult<&'a [u8], EventPacket<'a>> {
34    // We parse the title length and text length from `_e{<TITLE_UTF8_LENGTH>,<TEXT_UTF8_LENGTH>}:`
35    let (remaining, (title_len, text_len)) = delimited(
36        tag(EVENT_PREFIX),
37        separated_pair(parse_u32, tag(","), parse_u32),
38        tag("}:"),
39    )
40    .parse(input)?;
41
42    // Title and Text are the required fields of an event.
43    if title_len == 0 || text_len == 0 {
44        return Err(nom::Err::Error(Error::new(input, ErrorKind::Verify)));
45    }
46
47    let (remaining, (raw_title, raw_text)) =
48        separated_pair(take(title_len), tag("|"), take(text_len)).parse(remaining)?;
49
50    let title = match simdutf8::basic::from_utf8(raw_title) {
51        Ok(title) => title.replace("\\n", "\n"),
52        Err(_) => return Err(nom::Err::Error(Error::new(raw_title, ErrorKind::Verify))),
53    };
54
55    let text = match simdutf8::basic::from_utf8(raw_text) {
56        Ok(text) => text.replace("\\n", "\n"),
57        Err(_) => return Err(nom::Err::Error(Error::new(raw_text, ErrorKind::Verify))),
58    };
59
60    // At this point, we may have some of this additional data, and if so, we also then would have a pipe separator at
61    // the very front, which we'd want to consume before going further.
62    //
63    // After that, we simply split the remaining bytes by the pipe separator, and then try and parse each chunk to see
64    // if it's any of the protocol extensions we know of.
65    //
66    // Priority and Alert Type have default values
67    let mut maybe_priority = Some(Priority::Normal);
68    let mut maybe_alert_type = Some(AlertType::Info);
69    let mut maybe_timestamp = None;
70    let mut maybe_hostname = None;
71    let mut maybe_aggregation_key = None;
72    let mut maybe_source_type = None;
73    let mut maybe_tags = None;
74    let mut maybe_local_data = None;
75    let mut maybe_external_data = None;
76    let mut maybe_cardinality = None;
77
78    let remaining = if !remaining.is_empty() {
79        let (mut remaining, _) = tag("|")(remaining)?;
80        while let Some((chunk, tail)) = split_at_delimiter(remaining, b'|') {
81            if chunk.len() < 2 {
82                break;
83            }
84            match &chunk[..2] {
85                // Timestamp: client-provided timestamp for the event, relative to the Unix epoch, in seconds.
86                TIMESTAMP_PREFIX => {
87                    let (_, timestamp) = all_consuming(preceded(tag(TIMESTAMP_PREFIX), unix_timestamp)).parse(chunk)?;
88                    maybe_timestamp = Some(timestamp);
89                }
90                // Hostname: client-provided hostname for the host that this event originated from.
91                HOSTNAME_PREFIX if chunk != HOSTNAME_PREFIX => {
92                    let (_, hostname) =
93                        all_consuming(preceded(tag(HOSTNAME_PREFIX), ascii_alphanum_and_seps)).parse(chunk)?;
94                    maybe_hostname = Some(hostname);
95                }
96                // Aggregation key: key to be used to group this event with others that have the same key.
97                AGGREGATION_KEY_PREFIX if chunk != AGGREGATION_KEY_PREFIX => {
98                    let (_, aggregation_key) =
99                        all_consuming(preceded(tag(AGGREGATION_KEY_PREFIX), ascii_alphanum_and_seps)).parse(chunk)?;
100                    maybe_aggregation_key = Some(aggregation_key);
101                }
102                // Priority: client-provided priority of the event.
103                PRIORITY_PREFIX => {
104                    let (_, priority) =
105                        all_consuming(preceded(tag(PRIORITY_PREFIX), ascii_alphanum_and_seps)).parse(chunk)?;
106                    maybe_priority = Priority::try_from_string(priority);
107                }
108                // Source type name: client-provided source type name of the event.
109                SOURCE_TYPE_PREFIX if chunk != SOURCE_TYPE_PREFIX => {
110                    let (_, source_type) =
111                        all_consuming(preceded(tag(SOURCE_TYPE_PREFIX), ascii_alphanum_and_seps)).parse(chunk)?;
112                    maybe_source_type = Some(source_type);
113                }
114                // Alert type: client-provided alert type of the event.
115                ALERT_TYPE_PREFIX => {
116                    let (_, alert_type) =
117                        all_consuming(preceded(tag(ALERT_TYPE_PREFIX), ascii_alphanum_and_seps)).parse(chunk)?;
118                    maybe_alert_type = AlertType::try_from_string(alert_type);
119                }
120                // Local Data: client-provided data used for resolving the entity ID that this event originated from.
121                LOCAL_DATA_PREFIX if config.client_origin_detection && chunk != LOCAL_DATA_PREFIX => {
122                    let (_, local_data) = all_consuming(preceded(tag(LOCAL_DATA_PREFIX), local_data)).parse(chunk)?;
123                    maybe_local_data = Some(local_data);
124                }
125                // External Data: client-provided data used for resolving the entity ID that this event originated from.
126                EXTERNAL_DATA_PREFIX if config.client_origin_detection && chunk != EXTERNAL_DATA_PREFIX => {
127                    let (_, external_data) =
128                        all_consuming(preceded(tag(EXTERNAL_DATA_PREFIX), external_data)).parse(chunk)?;
129                    maybe_external_data = Some(external_data);
130                }
131                // Cardinality: client-provided cardinality for the event.
132                _ if chunk.starts_with(CARDINALITY_PREFIX)
133                    && config.client_origin_detection
134                    && chunk != CARDINALITY_PREFIX =>
135                {
136                    let (_, cardinality) = cardinality(chunk)?;
137                    maybe_cardinality = cardinality;
138                }
139                // Tags: additional tags to be added to the event.
140                _ if chunk.starts_with(TAGS_PREFIX) && chunk != TAGS_PREFIX => {
141                    let (_, tags) = all_consuming(preceded(tag("#"), tags(config))).parse(chunk)?;
142                    maybe_tags = Some(tags);
143                }
144                _ => {
145                    // We don't know what this is, so we just skip it.
146                    //
147                    // TODO: Should we throw an error, warn, or be silently permissive?
148                }
149            }
150            remaining = tail;
151        }
152        remaining
153    } else {
154        remaining
155    };
156
157    let tags = maybe_tags.unwrap_or_else(RawTags::empty);
158
159    let eventd = EventPacket {
160        title: title.into(),
161        text: text.into(),
162        tags,
163        timestamp: maybe_timestamp,
164        hostname: maybe_hostname,
165        aggregation_key: maybe_aggregation_key,
166        priority: maybe_priority,
167        alert_type: maybe_alert_type,
168        source_type_name: maybe_source_type,
169        local_data: maybe_local_data,
170        external_data: maybe_external_data,
171        cardinality: maybe_cardinality,
172    };
173    Ok((remaining, eventd))
174}
175
176#[cfg(test)]
177mod tests {
178    use nom::IResult;
179    use saluki_core::data_model::{
180        event::eventd::{AlertType, EventD, Priority},
181        origin::OriginTagCardinality,
182        tags::{SharedTagSet, Tag, TagSet},
183    };
184    use stringtheory::MetaString;
185
186    use super::{parse_dogstatsd_event, DogStatsDCodecConfiguration};
187
188    type NomResult<'input, T> = Result<T, nom::Err<nom::error::Error<&'input [u8]>>>;
189
190    fn parse_dsd_eventd(input: &[u8]) -> NomResult<'_, EventD> {
191        let default_config = DogStatsDCodecConfiguration::default();
192        parse_dsd_eventd_with_conf(input, &default_config)
193    }
194
195    fn parse_dsd_eventd_with_conf<'input>(
196        input: &'input [u8], config: &DogStatsDCodecConfiguration,
197    ) -> NomResult<'input, EventD> {
198        let (remaining, eventd) = parse_dsd_eventd_direct(input, config)?;
199        assert!(remaining.is_empty());
200        Ok(eventd)
201    }
202
203    fn parse_dsd_eventd_direct<'input>(
204        input: &'input [u8], config: &DogStatsDCodecConfiguration,
205    ) -> IResult<&'input [u8], EventD> {
206        let (remaining, packet) = parse_dogstatsd_event(input, config)?;
207        assert!(remaining.is_empty());
208
209        let mut event_tags = TagSet::default();
210        for tag in packet.tags.iter() {
211            event_tags.insert_tag(tag);
212        }
213
214        let eventd = EventD::new(packet.title, packet.text)
215            .with_timestamp(packet.timestamp)
216            .with_hostname(packet.hostname.map(|s| s.into()))
217            .with_aggregation_key(packet.aggregation_key.map(|s| s.into()))
218            .with_alert_type(packet.alert_type)
219            .with_priority(packet.priority)
220            .with_source_type_name(packet.source_type_name.map(|s| s.into()))
221            .with_alert_type(packet.alert_type)
222            .with_tags(event_tags);
223
224        Ok((remaining, eventd))
225    }
226
227    #[track_caller]
228    fn check_basic_eventd_eq(expected: EventD, actual: EventD) {
229        assert_eq!(expected.title(), actual.title());
230        assert_eq!(expected.text(), actual.text());
231        assert_eq!(expected.timestamp(), actual.timestamp());
232        assert_eq!(expected.hostname(), actual.hostname());
233        assert_eq!(expected.aggregation_key(), actual.aggregation_key());
234        assert_eq!(expected.priority(), actual.priority());
235        assert_eq!(expected.source_type_name(), actual.source_type_name());
236        assert_eq!(expected.alert_type(), actual.alert_type());
237        assert_eq!(expected.tags(), actual.tags());
238        assert_eq!(expected.origin_tags(), actual.origin_tags());
239    }
240
241    #[test]
242    fn basic_eventd() {
243        let event_title = "my event";
244        let event_text = "text";
245        let raw = format!(
246            "_e{{{},{}}}:{}|{}",
247            event_title.len(),
248            event_text.len(),
249            event_title,
250            event_text
251        );
252
253        let actual = parse_dsd_eventd(raw.as_bytes()).unwrap();
254        let expected = EventD::new(event_title, event_text);
255        check_basic_eventd_eq(expected, actual);
256    }
257
258    #[test]
259    fn eventd_tags() {
260        let event_title = "my event";
261        let event_text = "text";
262        let tags = ["tag1", "tag2"];
263        let shared_tag_set: SharedTagSet = tags.iter().map(|&s| Tag::from(s)).collect::<TagSet>().into_shared();
264        let raw = format!(
265            "_e{{{},{}}}:{}|{}|#{}",
266            event_title.len(),
267            event_text.len(),
268            event_title,
269            event_text,
270            tags.join(","),
271        );
272
273        let expected = EventD::new(event_title, event_text).with_tags(shared_tag_set);
274        let actual = parse_dsd_eventd(raw.as_bytes()).unwrap();
275        check_basic_eventd_eq(expected, actual);
276    }
277
278    #[test]
279    fn event_tags_with_invalid_utf8_are_normalized() {
280        let mut input = b"_e{5,4}:title|text|#env:prod,tag:".to_vec();
281        input.extend_from_slice(&[0xff, 0xfe]);
282
283        let tags = ["env:prod", "tag:\u{FFFD}"];
284        let expected = EventD::new("title", "text").with_tags(SharedTagSet::from(TagSet::from_iter(
285            tags.iter().map(|&tag| tag.into()),
286        )));
287        let actual = parse_dsd_eventd(&input).expect("event with invalid tag bytes should parse");
288
289        check_basic_eventd_eq(expected, actual);
290    }
291
292    #[test]
293    fn eventd_priority() {
294        let event_title = "my event";
295        let event_text = "text";
296        let event_priority = Priority::Low;
297        let raw = format!(
298            "_e{{{},{}}}:{}|{}|p:{}",
299            event_title.len(),
300            event_text.len(),
301            event_title,
302            event_text,
303            event_priority
304        );
305
306        let expected = EventD::new(event_title, event_text).with_priority(event_priority);
307        let actual = parse_dsd_eventd(raw.as_bytes()).unwrap();
308        check_basic_eventd_eq(expected, actual);
309    }
310
311    #[test]
312    fn eventd_alert_type() {
313        let event_title = "my event";
314        let event_text = "text";
315        let event_alert_type = AlertType::Warning;
316        let raw = format!(
317            "_e{{{},{}}}:{}|{}|t:{}",
318            event_title.len(),
319            event_text.len(),
320            event_title,
321            event_text,
322            event_alert_type
323        );
324
325        let expected = EventD::new(event_title, event_text).with_alert_type(event_alert_type);
326        let actual = parse_dsd_eventd(raw.as_bytes()).unwrap();
327        check_basic_eventd_eq(expected, actual);
328    }
329
330    #[test]
331    fn eventd_multiple_extensions() {
332        let event_title = "my event";
333        let event_text = "text";
334        let event_hostname = MetaString::from("testhost");
335        let event_aggregation_key = MetaString::from("testkey");
336        let event_priority = Priority::Low;
337        let event_source_type = MetaString::from("testsource");
338        let event_alert_type = AlertType::Success;
339        let event_timestamp = 1234567890;
340        let event_local_data = "abcdef123456";
341        let event_external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
342        let event_cardinality = "low";
343        let tags = ["tags1", "tags2"];
344        let shared_tag_set = SharedTagSet::from(TagSet::from_iter(tags.iter().map(|&s| s.into())));
345        let raw = format!(
346            "_e{{{},{}}}:{}|{}|h:{}|k:{}|p:{}|s:{}|t:{}|d:{}|c:{}|e:{}|card:{}|#{}",
347            event_title.len(),
348            event_text.len(),
349            event_title,
350            event_text,
351            event_hostname,
352            event_aggregation_key,
353            event_priority,
354            event_source_type,
355            event_alert_type,
356            event_timestamp,
357            event_local_data,
358            event_external_data,
359            event_cardinality,
360            tags.join(","),
361        );
362        let actual = parse_dsd_eventd(raw.as_bytes()).unwrap();
363        let expected = EventD::new(event_title, event_text)
364            .with_hostname(event_hostname)
365            .with_aggregation_key(event_aggregation_key)
366            .with_priority(event_priority)
367            .with_source_type_name(event_source_type)
368            .with_alert_type(event_alert_type)
369            .with_timestamp(event_timestamp)
370            .with_tags(shared_tag_set);
371        check_basic_eventd_eq(expected, actual);
372
373        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
374        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
375        let (_, packet) = parse_dogstatsd_event(raw.as_bytes(), &config).expect("should not fail to parse");
376        assert_eq!(packet.local_data, Some(event_local_data));
377        assert_eq!(packet.external_data, Some(event_external_data));
378        assert_eq!(packet.cardinality, Some(OriginTagCardinality::Low));
379    }
380
381    #[test]
382    fn event_rejects_empty_title_or_text() {
383        use nom::error::ErrorKind;
384
385        // Title and text are the two required fields of an event: a declared length of zero for either one is a
386        // structural error, so the parser rejects the frame with a `Verify` error rather than emitting an event with
387        // an empty title/text.
388        let config = DogStatsDCodecConfiguration::default();
389
390        match parse_dogstatsd_event(b"_e{0,4}:|text", &config) {
391            Err(nom::Err::Error(e)) => assert_eq!(e.code, ErrorKind::Verify),
392            Err(other) => panic!("expected Verify error for empty title, got {other:?}"),
393            Ok(_) => panic!("empty title must be rejected"),
394        }
395
396        match parse_dogstatsd_event(b"_e{5,0}:title|", &config) {
397            Err(nom::Err::Error(e)) => assert_eq!(e.code, ErrorKind::Verify),
398            Err(other) => panic!("expected Verify error for empty text, got {other:?}"),
399            Ok(_) => panic!("empty text must be rejected"),
400        }
401    }
402
403    #[test]
404    fn client_origin_fields_ignored_when_disabled() {
405        let local_data = "abcdef123456";
406        let external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
407        let raw = format!("_e{{5,4}}:title|text|c:{}|e:{}|card:low", local_data, external_data);
408        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(false);
409        let (_, packet) = parse_dogstatsd_event(raw.as_bytes(), &config).expect("should not fail to parse");
410        assert_eq!(packet.local_data, None);
411        assert_eq!(packet.external_data, None);
412        assert_eq!(packet.cardinality, None);
413    }
414
415    #[test]
416    fn empty_structured_fields_treated_as_missing() {
417        // All optional stringy fields are empty — should parse successfully and treat them as missing.
418        let raw = "_e{5,4}:title|text|h:|k:|s:|c:|e:|card:|#";
419        let config = DogStatsDCodecConfiguration::default();
420        let (_, packet) = parse_dogstatsd_event(raw.as_bytes(), &config).expect("should not fail to parse");
421        assert_eq!(packet.hostname, None);
422        assert_eq!(packet.aggregation_key, None);
423        assert_eq!(packet.source_type_name, None);
424        assert_eq!(packet.local_data, None);
425        assert_eq!(packet.external_data, None);
426        assert_eq!(packet.cardinality, None);
427        assert!(packet.tags.iter().next().is_none());
428    }
429}