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