use super::*;
#[test]
fn continuation_lines_at_file_boundary_must_not_inherit_previous_timestamp() {
use crate::TrackedChain;
let uuid_a = "aaaaaaaa-1111-2222-3333-444444444444";
let uuid_b = "bbbbbbbb-1111-2222-3333-444444444444";
let ts_old = "2025-01-15 23:58:03.000000";
let ts_new = "2025-01-16 08:37:12.000000";
let seg1: Vec<String> = vec![format!(
"{uuid_a} {ts_old} 95.00% [DEBUG] test.c:1 Last line in rotated file"
)];
let seg2: Vec<String> = vec![
format!("{uuid_b} CHANNEL_DATA:"),
format!("{uuid_b} Channel-State: [CS_EXECUTE]"),
format!("{uuid_b} {ts_new} 95.00% [DEBUG] test.c:1 First timestamped line in new file"),
];
let segments: Vec<(String, Box<dyn Iterator<Item = String>>)> = vec![
("rotated.log".to_string(), Box::new(seg1.into_iter())),
("freeswitch.log".to_string(), Box::new(seg2.into_iter())),
];
let (chain, _) = TrackedChain::new(segments);
let entries: Vec<_> = LogStream::new(chain).collect();
let b_entry = entries
.iter()
.find(|e| e.uuid.as_deref() == Some(uuid_b))
.expect("should find entry for uuid_b");
assert_ne!(
b_entry.timestamp, ts_old,
"continuation lines in a new file segment inherited timestamp \
'{ts_old}' from the previous segment — timestamps must not bleed \
across file boundaries"
);
}
#[test]
fn no_split_on_short_lines() {
let line = format!("variable_call_uuid: [{UUID2}]");
let lines = vec![full_line(UUID1, TS1, "CHANNEL_DATA:"), line];
let mut stream = LogStream::new(lines.into_iter());
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 1);
assert_eq!(stream.stats().lines_split, 0);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_splits_system_lines() {
let line = format!(
"{TS1} 98.03% [INFO] mod_event_socket.c:1752 Event Socket Command from ::1:42864: api sofia jsonstatus{TS2} 97.93% [INFO] mod_event_socket.c:1752 Event Socket Command from ::1:42898: api fsctl pause_check"
);
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 2);
assert_eq!(
entries[0].message,
"Event Socket Command from ::1:42864: api sofia jsonstatus"
);
assert_eq!(
entries[1].message,
"Event Socket Command from ::1:42898: api fsctl pause_check"
);
assert_eq!(stream.stats().lines_split, 1);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_splits_three_entries() {
let ts3 = "2025-01-15 10:30:47.345678";
let line = format!(
"{TS1} 95.00% [INFO] mod.c:1 first{TS2} 96.00% [INFO] mod.c:1 second{ts3} 97.00% [INFO] mod.c:1 third"
);
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 3);
assert_eq!(entries[0].message, "first");
assert_eq!(entries[1].message, "second");
assert_eq!(entries[2].message, "third");
assert_eq!(stream.stats().lines_split, 2);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_oversize_write_contention() {
let entry = |n: usize| {
format!(
"{TS1} 98.77% [INFO] mod_event_socket.c:1754 Event Socket Command from ::1:42864: api db select/ngcs_sip_call_id/entry-{n:04}"
)
};
let count: u64 = 20;
let line: String = (0..count).map(|n| entry(n as usize)).collect();
assert!(
line.len() > super::MAX_LINE_PAYLOAD,
"test fixture should exceed MAX_LINE_PAYLOAD, got {}",
line.len()
);
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len() as u64, count);
for (i, e) in entries.iter().enumerate() {
assert_eq!(
e.message,
format!(
"Event Socket Command from ::1:42864: api db select/ngcs_sip_call_id/entry-{i:04}"
)
);
}
assert_eq!(stream.stats().lines_split, count - 1);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_with_uuid_prefix() {
let line =
format!("{TS1} 95.00% [INFO] mod.c:1 first{UUID1} {TS2} 96.00% [DEBUG] sofia.c:100 second");
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].message, "first");
assert_eq!(entries[1].uuid.as_deref(), Some(UUID1));
assert_eq!(entries[1].message, "second");
assert_eq!(stream.stats().lines_split, 1);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_no_idle_pct_system() {
let line = format!(
"{TS1} [WARNING] sofia_presence.c:4546 Session does not exist, aborting REFER.{TS2} [WARNING] sofia_presence.c:4546 Session does not exist, aborting REFER."
);
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 2);
assert_eq!(
entries[0].message,
"Session does not exist, aborting REFER."
);
assert_eq!(
entries[1].message,
"Session does not exist, aborting REFER."
);
assert_eq!(stream.stats().lines_split, 1);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_no_idle_pct_uuid_suffix() {
let line = format!(
"{TS1} [WARNING] sofia_presence.c:4546 Session does not exist, aborting REFER.{UUID1} {TS2} [NOTICE] sofia.c:1114 Hangup sofia/internal/sos@192.0.2.10:5080 [CS_EXCHANGE_MEDIA] [NORMAL_CLEARING]"
);
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].uuid, None);
assert_eq!(
entries[0].message,
"Session does not exist, aborting REFER."
);
assert_eq!(entries[1].uuid.as_deref(), Some(UUID1));
assert_eq!(entries[1].level, Some(LogLevel::Notice));
assert_eq!(
entries[1].message,
"Hangup sofia/internal/sos@192.0.2.10:5080 [CS_EXCHANGE_MEDIA] [NORMAL_CLEARING]"
);
assert_eq!(stream.stats().lines_split, 1);
assert_accounting(&stream);
}
#[test]
fn timestamp_collision_no_idle_pct_run_on() {
let count: u64 = 15;
let line: String = (0..count)
.map(|n| {
format!(
"2024-04-02 10:31:{:02}.945614 [WARNING] sofia_presence.c:4546 Session does not exist, aborting REFER.",
n + 10
)
})
.collect();
let mut stream = LogStream::new(std::iter::once(line));
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len() as u64, count);
for e in &entries {
assert_eq!(e.message, "Session does not exist, aborting REFER.");
}
assert_eq!(stream.stats().lines_split, count - 1);
assert_accounting(&stream);
}
#[test]
fn truncated_collision_in_channel_data_variable() {
let padding = "x".repeat(2000);
let collision_line = format!(
"{UUID1} variable_long_xml: [{padding}{UUID1} EXECUTE [depth=0] sofia/internal/+15550001234@192.0.2.1 export(foo=bar)"
);
assert!(
collision_line.len() > super::MAX_LINE_PAYLOAD,
"test line must exceed buffer limit, got {}",
collision_line.len()
);
let lines = vec![
full_line(UUID1, TS1, "CHANNEL_DATA:"),
format!("{UUID1} Channel-Name: [sofia/internal/+15550001234@192.0.2.1]"),
format!("{UUID1} variable_direction: [inbound]"),
collision_line,
full_line(UUID1, TS2, "Next log entry"),
];
let entries: Vec<_> = LogStream::new(lines.into_iter()).collect();
assert_eq!(entries[0].message, "CHANNEL_DATA:");
let block = entries[0].block.as_ref().expect("should have block");
match block {
Block::ChannelData { fields, variables } => {
assert_eq!(fields.len(), 1, "should have Channel-Name field");
assert_eq!(fields[0].0, "Channel-Name");
assert_eq!(
variables.len(),
2,
"should have direction + unclosed long_xml"
);
assert_eq!(variables[0].0, "variable_direction");
assert_eq!(variables[0].1, "inbound");
assert_eq!(variables[1].0, "variable_long_xml");
}
other => panic!("expected ChannelData block, got {other:?}"),
}
assert!(
entries[0]
.warnings
.iter()
.any(|w| matches!(w, ParseWarning::OversizeLine { .. })),
"expected buffer overflow warning, got: {:?}",
entries[0].warnings
);
assert!(
entries[0]
.warnings
.iter()
.any(|w| matches!(w, ParseWarning::UnclosedVariable { .. })),
"expected unclosed variable warning, got: {:?}",
entries[0].warnings
);
assert_eq!(entries[1].uuid.as_deref(), Some(UUID1));
assert!(
entries[1].message.starts_with("EXECUTE "),
"split entry should be EXECUTE, got: {}",
entries[1].message
);
assert_eq!(entries.len(), 3);
assert_eq!(entries[2].message, "Next log entry");
}
#[test]
fn channel_data_uuid_drops_mid_block() {
let lines = vec![
full_line(UUID1, TS1, "CHANNEL_DATA:"),
format!("{UUID1} variable_max_forwards: [69]"),
format!("{UUID1} variable_presence_id: [1251@[2001:db8::10]]"),
format!("{UUID1} variable_sip_h_X-Custom-ID: [c4da84eb-88a7-40b2-b90d-e5bc2a0f634e]"),
"variable_sip_h_X-Call-Info: [<urn:test:callid:20260316>;purpose=emergency-CallId]"
.to_string(),
"variable_ep_codec_string: [mod_opus.opus@48000h@20i@2c]".to_string(),
"variable_remote_media_ip: [2001:db8::10]".to_string(),
"variable_remote_media_port: [9952]".to_string(),
"variable_rtp_use_codec_name: [opus]".to_string(),
full_line(UUID1, TS2, "Next entry"),
];
let mut stream = LogStream::new(lines.into_iter());
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].message, "CHANNEL_DATA:");
let block = entries[0].block.as_ref().expect("should have block");
match block {
Block::ChannelData { fields, variables } => {
assert_eq!(fields.len(), 0);
assert_eq!(variables.len(), 8);
assert_eq!(variables[0].0, "variable_max_forwards");
assert_eq!(variables[0].1, "69");
assert_eq!(variables[1].0, "variable_presence_id");
assert_eq!(variables[1].1, "1251@[2001:db8::10]");
assert_eq!(variables[2].0, "variable_sip_h_X-Custom-ID");
assert_eq!(variables[3].0, "variable_sip_h_X-Call-Info");
assert!(variables[3].1.contains("emergency-CallId"));
assert_eq!(variables[4].0, "variable_ep_codec_string");
assert_eq!(variables[7].0, "variable_rtp_use_codec_name");
assert_eq!(variables[7].1, "opus");
}
other => panic!("expected ChannelData block, got {other:?}"),
}
assert_eq!(entries[0].attached.len(), 8);
assert_eq!(entries[1].message, "Next entry");
assert_accounting(&stream);
}
#[test]
fn channel_data_uuid_drops_with_multiline_variable() {
let lines = vec![
full_line(UUID1, TS1, "CHANNEL_DATA:"),
format!("{UUID1} variable_max_forwards: [69]"),
format!("{UUID1} variable_sip_h_X-Custom-ID: [c4da84eb-88a7-40b2-b90d-e5bc2a0f634e]"),
"variable_switch_r_sdp: [v=0\r".to_string(),
"o=FreeSWITCH 1773663549 1773663550 IN IP6 2001:db8::10\r".to_string(),
"s=FreeSWITCH\r".to_string(),
"c=IN IP6 2001:db8::10\r".to_string(),
"t=0 0\r".to_string(),
"m=audio 9952 RTP/AVP 102 101 13\r".to_string(),
"a=rtpmap:102 opus/48000/2\r".to_string(),
"a=ptime:20\r".to_string(),
"]".to_string(),
"variable_ep_codec_string: [mod_opus.opus@48000h@20i@2c]".to_string(),
"variable_direction: [inbound]".to_string(),
full_line(UUID1, TS2, "Next entry"),
];
let mut stream = LogStream::new(lines.into_iter());
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 2);
let block = entries[0].block.as_ref().expect("should have block");
match block {
Block::ChannelData { fields, variables } => {
assert_eq!(fields.len(), 0);
assert_eq!(variables.len(), 5);
assert_eq!(variables[0].0, "variable_max_forwards");
assert_eq!(variables[1].0, "variable_sip_h_X-Custom-ID");
assert_eq!(variables[2].0, "variable_switch_r_sdp");
let sdp = &variables[2].1;
assert!(
sdp.starts_with("v=0\r\n"),
"SDP should start with v=0\\r\\n, got: {sdp:?}"
);
assert!(sdp.contains("m=audio 9952 RTP/AVP 102 101 13\r"));
assert!(sdp.contains("a=ptime:20\r"));
assert!(!sdp.ends_with(']'), "closing bracket should be stripped");
assert_eq!(variables[3].0, "variable_ep_codec_string");
assert_eq!(variables[4].0, "variable_direction");
assert_eq!(variables[4].1, "inbound");
}
other => panic!("expected ChannelData block, got {other:?}"),
}
assert_eq!(entries[0].attached.len(), 13);
assert_accounting(&stream);
}
#[test]
fn channel_data_bare_variable_collision_with_execute() {
let collision = format!(
"variable_call_uuid: {UUID1} EXECUTE [depth=0] \
sofia/internal-v6/1251@[2001:db8::10] export(nolocal:test_var=value)"
);
let lines = vec![
full_line(UUID1, TS1, "CHANNEL_DATA:"),
format!("{UUID1} variable_max_forwards: [69]"),
"variable_DP_MATCH: [ARRAY::create_conference|:create_conference]".to_string(),
collision,
full_line(
UUID1,
TS2,
"EXPORT (export_vars) (REMOTE ONLY) [test_var]=[value]",
),
];
let mut stream = LogStream::new(lines.into_iter());
let entries: Vec<_> = stream.by_ref().collect();
assert_eq!(entries.len(), 3);
let block = entries[0].block.as_ref().expect("should have block");
match block {
Block::ChannelData { fields, variables } => {
assert_eq!(fields.len(), 0);
assert_eq!(variables.len(), 2);
assert_eq!(variables[0].0, "variable_max_forwards");
assert_eq!(variables[1].0, "variable_DP_MATCH");
}
other => panic!("expected ChannelData block, got {other:?}"),
}
assert_eq!(entries[1].uuid.as_deref(), Some(UUID1));
assert_eq!(entries[1].kind, LineKind::Truncated);
assert!(
entries[1].message.starts_with("EXECUTE "),
"truncated line should yield EXECUTE, got: {}",
entries[1].message
);
assert_eq!(entries[2].message_kind.label(), "variable");
assert_accounting(&stream);
}
#[test]
fn oversize_primary_line_warns_on_its_own_entry() {
let long_args = "x".repeat(MAX_LINE_PAYLOAD);
let lines = vec![
full_line(UUID1, TS1, "an ordinary earlier entry"),
full_line(UUID2, TS2, &format!("Ring-Ready {long_args}")),
];
let entries: Vec<_> = LogStream::new(lines.into_iter()).collect();
assert_eq!(entries.len(), 2);
assert!(
entries[0].warnings.is_empty(),
"the earlier entry did not produce the oversize line, got: {:?}",
entries[0].warnings
);
assert!(
entries[1]
.warnings
.iter()
.any(|w| matches!(w, ParseWarning::OversizeLine { .. })),
"expected the oversize warning on the entry the line opened, got: {:?}",
entries[1].warnings
);
}