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#[configurable_component(source("syslog", "Collect logs sent via Syslog."))]
46#[derive(Clone, Debug)]
47pub struct SyslogConfig {
48 #[serde(flatten)]
49 mode: Mode,
50
51 #[serde(default = "crate::serde::default_max_length")]
55 #[configurable(metadata(docs::type_unit = "bytes"))]
56 max_length: usize,
57
58 host_key: Option<OptionalValuePath>,
67
68 #[configurable(metadata(docs::hidden))]
70 #[serde(default)]
71 pub log_namespace: Option<bool>,
72}
73
74#[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 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 #[configurable(metadata(docs::type_unit = "bytes"))]
99 receive_buffer_bytes: Option<usize>,
100
101 connection_limit: Option<u32>,
103 },
104
105 Udp {
107 #[configurable(derived)]
108 address: SocketListenAddr,
109
110 #[configurable(metadata(docs::type_unit = "bytes"))]
114 receive_buffer_bytes: Option<usize>,
115 },
116
117 #[cfg(unix)]
121 Unix {
122 #[configurable(metadata(docs::examples = "/path/to/socket"))]
126 path: PathBuf,
127
128 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_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 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 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 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_tcp(in_addr).await;
1167
1168 let output_events = CountReceiver::receive_events(rx);
1169
1170 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 sleep(Duration::from_secs(2)).await;
1182
1183 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")); e.as_mut_log().remove(event_path!("source_ip")); 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 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 sleep(Duration::from_millis(150)).await;
1230
1231 let output_events = CountReceiver::receive_events(rx);
1232
1233 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 sleep(Duration::from_secs(2)).await;
1248
1249 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")); e.as_mut_log().remove(event_path!("source_ip")); 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 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 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 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 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")); e.as_mut_log().remove(event_path!("source_ip")); 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 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_tcp(in_addr).await;
1378
1379 let output_events = CountReceiver::receive_events(rx);
1380
1381 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 sleep(Duration::from_secs(2)).await;
1404
1405 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")); e.as_mut_log().remove(event_path!("source_ip")); 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 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()) .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}