also allow json logs to be sent over a TCP stream, with multiple messages using json_lines ofrmat for framing

This commit is contained in:
Chris Beck
2025-12-05 14:49:21 -07:00
parent 0a5023ce6b
commit a42aad48f4
3 changed files with 403 additions and 55 deletions
+2 -2
View File
@@ -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
};
+204 -43
View File
@@ -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<R>(
reader: &mut BufReader<R>,
) -> std::io::Result<Vec<u8>>
where R: Reader,
) -> std::io::Result<Option<Vec<u8>>>
where
R: AsyncRead + Unpin,
{
let mut buf = Vec::<u8>::with_capacity(256);
@@ -17,55 +33,200 @@ pub async fn read_json_lines_value<R>(
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);
}
}
+197 -10
View File
@@ -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<Gateway>,
) -> 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<Gateway>,
) -> std::io::Result<tokio::task::JoinHandle<()>> {
@@ -58,6 +78,61 @@ impl UdpJsonConfig {
}
}))
}
async fn start_tcp_task(
&self,
gateway: Arc<Gateway>,
) -> std::io::Result<tokio::task::JoinHandle<()>> {
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<R: AsyncRead + Unpin>(
reader: &mut BufReader<R>,
) -> std::io::Result<Option<LogMessage>> {
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);
}
}