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
14pub 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_PREFIX => {
57 let (_, timestamp) = all_consuming(preceded(tag(TIMESTAMP_PREFIX), unix_timestamp)).parse(chunk)?;
58 maybe_timestamp = Some(timestamp);
59 }
60 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_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_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 _ 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 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 _ if chunk.starts_with(CARDINALITY_PREFIX) && config.client_origin_detection => {
89 let (_, cardinality) = cardinality(chunk)?;
90 maybe_cardinality = cardinality;
91 }
92 _ => {
93 }
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 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 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}