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