saluki_components/sources/checks_ipc/
mod.rs

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