Skip to main content

vector/sources/
syslog.rs

1#[cfg(unix)]
2use std::path::PathBuf;
3use std::{net::SocketAddr, time::Duration};
4
5use bytes::Bytes;
6use chrono::Utc;
7use futures::StreamExt;
8use listenfd::ListenFd;
9use smallvec::SmallVec;
10use tokio_util::udp::UdpFramed;
11use vector_lib::{
12    EstimatedJsonEncodedSizeOf,
13    codecs::{
14        BytesDecoder, OctetCountingDecoder, SyslogDeserializerConfig,
15        decoding::{Deserializer, Framer},
16    },
17    config::{LegacyKey, LogNamespace},
18    configurable::configurable_component,
19    internal_event::{ByteSize, BytesReceived, InternalEventHandle as _, Protocol},
20    ipallowlist::IpAllowlistConfig,
21    lookup::{OwnedValuePath, lookup_v2::OptionalValuePath, path},
22};
23use vrl::event_path;
24
25#[cfg(unix)]
26use crate::sources::util::build_unix_stream_source;
27use crate::{
28    SourceSender,
29    codecs::Decoder,
30    config::{
31        DataType, GenerateConfig, Resource, SourceConfig, SourceContext, SourceOutput, log_schema,
32    },
33    event::Event,
34    internal_events::{
35        SocketBindError, SocketEventsReceived, SocketMode, SocketReceiveError, StreamClosedError,
36    },
37    net,
38    shutdown::ShutdownSignal,
39    sources::util::net::{SocketListenAddr, TcpNullAcker, TcpSource, try_bind_udp_socket},
40    tcp::TcpKeepaliveConfig,
41    tls::{MaybeTlsSettings, TlsSourceConfig},
42};
43
44/// Configuration for the `syslog` source.
45#[configurable_component(source("syslog", "Collect logs sent via Syslog."))]
46#[derive(Clone, Debug)]
47pub struct SyslogConfig {
48    #[serde(flatten)]
49    mode: Mode,
50
51    /// The maximum buffer size of incoming messages, in bytes.
52    ///
53    /// Messages larger than this are truncated.
54    #[serde(default = "crate::serde::default_max_length")]
55    #[configurable(metadata(docs::type_unit = "bytes"))]
56    max_length: usize,
57
58    /// Overrides the name of the log field used to add the peer host to each event.
59    ///
60    /// If using TCP or UDP, the value is the peer host's address, including the port. For example, `1.2.3.4:9000`. If using
61    /// UDS, the value is the socket path itself.
62    ///
63    /// By default, the [global `log_schema.host_key` option][global_host_key] is used.
64    ///
65    /// [global_host_key]: https://vector.dev/docs/reference/configuration/global-options/#log_schema.host_key
66    host_key: Option<OptionalValuePath>,
67
68    /// The namespace to use for logs. This overrides the global setting.
69    #[configurable(metadata(docs::hidden))]
70    #[serde(default)]
71    pub log_namespace: Option<bool>,
72}
73
74/// Listener mode for the `syslog` source.
75#[configurable_component]
76#[derive(Clone, Debug)]
77#[serde(tag = "mode", rename_all = "snake_case")]
78#[configurable(metadata(docs::enum_tag_description = "The type of socket to use."))]
79#[allow(clippy::large_enum_variant)]
80pub enum Mode {
81    /// Listen on TCP.
82    Tcp {
83        #[configurable(derived)]
84        address: SocketListenAddr,
85
86        #[configurable(derived)]
87        keepalive: Option<TcpKeepaliveConfig>,
88
89        #[configurable(derived)]
90        permit_origin: Option<IpAllowlistConfig>,
91
92        #[configurable(derived)]
93        tls: Option<TlsSourceConfig>,
94
95        /// The size of the receive buffer used for each connection.
96        ///
97        /// This should not typically needed to be changed.
98        #[configurable(metadata(docs::type_unit = "bytes"))]
99        receive_buffer_bytes: Option<usize>,
100
101        /// The maximum number of TCP connections that are allowed at any given time.
102        connection_limit: Option<u32>,
103    },
104
105    /// Listen on UDP.
106    Udp {
107        #[configurable(derived)]
108        address: SocketListenAddr,
109
110        /// The size of the receive buffer used for the listening socket.
111        ///
112        /// This should not typically needed to be changed.
113        #[configurable(metadata(docs::type_unit = "bytes"))]
114        receive_buffer_bytes: Option<usize>,
115    },
116
117    /// Listen on UDS (Unix domain socket). This only supports Unix stream sockets.
118    ///
119    /// For Unix datagram sockets, use the `socket` source instead.
120    #[cfg(unix)]
121    Unix {
122        /// The Unix socket path.
123        ///
124        /// This should be an absolute path.
125        #[configurable(metadata(docs::examples = "/path/to/socket"))]
126        path: PathBuf,
127
128        /// Unix file mode bits to be applied to the unix socket file as its designated file permissions.
129        ///
130        /// The file mode value can be specified in any numeric format supported by your configuration
131        /// language, but it is most intuitive to use an octal number.
132        socket_file_mode: Option<u32>,
133    },
134}
135
136impl SyslogConfig {
137    #[cfg(test)]
138    pub fn from_mode(mode: Mode) -> Self {
139        Self {
140            mode,
141            host_key: None,
142            max_length: crate::serde::default_max_length(),
143            log_namespace: None,
144        }
145    }
146}
147
148impl Default for SyslogConfig {
149    fn default() -> Self {
150        Self {
151            mode: Mode::Tcp {
152                address: SocketListenAddr::SocketAddr("0.0.0.0:514".parse().unwrap()),
153                keepalive: None,
154                permit_origin: None,
155                tls: None,
156                receive_buffer_bytes: None,
157                connection_limit: None,
158            },
159            host_key: None,
160            max_length: crate::serde::default_max_length(),
161            log_namespace: None,
162        }
163    }
164}
165
166impl GenerateConfig for SyslogConfig {
167    fn generate_config() -> toml::Value {
168        toml::Value::try_from(SyslogConfig::default()).unwrap()
169    }
170}
171
172#[async_trait::async_trait]
173#[typetag::serde(name = "syslog")]
174impl SourceConfig for SyslogConfig {
175    async fn build(&self, cx: SourceContext) -> crate::Result<super::Source> {
176        let log_namespace = cx.log_namespace(self.log_namespace);
177        let host_key = self
178            .host_key
179            .clone()
180            .and_then(|k| k.path)
181            .or(log_schema().host_key().cloned());
182
183        match self.mode.clone() {
184            Mode::Tcp {
185                address,
186                keepalive,
187                permit_origin,
188                tls,
189                receive_buffer_bytes,
190                connection_limit,
191            } => {
192                let source = SyslogTcpSource {
193                    max_length: self.max_length,
194                    host_key,
195                    log_namespace,
196                };
197                let shutdown_secs = Duration::from_secs(30);
198                let tls_config = tls.as_ref().map(|tls| tls.tls_config.clone());
199                let tls_client_metadata_key = tls
200                    .as_ref()
201                    .and_then(|tls| tls.client_metadata_key.clone())
202                    .and_then(|k| k.path);
203                let tls = MaybeTlsSettings::from_config(tls_config.as_ref(), true)?;
204                source.run(
205                    address,
206                    keepalive,
207                    shutdown_secs,
208                    tls,
209                    None, // tls_reloader: not wired for this source
210                    tls_client_metadata_key,
211                    receive_buffer_bytes,
212                    None,
213                    cx,
214                    false.into(),
215                    connection_limit,
216                    permit_origin.map(Into::into),
217                    SyslogConfig::NAME,
218                    log_namespace,
219                )
220            }
221            Mode::Udp {
222                address,
223                receive_buffer_bytes,
224            } => Ok(udp(
225                address,
226                self.max_length,
227                host_key,
228                receive_buffer_bytes,
229                cx.shutdown,
230                log_namespace,
231                cx.out,
232            )),
233            #[cfg(unix)]
234            Mode::Unix {
235                path,
236                socket_file_mode,
237            } => {
238                let decoder = Decoder::new(
239                    Framer::OctetCounting(OctetCountingDecoder::new_with_max_length(
240                        self.max_length,
241                    )),
242                    Deserializer::Syslog(
243                        SyslogDeserializerConfig::from_source(SyslogConfig::NAME).build(),
244                    ),
245                );
246
247                build_unix_stream_source(
248                    path,
249                    socket_file_mode,
250                    decoder,
251                    move |events, host| handle_events(events, &host_key, host, log_namespace),
252                    cx.shutdown,
253                    cx.out,
254                )
255            }
256        }
257    }
258
259    fn outputs(&self, global_log_namespace: LogNamespace) -> Vec<SourceOutput> {
260        let log_namespace = global_log_namespace.merge(self.log_namespace);
261        let schema_definition = SyslogDeserializerConfig::from_source(SyslogConfig::NAME)
262            .schema_definition(log_namespace)
263            .with_standard_vector_source_metadata();
264
265        vec![SourceOutput::new_maybe_logs(
266            DataType::Log,
267            schema_definition,
268        )]
269    }
270
271    fn resources(&self) -> Vec<Resource> {
272        match self.mode.clone() {
273            Mode::Tcp { address, .. } => vec![address.as_tcp_resource()],
274            Mode::Udp { address, .. } => vec![address.as_udp_resource()],
275            #[cfg(unix)]
276            Mode::Unix { .. } => vec![],
277        }
278    }
279
280    fn can_acknowledge(&self) -> bool {
281        false
282    }
283}
284
285#[derive(Debug, Clone)]
286struct SyslogTcpSource {
287    max_length: usize,
288    host_key: Option<OwnedValuePath>,
289    log_namespace: LogNamespace,
290}
291
292impl TcpSource for SyslogTcpSource {
293    type Error = vector_lib::codecs::decoding::Error;
294    type Item = SmallVec<[Event; 1]>;
295    type Decoder = Decoder;
296    type Acker = TcpNullAcker;
297
298    fn decoder(&self) -> Self::Decoder {
299        Decoder::new(
300            Framer::OctetCounting(OctetCountingDecoder::new_with_max_length(self.max_length)),
301            Deserializer::Syslog(SyslogDeserializerConfig::from_source(SyslogConfig::NAME).build()),
302        )
303    }
304
305    fn handle_events(&self, events: &mut [Event], host: SocketAddr) {
306        handle_events(
307            events,
308            &self.host_key,
309            Some(host.ip().to_string().into()),
310            self.log_namespace,
311        );
312    }
313
314    fn build_acker(&self, _: &[Self::Item]) -> Self::Acker {
315        TcpNullAcker
316    }
317}
318
319pub fn udp(
320    addr: SocketListenAddr,
321    _max_length: usize,
322    host_key: Option<OwnedValuePath>,
323    receive_buffer_bytes: Option<usize>,
324    shutdown: ShutdownSignal,
325    log_namespace: LogNamespace,
326    mut out: SourceSender,
327) -> super::Source {
328    Box::pin(async move {
329        let listenfd = ListenFd::from_env();
330        let socket = try_bind_udp_socket(addr, listenfd).await.map_err(|error| {
331            emit!(SocketBindError {
332                mode: SocketMode::Udp,
333                error: &error,
334            })
335        })?;
336
337        if let Some(receive_buffer_bytes) = receive_buffer_bytes
338            && let Err(error) = net::set_receive_buffer_size(&socket, receive_buffer_bytes)
339        {
340            warn!(message = "Failed configuring receive buffer size on UDP socket.", %error);
341        }
342
343        info!(
344            message = "Listening.",
345            addr = %addr,
346            r#type = "udp"
347        );
348
349        let bytes_received = register!(BytesReceived::from(Protocol::UDP));
350
351        let mut stream = UdpFramed::new(
352            socket,
353            Decoder::new(
354                Framer::Bytes(BytesDecoder::new()),
355                Deserializer::Syslog(
356                    SyslogDeserializerConfig::from_source(SyslogConfig::NAME).build(),
357                ),
358            ),
359        )
360        .take_until(shutdown)
361        .filter_map(|frame| {
362            let host_key = host_key.clone();
363            let bytes_received = bytes_received.clone();
364            async move {
365                match frame {
366                    Ok(((mut events, byte_size), received_from)) => {
367                        let count = events.len();
368                        bytes_received.emit(ByteSize(byte_size));
369                        emit!(SocketEventsReceived {
370                            mode: SocketMode::Udp,
371                            byte_size: events.estimated_json_encoded_size_of(),
372                            count,
373                        });
374                        let received_from = received_from.ip().to_string().into();
375                        handle_events(&mut events, &host_key, Some(received_from), log_namespace);
376                        Some(events.remove(0))
377                    }
378                    Err(error) => {
379                        emit!(SocketReceiveError {
380                            mode: SocketMode::Udp,
381                            error: &error,
382                        });
383                        None
384                    }
385                }
386            }
387        })
388        .boxed();
389
390        match out.send_event_stream(&mut stream).await {
391            Ok(()) => {
392                debug!("Finished sending.");
393                Ok(())
394            }
395            Err(_) => {
396                let (count, _) = stream.size_hint();
397                emit!(StreamClosedError { count });
398                Err(())
399            }
400        }
401    })
402}
403
404fn handle_events(
405    events: &mut [Event],
406    host_key: &Option<OwnedValuePath>,
407    default_host: Option<Bytes>,
408    log_namespace: LogNamespace,
409) {
410    for event in events {
411        enrich_syslog_event(event, host_key, default_host.clone(), log_namespace);
412    }
413}
414
415fn enrich_syslog_event(
416    event: &mut Event,
417    host_key: &Option<OwnedValuePath>,
418    default_host: Option<Bytes>,
419    log_namespace: LogNamespace,
420) {
421    let log = event.as_mut_log();
422
423    if let Some(default_host) = &default_host {
424        log_namespace.insert_source_metadata(
425            SyslogConfig::NAME,
426            log,
427            Some(LegacyKey::Overwrite(path!("source_ip"))),
428            path!("source_ip"),
429            default_host.clone(),
430        );
431    }
432
433    let parsed_hostname = log
434        .get(event_path!("hostname"))
435        .map(|hostname| hostname.coerce_to_bytes());
436
437    if let Some(parsed_host) = parsed_hostname.or(default_host) {
438        let legacy_host_key = host_key.as_ref().map(LegacyKey::Overwrite);
439
440        log_namespace.insert_source_metadata(
441            SyslogConfig::NAME,
442            log,
443            legacy_host_key,
444            path!("host"),
445            parsed_host,
446        );
447    }
448
449    log_namespace.insert_standard_vector_source_metadata(log, SyslogConfig::NAME, Utc::now());
450
451    if log_namespace == LogNamespace::Legacy {
452        let timestamp = log
453            .get(event_path!("timestamp"))
454            .and_then(|timestamp| timestamp.as_timestamp().cloned())
455            .unwrap_or_else(Utc::now);
456        log.maybe_insert(log_schema().timestamp_key_target_path(), timestamp);
457    }
458
459    trace!(
460        message = "Processing one event.",
461        event = ?event
462    );
463}
464
465#[cfg(test)]
466mod test {
467    use std::{collections::HashMap, fmt, str::FromStr};
468
469    use chrono::prelude::*;
470    use indoc::indoc;
471    use rand::{Rng, rng};
472    use serde::Deserialize;
473    use tokio::time::{Duration, Instant, sleep};
474    use tokio_util::codec::BytesCodec;
475    use vector_lib::{
476        assert_event_data_eq,
477        codecs::decoding::format::Deserializer,
478        config::ComponentKey,
479        lookup::{OwnedTargetPath, PathPrefix, event_path, owned_value_path},
480        schema::Definition,
481    };
482    use vrl::value::{Kind, ObjectMap, Value, kind::Collection};
483
484    use super::*;
485    use crate::{
486        config::log_schema,
487        event::{Event, LogEvent},
488        test_util::{
489            CountReceiver,
490            addr::next_addr,
491            components::{SOCKET_PUSH_SOURCE_TAGS, assert_source_compliance},
492            random_maps, random_string, send_encodable, send_lines, wait_for_tcp,
493        },
494    };
495
496    fn event_from_bytes(
497        host_key: &str,
498        default_host: Option<Bytes>,
499        bytes: Bytes,
500        log_namespace: LogNamespace,
501    ) -> Option<Event> {
502        let parser = SyslogDeserializerConfig::from_source(SyslogConfig::NAME).build();
503        let mut events = parser.parse(bytes, LogNamespace::Legacy).ok()?;
504        handle_events(
505            &mut events,
506            &Some(owned_value_path!(host_key)),
507            default_host,
508            log_namespace,
509        );
510        Some(events.remove(0))
511    }
512
513    #[test]
514    fn generate_config() {
515        crate::test_util::test_generate_config::<SyslogConfig>();
516    }
517
518    #[test]
519    fn output_schema_definition_vector_namespace() {
520        let config = SyslogConfig {
521            log_namespace: Some(true),
522            ..Default::default()
523        };
524
525        let definitions = config
526            .outputs(LogNamespace::Vector)
527            .remove(0)
528            .schema_definition(true);
529
530        let expected_definition =
531            Definition::new_with_default_metadata(Kind::bytes(), [LogNamespace::Vector])
532                .with_meaning(OwnedTargetPath::event_root(), "message")
533                .with_metadata_field(
534                    &owned_value_path!("vector", "source_type"),
535                    Kind::bytes(),
536                    None,
537                )
538                .with_metadata_field(
539                    &owned_value_path!("vector", "ingest_timestamp"),
540                    Kind::timestamp(),
541                    None,
542                )
543                .with_metadata_field(
544                    &owned_value_path!("syslog", "timestamp"),
545                    Kind::timestamp(),
546                    Some("timestamp"),
547                )
548                .with_metadata_field(
549                    &owned_value_path!("syslog", "hostname"),
550                    Kind::bytes().or_undefined(),
551                    Some("host"),
552                )
553                .with_metadata_field(
554                    &owned_value_path!("syslog", "source_ip"),
555                    Kind::bytes().or_undefined(),
556                    None,
557                )
558                .with_metadata_field(
559                    &owned_value_path!("syslog", "severity"),
560                    Kind::bytes().or_undefined(),
561                    Some("severity"),
562                )
563                .with_metadata_field(
564                    &owned_value_path!("syslog", "facility"),
565                    Kind::bytes().or_undefined(),
566                    None,
567                )
568                .with_metadata_field(
569                    &owned_value_path!("syslog", "version"),
570                    Kind::integer().or_undefined(),
571                    None,
572                )
573                .with_metadata_field(
574                    &owned_value_path!("syslog", "appname"),
575                    Kind::bytes().or_undefined(),
576                    Some("service"),
577                )
578                .with_metadata_field(
579                    &owned_value_path!("syslog", "msgid"),
580                    Kind::bytes().or_undefined(),
581                    None,
582                )
583                .with_metadata_field(
584                    &owned_value_path!("syslog", "procid"),
585                    Kind::integer().or_bytes().or_undefined(),
586                    None,
587                )
588                .with_metadata_field(
589                    &owned_value_path!("syslog", "structured_data"),
590                    Kind::object(Collection::from_unknown(Kind::object(
591                        Collection::from_unknown(Kind::bytes()),
592                    ))),
593                    None,
594                )
595                .with_metadata_field(
596                    &owned_value_path!("syslog", "tls_client_metadata"),
597                    Kind::object(Collection::empty().with_unknown(Kind::bytes())).or_undefined(),
598                    None,
599                );
600
601        assert_eq!(definitions, Some(expected_definition));
602    }
603
604    #[test]
605    fn output_schema_definition_legacy_namespace() {
606        let config = SyslogConfig::default();
607
608        let definitions = config
609            .outputs(LogNamespace::Legacy)
610            .remove(0)
611            .schema_definition(true);
612
613        let expected_definition = Definition::new_with_default_metadata(
614            Kind::object(Collection::empty()),
615            [LogNamespace::Legacy],
616        )
617        .with_event_field(
618            &owned_value_path!("message"),
619            Kind::bytes(),
620            Some("message"),
621        )
622        .with_event_field(
623            &owned_value_path!("timestamp"),
624            Kind::timestamp(),
625            Some("timestamp"),
626        )
627        .with_event_field(
628            &owned_value_path!("hostname"),
629            Kind::bytes().or_undefined(),
630            Some("host"),
631        )
632        .with_event_field(
633            &owned_value_path!("source_ip"),
634            Kind::bytes().or_undefined(),
635            None,
636        )
637        .with_event_field(
638            &owned_value_path!("severity"),
639            Kind::bytes().or_undefined(),
640            Some("severity"),
641        )
642        .with_event_field(
643            &owned_value_path!("facility"),
644            Kind::bytes().or_undefined(),
645            None,
646        )
647        .with_event_field(
648            &owned_value_path!("version"),
649            Kind::integer().or_undefined(),
650            None,
651        )
652        .with_event_field(
653            &owned_value_path!("appname"),
654            Kind::bytes().or_undefined(),
655            Some("service"),
656        )
657        .with_event_field(
658            &owned_value_path!("msgid"),
659            Kind::bytes().or_undefined(),
660            None,
661        )
662        .with_event_field(
663            &owned_value_path!("procid"),
664            Kind::integer().or_bytes().or_undefined(),
665            None,
666        )
667        .unknown_fields(Kind::object(Collection::from_unknown(Kind::bytes())))
668        .with_standard_vector_source_metadata();
669
670        assert_eq!(definitions, Some(expected_definition));
671    }
672
673    #[test]
674    fn config_tcp() {
675        let config: SyslogConfig = serde_yaml::from_str(indoc! {
676            r#"
677            mode: tcp
678            address: "127.0.0.1:1235"
679            "#,
680        })
681        .unwrap();
682        assert!(matches!(config.mode, Mode::Tcp { .. }));
683    }
684
685    #[test]
686    fn config_tcp_with_receive_buffer_size() {
687        let config: SyslogConfig = serde_yaml::from_str(indoc! {
688            r#"
689            mode: tcp
690            address: "127.0.0.1:1235"
691            receive_buffer_bytes: 256
692            "#,
693        })
694        .unwrap();
695
696        let receive_buffer_bytes = match config.mode {
697            Mode::Tcp {
698                receive_buffer_bytes,
699                ..
700            } => receive_buffer_bytes,
701            _ => panic!("expected Mode::Tcp"),
702        };
703
704        assert_eq!(receive_buffer_bytes, Some(256));
705    }
706
707    #[test]
708    fn config_tcp_keepalive_empty() {
709        let config: SyslogConfig = serde_yaml::from_str(indoc! {
710            r#"
711            mode: tcp
712            address: "127.0.0.1:1235"
713            "#,
714        })
715        .unwrap();
716
717        let keepalive = match config.mode {
718            Mode::Tcp { keepalive, .. } => keepalive,
719            _ => panic!("expected Mode::Tcp"),
720        };
721
722        assert_eq!(keepalive, None);
723    }
724
725    #[test]
726    fn config_tcp_keepalive_full() {
727        let config: SyslogConfig = serde_yaml::from_str(indoc! {
728            r#"
729            mode: tcp
730            address: "127.0.0.1:1235"
731            keepalive:
732              time_secs: 7200
733            "#,
734        })
735        .unwrap();
736
737        let keepalive = match config.mode {
738            Mode::Tcp { keepalive, .. } => keepalive,
739            _ => panic!("expected Mode::Tcp"),
740        };
741
742        let keepalive = keepalive.expect("keepalive config not set");
743
744        assert_eq!(keepalive.time_secs, Some(7200));
745    }
746
747    #[test]
748    fn config_udp() {
749        let config: SyslogConfig = serde_yaml::from_str(indoc! {
750            r#"
751            mode: udp
752            address: "127.0.0.1:1235"
753            max_length: 32187
754            "#,
755        })
756        .unwrap();
757        assert!(matches!(config.mode, Mode::Udp { .. }));
758    }
759
760    #[test]
761    fn config_udp_with_receive_buffer_size() {
762        let config: SyslogConfig = serde_yaml::from_str(indoc! {
763            r#"
764            mode: udp
765            address: "127.0.0.1:1235"
766            max_length: 32187
767            receive_buffer_bytes: 256
768            "#,
769        })
770        .unwrap();
771
772        let receive_buffer_bytes = match config.mode {
773            Mode::Udp {
774                receive_buffer_bytes,
775                ..
776            } => receive_buffer_bytes,
777            _ => panic!("expected Mode::Udp"),
778        };
779
780        assert_eq!(receive_buffer_bytes, Some(256));
781    }
782
783    #[cfg(unix)]
784    #[test]
785    fn config_unix() {
786        let config: SyslogConfig = serde_yaml::from_str(indoc! {
787            r#"
788            mode: unix
789            path: "127.0.0.1:1235"
790            "#,
791        })
792        .unwrap();
793        assert!(matches!(config.mode, Mode::Unix { .. }));
794    }
795
796    #[cfg(unix)]
797    #[test]
798    fn config_unix_permissions() {
799        let config: SyslogConfig = serde_yaml::from_str(indoc! {
800            r#"
801            mode: unix
802            path: "127.0.0.1:1235"
803            socket_file_mode: 511
804            "#,
805        })
806        .unwrap();
807        let socket_file_mode = match config.mode {
808            Mode::Unix {
809                path: _,
810                socket_file_mode,
811            } => socket_file_mode,
812            _ => panic!("expected Mode::Unix"),
813        };
814
815        assert_eq!(socket_file_mode, Some(0o777));
816    }
817
818    #[test]
819    fn syslog_ng_network_syslog_protocol() {
820        // this should also match rsyslog omfwd with template=RSYSLOG_SyslogProtocol23Format
821        let msg = "i am foobar";
822        let raw = format!(
823            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {}{} {}"#,
824            r#"[meta sequenceId="1" sysUpTime="37" language="EN"]"#,
825            r#"[origin ip="192.168.0.1" software="test"]"#,
826            msg
827        );
828
829        let mut expected = Event::Log(LogEvent::from(msg));
830
831        {
832            let expected = expected.as_mut_log();
833            expected.insert(
834                (PathPrefix::Event, log_schema().timestamp_key().unwrap()),
835                Utc.with_ymd_and_hms(2019, 2, 13, 19, 48, 34)
836                    .single()
837                    .expect("invalid timestamp"),
838            );
839            expected.insert(
840                log_schema().source_type_key_target_path().unwrap(),
841                "syslog",
842            );
843            expected.insert(event_path!("host"), "74794bfb6795");
844            expected.insert(event_path!("hostname"), "74794bfb6795");
845
846            expected.insert(event_path!("meta", "sequenceId"), "1");
847            expected.insert(event_path!("meta", "sysUpTime"), "37");
848            expected.insert(event_path!("meta", "language"), "EN");
849            expected.insert(event_path!("origin", "software"), "test");
850            expected.insert(event_path!("origin", "ip"), "192.168.0.1");
851
852            expected.insert(event_path!("severity"), "notice");
853            expected.insert(event_path!("facility"), "user");
854            expected.insert(event_path!("version"), 1);
855            expected.insert(event_path!("appname"), "root");
856            expected.insert(event_path!("procid"), 8449);
857            expected.insert(event_path!("source_ip"), "192.168.0.254");
858        }
859
860        assert_event_data_eq!(
861            event_from_bytes(
862                "host",
863                Some(Bytes::from("192.168.0.254")),
864                raw.into(),
865                LogNamespace::Legacy
866            )
867            .unwrap(),
868            expected
869        );
870    }
871
872    #[test]
873    fn handles_incorrect_sd_element() {
874        let msg = "qwerty";
875        let raw = format!(
876            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {} {}"#,
877            r"[incorrect x]", msg
878        );
879
880        let mut expected = Event::Log(LogEvent::from(msg));
881        {
882            let expected = expected.as_mut_log();
883            expected.insert(
884                (PathPrefix::Event, log_schema().timestamp_key().unwrap()),
885                Utc.with_ymd_and_hms(2019, 2, 13, 19, 48, 34)
886                    .single()
887                    .expect("invalid timestamp"),
888            );
889            expected.insert(
890                (PathPrefix::Event, log_schema().host_key().unwrap()),
891                "74794bfb6795",
892            );
893            expected.insert(event_path!("hostname"), "74794bfb6795");
894            expected.insert(
895                log_schema().source_type_key_target_path().unwrap(),
896                "syslog",
897            );
898            expected.insert(event_path!("severity"), "notice");
899            expected.insert(event_path!("facility"), "user");
900            expected.insert(event_path!("version"), 1);
901            expected.insert(event_path!("appname"), "root");
902            expected.insert(event_path!("procid"), 8449);
903            expected.insert(event_path!("source_ip"), "192.168.0.254");
904        }
905
906        let event = event_from_bytes(
907            "host",
908            Some(Bytes::from("192.168.0.254")),
909            raw.into(),
910            LogNamespace::Legacy,
911        )
912        .unwrap();
913        assert_event_data_eq!(event, expected);
914
915        let raw = format!(
916            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {} {}"#,
917            r"[incorrect x=]", msg
918        );
919
920        let event = event_from_bytes(
921            "host",
922            Some(Bytes::from("192.168.0.254")),
923            raw.into(),
924            LogNamespace::Legacy,
925        )
926        .unwrap();
927        assert_event_data_eq!(event, expected);
928    }
929
930    #[test]
931    fn handles_empty_sd_element() {
932        fn there_is_map_called_empty(event: Event) -> bool {
933            event
934                .as_log()
935                .get(event_path!("empty"))
936                .expect("empty exists")
937                .is_object()
938        }
939
940        let msg = format!(
941            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {} qwerty"#,
942            r"[empty]"
943        );
944
945        let event = event_from_bytes("host", None, msg.into(), LogNamespace::Legacy).unwrap();
946        assert!(there_is_map_called_empty(event));
947
948        let msg = format!(
949            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {} qwerty"#,
950            r#"[non_empty x="1"][empty]"#
951        );
952
953        let event = event_from_bytes("host", None, msg.into(), LogNamespace::Legacy).unwrap();
954        assert!(there_is_map_called_empty(event));
955
956        let msg = format!(
957            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {} qwerty"#,
958            r#"[empty][non_empty x="1"]"#
959        );
960
961        let event = event_from_bytes("host", None, msg.into(), LogNamespace::Legacy).unwrap();
962        assert!(there_is_map_called_empty(event));
963
964        let msg = format!(
965            r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - {} qwerty"#,
966            r#"[empty not_really="testing the test"]"#
967        );
968
969        let event = event_from_bytes("host", None, msg.into(), LogNamespace::Legacy).unwrap();
970        assert!(there_is_map_called_empty(event));
971    }
972
973    #[test]
974    fn handles_weird_whitespace() {
975        // this should also match rsyslog omfwd with template=RSYSLOG_SyslogProtocol23Format
976        let raw = r#"
977            <13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - [meta sequenceId="1"] i am foobar
978            "#;
979        let cleaned = r#"<13>1 2019-02-13T19:48:34+00:00 74794bfb6795 root 8449 - [meta sequenceId="1"] i am foobar"#;
980
981        assert_event_data_eq!(
982            event_from_bytes("host", None, raw.to_owned().into(), LogNamespace::Legacy).unwrap(),
983            event_from_bytes(
984                "host",
985                None,
986                cleaned.to_owned().into(),
987                LogNamespace::Legacy
988            )
989            .unwrap()
990        );
991    }
992
993    #[test]
994    fn handles_dots_in_sdata() {
995        let raw =
996            r#"<190>Feb 13 21:31:56 74794bfb6795 liblogging-stdlog:  [origin foo.bar="baz"] hello"#;
997        let event =
998            event_from_bytes("host", None, raw.to_owned().into(), LogNamespace::Legacy).unwrap();
999        assert_eq!(
1000            event.as_log().get(event_path!("origin", "foo.bar")),
1001            Some(&Value::from("baz"))
1002        );
1003    }
1004
1005    #[test]
1006    fn syslog_ng_default_network() {
1007        let msg = "i am foobar";
1008        let raw = format!(r#"<13>Feb 13 20:07:26 74794bfb6795 root[8539]: {msg}"#);
1009        let event = event_from_bytes(
1010            "host",
1011            Some(Bytes::from("192.168.0.254")),
1012            raw.into(),
1013            LogNamespace::Legacy,
1014        )
1015        .unwrap();
1016
1017        let mut expected = Event::Log(LogEvent::from(msg));
1018        {
1019            let value = event.as_log().get(event_path!("timestamp")).unwrap();
1020            let year = value.as_timestamp().unwrap().naive_local().year();
1021
1022            let expected = expected.as_mut_log();
1023            let expected_date: DateTime<Utc> = Local
1024                .with_ymd_and_hms(year, 2, 13, 20, 7, 26)
1025                .single()
1026                .expect("invalid timestamp")
1027                .into();
1028
1029            expected.insert(
1030                (PathPrefix::Event, log_schema().timestamp_key().unwrap()),
1031                expected_date,
1032            );
1033            expected.insert(
1034                (PathPrefix::Event, log_schema().host_key().unwrap()),
1035                "74794bfb6795",
1036            );
1037            expected.insert(
1038                log_schema().source_type_key_target_path().unwrap(),
1039                "syslog",
1040            );
1041            expected.insert(event_path!("hostname"), "74794bfb6795");
1042            expected.insert(event_path!("severity"), "notice");
1043            expected.insert(event_path!("facility"), "user");
1044            expected.insert(event_path!("appname"), "root");
1045            expected.insert(event_path!("procid"), 8539);
1046            expected.insert(event_path!("source_ip"), "192.168.0.254");
1047        }
1048
1049        assert_event_data_eq!(event, expected);
1050    }
1051
1052    #[test]
1053    fn rsyslog_omfwd_tcp_default() {
1054        let msg = "start";
1055        let raw = format!(
1056            r#"<190>Feb 13 21:31:56 74794bfb6795 liblogging-stdlog:  [origin software="rsyslogd" swVersion="8.24.0" x-pid="8979" x-info="http://www.rsyslog.com"] {msg}"#
1057        );
1058        let event = event_from_bytes(
1059            "host",
1060            Some(Bytes::from("192.168.0.254")),
1061            raw.into(),
1062            LogNamespace::Legacy,
1063        )
1064        .unwrap();
1065
1066        let mut expected = Event::Log(LogEvent::from(msg));
1067        {
1068            let value = event.as_log().get(event_path!("timestamp")).unwrap();
1069            let year = value.as_timestamp().unwrap().naive_local().year();
1070
1071            let expected = expected.as_mut_log();
1072            let expected_date: DateTime<Utc> = Local
1073                .with_ymd_and_hms(year, 2, 13, 21, 31, 56)
1074                .single()
1075                .expect("invalid timestamp")
1076                .into();
1077            expected.insert(
1078                (PathPrefix::Event, log_schema().timestamp_key().unwrap()),
1079                expected_date,
1080            );
1081            expected.insert(
1082                log_schema().source_type_key_target_path().unwrap(),
1083                "syslog",
1084            );
1085            expected.insert(event_path!("host"), "74794bfb6795");
1086            expected.insert(event_path!("hostname"), "74794bfb6795");
1087            expected.insert(event_path!("severity"), "info");
1088            expected.insert(event_path!("facility"), "local7");
1089            expected.insert(event_path!("appname"), "liblogging-stdlog");
1090            expected.insert(event_path!("origin", "software"), "rsyslogd");
1091            expected.insert(event_path!("origin", "swVersion"), "8.24.0");
1092            expected.insert(event_path!("source_ip"), "192.168.0.254");
1093            expected.insert(event_path!("origin", "x-pid"), "8979");
1094            expected.insert(event_path!("origin", "x-info"), "http://www.rsyslog.com");
1095        }
1096
1097        assert_event_data_eq!(event, expected);
1098    }
1099
1100    #[test]
1101    fn rsyslog_omfwd_tcp_forward_format() {
1102        let msg = "start";
1103        let raw = format!(
1104            r#"<190>2019-02-13T21:53:30.605850+00:00 74794bfb6795 liblogging-stdlog:  [origin software="rsyslogd" swVersion="8.24.0" x-pid="9043" x-info="http://www.rsyslog.com"] {msg}"#
1105        );
1106
1107        let mut expected = Event::Log(LogEvent::from(msg));
1108        {
1109            let expected = expected.as_mut_log();
1110            expected.insert(
1111                (PathPrefix::Event, log_schema().timestamp_key().unwrap()),
1112                Utc.with_ymd_and_hms(2019, 2, 13, 21, 53, 30)
1113                    .single()
1114                    .and_then(|t| t.with_nanosecond(605_850 * 1000))
1115                    .expect("invalid timestamp"),
1116            );
1117            expected.insert(
1118                log_schema().source_type_key_target_path().unwrap(),
1119                "syslog",
1120            );
1121            expected.insert(event_path!("host"), "74794bfb6795");
1122            expected.insert(event_path!("hostname"), "74794bfb6795");
1123            expected.insert(event_path!("severity"), "info");
1124            expected.insert(event_path!("facility"), "local7");
1125            expected.insert(event_path!("appname"), "liblogging-stdlog");
1126            expected.insert(event_path!("origin", "software"), "rsyslogd");
1127            expected.insert(event_path!("origin", "swVersion"), "8.24.0");
1128            expected.insert(event_path!("origin", "x-pid"), "9043");
1129            expected.insert(event_path!("origin", "x-info"), "http://www.rsyslog.com");
1130        }
1131
1132        assert_event_data_eq!(
1133            event_from_bytes("host", None, raw.into(), LogNamespace::Legacy).unwrap(),
1134            expected
1135        );
1136    }
1137
1138    #[tokio::test]
1139    async fn test_tcp_syslog() {
1140        assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
1141            let num_messages: usize = 10000;
1142            let (_guard, in_addr) = next_addr();
1143
1144            // Create and spawn the source.
1145            let config = SyslogConfig::from_mode(Mode::Tcp {
1146                address: in_addr.into(),
1147                permit_origin: None,
1148                keepalive: None,
1149                tls: None,
1150                receive_buffer_bytes: None,
1151                connection_limit: None,
1152            });
1153
1154            let key = ComponentKey::from("in");
1155            let (tx, rx) = SourceSender::new_test();
1156            let (context, shutdown) = SourceContext::new_shutdown(&key, tx);
1157            let shutdown_complete = shutdown.shutdown_tripwire();
1158
1159            let source = config
1160                .build(context)
1161                .await
1162                .expect("source should not fail to build");
1163            tokio::spawn(source);
1164
1165            // Wait for source to become ready to accept traffic.
1166            wait_for_tcp(in_addr).await;
1167
1168            let output_events = CountReceiver::receive_events(rx);
1169
1170            // Now craft and send syslog messages to the source, and collect them on the other side.
1171            let input_messages: Vec<SyslogMessageRfc5424> = (0..num_messages)
1172                .map(|i| SyslogMessageRfc5424::random(i, 30, 4, 3, 3))
1173                .collect();
1174
1175            let input_lines: Vec<String> =
1176                input_messages.iter().map(|msg| msg.to_string()).collect();
1177
1178            send_lines(in_addr, input_lines).await.unwrap();
1179
1180            // Wait a short period of time to ensure the messages get sent.
1181            sleep(Duration::from_secs(2)).await;
1182
1183            // Shutdown the source, and make sure we've got all the messages we sent in.
1184            shutdown
1185                .shutdown_all(Some(Instant::now() + Duration::from_millis(100)))
1186                .await;
1187            shutdown_complete.await;
1188
1189            let output_events = output_events.await;
1190            assert_eq!(output_events.len(), num_messages);
1191
1192            let output_messages: Vec<SyslogMessageRfc5424> = output_events
1193                .into_iter()
1194                .map(|mut e| {
1195                    e.as_mut_log().remove(event_path!("hostname")); // Vector adds this field which will cause a parse error.
1196                    e.as_mut_log().remove(event_path!("source_ip")); // Vector adds this field which will cause a parse error.
1197                    e.into()
1198                })
1199                .collect();
1200            assert_eq!(output_messages, input_messages);
1201        })
1202        .await;
1203    }
1204
1205    #[tokio::test]
1206    async fn test_udp_syslog() {
1207        assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
1208            let num_messages: usize = 1000;
1209            let (_guard, in_addr) = next_addr();
1210
1211            // Create and spawn the source.
1212            let config = SyslogConfig::from_mode(Mode::Udp {
1213                address: in_addr.into(),
1214                receive_buffer_bytes: Some(4 * 1024 * 1024),
1215            });
1216
1217            let key = ComponentKey::from("in");
1218            let (tx, rx) = SourceSender::new_test();
1219            let (context, shutdown) = SourceContext::new_shutdown(&key, tx);
1220            let shutdown_complete = shutdown.shutdown_tripwire();
1221
1222            let source = config
1223                .build(context)
1224                .await
1225                .expect("source should not fail to build");
1226            tokio::spawn(source);
1227
1228            // Give UDP a brief moment to start listening.
1229            sleep(Duration::from_millis(150)).await;
1230
1231            let output_events = CountReceiver::receive_events(rx);
1232
1233            // Craft and send syslog messages as individual UDP datagrams.
1234            let input_messages: Vec<SyslogMessageRfc5424> = (0..num_messages)
1235                .map(|i| SyslogMessageRfc5424::random(i, 30, 4, 3, 3))
1236                .collect();
1237
1238            let input_lines: Vec<String> =
1239                input_messages.iter().map(|msg| msg.to_string()).collect();
1240
1241            let socket = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
1242            for line in input_lines {
1243                socket.send_to(line.as_bytes(), in_addr).await.unwrap();
1244            }
1245
1246            // Wait a short period of time to ensure the messages get sent.
1247            sleep(Duration::from_secs(2)).await;
1248
1249            // Shutdown the source, and make sure we've got all the messages we sent in.
1250            shutdown
1251                .shutdown_all(Some(Instant::now() + Duration::from_millis(100)))
1252                .await;
1253            shutdown_complete.await;
1254
1255            let output_events = output_events.await;
1256            assert_eq!(output_events.len(), num_messages);
1257
1258            let output_messages: Vec<SyslogMessageRfc5424> = output_events
1259                .into_iter()
1260                .map(|mut e| {
1261                    e.as_mut_log().remove(event_path!("hostname")); // Vector adds this field which will cause a parse error.
1262                    e.as_mut_log().remove(event_path!("source_ip")); // Vector adds this field which will cause a parse error.
1263                    e.into()
1264                })
1265                .collect();
1266            assert_eq!(output_messages, input_messages);
1267        })
1268        .await;
1269    }
1270
1271    #[cfg(unix)]
1272    #[tokio::test]
1273    async fn test_unix_stream_syslog() {
1274        use std::os::unix::net::UnixStream as StdUnixStream;
1275
1276        use futures_util::{SinkExt, stream};
1277        use tokio::{io::AsyncWriteExt, net::UnixStream};
1278        use tokio_util::codec::{FramedWrite, LinesCodec};
1279
1280        use crate::test_util::components::SOCKET_PUSH_SOURCE_TAGS;
1281
1282        assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
1283            let num_messages: usize = 1;
1284            let in_path = tempfile::tempdir().unwrap().keep().join("stream_test");
1285
1286            // Create and spawn the source.
1287            let config = SyslogConfig::from_mode(Mode::Unix {
1288                path: in_path.clone(),
1289                socket_file_mode: None,
1290            });
1291
1292            let key = ComponentKey::from("in");
1293            let (tx, rx) = SourceSender::new_test();
1294            let (context, shutdown) = SourceContext::new_shutdown(&key, tx);
1295            let shutdown_complete = shutdown.shutdown_tripwire();
1296
1297            let source = config
1298                .build(context)
1299                .await
1300                .expect("source should not fail to build");
1301            tokio::spawn(source);
1302
1303            // Wait for source to become ready to accept traffic.
1304            while StdUnixStream::connect(&in_path).is_err() {
1305                tokio::task::yield_now().await;
1306            }
1307
1308            let output_events = CountReceiver::receive_events(rx);
1309
1310            // Now craft and send syslog messages to the source, and collect them on the other side.
1311            let input_messages: Vec<SyslogMessageRfc5424> = (0..num_messages)
1312                .map(|i| SyslogMessageRfc5424::random(i, 30, 4, 3, 3))
1313                .collect();
1314
1315            let stream = UnixStream::connect(&in_path).await.unwrap();
1316            let mut sink = FramedWrite::new(stream, LinesCodec::new());
1317
1318            let lines: Vec<String> = input_messages.iter().map(|msg| msg.to_string()).collect();
1319            let mut lines = stream::iter(lines).map(Ok);
1320            sink.send_all(&mut lines).await.unwrap();
1321
1322            let stream = sink.get_mut();
1323            stream.shutdown().await.unwrap();
1324
1325            // Wait a short period of time to ensure the messages get sent.
1326            sleep(Duration::from_secs(1)).await;
1327
1328            shutdown
1329                .shutdown_all(Some(Instant::now() + Duration::from_millis(100)))
1330                .await;
1331            shutdown_complete.await;
1332
1333            let output_events = output_events.await;
1334            assert_eq!(output_events.len(), num_messages);
1335
1336            let output_messages: Vec<SyslogMessageRfc5424> = output_events
1337                .into_iter()
1338                .map(|mut e| {
1339                    e.as_mut_log().remove(event_path!("hostname")); // Vector adds this field which will cause a parse error.
1340                    e.as_mut_log().remove(event_path!("source_ip")); // Vector adds this field which will cause a parse error.
1341                    e.into()
1342                })
1343                .collect();
1344            assert_eq!(output_messages, input_messages);
1345        })
1346        .await;
1347    }
1348
1349    #[tokio::test]
1350    async fn test_octet_counting_syslog() {
1351        assert_source_compliance(&SOCKET_PUSH_SOURCE_TAGS, async {
1352            let num_messages: usize = 10000;
1353            let (_guard, in_addr) = next_addr();
1354
1355            // Create and spawn the source.
1356            let config = SyslogConfig::from_mode(Mode::Tcp {
1357                address: in_addr.into(),
1358                permit_origin: None,
1359                keepalive: None,
1360                tls: None,
1361                receive_buffer_bytes: None,
1362                connection_limit: None,
1363            });
1364
1365            let key = ComponentKey::from("in");
1366            let (tx, rx) = SourceSender::new_test();
1367            let (context, shutdown) = SourceContext::new_shutdown(&key, tx);
1368            let shutdown_complete = shutdown.shutdown_tripwire();
1369
1370            let source = config
1371                .build(context)
1372                .await
1373                .expect("source should not fail to build");
1374            tokio::spawn(source);
1375
1376            // Wait for source to become ready to accept traffic.
1377            wait_for_tcp(in_addr).await;
1378
1379            let output_events = CountReceiver::receive_events(rx);
1380
1381            // Now craft and send syslog messages to the source, and collect them on the other side.
1382            let input_messages: Vec<SyslogMessageRfc5424> = (0..num_messages)
1383                .map(|i| {
1384                    let mut msg = SyslogMessageRfc5424::random(i, 30, 4, 3, 3);
1385                    msg.message.push('\n');
1386                    msg.message.push_str(&random_string(30));
1387                    msg
1388                })
1389                .collect();
1390
1391            let codec = BytesCodec::new();
1392            let input_lines: Vec<Bytes> = input_messages
1393                .iter()
1394                .map(|msg| {
1395                    let s = msg.to_string();
1396                    format!("{} {}", s.len(), s).into()
1397                })
1398                .collect();
1399
1400            send_encodable(in_addr, codec, input_lines).await.unwrap();
1401
1402            // Wait a short period of time to ensure the messages get sent.
1403            sleep(Duration::from_secs(2)).await;
1404
1405            // Shutdown the source, and make sure we've got all the messages we sent in.
1406            shutdown
1407                .shutdown_all(Some(Instant::now() + Duration::from_millis(100)))
1408                .await;
1409            shutdown_complete.await;
1410
1411            let output_events = output_events.await;
1412            assert_eq!(output_events.len(), num_messages);
1413
1414            let output_messages: Vec<SyslogMessageRfc5424> = output_events
1415                .into_iter()
1416                .map(|mut e| {
1417                    e.as_mut_log().remove(event_path!("hostname")); // Vector adds this field which will cause a parse error.
1418                    e.as_mut_log().remove(event_path!("source_ip")); // Vector adds this field which will cause a parse error.
1419                    e.into()
1420                })
1421                .collect();
1422            assert_eq!(output_messages, input_messages);
1423        })
1424        .await;
1425    }
1426
1427    #[derive(Deserialize, PartialEq, Clone, Debug)]
1428    struct SyslogMessageRfc5424 {
1429        msgid: String,
1430        severity: Severity,
1431        facility: Facility,
1432        version: u8,
1433        timestamp: String,
1434        host: String,
1435        source_type: String,
1436        appname: String,
1437        procid: usize,
1438        message: String,
1439        #[serde(flatten)]
1440        structured_data: StructuredData,
1441    }
1442
1443    impl SyslogMessageRfc5424 {
1444        fn random(
1445            id: usize,
1446            msg_len: usize,
1447            field_len: usize,
1448            max_map_size: usize,
1449            max_children: usize,
1450        ) -> Self {
1451            let msg = random_string(msg_len);
1452            let structured_data = random_structured_data(max_map_size, max_children, field_len);
1453
1454            let timestamp = Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true);
1455            //"secfrac" can contain up to 6 digits, but TCP sinks uses `AutoSi`
1456
1457            Self {
1458                msgid: format!("test{id}"),
1459                severity: Severity::LOG_INFO,
1460                facility: Facility::LOG_USER,
1461                version: 1,
1462                timestamp,
1463                host: "hogwarts".to_owned(),
1464                source_type: "syslog".to_owned(),
1465                appname: "harry".to_owned(),
1466                procid: rng().random_range(0..32768),
1467                structured_data,
1468                message: msg,
1469            }
1470        }
1471    }
1472
1473    impl fmt::Display for SyslogMessageRfc5424 {
1474        fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
1475            write!(
1476                f,
1477                "<{}>{} {} {} {} {} {} {} {}",
1478                encode_priority(self.severity, self.facility),
1479                self.version,
1480                self.timestamp,
1481                self.host,
1482                self.appname,
1483                self.procid,
1484                self.msgid,
1485                format_structured_data_rfc5424(&self.structured_data),
1486                self.message
1487            )
1488        }
1489    }
1490
1491    impl From<Event> for SyslogMessageRfc5424 {
1492        fn from(e: Event) -> Self {
1493            let (value, _) = e.into_log().into_parts();
1494            let mut fields = value.into_object().unwrap();
1495
1496            Self {
1497                msgid: fields.remove("msgid").map(value_to_string).unwrap(),
1498                severity: fields
1499                    .remove("severity")
1500                    .map(value_to_string)
1501                    .and_then(|s| Severity::from_str(s.as_str()))
1502                    .unwrap(),
1503                facility: fields
1504                    .remove("facility")
1505                    .map(value_to_string)
1506                    .and_then(|s| Facility::from_str(s.as_str()))
1507                    .unwrap(),
1508                version: fields
1509                    .remove("version")
1510                    .map(value_to_string)
1511                    .map(|s| u8::from_str(s.as_str()).unwrap())
1512                    .unwrap(),
1513                timestamp: fields.remove("timestamp").map(value_to_string).unwrap(),
1514                host: fields.remove("host").map(value_to_string).unwrap(),
1515                source_type: fields.remove("source_type").map(value_to_string).unwrap(),
1516                appname: fields.remove("appname").map(value_to_string).unwrap(),
1517                procid: fields
1518                    .remove("procid")
1519                    .map(value_to_string)
1520                    .map(|s| usize::from_str(s.as_str()).unwrap())
1521                    .unwrap(),
1522                message: fields.remove("message").map(value_to_string).unwrap(),
1523                structured_data: structured_data_from_fields(fields),
1524            }
1525        }
1526    }
1527
1528    fn structured_data_from_fields(fields: ObjectMap) -> StructuredData {
1529        let mut structured_data = StructuredData::default();
1530
1531        for (key, value) in fields.into_iter() {
1532            let subfields = value
1533                .into_object()
1534                .unwrap()
1535                .into_iter()
1536                .map(|(k, v)| (k.into(), value_to_string(v)))
1537                .collect();
1538
1539            structured_data.insert(key.into(), subfields);
1540        }
1541
1542        structured_data
1543    }
1544
1545    #[allow(non_camel_case_types, clippy::upper_case_acronyms)]
1546    #[derive(Copy, Clone, Deserialize, PartialEq, Eq, Debug)]
1547    pub enum Severity {
1548        #[serde(rename(deserialize = "emergency"))]
1549        LOG_EMERG,
1550        #[serde(rename(deserialize = "alert"))]
1551        LOG_ALERT,
1552        #[serde(rename(deserialize = "critical"))]
1553        LOG_CRIT,
1554        #[serde(rename(deserialize = "error"))]
1555        LOG_ERR,
1556        #[serde(rename(deserialize = "warn"))]
1557        LOG_WARNING,
1558        #[serde(rename(deserialize = "notice"))]
1559        LOG_NOTICE,
1560        #[serde(rename(deserialize = "info"))]
1561        LOG_INFO,
1562        #[serde(rename(deserialize = "debug"))]
1563        LOG_DEBUG,
1564    }
1565
1566    impl Severity {
1567        fn from_str(s: &str) -> Option<Self> {
1568            match s {
1569                "emergency" => Some(Self::LOG_EMERG),
1570                "alert" => Some(Self::LOG_ALERT),
1571                "critical" => Some(Self::LOG_CRIT),
1572                "error" => Some(Self::LOG_ERR),
1573                "warn" => Some(Self::LOG_WARNING),
1574                "notice" => Some(Self::LOG_NOTICE),
1575                "info" => Some(Self::LOG_INFO),
1576                "debug" => Some(Self::LOG_DEBUG),
1577
1578                x => {
1579                    #[allow(clippy::print_stdout)]
1580                    {
1581                        println!("converting severity str, got {x}");
1582                    }
1583                    None
1584                }
1585            }
1586        }
1587    }
1588
1589    #[allow(non_camel_case_types, clippy::upper_case_acronyms)]
1590    #[derive(Copy, Clone, PartialEq, Eq, Deserialize, Debug)]
1591    pub enum Facility {
1592        #[serde(rename(deserialize = "kernel"))]
1593        LOG_KERN = 0 << 3,
1594        #[serde(rename(deserialize = "user"))]
1595        LOG_USER = 1 << 3,
1596        #[serde(rename(deserialize = "mail"))]
1597        LOG_MAIL = 2 << 3,
1598        #[serde(rename(deserialize = "daemon"))]
1599        LOG_DAEMON = 3 << 3,
1600        #[serde(rename(deserialize = "auth"))]
1601        LOG_AUTH = 4 << 3,
1602        #[serde(rename(deserialize = "syslog"))]
1603        LOG_SYSLOG = 5 << 3,
1604    }
1605
1606    impl Facility {
1607        fn from_str(s: &str) -> Option<Self> {
1608            match s {
1609                "kernel" => Some(Self::LOG_KERN),
1610                "user" => Some(Self::LOG_USER),
1611                "mail" => Some(Self::LOG_MAIL),
1612                "daemon" => Some(Self::LOG_DAEMON),
1613                "auth" => Some(Self::LOG_AUTH),
1614                "syslog" => Some(Self::LOG_SYSLOG),
1615                _ => None,
1616            }
1617        }
1618    }
1619
1620    type StructuredData = HashMap<String, HashMap<String, String>>;
1621
1622    fn random_structured_data(
1623        max_map_size: usize,
1624        max_children: usize,
1625        field_len: usize,
1626    ) -> StructuredData {
1627        let amount = rng().random_range(0..max_children);
1628
1629        random_maps(max_map_size, field_len)
1630            .filter(|m| !m.is_empty()) //syslog_rfc5424 ignores empty maps, tested separately
1631            .take(amount)
1632            .enumerate()
1633            .map(|(i, map)| (format!("id{i}"), map))
1634            .collect()
1635    }
1636
1637    fn format_structured_data_rfc5424(data: &StructuredData) -> String {
1638        if data.is_empty() {
1639            "-".to_string()
1640        } else {
1641            let mut res = String::new();
1642            for (id, params) in data {
1643                res = res + "[" + id;
1644                for (name, value) in params {
1645                    res = res + " " + name + "=\"" + value + "\"";
1646                }
1647                res += "]";
1648            }
1649
1650            res
1651        }
1652    }
1653
1654    const fn encode_priority(severity: Severity, facility: Facility) -> u8 {
1655        facility as u8 | severity as u8
1656    }
1657
1658    fn value_to_string(v: Value) -> String {
1659        if v.is_bytes() {
1660            let buf = v.as_bytes().unwrap();
1661            String::from_utf8_lossy(buf).to_string()
1662        } else if v.is_timestamp() {
1663            let ts = v.as_timestamp().unwrap();
1664            ts.to_rfc3339_opts(SecondsFormat::AutoSi, true)
1665        } else {
1666            v.to_string()
1667        }
1668    }
1669}