saluki_io/deser/codec/dogstatsd/
service_check.rs

1use nom::{
2    bytes::complete::tag,
3    character::complete::u8 as parse_u8,
4    combinator::all_consuming,
5    error::{Error, ErrorKind},
6    sequence::{preceded, separated_pair},
7    IResult, Parser as _,
8};
9use saluki_core::data_model::{event::service_check::*, origin::OriginTagCardinality, tags::RawTags};
10use stringtheory::MetaString;
11
12use super::{helpers::*, DogStatsDCodecConfiguration};
13
14/// A DogStatsD service check packet.
15pub struct ServiceCheckPacket<'a> {
16    pub name: MetaString,
17    pub status: CheckStatus,
18    pub timestamp: Option<u64>,
19    pub hostname: Option<&'a str>,
20    pub message: Option<&'a str>,
21    pub tags: RawTags<'a>,
22    pub local_data: Option<&'a str>,
23    pub external_data: Option<&'a str>,
24    pub cardinality: Option<OriginTagCardinality>,
25}
26
27pub fn parse_dogstatsd_service_check<'a>(
28    input: &'a [u8], config: &DogStatsDCodecConfiguration,
29) -> IResult<&'a [u8], ServiceCheckPacket<'a>> {
30    let (remaining, (name, raw_check_status)) = preceded(
31        tag(SERVICE_CHECK_PREFIX),
32        separated_pair(ascii_alphanum_and_seps, tag("|"), parse_u8),
33    )
34    .parse(input)?;
35
36    let check_status =
37        CheckStatus::try_from(raw_check_status).map_err(|_| nom::Err::Error(Error::new(input, ErrorKind::Verify)))?;
38
39    let mut maybe_timestamp = None;
40    let mut maybe_hostname = None;
41    let mut maybe_tags = None;
42    let mut maybe_message = None;
43    let mut maybe_local_data = None;
44    let mut maybe_external_data = None;
45    let mut maybe_cardinality = None;
46
47    let remaining = if !remaining.is_empty() {
48        let (mut remaining, _) = tag("|")(remaining)?;
49        while let Some((chunk, tail)) = split_at_delimiter(remaining, b'|') {
50            if chunk.len() < 2 {
51                break;
52            }
53
54            match &chunk[..2] {
55                // Timestamp: client-provided timestamp for the event, relative to the Unix epoch, in seconds.
56                TIMESTAMP_PREFIX => {
57                    let (_, timestamp) = all_consuming(preceded(tag(TIMESTAMP_PREFIX), unix_timestamp)).parse(chunk)?;
58                    maybe_timestamp = Some(timestamp);
59                }
60                // Hostname: client-provided hostname for the host that this service check originated from.
61                HOSTNAME_PREFIX => {
62                    let (_, hostname) =
63                        all_consuming(preceded(tag(HOSTNAME_PREFIX), ascii_alphanum_and_seps)).parse(chunk)?;
64                    maybe_hostname = Some(hostname);
65                }
66                // Local Data: client-provided data used for resolving the entity ID that this service check originated from.
67                LOCAL_DATA_PREFIX if config.client_origin_detection => {
68                    let (_, local_data) = all_consuming(preceded(tag(LOCAL_DATA_PREFIX), local_data)).parse(chunk)?;
69                    maybe_local_data = Some(local_data);
70                }
71                // External Data: client-provided data used for resolving the entity ID that this service check originated from.
72                EXTERNAL_DATA_PREFIX if config.client_origin_detection => {
73                    let (_, external_data) =
74                        all_consuming(preceded(tag(EXTERNAL_DATA_PREFIX), external_data)).parse(chunk)?;
75                    maybe_external_data = Some(external_data);
76                }
77                // Tags: additional tags to be added to the service check.
78                _ if chunk.starts_with(TAGS_PREFIX) => {
79                    let (_, tags) = all_consuming(preceded(tag(TAGS_PREFIX), tags(config))).parse(chunk)?;
80                    maybe_tags = Some(tags);
81                }
82                // Message: A message describing the current state of the service check.
83                SERVICE_CHECK_MESSAGE_PREFIX => {
84                    let (_, message) = all_consuming(preceded(tag(SERVICE_CHECK_MESSAGE_PREFIX), utf8)).parse(chunk)?;
85                    maybe_message = Some(message);
86                }
87                // Cardinality: client-provided cardinality for the service check.
88                _ if chunk.starts_with(CARDINALITY_PREFIX) && config.client_origin_detection => {
89                    let (_, cardinality) = cardinality(chunk)?;
90                    maybe_cardinality = cardinality;
91                }
92                _ => {
93                    // We don't know what this is, so we just skip it.
94                    //
95                    // TODO: Should we throw an error, warn, or be silently permissive?
96                }
97            }
98            remaining = tail;
99        }
100        remaining
101    } else {
102        remaining
103    };
104
105    let tags = maybe_tags.unwrap_or_else(RawTags::empty);
106
107    let service_check_packet = ServiceCheckPacket {
108        name: name.into(),
109        status: check_status,
110        tags,
111        timestamp: maybe_timestamp,
112        hostname: maybe_hostname,
113        message: maybe_message,
114        local_data: maybe_local_data,
115        external_data: maybe_external_data,
116        cardinality: maybe_cardinality,
117    };
118    Ok((remaining, service_check_packet))
119}
120
121#[cfg(test)]
122mod tests {
123    use nom::IResult;
124    use saluki_core::data_model::{
125        event::service_check::{CheckStatus, ServiceCheck},
126        origin::OriginTagCardinality,
127        tags::{SharedTagSet, TagSet},
128    };
129    use stringtheory::MetaString;
130
131    use super::{parse_dogstatsd_service_check, DogStatsDCodecConfiguration};
132
133    type NomResult<'input, T> = Result<T, nom::Err<nom::error::Error<&'input [u8]>>>;
134
135    fn parse_dsd_service_check(input: &[u8]) -> NomResult<'_, ServiceCheck> {
136        let default_config = DogStatsDCodecConfiguration::default();
137        parse_dsd_service_check_with_conf(input, &default_config)
138    }
139
140    fn parse_dsd_service_check_with_conf<'input>(
141        input: &'input [u8], config: &DogStatsDCodecConfiguration,
142    ) -> NomResult<'input, ServiceCheck> {
143        let (remaining, service_check) = parse_dsd_service_check_direct(input, config)?;
144        assert!(remaining.is_empty());
145
146        Ok(service_check)
147    }
148
149    fn parse_dsd_service_check_direct<'input>(
150        input: &'input [u8], config: &DogStatsDCodecConfiguration,
151    ) -> IResult<&'input [u8], ServiceCheck> {
152        let (remaining, packet) = parse_dogstatsd_service_check(input, config)?;
153        assert!(remaining.is_empty());
154
155        let mut service_check_tags = TagSet::default();
156        for tag in packet.tags.iter() {
157            service_check_tags.insert_tag(tag);
158        }
159
160        let service_check = ServiceCheck::new(packet.name, packet.status)
161            .with_timestamp(packet.timestamp)
162            .with_hostname(packet.hostname.map(|s| s.into()))
163            .with_tags(service_check_tags)
164            .with_message(packet.message.map(|s| s.into()));
165
166        Ok((remaining, service_check))
167    }
168
169    #[track_caller]
170    fn check_basic_service_check_eq(expected: ServiceCheck, actual: ServiceCheck) {
171        assert_eq!(expected.name(), actual.name());
172        assert_eq!(expected.status(), actual.status());
173        assert_eq!(expected.timestamp(), actual.timestamp());
174        assert_eq!(expected.hostname(), actual.hostname());
175        assert_eq!(expected.tags(), actual.tags());
176        assert_eq!(expected.message(), actual.message());
177        assert_eq!(expected.origin_tags(), actual.origin_tags());
178    }
179
180    #[test]
181    fn basic_service_checks() {
182        let name = "testsvc";
183        let sc_status = CheckStatus::Warning;
184        let raw = format!("_sc|{}|{}", name, sc_status.as_u8());
185        let actual = parse_dsd_service_check(raw.as_bytes()).unwrap();
186        let expected = ServiceCheck::new(name, sc_status);
187        check_basic_service_check_eq(expected, actual);
188    }
189
190    #[test]
191    fn service_check_timestamp() {
192        let name = "testsvc";
193        let sc_status = CheckStatus::Warning;
194        let sc_timestamp = 1234567890;
195        let raw = format!("_sc|{}|{}|d:{}", name, sc_status.as_u8(), sc_timestamp);
196        let actual = parse_dsd_service_check(raw.as_bytes()).unwrap();
197        let expected = ServiceCheck::new(name, sc_status).with_timestamp(sc_timestamp);
198        check_basic_service_check_eq(expected, actual);
199    }
200
201    #[test]
202    fn service_check_tags() {
203        let name = "testsvc";
204        let sc_status = CheckStatus::Warning;
205        let tags = ["tag1", "tag2"];
206        let raw = format!("_sc|{}|{}|#{}", name, sc_status.as_u8(), tags.join(","));
207        let actual = parse_dsd_service_check(raw.as_bytes()).unwrap();
208        let shared_tag_set = SharedTagSet::from(TagSet::from_iter(tags.iter().map(|&s| s.into())));
209        let expected = ServiceCheck::new(name, sc_status).with_tags(shared_tag_set);
210        check_basic_service_check_eq(expected, actual);
211    }
212
213    #[test]
214    fn service_check_tags_with_invalid_utf8_are_normalized() {
215        let mut input = b"_sc|testsvc|0|#env:prod,tag:".to_vec();
216        input.extend_from_slice(&[0xff, 0xfe]);
217
218        let tags = ["env:prod", "tag:\u{FFFD}"];
219        let expected = ServiceCheck::new("testsvc", CheckStatus::Ok).with_tags(SharedTagSet::from(TagSet::from_iter(
220            tags.iter().map(|&tag| tag.into()),
221        )));
222        let actual = parse_dsd_service_check(&input).expect("service check with invalid tag bytes should parse");
223
224        check_basic_service_check_eq(expected, actual);
225    }
226
227    #[test]
228    fn service_check_message() {
229        let name = "testsvc";
230        let sc_status = CheckStatus::Ok;
231        let sc_message = MetaString::from("service running properly");
232        let raw = format!("_sc|{}|{}|m:{}", name, sc_status.as_u8(), sc_message);
233        let actual = parse_dsd_service_check(raw.as_bytes()).unwrap();
234        let expected = ServiceCheck::new(name, sc_status).with_message(sc_message);
235        check_basic_service_check_eq(expected, actual);
236    }
237
238    #[test]
239    fn service_check_fields_after_message() {
240        let name = "testsvc";
241        let sc_status = CheckStatus::Ok;
242        let sc_message = MetaString::from("service running properly");
243        let sc_local_data = "ci-1234567890";
244        let sc_external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
245        let raw = format!(
246            "_sc|{}|{}|#tag1,tag2|m:{}|c:{}|e:{}",
247            name,
248            sc_status.as_u8(),
249            sc_message,
250            sc_local_data,
251            sc_external_data,
252        );
253        let result = parse_dsd_service_check(raw.as_bytes()).unwrap();
254        let tags = ["tag1", "tag2"];
255        let shared_tag_set = SharedTagSet::from(TagSet::from_iter(tags.iter().map(|&s| s.into())));
256        let expected = ServiceCheck::new(name, sc_status)
257            .with_message(sc_message)
258            .with_tags(shared_tag_set);
259        check_basic_service_check_eq(expected, result);
260
261        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
262        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
263        let (_, packet) = parse_dogstatsd_service_check(raw.as_bytes(), &config).expect("should not fail to parse");
264        assert_eq!(packet.local_data, Some(sc_local_data));
265        assert_eq!(packet.external_data, Some(sc_external_data));
266    }
267
268    #[test]
269    fn service_check_multiple_extensions() {
270        let name = "testsvc";
271        let sc_status = CheckStatus::Unknown;
272        let sc_timestamp = 1234567890;
273        let sc_hostname = MetaString::from("myhost");
274        let sc_local_data = "abcdef123456";
275        let sc_external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
276        let tags = ["tag1", "tag2"];
277        let sc_message = MetaString::from("service status unknown");
278        let sc_cardinality = "none";
279        let raw = format!(
280            "_sc|{}|{}|d:{}|h:{}|c:{}|e:{}|card:{}|#{}|m:{}",
281            name,
282            sc_status.as_u8(),
283            sc_timestamp,
284            sc_hostname,
285            sc_local_data,
286            sc_external_data,
287            sc_cardinality,
288            tags.join(","),
289            sc_message
290        );
291        let shared_tag_set = SharedTagSet::from(TagSet::from_iter(tags.iter().map(|&s| s.into())));
292        let actual = parse_dsd_service_check(raw.as_bytes()).unwrap();
293        let expected = ServiceCheck::new(name, sc_status)
294            .with_timestamp(sc_timestamp)
295            .with_hostname(sc_hostname)
296            .with_tags(shared_tag_set)
297            .with_message(sc_message);
298        check_basic_service_check_eq(expected, actual);
299
300        // We need client_origin_detection on in order to parse local_data, external_data, and cardinality fields
301        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(true);
302        let (_, packet) = parse_dogstatsd_service_check(raw.as_bytes(), &config).expect("should not fail to parse");
303        assert_eq!(packet.local_data, Some(sc_local_data));
304        assert_eq!(packet.external_data, Some(sc_external_data));
305        assert_eq!(packet.cardinality, Some(OriginTagCardinality::None));
306    }
307
308    #[test]
309    fn client_origin_fields_ignored_when_disabled() {
310        let local_data = "abcdef123456";
311        let external_data = "it-false,cn-redis,pu-810fe89d-da47-410b-8979-9154a40f8183";
312        let raw = format!("_sc|testsvc|0|c:{}|e:{}|card:low", local_data, external_data);
313        let config = DogStatsDCodecConfiguration::default().with_client_origin_detection(false);
314        let (_, packet) = parse_dogstatsd_service_check(raw.as_bytes(), &config).expect("should not fail to parse");
315        assert_eq!(packet.local_data, None);
316        assert_eq!(packet.external_data, None);
317        assert_eq!(packet.cardinality, None);
318    }
319
320    #[test]
321    fn service_check_semi_real_payload_kafka() {
322        let raw_payload = "_sc|kafka.can_connect|2|#env:staging,service:datadog-agent,dd.internal.entity_id:none,dd.internal.card:none,instance:kafka-127.0.0.1-9999,jmx_server:127.0.0.1|m:Unable to instantiate or initialize instance 127.0.0.1:9999. Is the target JMX Server or JVM running? Failed to retrieve RMIServer stub: javax.naming.ServiceUnavailableException [Root exception is java.rmi.ConnectException: Connection refused to host: 127.0.0.1; nested exception is: \\n\tjava.net.ConnectException: Connection refused (Connection refused)]";
323        let _ = parse_dsd_service_check(raw_payload.as_bytes()).unwrap();
324    }
325}