#[cfg(test)]
mod tests {
use bytes::BytesMut;
use tokio::sync::mpsc::channel;
use crate::sources::syslog::{TcpSyslogSource, constants::Message};
#[tokio::test]
async fn test_framing_octet_counting() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "<165>1 2023-01-01T00:00:00Z msg";
let frame = format!("{} {}", msg.len(), msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (ip, received) = rx.recv().await.unwrap();
assert_eq!(ip.as_ref(), "127.0.0.1");
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_framing_newline_delimiter() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "<165>1 2023-01-01T00:00:00Z simple message";
let frame = format!("{}\n", msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "192.168.1.1", &tx)
.await
.unwrap();
let (ip, received) = rx.recv().await.unwrap();
assert_eq!(ip.as_ref(), "192.168.1.1");
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_framing_message_with_newlines() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "<165>1 2023-01-01T00:00:00Z - - - msg with\nnewlines\n";
let frame = format!("{} {}", msg.len(), msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "10.0.0.1", &tx)
.await
.unwrap();
let (ip, received) = rx.recv().await.unwrap();
assert_eq!(ip.as_ref(), "10.0.0.1");
assert_eq!(received, msg);
assert!(received.contains(&b'\n'));
}
#[tokio::test]
async fn test_framing_partial_messages() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "<165>1 2023-01-01T00:00:00Z complete";
let frame = format!("{} {}", msg.len(), msg);
let (part1, part2) = frame.split_at(15);
TcpSyslogSource::process_buffer(&mut buffer, part1.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
TcpSyslogSource::process_buffer(&mut buffer, part2.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_framing_multiple_messages() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg1 = "<165>1 2023-01-01T00:00:00Z first";
let msg2 = "<165>1 2023-01-01T00:00:01Z second";
let frame = format!("{} {}{} {}", msg1.len(), msg1, msg2.len(), msg2);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received1) = rx.recv().await.unwrap();
assert_eq!(received1, msg1);
let (_, received2) = rx.recv().await.unwrap();
assert_eq!(received2, msg2);
}
#[tokio::test]
async fn test_framing_empty_lines_newline_delimiter() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "\n\n\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
}
#[tokio::test]
async fn test_framing_whitespace_only_lines() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = " \n\t\t\n \n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
}
#[tokio::test]
async fn test_framing_carriage_return_handling() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "<165>1 2023-01-01T00:00:00Z test";
let frame = format!("{}\r\n", msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_framing_zero_length_message() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "0 ";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
}
#[tokio::test]
async fn test_framing_large_valid_message() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "X".repeat(1_000_000);
let frame = format!("{} {}", msg.len(), msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received.len(), 1_000_000);
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_framing_mixed_octet_and_newline() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg1 = "<165>1 2023-01-01T00:00:00Z octet";
let msg2 = "<165>1 2023-01-01T00:00:01Z newline";
let frame = format!("{} {}{}\n", msg1.len(), msg1, msg2);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received1) = rx.recv().await.unwrap();
assert_eq!(received1, msg1);
let (_, received2) = rx.recv().await.unwrap();
assert_eq!(received2, msg2);
}
#[tokio::test]
async fn test_framing_newline_then_octet() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg1 = "<165>1 2023-01-01T00:00:00Z newline";
let msg2 = "<165>1 2023-01-01T00:00:01Z octet";
let frame = format!("{}\n{} {}", msg1, msg2.len(), msg2);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received1) = rx.recv().await.unwrap();
assert_eq!(received1, msg1);
let (_, received2) = rx.recv().await.unwrap();
assert_eq!(received2, msg2);
}
#[tokio::test]
async fn test_framing_alternating_methods() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg1 = "msg1";
let msg2 = "msg2";
let msg3 = "msg3";
let frame = format!("{} {}{}\n{} {}", msg1.len(), msg1, msg2, msg3.len(), msg3);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert_eq!(rx.recv().await.unwrap().1, msg1);
assert_eq!(rx.recv().await.unwrap().1, msg2);
assert_eq!(rx.recv().await.unwrap().1, msg3);
}
#[tokio::test]
async fn test_framing_invalid_length_prefix_non_numeric() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "12a <165>1 2023-01-01T00:00:00Z msg\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "12a <165>1 2023-01-01T00:00:00Z msg");
}
#[tokio::test]
async fn test_framing_invalid_length_prefix_negative() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "-10 test\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "-10 test");
}
#[tokio::test]
async fn test_framing_length_mismatch() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "100 short\nactual message\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err()); }
#[tokio::test]
async fn test_framing_extremely_large_length_prefix() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "99999999 msg\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "99999999 msg");
}
#[tokio::test]
async fn test_framing_buffer_overflow_protection() {
let (tx, _rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let large_data = vec![b'X'; 11_000_000];
let result =
TcpSyslogSource::process_buffer(&mut buffer, &large_data, "127.0.0.1", &tx).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_framing_no_space_in_length_prefix() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "12345678901234567890\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "12345678901234567890");
}
#[tokio::test]
async fn test_octet_in_progress_does_not_fallback_newline() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let first = b"12 short\n"; TcpSyslogSource::process_buffer(&mut buffer, first, "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
let second = b"remain"; TcpSyslogSource::process_buffer(&mut buffer, second, "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "short\nremain");
}
#[tokio::test]
async fn test_framing_length_with_leading_zeros() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "0005 hello"; TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "hello");
}
#[tokio::test]
async fn test_zero_length_prefix_fallback_newline() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "0 test\n"; TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "0 test");
}
#[tokio::test]
async fn test_double_space_after_prefix() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "6 hello"; TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, " hello");
}
#[tokio::test]
async fn test_prefix_longer_than_search_window_fallback_newline() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "12345678901 msg\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "12345678901 msg");
}
#[tokio::test]
async fn test_mixed_octet_then_newline_one_buffer() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg1 = "hello";
let msg2 = "line_after";
let frame = format!("{} {}\n{}\n", msg1.len(), msg1, msg2);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, r1) = rx.recv().await.unwrap();
let (_, r2) = rx.recv().await.unwrap();
assert_eq!(r1, msg1);
assert_eq!(r2, msg2);
}
#[tokio::test]
async fn test_mixed_newline_then_octet_one_buffer() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let m1 = "first_line";
let m2 = "second";
let frame = format!("{}\n{} {}", m1, m2.len(), m2);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, r1) = rx.recv().await.unwrap();
let (_, r2) = rx.recv().await.unwrap();
assert_eq!(r1, m1);
assert_eq!(r2, m2);
}
#[tokio::test]
async fn test_partial_header_across_chunks_then_newline() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let m_oct = "hello world!"; let m_nl = "after_octet";
TcpSyslogSource::process_buffer(&mut buffer, b"1", "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
TcpSyslogSource::process_buffer(&mut buffer, b"2 ", "127.0.0.1", &tx)
.await
.unwrap();
assert!(rx.try_recv().is_err());
let chunk3 = format!("{}\n{}\n", m_oct, m_nl);
TcpSyslogSource::process_buffer(&mut buffer, chunk3.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, r1) = rx.recv().await.unwrap();
assert_eq!(r1, m_oct);
let (_, r2) = rx.recv().await.unwrap();
assert_eq!(r2, m_nl);
}
#[tokio::test]
async fn test_invalid_large_then_newline_then_valid_octet() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "99999999 too_big\n7 welcome"; TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, r1) = rx.recv().await.unwrap();
assert_eq!(r1, "99999999 too_big");
let (_, r2) = rx.recv().await.unwrap();
assert_eq!(r2, "welcome");
}
#[tokio::test]
async fn test_many_mixed_messages_in_stream() {
let (tx, mut rx) = channel::<Message>(50);
let mut buffer = BytesMut::new();
let m1 = "alpha"; let m2 = "bravo"; let m3 = "charlie"; let m4 = "delta"; let frame = format!("{}\n{} {}\n{}\n{} {}", m1, m2.len(), m2, m3, m4.len(), m4);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let mut outs = Vec::new();
for _ in 0..4 {
outs.push(rx.recv().await.unwrap().1);
}
assert_eq!(
outs,
vec![
m1.to_string(),
m2.to_string(),
m3.to_string(),
m4.to_string()
]
);
}
#[tokio::test]
async fn test_framing_space_at_position_10() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let frame = "1234567890 msg\n";
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, "1234567890 msg");
}
#[tokio::test]
async fn test_framing_many_small_messages() {
let (tx, mut rx) = channel::<Message>(1000);
let mut buffer = BytesMut::new();
let mut frame = String::new();
for i in 0..100 {
let msg = format!("message_{}", i);
frame.push_str(&format!("{} {}", msg.len(), msg));
}
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
for i in 0..100 {
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, format!("message_{}", i));
}
}
#[tokio::test]
async fn test_framing_many_newline_messages() {
let (tx, mut rx) = channel::<Message>(1000);
let mut buffer = BytesMut::new();
let mut frame = String::new();
for i in 0..100 {
frame.push_str(&format!("message_{}\n", i));
}
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
for i in 0..100 {
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, format!("message_{}", i));
}
}
#[tokio::test]
async fn test_framing_incremental_large_message() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "A".repeat(100_000);
let frame = format!("{} {}", msg.len(), msg);
for chunk in frame.as_bytes().chunks(1024) {
TcpSyslogSource::process_buffer(&mut buffer, chunk, "127.0.0.1", &tx)
.await
.unwrap();
}
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received.len(), 100_000);
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_rfc6587_example_octet_counting() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "<34>1 2003-10-11T22:14:15.003Z mymachine.example.com su - ID47 - BOM'su root' failed for lonvick on /dev/pts/8";
let frame = format!("{} {}", msg.len(), msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, msg);
}
#[tokio::test]
async fn test_rfc6587_octet_counting_preserves_all_data() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg = "line1\nline2\r\nline3\tspaces multiple";
let frame = format!("{} {}", msg.len(), msg);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
let (_, received) = rx.recv().await.unwrap();
assert_eq!(received, msg);
assert!(received.contains(&b'\n'));
assert!(received.contains(&b'\t'));
}
#[tokio::test]
async fn test_rfc6587_trailer_behavior() {
let (tx, mut rx) = channel::<Message>(10);
let mut buffer = BytesMut::new();
let msg1 = "first";
let msg2 = "second";
let frame = format!("{} {}{} {}", msg1.len(), msg1, msg2.len(), msg2);
TcpSyslogSource::process_buffer(&mut buffer, frame.as_bytes(), "127.0.0.1", &tx)
.await
.unwrap();
assert_eq!(rx.recv().await.unwrap().1, msg1);
assert_eq!(rx.recv().await.unwrap().1, msg2);
}
}