diff --git a/signal-gateway-bin/src/main.rs b/signal-gateway-bin/src/main.rs index 4186713..efe1856 100644 --- a/signal-gateway-bin/src/main.rs +++ b/signal-gateway-bin/src/main.rs @@ -124,8 +124,8 @@ async fn main() { } else { None }; - let _udp_json_task = if let Some(udp_json) = &config.udp_json { - Some(udp_json.start_udp_task(gateway.clone()).await.unwrap()) + let _udp_json_tasks = if let Some(udp_json) = &config.udp_json { + Some(udp_json.start_tasks(gateway.clone()).await.unwrap()) } else { None }; diff --git a/signal-gateway-bin/src/udp_json/json_lines.rs b/signal-gateway-bin/src/udp_json/json_lines.rs index 85e38b9..f478b1c 100644 --- a/signal-gateway-bin/src/udp_json/json_lines.rs +++ b/signal-gateway-bin/src/udp_json/json_lines.rs @@ -1,11 +1,27 @@ -use tokio::{ - io::{Reader, BufReader}, -}; +//! Relaxed JSON Lines reader that allows newlines within JSON objects. +//! +//! Standard JSON Lines requires each JSON value to be on a single line. +//! This implementation relaxes that by tracking brace/bracket depth and +//! only treating newlines as delimiters when outside of JSON structures. +use tokio::io::{AsyncRead, AsyncReadExt, BufReader}; + +/// Read a single JSON value from a relaxed JSON Lines stream. +/// +/// Returns: +/// - `Ok(Some(bytes))` - A complete JSON value was read +/// - `Ok(None)` - EOF reached (no more values) +/// - `Err(InvalidData)` - Malformed JSON structure (imbalanced braces/brackets) +/// - `Err(other)` - I/O error +/// +/// This parser tracks brace and bracket depth to allow newlines within +/// JSON objects and arrays. A newline only ends the value when we're at +/// depth 0 (outside any object or array). pub async fn read_json_lines_value( reader: &mut BufReader, -) -> std::io::Result> - where R: Reader, +) -> std::io::Result>> +where + R: AsyncRead + Unpin, { let mut buf = Vec::::with_capacity(256); @@ -17,55 +33,200 @@ pub async fn read_json_lines_value( loop { let b = match reader.read_u8().await { Ok(b) => b, - Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => break, // EOF - Err(e) => return e; + Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => { + // EOF - return what we have if anything + if buf.is_empty() { + return Ok(None); + } + break; + } + Err(e) => return Err(e), }; if escaped { escaped = false; - } else if quoted { + buf.push(b); + continue; + } + + if quoted { if b == b'\\' { escaped = true; } else if b == b'"' { quoted = false; } - } else { - match b { - b'\\' => { escaped = true; }, - b'"' => { quoted = true; }, - b'{' => { brace_depth += 1; }, - b'[' => { bracket_depth += 1; }, - b'}' => { - if brace_depth == 0 { - let rendered_buf = String::from_utf8_lossy(&buf); - return Err(std::io::Error::new( - std::io::ErrorKind::InvalidData, - format!("imbalanced braces: {rendered_buf}\}"), - )); - } - brace_depth -= 1; - }, - b']' => { - if bracket_depth == 0 { - let rendered_buf = String::from_utf8_lossy(&buf); - return Err(std::io::Error::new( - std::io::ErrorKind::InvalidData, - format!("imbalanced brackets: {rendered_buf}]"), - )); - } - bracket_depth -= 1; - }, - b'\n' => { - if !escaped && !quoted && brace_depth == 0 && bracket_depth == 0 { - // Clean line break - buf.push('\n'); - break; - } + buf.push(b); + continue; + } + + // Not escaped, not quoted + match b { + b'"' => { + quoted = true; + buf.push(b); + } + b'{' => { + brace_depth += 1; + buf.push(b); + } + b'[' => { + bracket_depth += 1; + buf.push(b); + } + b'}' => { + if brace_depth == 0 { + let rendered_buf = String::from_utf8_lossy(&buf); + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("imbalanced braces: {rendered_buf}}}"), + )); } - _ => {}, + brace_depth -= 1; + buf.push(b); + } + b']' => { + if bracket_depth == 0 { + let rendered_buf = String::from_utf8_lossy(&buf); + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("imbalanced brackets: {rendered_buf}]"), + )); + } + bracket_depth -= 1; + buf.push(b); + } + b'\n' => { + if brace_depth == 0 && bracket_depth == 0 { + // Clean line break at depth 0 - end of value + if buf.is_empty() { + // Skip empty lines + continue; + } + break; + } + // Newline inside object/array - keep it + buf.push(b); + } + _ => { + buf.push(b); } } - buf.push(b); } - Ok(buf) + + Ok(Some(buf)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_simple_json_object() { + let data: &[u8] = b"{\"msg\": \"hello\"}\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"{\"msg\": \"hello\"}".to_vec())); + } + + #[tokio::test] + async fn test_json_with_internal_newlines() { + let data: &[u8] = b"{\n \"msg\": \"hello\",\n \"level\": \"info\"\n}\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!( + result, + Some(b"{\n \"msg\": \"hello\",\n \"level\": \"info\"\n}".to_vec()) + ); + } + + #[tokio::test] + async fn test_eof_returns_none() { + let data: &[u8] = b""; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, None); + } + + #[tokio::test] + async fn test_eof_returns_partial_buffer() { + // No trailing newline + let data: &[u8] = b"{\"msg\": \"hello\"}"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"{\"msg\": \"hello\"}".to_vec())); + } + + #[tokio::test] + async fn test_multiple_values() { + let data: &[u8] = b"{\"a\": 1}\n{\"b\": 2}\n"; + let mut reader = BufReader::new(data); + + let r1 = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(r1, Some(b"{\"a\": 1}".to_vec())); + + let r2 = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(r2, Some(b"{\"b\": 2}".to_vec())); + + let r3 = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(r3, None); + } + + #[tokio::test] + async fn test_skips_empty_lines() { + let data: &[u8] = b"\n\n{\"msg\": \"hello\"}\n\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"{\"msg\": \"hello\"}".to_vec())); + } + + #[tokio::test] + async fn test_nested_braces() { + let data: &[u8] = b"{\"nested\": {\"deep\": {}}}\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"{\"nested\": {\"deep\": {}}}".to_vec())); + } + + #[tokio::test] + async fn test_array_value() { + let data: &[u8] = b"[1, 2, 3]\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"[1, 2, 3]".to_vec())); + } + + #[tokio::test] + async fn test_string_with_escaped_quote() { + let data: &[u8] = b"{\"msg\": \"hello \\\"world\\\"\"}\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"{\"msg\": \"hello \\\"world\\\"\"}".to_vec())); + } + + #[tokio::test] + async fn test_string_with_newline() { + // Newline inside a string should be preserved (though this is technically invalid JSON) + let data: &[u8] = b"{\"msg\": \"line1\nline2\"}\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await.unwrap(); + assert_eq!(result, Some(b"{\"msg\": \"line1\nline2\"}".to_vec())); + } + + #[tokio::test] + async fn test_imbalanced_brace_error() { + let data: &[u8] = b"{\"msg\": \"hello\"}}\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await; + assert!(result.is_err()); + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::InvalidData); + } + + #[tokio::test] + async fn test_imbalanced_bracket_error() { + let data: &[u8] = b"[1, 2]]\n"; + let mut reader = BufReader::new(data); + let result = read_json_lines_value(&mut reader).await; + assert!(result.is_err()); + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::InvalidData); + } } diff --git a/signal-gateway-bin/src/udp_json/mod.rs b/signal-gateway-bin/src/udp_json/mod.rs index 6127375..e11a19e 100644 --- a/signal-gateway-bin/src/udp_json/mod.rs +++ b/signal-gateway-bin/src/udp_json/mod.rs @@ -1,29 +1,49 @@ -//! JSON UDP listener for receiving log messages in a logstash-compatible format +//! JSON listener for receiving log messages in a logstash-compatible format. //! -//! This module accepts JSON log messages over UDP and converts them to LogMessage. +//! This module accepts JSON log messages over UDP and TCP, converting them to LogMessage. //! It's designed to be flexible and accept various common formats. +//! +//! TCP connections use a relaxed JSON Lines format that allows newlines within +//! JSON objects (standard JSON Lines requires each value on a single line). use chrono::{DateTime, TimeZone, Utc}; use conf::Conf; use serde::Deserialize; use signal_gateway::{Gateway, Level, LogMessage}; use std::{net::SocketAddr, sync::Arc}; -use tokio::net::UdpSocket; -use tracing::{error, info}; +use tokio::{ + io::BufReader, + net::{TcpListener, TcpStream, UdpSocket}, +}; +use tracing::{error, info, trace}; -/// Configuration for the JSON UDP listener -#[derive(Conf, Debug)] +mod json_lines; +use json_lines::read_json_lines_value; + +/// Configuration for the JSON listener (UDP and TCP). +#[derive(Clone, Conf, Debug)] pub struct UdpJsonConfig { - /// Socket to listen for UDP messages in JSON format + /// Socket to listen for JSON log messages. + /// Both UDP and TCP listeners are started on this address. + /// TCP uses relaxed JSON Lines format (newlines allowed within objects). #[conf(long, env)] pub listen_addr: SocketAddr, } impl UdpJsonConfig { - /// Bind a UDP socket and start a background task to handle incoming JSON log messages. + /// Bind UDP and TCP sockets and start background tasks to handle incoming JSON log messages. /// - /// Returns a join handle for the background task. - pub async fn start_udp_task( + /// Returns join handles for the background tasks. + pub async fn start_tasks( + &self, + gateway: Arc, + ) -> std::io::Result<(tokio::task::JoinHandle<()>, tokio::task::JoinHandle<()>)> { + let udp_handle = self.start_udp_task(gateway.clone()).await?; + let tcp_handle = self.start_tcp_task(gateway).await?; + Ok((udp_handle, tcp_handle)) + } + + async fn start_udp_task( &self, gateway: Arc, ) -> std::io::Result> { @@ -58,6 +78,61 @@ impl UdpJsonConfig { } })) } + + async fn start_tcp_task( + &self, + gateway: Arc, + ) -> std::io::Result> { + let tcp_listener = TcpListener::bind(self.listen_addr).await?; + info!("Listening for JSON TCP on {}", self.listen_addr); + + Ok(tokio::task::spawn(async move { + loop { + let Ok((stream, addr)) = tcp_listener + .accept() + .await + .inspect_err(|err| error!("Error accepting JSON TCP connection: {err}")) + else { + continue; + }; + + trace!("Accepted JSON TCP connection from {addr}"); + + let gateway = gateway.clone(); + + // Spawn a task for each connection + tokio::spawn(async move { + if let Err(err) = handle_tcp_connection(stream, &gateway).await { + error!("JSON TCP connection from {addr} error: {err}"); + } else { + trace!("JSON TCP connection from {addr} closed"); + } + }); + } + })) + } +} + +/// Handle a single TCP connection using relaxed JSON Lines framing. +async fn handle_tcp_connection(stream: TcpStream, gateway: &Gateway) -> std::io::Result<()> { + let mut reader = BufReader::new(stream); + + loop { + let Some(msg_bytes) = read_json_lines_value(&mut reader).await? else { + return Ok(()); // Clean EOF + }; + + let text = std::str::from_utf8(&msg_bytes).map_err(|err| { + std::io::Error::new(std::io::ErrorKind::InvalidData, err) + })?; + + let json_msg: JsonLogMessage = serde_json::from_str(text).map_err(|err| { + std::io::Error::new(std::io::ErrorKind::InvalidData, err) + })?; + + let log_msg = json_msg.into_log_message(); + gateway.handle_log_message(log_msg).await; + } } /// A flexible JSON log message format compatible with logstash and similar systems. @@ -105,6 +180,8 @@ pub struct JsonLogMessage { impl JsonLogMessage { /// Convert to a LogMessage + /// + /// TODO: Allow this to take configuration options to customize how fields are mapped pub fn into_log_message(self) -> LogMessage { let level = self.level.unwrap_or(Level::INFO); let mut builder = LogMessage::builder(level, self.message); @@ -469,4 +546,114 @@ mod tests { let msg: JsonLogMessage = serde_json::from_str(json).unwrap(); assert_eq!(msg.message, "test"); } + + // Tests for TCP stream parsing (read_json_lines_value + JSON parsing) + + use tokio::io::AsyncRead; + + /// Test helper that mimics handle_tcp_connection's parsing logic + async fn parse_next_message( + reader: &mut BufReader, + ) -> std::io::Result> { + let Some(msg_bytes) = read_json_lines_value(reader).await? else { + return Ok(None); + }; + + let text = std::str::from_utf8(&msg_bytes) + .map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err))?; + + let json_msg: JsonLogMessage = serde_json::from_str(text) + .map_err(|err| std::io::Error::new(std::io::ErrorKind::InvalidData, err))?; + + Ok(Some(json_msg.into_log_message())) + } + + #[tokio::test] + async fn test_tcp_valid_message() { + let data: &[u8] = b"{\"message\": \"hello\"}\n"; + let mut reader = BufReader::new(data); + let result = parse_next_message(&mut reader).await.unwrap(); + assert!(result.is_some()); + let msg = result.unwrap(); + assert_eq!(&*msg.msg, "hello"); + } + + #[tokio::test] + async fn test_tcp_multiple_messages() { + let data: &[u8] = b"{\"message\": \"first\"}\n{\"message\": \"second\"}\n"; + let mut reader = BufReader::new(data); + + let msg1 = parse_next_message(&mut reader).await.unwrap().unwrap(); + assert_eq!(&*msg1.msg, "first"); + + let msg2 = parse_next_message(&mut reader).await.unwrap().unwrap(); + assert_eq!(&*msg2.msg, "second"); + + let msg3 = parse_next_message(&mut reader).await.unwrap(); + assert!(msg3.is_none()); // EOF + } + + #[tokio::test] + async fn test_tcp_eof_returns_none() { + let data: &[u8] = b""; + let mut reader = BufReader::new(data); + let result = parse_next_message(&mut reader).await.unwrap(); + assert!(result.is_none()); + } + + #[tokio::test] + async fn test_tcp_invalid_utf8_returns_error() { + // Invalid UTF-8 bytes inside a "JSON" structure + let data: &[u8] = b"{\"\xff\xfe\": \"bad\"}\n"; + let mut reader = BufReader::new(data); + let result = parse_next_message(&mut reader).await; + assert!(result.is_err()); + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::InvalidData); + } + + #[tokio::test] + async fn test_tcp_invalid_json_returns_error() { + let data: &[u8] = b"{not valid json}\n"; + let mut reader = BufReader::new(data); + let result = parse_next_message(&mut reader).await; + assert!(result.is_err()); + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::InvalidData); + } + + #[tokio::test] + async fn test_tcp_missing_required_field_returns_error() { + // Valid JSON but missing required "message" field + let data: &[u8] = b"{\"level\": \"info\"}\n"; + let mut reader = BufReader::new(data); + let result = parse_next_message(&mut reader).await; + assert!(result.is_err()); + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::InvalidData); + } + + #[tokio::test] + async fn test_tcp_error_stops_processing() { + // First message valid, second invalid - error should stop further processing + let data: &[u8] = b"{\"message\": \"ok\"}\n{invalid}\n{\"message\": \"never reached\"}\n"; + let mut reader = BufReader::new(data); + + // First message succeeds + let msg1 = parse_next_message(&mut reader).await.unwrap(); + assert!(msg1.is_some()); + + // Second message fails with error + let result = parse_next_message(&mut reader).await; + assert!(result.is_err()); + // After error, caller should close connection - no third read attempted + } + + #[tokio::test] + async fn test_tcp_multiline_json() { + let data: &[u8] = b"{\n \"message\": \"hello\",\n \"level\": \"error\"\n}\n"; + let mut reader = BufReader::new(data); + let result = parse_next_message(&mut reader).await.unwrap(); + assert!(result.is_some()); + let msg = result.unwrap(); + assert_eq!(&*msg.msg, "hello"); + assert_eq!(msg.level, Level::ERROR); + } }