From e564027a927f917c1a3e7d1c25038d8251c25494 Mon Sep 17 00:00:00 2001 From: Donemmanuelo Date: Wed, 12 Aug 2026 13:34:43 +0100 Subject: [PATCH] test(encoding): add direct unit tests for sink encodings (#33) Covers the framing and shaping logic that previously had no direct coverage: RFC 3164 and RFC 5424 frame structure, Loki stream grouping by label set with nanosecond timestamps, account_id defaulting, timestamp parsing/fallback robustness, header parsing edge cases, encoding name round-trips, and the cross-transport validation errors. The targeted functions are exposed as pub(crate) so the existing single test module can exercise them directly. --- src/config.rs | 2 +- src/sinks/encoding.rs | 12 +-- src/tests.rs | 187 +++++++++++++++++++++++++++++++++++++++++- 3 files changed, 192 insertions(+), 9 deletions(-) diff --git a/src/config.rs b/src/config.rs index b92f760..700cabb 100644 --- a/src/config.rs +++ b/src/config.rs @@ -208,7 +208,7 @@ fn parse_sink_spec(name: &str) -> Result { } // "Key1:Value1,Key2:Value2" -> [(Key1, Value1), (Key2, Value2)] -fn parse_headers(raw: &str) -> Result> { +pub(crate) fn parse_headers(raw: &str) -> Result> { raw.split(',') .map(str::trim) .filter(|s| !s.is_empty()) diff --git a/src/sinks/encoding.rs b/src/sinks/encoding.rs index b89124d..bd27916 100644 --- a/src/sinks/encoding.rs +++ b/src/sinks/encoding.rs @@ -96,7 +96,7 @@ fn event_json(event: &Event) -> serde_json::Value { // syslog collectors expect for application-level logs. const SYSLOG_PRI: u32 = 134; -fn format_rfc3164(event: &Event) -> String { +pub(crate) fn format_rfc3164(event: &Event) -> String { // TIMESTAMP HOSTNAME TAG: MSG. RFC 3164 timestamps are notionally // "Mmm dd hh:mm:ss" local time, but most collectors (including syslog-ng // and rsyslog defaults) accept ISO 8601 in the field without complaint, @@ -109,7 +109,7 @@ fn format_rfc3164(event: &Event) -> String { ) } -fn format_rfc5424(event: &Event) -> String { +pub(crate) fn format_rfc5424(event: &Event) -> String { // VERSION TIMESTAMP HOSTNAME APP-NAME PROCID MSGID STRUCTURED-DATA MSG format!( "<{}>1 {} auditbridge netbird-audit - AUDIT - {}\n", @@ -120,17 +120,17 @@ fn format_rfc5424(event: &Event) -> String { } #[derive(Debug, Serialize)] -struct LokiStream { +pub(crate) struct LokiStream { stream: HashMap, values: Vec<(String, String)>, } #[derive(Debug, Serialize)] -struct LokiPushRequest { +pub(crate) struct LokiPushRequest { streams: Vec, } -fn build_loki_push_request(events: &[Event]) -> LokiPushRequest { +pub(crate) fn build_loki_push_request(events: &[Event]) -> LokiPushRequest { let mut streams: HashMap = HashMap::new(); for event in events { @@ -176,7 +176,7 @@ fn build_loki_push_request(events: &[Event]) -> LokiPushRequest { } } -fn timestamp_to_nanoseconds(timestamp: &str) -> String { +pub(crate) fn timestamp_to_nanoseconds(timestamp: &str) -> String { DateTime::parse_from_rfc3339(timestamp) .or_else(|_| { let ts = timestamp.trim_end_matches('Z'); diff --git a/src/tests.rs b/src/tests.rs index bdcc869..a03fa62 100644 --- a/src/tests.rs +++ b/src/tests.rs @@ -1,10 +1,12 @@ -use crate::config::Config; +use crate::config::{parse_headers, Config}; use crate::cursor; use crate::metrics::{router, Metrics}; use crate::models::Event; use crate::netbird::NetbirdClient; use crate::retry::{with_retry, RetryConfig}; -use crate::sinks::encoding::Encoding; +use crate::sinks::encoding::{ + build_loki_push_request, format_rfc3164, format_rfc5424, timestamp_to_nanoseconds, Encoding, +}; use crate::sinks::http::HttpSink; use crate::sinks::syslog::{SyslogProtocol, SyslogSink}; use crate::sinks::Sink; @@ -812,3 +814,184 @@ async fn test_run_lets_in_flight_cycle_finish_before_exiting() { "the in-flight cycle must finish and its result apply even though shutdown fired mid-fetch" ); } + +// ---- direct unit tests for the sink encoding layer ---- + +#[test] +fn test_format_rfc3164_frame_structure() { + let event = sample_event("1"); + let frame = format_rfc3164(&event); + + assert!( + frame.starts_with("<134>"), + "expected the PRI prefix, got: {}", + frame + ); + assert!( + frame.contains("2023-01-01T00:00:00Z"), + "timestamp must pass through" + ); + assert!( + frame.contains("auditbridge netbird-audit"), + "expected hostname and tag" + ); + assert!( + frame.contains("test.activity"), + "MSG must carry the event JSON" + ); + assert!(frame.contains("\"event_id\":\"1\"")); + assert!( + frame.ends_with('\n'), + "each frame must be newline-terminated" + ); +} + +#[test] +fn test_format_rfc5424_frame_structure() { + let event = sample_event("1"); + let frame = format_rfc5424(&event); + + assert!( + frame.starts_with("<134>1 "), + "expected PRI and VERSION, got: {}", + frame + ); + assert!(frame.contains("2023-01-01T00:00:00Z")); + assert!(frame.contains("auditbridge netbird-audit - AUDIT")); + assert!(frame.contains("test.activity")); + assert!(frame.ends_with('\n')); +} + +#[test] +fn test_loki_push_request_groups_streams_by_labels() { + let ev1 = sample_event("1"); + let mut ev2 = sample_event("2"); + ev2.activity = "different_activity".to_string(); + + let request = serde_json::to_value(build_loki_push_request(&[ev1, ev2])).unwrap(); + let streams = request["streams"].as_array().unwrap(); + assert_eq!( + streams.len(), + 2, + "events with different labels must split streams" + ); + + let mut activities = std::collections::BTreeSet::new(); + for stream in streams { + let labels = stream["stream"].as_object().unwrap(); + assert_eq!(labels["job"], "netbird-events"); + assert_eq!(labels["account_id"], "acc1"); + activities.insert(labels["activity"].as_str().unwrap()); + + let values = stream["values"].as_array().unwrap(); + assert_eq!(values.len(), 1); + assert_eq!(values[0][0], "1672531200000000000", "nanosecond timestamp"); + assert!(values[0][1].as_str().unwrap().contains("test.activity")); + } + assert_eq!(activities.len(), 2); +} + +#[test] +fn test_loki_push_request_defaults_unknown_account_id() { + let mut event = sample_event("1"); + event.account_id = None; + + let request = serde_json::to_value(build_loki_push_request(&[event])).unwrap(); + let labels = request["streams"][0]["stream"].as_object().unwrap(); + assert_eq!(labels["account_id"], "unknown"); +} + +#[test] +fn test_timestamp_to_nanoseconds_parses_rfc3339() { + assert_eq!( + timestamp_to_nanoseconds("2023-01-01T00:00:00Z"), + "1672531200000000000" + ); +} + +#[test] +fn test_timestamp_to_nanoseconds_handles_offsets_without_z() { + assert_eq!( + timestamp_to_nanoseconds("2023-01-01T00:00:00+00:00"), + "1672531200000000000" + ); + assert_eq!( + timestamp_to_nanoseconds("2023-01-01T01:00:00+01:00"), + "1672531200000000000", + "offsets must be normalized to UTC" + ); +} + +#[test] +fn test_timestamp_to_nanoseconds_handles_unparseable_input() { + // Robustness contract: a malformed timestamp must not panic or fail the + // whole batch encode; it resolves to a numeric fallback. The exact + // fallback value is pinned by a test in #35. + let result = timestamp_to_nanoseconds("not-a-timestamp"); + assert!( + result.chars().all(|c| c.is_ascii_digit()), + "fallback must be numeric, got '{}'", + result + ); +} + +#[test] +fn test_parse_headers_parses_and_trims_entries() { + let headers = parse_headers("Key1:Value1, Key2: Value2").unwrap(); + assert_eq!( + headers, + vec![ + ("Key1".to_string(), "Value1".to_string()), + ("Key2".to_string(), "Value2".to_string()), + ] + ); +} + +#[test] +fn test_parse_headers_ignores_empty_entries() { + let headers = parse_headers("A:1,,B:2,").unwrap(); + assert_eq!( + headers, + vec![ + ("A".to_string(), "1".to_string()), + ("B".to_string(), "2".to_string()) + ] + ); +} + +#[test] +fn test_parse_headers_rejects_entry_without_colon() { + assert!(parse_headers("KeyValue").is_err()); +} + +#[test] +fn test_encoding_parse_round_trips_all_variants() { + for (name, expected) in [ + ("json", Encoding::Json), + ("ndjson", Encoding::Ndjson), + ("loki", Encoding::Loki), + ("syslog3164", Encoding::Syslog3164), + ("syslog5424", Encoding::Syslog5424), + ] { + assert_eq!( + Encoding::parse(name).unwrap(), + expected, + "encoding '{}'", + name + ); + } + assert!(Encoding::parse("unknown").is_err()); +} + +#[test] +fn test_encode_http_body_rejects_syslog_encodings() { + assert!(crate::sinks::encoding::encode_http_body(Encoding::Syslog3164, &[]).is_err()); + assert!(crate::sinks::encoding::encode_http_body(Encoding::Syslog5424, &[]).is_err()); +} + +#[test] +fn test_encode_syslog_frames_rejects_http_encodings() { + assert!(crate::sinks::encoding::encode_syslog_frames(Encoding::Json, &[]).is_err()); + assert!(crate::sinks::encoding::encode_syslog_frames(Encoding::Ndjson, &[]).is_err()); + assert!(crate::sinks::encoding::encode_syslog_frames(Encoding::Loki, &[]).is_err()); +}