diff --git a/signal-gateway-bin/src/main.rs b/signal-gateway-bin/src/main.rs index fafabae..dfe0910 100644 --- a/signal-gateway-bin/src/main.rs +++ b/signal-gateway-bin/src/main.rs @@ -2,61 +2,15 @@ use conf::Conf; use hyper::service::service_fn; use hyper_util::rt::TokioIo; use hyper_util::server::conn::auto; -use signal_gateway::{Gateway, GatewayConfig, Level, LogMessage}; -use std::{net::SocketAddr, str::FromStr, sync::Arc, time::Duration}; -use syslog_rfc5424::{SyslogMessage, SyslogSeverity}; -use tokio::net::{TcpListener, UdpSocket}; +use signal_gateway::{Gateway, GatewayConfig}; +use std::{net::SocketAddr, sync::Arc, time::Duration}; +use tokio::net::TcpListener; use tokio_util::sync::CancellationToken; use tracing::{error, info, warn}; use tracing_subscriber::EnvFilter; -/// Convert SyslogSeverity to our Level enum -fn severity_to_level(sev: SyslogSeverity) -> Level { - match sev { - SyslogSeverity::SEV_EMERG => Level::EMERGENCY, - SyslogSeverity::SEV_ALERT => Level::ALERT, - SyslogSeverity::SEV_CRIT => Level::CRITICAL, - SyslogSeverity::SEV_ERR => Level::ERROR, - SyslogSeverity::SEV_WARNING => Level::WARNING, - SyslogSeverity::SEV_NOTICE => Level::NOTICE, - SyslogSeverity::SEV_INFO => Level::INFO, - SyslogSeverity::SEV_DEBUG => Level::DEBUG, - } -} - -/// Convert a SyslogMessage to a LogMessage, extracting structured data for tracing metadata -fn syslog_to_log_message(msg: SyslogMessage, sd_id: &str) -> LogMessage { - let level = severity_to_level(msg.severity); - let mut builder = LogMessage::builder(level, msg.msg); - - if let Some(ts) = msg.timestamp { - builder = builder.timestamp(ts); - } - if let Some(nanos) = msg.timestamp_nanos { - builder = builder.timestamp_nanos(nanos); - } - if let Some(hostname) = msg.hostname { - builder = builder.hostname(hostname); - } - if let Some(appname) = msg.appname { - builder = builder.appname(appname); - } - - // Extract tracing metadata from structured data - if let Some(sd_element) = msg.sd.find_sdid(sd_id) { - if let Some(module) = sd_element.get("module") { - builder = builder.module_path(module.clone()); - } - if let Some(file) = sd_element.get("file") { - builder = builder.file(file.clone()); - } - if let Some(line) = sd_element.get("line") { - builder = builder.line(line.clone()); - } - } - - builder.build() -} +mod syslog; +use syslog::SyslogUdpConfig; #[derive(Conf, Debug)] struct Config { @@ -66,12 +20,8 @@ struct Config { /// Socket to listen for HTTP requests (GET /health, POST /alert) #[conf(long, env, default_value = "0.0.0.0:8000")] http_listen_addr: SocketAddr, - /// Socket to listen for UDP messages, in syslog RFC 5424 format - #[conf(long, env, default_value = "0.0.0.0:5424")] - udp_listen_addr: SocketAddr, - /// Structured data ID for tracing metadata (module, file, line) in syslog messages - #[conf(long, env, default_value = "tracing-meta@64700")] - sd_id: String, + #[conf(flatten, prefix)] + syslog_udp: Option, #[conf(flatten)] gateway: GatewayConfig, } @@ -123,9 +73,6 @@ async fn main() { let listener = TcpListener::bind(config.http_listen_addr).await.unwrap(); info!("Listening for http on {}", config.http_listen_addr); - let udp_socket = UdpSocket::bind(config.udp_listen_addr).await.unwrap(); - info!("Listening for udp on {}", config.udp_listen_addr); - // Listen for ctrl-c let thread_token = token.clone(); tokio::task::spawn(async move { @@ -136,7 +83,11 @@ async fn main() { // Start the two server tasks let _http_task = start_http_task(listener, gateway.clone()); - let _udp_task = start_udp_task(udp_socket, gateway.clone(), config.sd_id); + let _udp_task = if let Some(syslog_udp) = &config.syslog_udp { + Some(syslog_udp.start_udp_task(gateway.clone()).await.unwrap()) + } else { + None + }; // Run gateway task and block on it returning. Note that it exits if the token is canceled. gateway.run().await; @@ -178,59 +129,3 @@ fn start_http_task(listener: TcpListener, gateway: Arc) -> tokio::task: } }) } - -fn start_udp_task( - udp_socket: UdpSocket, - gateway: Arc, - sd_id: String, -) -> tokio::task::JoinHandle<()> { - // Loop waiting for UDP syslog messages - tokio::task::spawn(async move { - let mut buf = vec![0u8; 8192]; - loop { - let Ok((len, _addr)) = udp_socket - .recv_from(&mut buf) - .await - .inspect_err(|err| error!("Error receiving UDP packet: {err}")) - else { - continue; - }; - - let Ok(text) = str::from_utf8(&buf[0..len]) - .inspect_err(|err| error!("UDP packet was not utf8: {err}")) - else { - continue; - }; - - let Ok(syslog_msg) = SyslogMessage::from_str(text) - .inspect_err(|err| error!("UDP packet was not valid syslog: {err}:\n{text}")) - else { - continue; - }; - - let log_msg = syslog_to_log_message(syslog_msg, &sd_id); - gateway.handle_log_message(log_msg).await; - } - }) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_syslog_parsing() { - SyslogMessage::from_str( - "<12>1 2025-11-08T02:24:10.815221698+00:00 ip-172-31-5-8 app 92748 - - Dropped 3/4 reports", - ) - .unwrap(); - SyslogMessage::from_str( - "<12>1 2025-11-08T02:24:10.815221698+00:00 ip-172-31-5-8 app 92748 - - Dropped 3/4 reports", - ) - .unwrap(); - SyslogMessage::from_str( - "<12>1 2025-11-08T02:24:10.815+00:00 ip-172-31-5-8 app 92748 - - Dropped 3/4 reports due to staleness" - ) - .unwrap(); - } -} diff --git a/signal-gateway-bin/src/syslog/mod.rs b/signal-gateway-bin/src/syslog/mod.rs new file mode 100644 index 0000000..819df08 --- /dev/null +++ b/signal-gateway-bin/src/syslog/mod.rs @@ -0,0 +1,131 @@ +//! Syslog UDP listener for receiving RFC 5424 syslog messages + +use conf::Conf; +use signal_gateway::{Gateway, Level, LogMessage}; +use std::{net::SocketAddr, str::FromStr, sync::Arc}; +use syslog_rfc5424::{SyslogMessage, SyslogSeverity}; +use tokio::net::UdpSocket; +use tracing::{error, info}; + +/// Configuration for the syslog UDP listener +#[derive(Conf, Debug)] +pub struct SyslogUdpConfig { + /// Socket to listen for UDP messages, in syslog RFC 5424 format + #[conf(long, env)] + pub listen_addr: SocketAddr, + /// Structured data ID for tracing metadata (module, file, line) in syslog messages + #[conf(long, env, default_value = "tracing-meta@64700")] + pub sd_id: String, +} + +impl SyslogUdpConfig { + /// Bind a UDP socket and start a background task to handle incoming syslog messages. + /// + /// Returns a join handle for the background task. + pub async fn start_udp_task( + &self, + gateway: Arc, + ) -> std::io::Result> { + let udp_socket = UdpSocket::bind(self.listen_addr).await?; + info!("Listening for syslog UDP on {}", self.listen_addr); + + let sd_id = self.sd_id.clone(); + + Ok(tokio::task::spawn(async move { + let mut buf = vec![0u8; 8192]; + loop { + let Ok((len, _addr)) = udp_socket + .recv_from(&mut buf) + .await + .inspect_err(|err| error!("Error receiving UDP packet: {err}")) + else { + continue; + }; + + let Ok(text) = std::str::from_utf8(&buf[0..len]) + .inspect_err(|err| error!("UDP packet was not utf8: {err}")) + else { + continue; + }; + + let Ok(syslog_msg) = SyslogMessage::from_str(text) + .inspect_err(|err| error!("UDP packet was not valid syslog: {err}:\n{text}")) + else { + continue; + }; + + let log_msg = syslog_to_log_message(syslog_msg, &sd_id); + gateway.handle_log_message(log_msg).await; + } + })) + } +} + +/// Convert SyslogSeverity to our Level enum +fn severity_to_level(sev: SyslogSeverity) -> Level { + match sev { + SyslogSeverity::SEV_EMERG => Level::EMERGENCY, + SyslogSeverity::SEV_ALERT => Level::ALERT, + SyslogSeverity::SEV_CRIT => Level::CRITICAL, + SyslogSeverity::SEV_ERR => Level::ERROR, + SyslogSeverity::SEV_WARNING => Level::WARNING, + SyslogSeverity::SEV_NOTICE => Level::NOTICE, + SyslogSeverity::SEV_INFO => Level::INFO, + SyslogSeverity::SEV_DEBUG => Level::DEBUG, + } +} + +/// Convert a SyslogMessage to a LogMessage, extracting structured data for tracing metadata +fn syslog_to_log_message(msg: SyslogMessage, sd_id: &str) -> LogMessage { + let level = severity_to_level(msg.severity); + let mut builder = LogMessage::builder(level, msg.msg); + + if let Some(ts) = msg.timestamp { + builder = builder.timestamp(ts); + } + if let Some(nanos) = msg.timestamp_nanos { + builder = builder.timestamp_nanos(nanos); + } + if let Some(hostname) = msg.hostname { + builder = builder.hostname(hostname); + } + if let Some(appname) = msg.appname { + builder = builder.appname(appname); + } + + // Extract tracing metadata from structured data + if let Some(sd_element) = msg.sd.find_sdid(sd_id) { + if let Some(module) = sd_element.get("module") { + builder = builder.module_path(module.clone()); + } + if let Some(file) = sd_element.get("file") { + builder = builder.file(file.clone()); + } + if let Some(line) = sd_element.get("line") { + builder = builder.line(line.clone()); + } + } + + builder.build() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_syslog_parsing() { + SyslogMessage::from_str( + "<12>1 2025-11-08T02:24:10.815221698+00:00 ip-172-31-5-8 app 92748 - - Dropped 3/4 reports", + ) + .unwrap(); + SyslogMessage::from_str( + "<12>1 2025-11-08T02:24:10.815221698+00:00 ip-172-31-5-8 app 92748 - - Dropped 3/4 reports", + ) + .unwrap(); + SyslogMessage::from_str( + "<12>1 2025-11-08T02:24:10.815+00:00 ip-172-31-5-8 app 92748 - - Dropped 3/4 reports due to staleness" + ) + .unwrap(); + } +}