mod block;
mod collision;
mod entry;
mod stats;
#[cfg(test)]
mod tests;
use std::collections::VecDeque;
use crate::attached::AttachedLines;
use crate::line::{
is_date_at, is_log_header_at, is_uuid_at, parse_line, LineKind, RawLine, UUID_PREFIX_LEN,
};
use crate::message::{classify_message, MessageKind};
use block::BlockBuilder;
use collision::{COLLISION_SCAN_SLACK, MAX_LINE_PAYLOAD};
pub use entry::{Block, LogEntry, ParseWarning};
pub use stats::{ParseStats, UnclassifiedLine, UnclassifiedReason, UnclassifiedTracking};
struct Pending {
entry: LogEntry,
block: BlockBuilder,
}
pub struct LogStream<I> {
lines: I,
last_uuid: String,
last_timestamp: String,
pending: Option<Pending>,
stats: ParseStats,
tracking: UnclassifiedTracking,
line_number: u64,
split_pending: VecDeque<String>,
deferred_warning: Option<ParseWarning>,
}
impl<I: Iterator<Item = String>> LogStream<I> {
pub fn new(lines: I) -> Self {
LogStream {
lines,
last_uuid: String::new(),
last_timestamp: String::new(),
pending: None,
stats: ParseStats::default(),
tracking: UnclassifiedTracking::CountOnly,
line_number: 0,
split_pending: VecDeque::new(),
deferred_warning: None,
}
}
pub fn unclassified_tracking(mut self, level: UnclassifiedTracking) -> Self {
self.tracking = level;
self
}
pub fn stats(&self) -> &ParseStats {
&self.stats
}
pub fn drain_unclassified(&mut self) -> Vec<UnclassifiedLine> {
std::mem::take(&mut self.stats.unclassified_lines)
}
fn record_unclassified(&mut self, reason: UnclassifiedReason, data: Option<&str>) {
self.stats.lines_unclassified += 1;
match self.tracking {
UnclassifiedTracking::CountOnly => {}
UnclassifiedTracking::TrackLines => {
self.stats.unclassified_lines.push(UnclassifiedLine {
line_number: self.line_number,
reason,
data: None,
});
}
UnclassifiedTracking::CaptureData => {
self.stats.unclassified_lines.push(UnclassifiedLine {
line_number: self.line_number,
reason,
data: data.map(|s| s.to_string()),
});
}
}
}
fn take_pending(&mut self) -> Option<LogEntry> {
let mut pending = self.pending.take()?;
let (block, warnings) = pending.block.finish();
pending.entry.block = block;
pending.entry.warnings.extend(warnings);
self.stats.lines_in_entries += 1 + pending.entry.attached.len() as u64;
Some(pending.entry)
}
fn warn(&mut self, warning: ParseWarning) {
match self.pending {
Some(ref mut pending) => pending.entry.warnings.push(warning),
None => self.deferred_warning = Some(warning),
}
}
fn merge_codec_run(&mut self, parsed: &RawLine<'_>, uuid: &str, line: &str) -> bool {
let Some(pending) = self.pending.as_mut() else {
return false;
};
let Some(open_media) = pending.block.codec_media() else {
return false;
};
if pending.entry.uuid != uuid {
return false;
}
let MessageKind::CodecNegotiation { media } = classify_message(parsed.message) else {
return false;
};
if media != open_media {
return false;
}
let warning = pending.block.push_codec_trace(parsed.message);
pending.entry.warnings.extend(warning);
pending.entry.attached.push(line);
true
}
fn accumulate_continuation(&mut self, msg: &str, line: &str) {
let Some(pending) = self.pending.as_mut() else {
return;
};
let warning = pending.block.push_continuation(msg);
pending.entry.warnings.extend(warning);
pending.entry.attached.push(line);
}
fn open_entry(&mut self, parsed: &RawLine<'_>, uuid: String, timestamp: String) {
let message_kind = classify_message(parsed.message);
if !uuid.is_empty() {
self.last_uuid = uuid.clone();
}
if parsed.timestamp.is_some() {
self.last_timestamp = timestamp.clone();
}
let mut block = BlockBuilder::open(&message_kind);
let opening_warning = block.push_codec_trace(parsed.message);
let entry = LogEntry {
uuid,
timestamp,
message: parsed.message.to_string(),
kind: parsed.kind,
message_kind,
level: parsed.level,
idle_pct: parsed.idle_pct.map(|s| s.to_string()),
source: parsed.source.map(|s| s.to_string()),
block: None,
attached: AttachedLines::new(),
line_number: self.line_number,
warnings: self
.deferred_warning
.take()
.into_iter()
.chain(opening_warning)
.collect(),
};
self.pending = Some(Pending { entry, block });
}
}
impl<I: Iterator<Item = String>> LogStream<I> {
fn detect_collision(&mut self, line: String) -> String {
if line.len() > MAX_LINE_PAYLOAD {
self.warn(ParseWarning::OversizeLine {
bytes: line.len() + UUID_PREFIX_LEN + 1,
});
}
let bytes = line.as_bytes();
let min_scan = if is_uuid_at(bytes, 0) {
if bytes.len() > UUID_PREFIX_LEN && bytes[UUID_PREFIX_LEN].is_ascii_digit() {
64 } else {
UUID_PREFIX_LEN }
} else if is_date_at(bytes, 0) {
27 } else {
0
};
let end = bytes.len().saturating_sub(28);
let oversize = bytes.len() > MAX_LINE_PAYLOAD;
let mut splits: Vec<usize> = Vec::new();
let mut chunk_start = 0usize;
let mut offset = min_scan;
while offset <= end {
if is_log_header_at(bytes, offset) {
let split_at = if offset >= chunk_start + UUID_PREFIX_LEN
&& is_uuid_at(bytes, offset - UUID_PREFIX_LEN)
{
offset - UUID_PREFIX_LEN
} else {
offset
};
if split_at > chunk_start {
splits.push(split_at);
chunk_start = split_at;
offset += 27;
} else {
offset = (offset + 27).max(offset + 1);
}
continue;
}
if oversize {
let boundary = chunk_start + MAX_LINE_PAYLOAD;
if offset + COLLISION_SCAN_SLACK >= boundary
&& offset <= boundary + COLLISION_SCAN_SLACK
&& is_uuid_at(bytes, offset)
{
splits.push(offset);
chunk_start = offset;
offset += UUID_PREFIX_LEN;
continue;
}
}
offset += 1;
}
if splits.is_empty() {
return line;
}
let mut tail = line;
let mut chunks: Vec<String> = Vec::with_capacity(splits.len());
for &at in splits.iter().rev() {
chunks.push(tail.split_off(at));
}
chunks.reverse();
self.split_pending.extend(chunks);
tail
}
}
impl<I: Iterator<Item = String>> Iterator for LogStream<I> {
type Item = LogEntry;
fn next(&mut self) -> Option<LogEntry> {
loop {
let line = if let Some(split) = self.split_pending.pop_front() {
self.stats.lines_split += 1;
split
} else {
let Some(line) = self.lines.next() else {
return self.take_pending();
};
if line.starts_with('\x00') {
let yielded = self.take_pending();
self.last_uuid.clear();
self.last_timestamp.clear();
if yielded.is_some() {
return yielded;
}
continue;
}
self.line_number += 1;
self.stats.lines_processed += 1;
self.detect_collision(line)
};
let parsed = parse_line(&line);
match parsed.kind {
LineKind::Full | LineKind::System | LineKind::Truncated => {
let uuid = parsed.uuid.unwrap_or("").to_string();
if self.merge_codec_run(&parsed, &uuid, &line) {
continue;
}
let yielded = self.take_pending();
let timestamp = parsed
.timestamp
.map(|t| t.to_string())
.unwrap_or_else(|| self.last_timestamp.clone());
self.open_entry(&parsed, uuid, timestamp);
if yielded.is_some() {
return yielded;
}
}
LineKind::UuidContinuation => {
let uuid = parsed.uuid.unwrap_or("").to_string();
let continues = !parsed.message.starts_with("EXECUTE ")
&& self.pending.as_ref().is_some_and(|p| p.entry.uuid == uuid);
if continues {
self.accumulate_continuation(parsed.message, &line);
} else {
let yielded = self.take_pending();
self.open_entry(&parsed, uuid, self.last_timestamp.clone());
if yielded.is_some() {
return yielded;
}
}
}
LineKind::BareContinuation => {
if self.pending.is_some() {
self.accumulate_continuation(parsed.message, &line);
} else {
self.record_unclassified(
UnclassifiedReason::OrphanContinuation,
Some(&line),
);
let (uuid, timestamp) =
(self.last_uuid.clone(), self.last_timestamp.clone());
self.open_entry(&parsed, uuid, timestamp);
}
}
LineKind::Empty => {
if let Some(ref mut pending) = self.pending {
pending.entry.attached.push(&line);
} else {
self.stats.lines_empty_orphan += 1;
}
}
}
}
}
}