1use std::sync::LazyLock;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use datadog_protos::checks::{
6 check_data::Data,
7 checks_server::{Checks, ChecksServer},
8 event::{AlertType as ProtoAlertType, Event as ProtoEvent, Priority as ProtoPriority},
9 log::{Log as ProtoLog, LogLevel},
10 metric::{Metric as ProtoMetric, MetricType},
11 service_check::{ServiceCheck as ProtoServiceCheck, Status as ServiceCheckStatus},
12 SendCheckPayloadRequest, SendCheckPayloadResponse,
13};
14use saluki_core::{
15 accounting::{MemoryBounds, MemoryBoundsBuilder},
16 components::{sources::*, BuildContext},
17 data_model::{
18 event::{
19 eventd::{AlertType, EventD, Priority},
20 log::{Log, LogStatus},
21 metric::{context::Context, Metric},
22 service_check::{CheckStatus, ServiceCheck},
23 Event, EventType,
24 },
25 tags::{Tag, TagSet},
26 },
27 runtime,
28 topology::OutputDefinition,
29};
30use saluki_error::{generic_error, GenericError};
31use saluki_io::net::{
32 server::http::{Http2Config, HttpServer},
33 ListenAddress,
34};
35use stringtheory::MetaString;
36use tokio::sync::mpsc;
37use tokio::{pin, select};
38use tonic::{Response, Status};
39use tracing::{debug, trace, warn};
40
41#[derive(Debug)]
43pub struct ChecksIPCConfiguration {
44 default_hostname: MetaString,
45 grpc_endpoint: ListenAddress,
46}
47
48impl ChecksIPCConfiguration {
49 pub fn new(grpc_endpoint: ListenAddress, default_hostname: impl Into<MetaString>) -> Self {
51 Self {
52 default_hostname: default_hostname.into(),
53 grpc_endpoint,
54 }
55 }
56}
57
58#[async_trait]
59impl SourceBuilder for ChecksIPCConfiguration {
60 fn outputs(&self) -> &[OutputDefinition<EventType>] {
61 static OUTPUTS: LazyLock<Vec<OutputDefinition<EventType>>> = LazyLock::new(|| {
62 vec![
63 OutputDefinition::named_output("metrics", EventType::Metric),
64 OutputDefinition::named_output("logs", EventType::Log),
65 OutputDefinition::named_output("events", EventType::EventD),
66 OutputDefinition::named_output("service_checks", EventType::ServiceCheck),
67 ]
68 });
69
70 &OUTPUTS
71 }
72
73 async fn build(&self, _context: BuildContext) -> Result<Box<dyn Source + Send>, GenericError> {
74 Ok(Box::new(ChecksIPC {
75 grpc_endpoint: self.grpc_endpoint.clone(),
76 default_hostname: self.default_hostname.clone(),
77 }))
78 }
79}
80
81impl MemoryBounds for ChecksIPCConfiguration {
82 fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
83 builder.minimum().with_single_value::<ChecksIPC>("checks_ipc");
85 }
86}
87
88struct ChecksIPC {
89 grpc_endpoint: ListenAddress,
90 default_hostname: MetaString,
91}
92
93#[async_trait]
94impl Source for ChecksIPC {
95 async fn run(self: Box<Self>, mut context: SourceContext) -> Result<(), GenericError> {
96 let ChecksIPC {
97 grpc_endpoint,
98 default_hostname,
99 } = *self;
100
101 let global_shutdown = context.take_shutdown_handle();
102 pin!(global_shutdown);
103
104 let mut health = context.take_health_handle();
105
106 let (events_tx, mut events_rx) = mpsc::channel(16);
107
108 let ListenAddress::Tcp(grpc_socket_addr) = grpc_endpoint else {
109 return Err(generic_error!("Checks IPC gRPC endpoint must be a TCP address."));
110 };
111
112 let grpc_server = HttpServer::from_listen_address(ListenAddress::Tcp(grpc_socket_addr))
115 .add_grpc_service(ChecksServer::new(ChecksService {
116 events_tx,
117 default_hostname,
118 }))
119 .with_http2_only()
120 .with_http2_config(Http2Config::grpc_defaults())
121 .with_worker_pool(context.topology_context().global_thread_pool().clone());
122
123 runtime::nested_supervisor(grpc_server.into_supervisor()).spawn();
124
125 health.mark_ready();
126 debug!("Checks IPC source started.");
127
128 loop {
129 select! {
130 _ = &mut global_shutdown => {
131 debug!("Received shutdown signal.");
132 break;
133 },
134 _ = health.live() => continue,
135 Some(event) = events_rx.recv() => {
136 let output_name = match &event {
137 Event::Metric(_) => "metrics",
138 Event::Log(_) => "logs",
139 Event::EventD(_) => "events",
140 Event::ServiceCheck(_) => "service_checks",
141 _ => continue,
142 };
143
144 if let Err(e) = context.dispatcher().dispatch_one_named(output_name, event).await {
145 warn!("Failed to dispatch {output_name} event: {:?}", e);
146 }
147 },
148 }
149 }
150
151 debug!("Checks IPC source stopped.");
152 Ok(())
153 }
154}
155
156struct ChecksService {
157 events_tx: mpsc::Sender<Event>,
158 default_hostname: MetaString,
159}
160
161#[async_trait]
162impl Checks for ChecksService {
163 async fn send_check_payload(
164 &self, request: tonic::Request<SendCheckPayloadRequest>,
165 ) -> Result<Response<SendCheckPayloadResponse>, Status> {
166 trace!("Received check payload.");
167
168 let payload = request.into_inner();
169 for check_data in payload.data.into_iter().filter_map(|data| data.data) {
170 let Some(event) = check_data_to_event(check_data, &self.default_hostname) else {
171 continue;
172 };
173
174 if let Err(e) = self.events_tx.send(event).await {
175 warn!("Failed to send check event: {:?}", e);
176 }
177 }
178
179 Ok(Response::new(SendCheckPayloadResponse {}))
180 }
181}
182
183fn check_data_to_event(check_data: Data, default_hostname: &MetaString) -> Option<Event> {
184 match check_data {
187 Data::Metric(metric) => {
188 let ProtoMetric {
189 r#type,
190 name,
191 value,
192 timestamp,
193 tags,
194 hostname,
195 interval_secs,
196 } = metric;
197
198 let metric_type = MetricType::try_from(r#type).ok()?;
199
200 let tags = tags.into_iter().map(Tag::from).collect::<TagSet>();
201 let mut context = Context::from_parts(name, tags.into_shared());
202 let hostname = if hostname.is_empty() {
203 default_hostname.clone()
204 } else {
205 MetaString::from(hostname)
206 };
207 context = context.with_host(Some(hostname));
208 let metric = match metric_type {
209 MetricType::Counter => Metric::counter(context, (timestamp, value)),
210 MetricType::Gauge => Metric::gauge(context, (timestamp, value)),
211 MetricType::Rate => {
212 if interval_secs == 0 {
213 warn!("Received rate metric from check with interval of zero. Skipping.");
214 return None;
215 }
216 Metric::rate(context, (timestamp, value), Duration::from_secs(interval_secs))
217 }
218 MetricType::Histogram => Metric::histogram(context, (timestamp, value)),
219 MetricType::Unspecified => {
220 warn!("Received metric with unspecified type. Skipping.");
221 return None;
222 }
223 };
224 Some(Event::Metric(metric))
225 }
226 Data::Log(log) => {
227 let ProtoLog { message, level } = log;
228
229 let level = LogLevel::try_from(level).ok()?;
230 let status = log_level_to_log_status(level);
231
232 Some(Event::Log(Log::new(message).with_status(status)))
233 }
234 Data::Event(event) => {
235 let ProtoEvent {
236 title,
237 text,
238 priority,
239 hostname,
240 tags,
241 alert_type,
242 aggregation_key,
243 source_type_name,
244 timestamp,
245 } = event;
246
247 let tags = tags.into_iter().map(Tag::from).collect::<TagSet>();
248 let mut eventd = EventD::new(title, text)
249 .with_timestamp(timestamp)
250 .with_tags(tags.into_shared());
251
252 if !hostname.is_empty() {
253 eventd.set_hostname(MetaString::from(hostname));
254 }
255 if !aggregation_key.is_empty() {
256 eventd.set_aggregation_key(MetaString::from(aggregation_key));
257 }
258 if !source_type_name.is_empty() {
259 eventd.set_source_type_name(MetaString::from(source_type_name));
260 }
261 if let Some(p) = ProtoPriority::try_from(priority)
262 .ok()
263 .and_then(proto_priority_to_priority)
264 {
265 eventd.set_priority(p);
266 }
267 if let Some(a) = ProtoAlertType::try_from(alert_type)
268 .ok()
269 .and_then(proto_alert_type_to_alert_type)
270 {
271 eventd.set_alert_type(a);
272 }
273 Some(Event::EventD(eventd))
274 }
275 Data::ServiceCheck(sc) => {
276 let ProtoServiceCheck {
277 status,
278 name,
279 message,
280 tags,
281 hostname,
282 } = sc;
283
284 let Some(status) = ServiceCheckStatus::try_from(status)
285 .ok()
286 .and_then(service_check_status_to_check_status)
287 else {
288 warn!(
289 "Received service check with unspecified or invalid status: {}. Skipping.",
290 status
291 );
292 return None;
293 };
294 let tags = tags.into_iter().map(Tag::from).collect::<TagSet>();
295 let mut service_check = ServiceCheck::new(name, status)
296 .with_message(MetaString::from(message))
297 .with_tags(tags.into_shared());
298 if !hostname.is_empty() {
299 service_check.set_hostname(MetaString::from(hostname));
300 }
301 Some(Event::ServiceCheck(service_check))
302 }
303 }
304}
305
306fn log_level_to_log_status(log_level: LogLevel) -> LogStatus {
307 match log_level {
308 LogLevel::Trace => LogStatus::Trace,
309 LogLevel::Debug => LogStatus::Debug,
310 LogLevel::Info => LogStatus::Info,
311 LogLevel::Warning => LogStatus::Warning,
312 LogLevel::Error => LogStatus::Error,
313 LogLevel::Critical => LogStatus::Emergency,
314 _ => LogStatus::Info,
315 }
316}
317
318fn service_check_status_to_check_status(status: ServiceCheckStatus) -> Option<CheckStatus> {
319 match status {
320 ServiceCheckStatus::Ok => Some(CheckStatus::Ok),
321 ServiceCheckStatus::Warning => Some(CheckStatus::Warning),
322 ServiceCheckStatus::Critical => Some(CheckStatus::Critical),
323 ServiceCheckStatus::Unknown => Some(CheckStatus::Unknown),
324 ServiceCheckStatus::Unspecified => None,
325 }
326}
327
328fn proto_priority_to_priority(priority: ProtoPriority) -> Option<Priority> {
329 match priority {
330 ProtoPriority::Normal => Some(Priority::Normal),
331 ProtoPriority::Low => Some(Priority::Low),
332 ProtoPriority::Unspecified => None,
333 }
334}
335
336fn proto_alert_type_to_alert_type(alert_type: ProtoAlertType) -> Option<AlertType> {
337 match alert_type {
338 ProtoAlertType::Info => Some(AlertType::Info),
339 ProtoAlertType::Error => Some(AlertType::Error),
340 ProtoAlertType::Warning => Some(AlertType::Warning),
341 ProtoAlertType::Success => Some(AlertType::Success),
342 ProtoAlertType::Unspecified => None,
343 }
344}
345
346#[cfg(test)]
347mod tests {
348 use datadog_protos::checks::{
349 check_data::Data,
350 event::Event as ProtoEvent,
351 log::Log as ProtoLog,
352 metric::{Metric as ProtoMetric, MetricType as ProtoMetricType},
353 service_check::{ServiceCheck as ProtoServiceCheck, Status as ProtoServiceCheckStatus},
354 };
355 use saluki_core::data_model::event::metric::MetricValues;
356
357 use super::*;
358
359 fn metric_data(
360 r#type: i32, name: &str, value: f64, timestamp: u64, interval_secs: u64, tags: &[&str], hostname: &str,
361 ) -> Data {
362 Data::Metric(ProtoMetric {
363 r#type,
364 name: name.to_string(),
365 value,
366 timestamp,
367 tags: tags.iter().map(|t| (*t).to_string()).collect(),
368 hostname: hostname.to_string(),
369 interval_secs,
370 })
371 }
372
373 fn log_data(level: i32, message: &str) -> Data {
374 Data::Log(ProtoLog {
375 message: message.to_string(),
376 level,
377 })
378 }
379
380 fn event_data(title: &str, text: &str, timestamp: u64, tags: &[&str], hostname: &str) -> Data {
381 Data::Event(ProtoEvent {
382 title: title.to_string(),
383 text: text.to_string(),
384 priority: 0,
385 hostname: hostname.to_string(),
386 tags: tags.iter().map(|t| (*t).to_string()).collect(),
387 alert_type: 0,
388 aggregation_key: String::new(),
389 source_type_name: String::new(),
390 timestamp,
391 })
392 }
393
394 fn service_check_data(status: i32, name: &str, message: &str, tags: &[&str], hostname: &str) -> Data {
395 Data::ServiceCheck(ProtoServiceCheck {
396 status,
397 name: name.to_string(),
398 message: message.to_string(),
399 tags: tags.iter().map(|t| (*t).to_string()).collect(),
400 hostname: hostname.to_string(),
401 })
402 }
403
404 fn check_data_to_event_for_tests(check_data: Data) -> Option<Event> {
405 check_data_to_event(check_data, &MetaString::from_static("default-host"))
406 }
407
408 #[test]
409 fn metric_counter_conversion() {
410 let event = check_data_to_event_for_tests(metric_data(
411 ProtoMetricType::Counter as i32,
412 "my_counter",
413 1.0,
414 1234,
415 0,
416 &["tag1:value1", "tag2:value2"],
417 "",
418 ))
419 .expect("counter should convert");
420
421 let Event::Metric(metric) = event else {
422 panic!("expected Metric event");
423 };
424 assert_eq!(metric.context().name().as_ref(), "my_counter");
425 assert!(metric.context().tags().has_tag("tag1:value1"));
426 assert!(metric.context().tags().has_tag("tag2:value2"));
427 assert!(matches!(metric.values(), MetricValues::Counter(_)));
428 }
429
430 #[test]
431 fn metric_gauge_conversion() {
432 let event = check_data_to_event_for_tests(metric_data(
433 ProtoMetricType::Gauge as i32,
434 "my_gauge",
435 42.0,
436 1234,
437 0,
438 &[],
439 "",
440 ))
441 .expect("gauge should convert");
442 let Event::Metric(metric) = event else {
443 panic!("expected Metric event");
444 };
445 assert!(matches!(metric.values(), MetricValues::Gauge(_)));
446 }
447
448 #[test]
449 fn metric_histogram_conversion() {
450 let event = check_data_to_event_for_tests(metric_data(
451 ProtoMetricType::Histogram as i32,
452 "my_hist",
453 1.0,
454 1234,
455 0,
456 &[],
457 "",
458 ))
459 .expect("histogram should convert");
460 let Event::Metric(metric) = event else {
461 panic!("expected Metric event");
462 };
463 assert!(matches!(metric.values(), MetricValues::Histogram(_)));
464 }
465
466 #[test]
467 fn metric_rate_conversion_uses_interval() {
468 let event = check_data_to_event_for_tests(metric_data(
469 ProtoMetricType::Rate as i32,
470 "my_rate",
471 10.0,
472 1234,
473 60,
474 &[],
475 "",
476 ))
477 .expect("rate should convert");
478 let Event::Metric(metric) = event else {
479 panic!("expected Metric event");
480 };
481 match metric.values() {
482 MetricValues::Rate(_, interval) => assert_eq!(*interval, Duration::from_secs(60)),
483 other => panic!("expected Rate values, got {other:?}"),
484 }
485 }
486
487 #[test]
488 fn metric_rate_with_zero_interval_is_skipped() {
489 let event = check_data_to_event_for_tests(metric_data(
490 ProtoMetricType::Rate as i32,
491 "my_rate",
492 10.0,
493 1234,
494 0,
495 &[],
496 "",
497 ));
498 assert!(event.is_none(), "rate with zero interval must be skipped");
499 }
500
501 #[test]
502 fn metric_unspecified_type_is_skipped() {
503 let event = check_data_to_event_for_tests(metric_data(
504 ProtoMetricType::Unspecified as i32,
505 "x",
506 1.0,
507 1234,
508 0,
509 &[],
510 "",
511 ));
512 assert!(event.is_none(), "unspecified metric type must be skipped");
513 }
514
515 #[test]
516 fn metric_unknown_type_is_skipped() {
517 let event = check_data_to_event_for_tests(metric_data(99, "x", 1.0, 1234, 0, &[], ""));
519 assert!(event.is_none(), "unknown metric type must be skipped");
520 }
521
522 #[test]
523 fn log_unknown_level_is_skipped() {
524 let event = check_data_to_event_for_tests(log_data(99, "hello"));
526 assert!(event.is_none(), "unknown log level must be skipped");
527 }
528
529 #[test]
530 fn event_conversion_preserves_fields() {
531 let event = check_data_to_event_for_tests(event_data("title", "body", 1234, &["env:prod", "team:foo"], ""))
532 .expect("event should convert");
533 let Event::EventD(ev) = event else {
534 panic!("expected EventD event");
535 };
536 assert_eq!(ev.title(), "title");
537 assert_eq!(ev.text(), "body");
538 assert_eq!(ev.timestamp(), Some(1234));
539 assert!(ev.tags().has_tag("env:prod"));
540 assert!(ev.tags().has_tag("team:foo"));
541 }
542
543 #[test]
544 fn service_check_status_mapping() {
545 let cases = [
546 (ProtoServiceCheckStatus::Ok, CheckStatus::Ok),
547 (ProtoServiceCheckStatus::Warning, CheckStatus::Warning),
548 (ProtoServiceCheckStatus::Critical, CheckStatus::Critical),
549 (ProtoServiceCheckStatus::Unknown, CheckStatus::Unknown),
550 ];
551
552 for (proto_status, expected) in cases {
553 let event = check_data_to_event_for_tests(service_check_data(proto_status as i32, "n", "m", &[], ""))
554 .unwrap_or_else(|| panic!("status {proto_status:?} should convert"));
555 let Event::ServiceCheck(sc) = event else {
556 panic!("expected ServiceCheck event for {proto_status:?}");
557 };
558 assert_eq!(sc.status(), expected, "status {proto_status:?}");
559 }
560 }
561
562 #[test]
563 fn service_check_unspecified_status_is_skipped() {
564 let event = check_data_to_event_for_tests(service_check_data(
565 ProtoServiceCheckStatus::Unspecified as i32,
566 "n",
567 "m",
568 &[],
569 "",
570 ));
571 assert!(event.is_none(), "service check with unspecified status must be skipped");
572 }
573
574 #[test]
575 fn service_check_unknown_status_value_is_skipped() {
576 let event = check_data_to_event_for_tests(service_check_data(99, "n", "m", &[], ""));
578 assert!(
579 event.is_none(),
580 "service check with out-of-range status must be skipped"
581 );
582 }
583
584 #[test]
585 fn service_check_preserves_name_message_and_tags() {
586 let event = check_data_to_event_for_tests(service_check_data(
587 ProtoServiceCheckStatus::Ok as i32,
588 "my.check",
589 "all good",
590 &["env:prod"],
591 "",
592 ))
593 .expect("service check should convert");
594 let Event::ServiceCheck(sc) = event else {
595 panic!("expected ServiceCheck event");
596 };
597 assert_eq!(sc.name(), "my.check");
598 assert_eq!(sc.status(), CheckStatus::Ok);
599 assert_eq!(sc.message(), Some("all good"));
600 assert!(sc.tags().has_tag("env:prod"));
601 }
602
603 #[test]
604 fn metric_hostname_propagates() {
605 let event = check_data_to_event_for_tests(metric_data(
606 ProtoMetricType::Counter as i32,
607 "n",
608 1.0,
609 0,
610 0,
611 &[],
612 "host-a",
613 ))
614 .expect("metric should convert");
615 let Event::Metric(m) = event else {
616 panic!("expected Metric event");
617 };
618 assert_eq!(m.context().host(), Some("host-a"));
619 }
620
621 #[test]
622 fn metric_empty_hostname_uses_default_host() {
623 let event =
624 check_data_to_event_for_tests(metric_data(ProtoMetricType::Counter as i32, "n", 1.0, 0, 0, &[], ""))
625 .expect("metric should convert");
626 let Event::Metric(m) = event else {
627 panic!("expected Metric event");
628 };
629 assert_eq!(m.context().host(), Some("default-host"));
630 }
631
632 #[test]
633 fn eventd_hostname_propagates() {
634 let event =
635 check_data_to_event_for_tests(event_data("title", "body", 0, &[], "host-b")).expect("event should convert");
636 let Event::EventD(ev) = event else {
637 panic!("expected EventD event");
638 };
639 assert_eq!(ev.hostname(), Some("host-b"));
640 }
641
642 #[test]
643 fn eventd_empty_hostname_stays_unset() {
644 let event =
645 check_data_to_event_for_tests(event_data("title", "body", 0, &[], "")).expect("event should convert");
646 let Event::EventD(ev) = event else {
647 panic!("expected EventD event");
648 };
649 assert_eq!(ev.hostname(), None);
650 }
651
652 #[test]
653 fn service_check_hostname_propagates() {
654 let event = check_data_to_event_for_tests(service_check_data(
655 ProtoServiceCheckStatus::Ok as i32,
656 "n",
657 "m",
658 &[],
659 "host-c",
660 ))
661 .expect("service check should convert");
662 let Event::ServiceCheck(sc) = event else {
663 panic!("expected ServiceCheck event");
664 };
665 assert_eq!(sc.hostname(), Some("host-c"));
666 }
667
668 #[test]
669 fn service_check_empty_hostname_stays_unset() {
670 let event = check_data_to_event_for_tests(service_check_data(
671 ProtoServiceCheckStatus::Ok as i32,
672 "n",
673 "m",
674 &[],
675 "",
676 ))
677 .expect("service check should convert");
678 let Event::ServiceCheck(sc) = event else {
679 panic!("expected ServiceCheck event");
680 };
681 assert_eq!(sc.hostname(), None);
682 }
683
684 #[test]
685 fn eventd_priority_propagates() {
686 let event = check_data_to_event_for_tests(Data::Event(ProtoEvent {
687 priority: ProtoPriority::Low as i32,
688 ..Default::default()
689 }))
690 .expect("event should convert");
691 let Event::EventD(ev) = event else {
692 panic!("expected EventD event");
693 };
694 assert_eq!(ev.priority(), Some(Priority::Low));
695 }
696
697 #[test]
698 fn eventd_alert_type_propagates() {
699 let event = check_data_to_event_for_tests(Data::Event(ProtoEvent {
700 alert_type: ProtoAlertType::Warning as i32,
701 ..Default::default()
702 }))
703 .expect("event should convert");
704 let Event::EventD(ev) = event else {
705 panic!("expected EventD event");
706 };
707 assert_eq!(ev.alert_type(), Some(AlertType::Warning));
708 }
709
710 #[test]
711 fn eventd_aggregation_key_propagates() {
712 let event = check_data_to_event_for_tests(Data::Event(ProtoEvent {
713 aggregation_key: "agg-key-1".to_string(),
714 ..Default::default()
715 }))
716 .expect("event should convert");
717 let Event::EventD(ev) = event else {
718 panic!("expected EventD event");
719 };
720 assert_eq!(ev.aggregation_key(), Some("agg-key-1"));
721 }
722
723 #[test]
724 fn eventd_source_type_name_propagates() {
725 let event = check_data_to_event_for_tests(Data::Event(ProtoEvent {
726 source_type_name: "my-source".to_string(),
727 ..Default::default()
728 }))
729 .expect("event should convert");
730 let Event::EventD(ev) = event else {
731 panic!("expected EventD event");
732 };
733 assert_eq!(ev.source_type_name(), Some("my-source"));
734 }
735
736 #[test]
737 fn eventd_unspecified_proto_keeps_saluki_defaults() {
738 let event = check_data_to_event_for_tests(Data::Event(ProtoEvent::default())).expect("event should convert");
743 let Event::EventD(ev) = event else {
744 panic!("expected EventD event");
745 };
746 assert_eq!(ev.priority(), Some(Priority::Normal));
747 assert_eq!(ev.alert_type(), Some(AlertType::Info));
748 assert_eq!(ev.aggregation_key(), None);
749 assert_eq!(ev.source_type_name(), None);
750 assert_eq!(ev.hostname(), None);
751 }
752}