Skip to main content

freeswitch_sofia_trace_parser/
frame.rs

1use std::fmt;
2use std::io::Read;
3
4use memchr::memmem;
5use tracing::{debug, trace};
6
7use crate::finders::BOUNDARY;
8use crate::types::{
9    Direction, Frame, ParseStats, SkipReason, SkipTracking, Timestamp, Transport, UnparsedRegion,
10};
11
12const RECV_PREFIX: &[u8] = b"recv ";
13const SENT_PREFIX: &[u8] = b"sent ";
14/// Maximum skip size classified as a partial first frame.
15/// Based on IP max datagram size (65535) plus the `\x0B\n` boundary (2 bytes).
16const MAX_PARTIAL_FRAME: usize = 65537;
17
18/// Errors produced during frame parsing (Level 1) and SIP parsing (Level 3).
19///
20/// Returned as `Iterator::Item = Result<T, ParseError>`. The caller decides
21/// whether to skip, log, or fail on each error.
22#[derive(Debug)]
23#[non_exhaustive]
24pub enum ParseError {
25    /// Frame header is malformed (e.g., missing colon, bad timestamp).
26    InvalidHeader(String),
27    /// Reassembled content is not a valid SIP message.
28    InvalidMessage(String),
29    /// Whitespace-only content from TLS/TCP keep-alive probes (RFC 5626).
30    /// Not a parse failure; can be safely ignored.
31    TransportNoise {
32        /// Number of whitespace bytes.
33        bytes: usize,
34        /// Transport of the connection that produced the noise.
35        transport: Transport,
36        /// Remote address of the connection.
37        address: String,
38    },
39    /// Underlying reader returned an I/O error.
40    Io(std::io::Error),
41}
42
43impl fmt::Display for ParseError {
44    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
45        match self {
46            ParseError::InvalidHeader(msg) => write!(f, "invalid frame header: {msg}"),
47            ParseError::InvalidMessage(msg) => write!(f, "invalid SIP message: {msg}"),
48            ParseError::TransportNoise {
49                bytes,
50                transport,
51                address,
52            } => write!(
53                f,
54                "transport noise: {bytes} bytes of non-SIP data from {transport}/{address}"
55            ),
56            ParseError::Io(e) => write!(f, "I/O error: {e}"),
57        }
58    }
59}
60
61impl std::error::Error for ParseError {
62    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
63        match self {
64            ParseError::Io(e) => Some(e),
65            _ => None,
66        }
67    }
68}
69
70impl From<std::io::Error> for ParseError {
71    fn from(e: std::io::Error) -> Self {
72        ParseError::Io(e)
73    }
74}
75
76fn digit(b: u8) -> Option<u8> {
77    match b {
78        b'0'..=b'9' => Some(b - b'0'),
79        _ => None,
80    }
81}
82
83/// Parse up to `MAX_DIGITS` ASCII digits; a value the target type cannot hold
84/// is `None`, as is any non-digit byte.
85fn parse_digits<const MAX_DIGITS: usize, T: TryFrom<u64>>(bytes: &[u8]) -> Option<T> {
86    if bytes.is_empty() || bytes.len() > MAX_DIGITS {
87        return None;
88    }
89    let mut val: u64 = 0;
90    for &b in bytes {
91        val = val.checked_mul(10)?.checked_add(u64::from(digit(b)?))?;
92    }
93    T::try_from(val).ok()
94}
95
96/// Parse timestamp from bytes: either `HH:MM:SS.usec` or `YYYY-MM-DD HH:MM:SS.usec`
97fn parse_timestamp(bytes: &[u8]) -> Option<Timestamp> {
98    // Try full datetime first: YYYY-MM-DD HH:MM:SS.usec (min 26 bytes)
99    if bytes.len() >= 26 && bytes[4] == b'-' && bytes[7] == b'-' && bytes[10] == b' ' {
100        let year = parse_digits::<5, u16>(&bytes[0..4])?;
101        let month = parse_digits::<3, u8>(&bytes[5..7])?;
102        let day = parse_digits::<3, u8>(&bytes[8..10])?;
103        let ts = parse_time_part(&bytes[11..])?;
104        return Some(Timestamp::DateTime {
105            year,
106            month,
107            day,
108            hour: ts.0,
109            min: ts.1,
110            sec: ts.2,
111            usec: ts.3,
112        });
113    }
114    // Time-only: HH:MM:SS.usec (min 15 bytes)
115    let (hour, min, sec, usec) = parse_time_part(bytes)?;
116    Some(Timestamp::TimeOnly {
117        hour,
118        min,
119        sec,
120        usec,
121    })
122}
123
124/// Parse `HH:MM:SS.usec` from bytes, returns (hour, min, sec, usec)
125fn parse_time_part(bytes: &[u8]) -> Option<(u8, u8, u8, u32)> {
126    if bytes.len() < 15 {
127        return None;
128    }
129    if bytes[2] != b':' || bytes[5] != b':' || bytes[8] != b'.' {
130        return None;
131    }
132    let hour = parse_digits::<3, u8>(&bytes[0..2])?;
133    let min = parse_digits::<3, u8>(&bytes[3..5])?;
134    let sec = parse_digits::<3, u8>(&bytes[6..8])?;
135    let usec = parse_digits::<10, u32>(&bytes[9..15])?;
136    Some((hour, min, sec, usec))
137}
138
139/// Fields of one frame header line.
140#[derive(Debug)]
141pub struct FrameHeader {
142    /// Whether FreeSWITCH received or sent the frame.
143    pub direction: Direction,
144    /// Content length declared by the header.
145    pub byte_count: usize,
146    /// Transport the frame travelled over.
147    pub transport: Transport,
148    /// Remote address, as written in the header.
149    pub address: String,
150    /// Time the frame was written.
151    pub timestamp: Timestamp,
152    /// Header length in bytes, including the trailing `\n`.
153    pub header_len: usize,
154}
155
156/// Parse a frame header line from `&[u8]`.
157///
158/// Expected format:
159/// `(recv|sent) <N> bytes (from|to) <transport>/<address> at <timestamp>:\n`
160pub fn parse_frame_header(data: &[u8]) -> Result<FrameHeader, ParseError> {
161    let newline_pos = memchr::memchr(b'\n', data)
162        .ok_or_else(|| ParseError::InvalidHeader("no newline in header".into()))?;
163    let line = &data[..newline_pos];
164    // Strip trailing \r if present
165    let line = line.strip_suffix(b"\r").unwrap_or(line);
166    // Must end with ':'
167    let line = line
168        .strip_suffix(b":")
169        .ok_or_else(|| ParseError::InvalidHeader("header does not end with ':'".into()))?;
170
171    // Direction — both "recv " and "sent " are 5 bytes
172    let direction = if line.starts_with(RECV_PREFIX) {
173        Direction::Recv
174    } else if line.starts_with(SENT_PREFIX) {
175        Direction::Sent
176    } else {
177        return Err(ParseError::InvalidHeader(
178            "expected 'recv' or 'sent'".into(),
179        ));
180    };
181    let mut pos = 5;
182
183    // Byte count: digits until ' '
184    let space = memchr::memchr(b' ', &line[pos..])
185        .ok_or_else(|| ParseError::InvalidHeader("no space after byte count".into()))?;
186    let byte_count = parse_digits::<10, usize>(&line[pos..pos + space])
187        .ok_or_else(|| ParseError::InvalidHeader("invalid byte count".into()))?;
188    pos += space + 1;
189
190    // " bytes from/to "
191    let expected_recv = b"bytes from ";
192    let expected_sent = b"bytes to ";
193    if direction == Direction::Recv {
194        if !line[pos..].starts_with(expected_recv) {
195            return Err(ParseError::InvalidHeader("expected 'bytes from '".into()));
196        }
197        pos += expected_recv.len();
198    } else {
199        if !line[pos..].starts_with(expected_sent) {
200            return Err(ParseError::InvalidHeader("expected 'bytes to '".into()));
201        }
202        pos += expected_sent.len();
203    }
204
205    // Transport: tcp/ udp/ tls/ wss/
206    let transport = if line[pos..].starts_with(b"tcp/") {
207        pos += 4;
208        Transport::Tcp
209    } else if line[pos..].starts_with(b"udp/") {
210        pos += 4;
211        Transport::Udp
212    } else if line[pos..].starts_with(b"tls/") {
213        pos += 4;
214        Transport::Tls
215    } else if line[pos..].starts_with(b"wss/") {
216        pos += 4;
217        Transport::Wss
218    } else {
219        return Err(ParseError::InvalidHeader("unknown transport".into()));
220    };
221
222    // Address: until " at "
223    let at_marker = b" at ";
224    let at_pos = memmem::find(&line[pos..], at_marker)
225        .ok_or_else(|| ParseError::InvalidHeader("no ' at ' in header".into()))?;
226    let address = String::from_utf8_lossy(&line[pos..pos + at_pos]).into_owned();
227    pos += at_pos + at_marker.len();
228
229    // Timestamp: rest of line (after stripping trailing ':' already done)
230    let timestamp = parse_timestamp(&line[pos..])
231        .ok_or_else(|| ParseError::InvalidHeader("invalid timestamp".into()))?;
232
233    Ok(FrameHeader {
234        direction,
235        byte_count,
236        transport,
237        address,
238        timestamp,
239        header_len: newline_pos + 1,
240    })
241}
242
243enum HeaderParse {
244    Ok(FrameHeader),
245    NeedMore,
246    Invalid(ParseError),
247}
248
249enum SyncStep {
250    Ready,
251    End,
252    Failed(ParseError),
253}
254
255enum HeaderStep {
256    Got(FrameHeader),
257    Restart,
258    End,
259    Failed(ParseError),
260}
261
262/// Distinguish a header still arriving from one the parser rejects.
263fn classify_header(data: &[u8]) -> HeaderParse {
264    match parse_frame_header(data) {
265        Ok(header) => HeaderParse::Ok(header),
266        Err(e) => {
267            if memchr::memchr(b'\n', data).is_none() {
268                HeaderParse::NeedMore
269            } else {
270                HeaderParse::Invalid(e)
271            }
272        }
273    }
274}
275
276/// Check if data at given position looks like a valid frame header start.
277/// Used to validate `\x0B\n` boundaries.
278pub fn is_frame_header(data: &[u8]) -> bool {
279    if data.len() < 20 {
280        return false;
281    }
282    let starts_valid = data.starts_with(RECV_PREFIX) || data.starts_with(SENT_PREFIX);
283    if !starts_valid {
284        return false;
285    }
286    // Check that after direction there are digits followed by " bytes "
287    let rest = &data[5..];
288    let space = match memchr::memchr(b' ', rest) {
289        Some(p) => p,
290        None => return false,
291    };
292    if space == 0 || space > 10 {
293        return false;
294    }
295    for &b in &rest[..space] {
296        if !b.is_ascii_digit() {
297            return false;
298        }
299    }
300    rest[space..].starts_with(b" bytes ")
301}
302
303const READ_BUF_SIZE: usize = 32 * 1024;
304
305/// Level 1 streaming parser: splits raw dump bytes into [`Frame`]s.
306///
307/// Reads from any [`Read`] source and yields frames delimited by `\x0B\n`
308/// boundaries. Handles truncated first/last frames, file concatenation,
309/// and garbage recovery.
310///
311/// # Example
312///
313/// ```no_run
314/// use std::fs::File;
315/// use freeswitch_sofia_trace_parser::FrameIterator;
316///
317/// let file = File::open("profile.dump").unwrap();
318/// for frame in FrameIterator::new(file) {
319///     let frame = frame.unwrap();
320///     println!("{} {} bytes {} {}",
321///         frame.timestamp, frame.byte_count, frame.direction, frame.address);
322/// }
323/// ```
324pub struct FrameIterator<R> {
325    reader: R,
326    buf: Vec<u8>,
327    eof: bool,
328    frame_count: u64,
329    offset: u64,
330    stats: ParseStats,
331    skip_tracking: SkipTracking,
332}
333
334impl<R: Read> FrameIterator<R> {
335    /// Create a new frame iterator reading from the given source.
336    pub fn new(reader: R) -> Self {
337        FrameIterator {
338            reader,
339            buf: Vec::with_capacity(READ_BUF_SIZE * 2),
340            eof: false,
341            frame_count: 0,
342            offset: 0,
343            stats: ParseStats::default(),
344            skip_tracking: SkipTracking::CountOnly,
345        }
346    }
347
348    /// Enable capturing of skipped bytes (shorthand for
349    /// [`SkipTracking::CaptureData`]); `false` selects
350    /// [`SkipTracking::CountOnly`]. Whichever of this and
351    /// [`skip_tracking`](Self::skip_tracking) is called last wins.
352    pub fn capture_skipped(mut self, enable: bool) -> Self {
353        self.skip_tracking = if enable {
354            SkipTracking::CaptureData
355        } else {
356            SkipTracking::CountOnly
357        };
358        self
359    }
360
361    /// Set the level of detail for unparsed region tracking.
362    pub fn skip_tracking(mut self, tracking: SkipTracking) -> Self {
363        self.skip_tracking = tracking;
364        self
365    }
366
367    /// Borrow the accumulated parse statistics.
368    pub fn stats(&self) -> &ParseStats {
369        &self.stats
370    }
371
372    /// Take all accumulated unparsed regions, leaving the list empty.
373    pub fn drain_unparsed(&mut self) -> Vec<UnparsedRegion> {
374        self.stats.drain_regions()
375    }
376
377    fn consume(&mut self, n: usize) {
378        self.buf.drain(..n);
379        self.offset += n as u64;
380    }
381
382    fn consume_skipped(&mut self, n: usize, reason: SkipReason) {
383        if self.skip_tracking != SkipTracking::CountOnly {
384            let data = if self.skip_tracking == SkipTracking::CaptureData {
385                Some(self.buf[..n].to_vec())
386            } else {
387                None
388            };
389            self.stats.unparsed_regions.push(UnparsedRegion {
390                offset: self.offset,
391                length: n as u64,
392                reason,
393                data,
394            });
395        }
396        self.stats.bytes_skipped += n as u64;
397        self.consume(n);
398    }
399
400    fn fill_buf(&mut self) -> Result<bool, std::io::Error> {
401        if self.eof {
402            return Ok(false);
403        }
404        let old_len = self.buf.len();
405        self.buf.resize(old_len + READ_BUF_SIZE, 0);
406        let n = self.reader.read(&mut self.buf[old_len..])?;
407        self.buf.truncate(old_len + n);
408        if n == 0 {
409            self.eof = true;
410            return Ok(false);
411        }
412        self.stats.bytes_read += n as u64;
413        Ok(true)
414    }
415
416    /// A SIP header terminator right before a frame boundary is the tail of a
417    /// frame logrotate copied into both files.
418    fn is_replay(&self, skipped: &[u8]) -> bool {
419        if self.frame_count == 0 {
420            return false;
421        }
422        skipped.ends_with(b"\r\n\r\n\x0B\n")
423    }
424
425    /// Find the next `\x0B\n` boundary that is followed by a valid frame header.
426    fn find_boundary(&self, start: usize) -> Option<usize> {
427        let mut search_from = start;
428        loop {
429            let pos = BOUNDARY.find(&self.buf[search_from..])?;
430            let abs_pos = search_from + pos;
431            let after = abs_pos + 2;
432            if after >= self.buf.len() {
433                // Boundary at very end — could be real, but we can't validate header yet
434                // If EOF, accept it as boundary (content ends at \x0B)
435                if self.eof {
436                    return Some(abs_pos);
437                }
438                return None; // Need more data
439            }
440            if is_frame_header(&self.buf[after..]) {
441                return Some(abs_pos);
442            }
443            // \x0B\n in content, not a boundary — skip past it
444            trace!(
445                offset = abs_pos,
446                "found \\x0B\\n in content (not a boundary), skipping"
447            );
448            search_from = abs_pos + 2;
449        }
450    }
451
452    /// Skip to the first valid frame header in the buffer (for partial first frames).
453    fn skip_to_first_header(&mut self) -> Option<usize> {
454        if is_frame_header(&self.buf) {
455            return Some(0);
456        }
457        // Look for \x0B\n followed by a valid header
458        let mut search_from = 0;
459        loop {
460            let pos = BOUNDARY.find(&self.buf[search_from..])?;
461            let abs_pos = search_from + pos;
462            let after = abs_pos + 2;
463            if after < self.buf.len() && is_frame_header(&self.buf[after..]) {
464                debug!(skipped_bytes = after, "skipped partial first frame");
465                return Some(after);
466            }
467            search_from = abs_pos + 2;
468        }
469    }
470
471    /// Drop a partial first frame so the buffer starts at a frame header.
472    fn sync_to_first_header(&mut self) -> SyncStep {
473        loop {
474            match self.skip_to_first_header() {
475                Some(offset) => {
476                    if offset > 0 {
477                        let reason = if offset <= MAX_PARTIAL_FRAME {
478                            SkipReason::PartialFirstFrame
479                        } else {
480                            SkipReason::OversizedFrame
481                        };
482                        self.consume_skipped(offset, reason);
483                    }
484                    return SyncStep::Ready;
485                }
486                None => {
487                    if self.eof {
488                        debug!("no valid frame header found in entire input");
489                        let remaining = self.buf.len();
490                        if remaining > 0 {
491                            self.consume_skipped(remaining, SkipReason::InvalidHeader);
492                        }
493                        return SyncStep::End;
494                    }
495                    if let Err(e) = self.fill_buf() {
496                        return SyncStep::Failed(ParseError::Io(e));
497                    }
498                }
499            }
500        }
501    }
502
503    /// Drop `\n` and `\r\n` padding between frames.
504    fn strip_padding(&mut self) {
505        let mut strip = 0;
506        while strip < self.buf.len() {
507            if self.buf[strip] == b'\n' {
508                strip += 1;
509            } else if strip + 1 < self.buf.len()
510                && self.buf[strip] == b'\r'
511                && self.buf[strip + 1] == b'\n'
512            {
513                strip += 2;
514            } else {
515                break;
516            }
517        }
518        if strip > 0 {
519            self.consume(strip);
520        }
521    }
522
523    /// Parse the header at the head of the buffer, reading more as it needs.
524    fn read_header(&mut self) -> HeaderStep {
525        loop {
526            match classify_header(&self.buf) {
527                HeaderParse::Ok(h) => return HeaderStep::Got(h),
528                HeaderParse::NeedMore => {
529                    if self.eof {
530                        debug!("truncated frame header at EOF");
531                        let remaining = self.buf.len();
532                        if remaining > 0 {
533                            self.consume_skipped(remaining, SkipReason::InvalidHeader);
534                        }
535                        return HeaderStep::End;
536                    }
537                    if let Err(e) = self.fill_buf() {
538                        return HeaderStep::Failed(ParseError::Io(e));
539                    }
540                }
541                HeaderParse::Invalid(e) => {
542                    let header_preview: String = self
543                        .buf
544                        .iter()
545                        .take(200)
546                        .take_while(|&&b| b != b'\n')
547                        .map(|&b| {
548                            if b.is_ascii_graphic() || b == b' ' {
549                                b as char
550                            } else {
551                                '.'
552                            }
553                        })
554                        .collect();
555                    if header_preview.starts_with("dump started at ") {
556                        let skip = memchr::memchr(b'\n', &self.buf)
557                            .map(|p| {
558                                let mut end = p + 1;
559                                while end < self.buf.len() && self.buf[end] == b'\n' {
560                                    end += 1;
561                                }
562                                end
563                            })
564                            .unwrap_or(self.buf.len());
565                        debug!(
566                            header = %header_preview,
567                            skipped_bytes = skip,
568                            "skipped dump restart marker",
569                        );
570                        self.consume(skip);
571                        return HeaderStep::Restart;
572                    }
573                    let skip = if let Some(b) = self.find_boundary(0) {
574                        b + 2
575                    } else {
576                        memchr::memchr(b'\n', &self.buf)
577                            .map(|p| p + 1)
578                            .unwrap_or(self.buf.len())
579                    };
580                    let reason = self.classify_skip(skip);
581                    self.consume_skipped(skip, reason);
582                    return HeaderStep::Failed(e);
583                }
584            }
585        }
586    }
587
588    /// Reason for dropping `skip` bytes that failed header parsing.
589    fn classify_skip(&self, skip: usize) -> SkipReason {
590        if self.buf.starts_with(RECV_PREFIX) || self.buf.starts_with(SENT_PREFIX) {
591            SkipReason::InvalidHeader
592        } else if skip > MAX_PARTIAL_FRAME {
593            SkipReason::OversizedFrame
594        } else if self.frame_count == 0 {
595            SkipReason::PartialFirstFrame
596        } else if self.is_replay(&self.buf[..skip]) {
597            SkipReason::ReplayedFrame
598        } else {
599            SkipReason::MidStreamSkip
600        }
601    }
602
603    /// Read one frame's content, from the header's end to the frame boundary.
604    fn read_content(&mut self, header: FrameHeader) -> Option<Result<Frame, ParseError>> {
605        let FrameHeader {
606            direction,
607            byte_count,
608            transport,
609            address,
610            timestamp,
611            header_len,
612        } = header;
613        let content_start = header_len;
614        let expected_end = content_start + byte_count;
615
616        // Find the boundary for this frame.
617        // Strategy: first check at the expected position (content_start + byte_count),
618        // then fall back to scanning. This handles file concatenation where \x0B\n
619        // is followed by garbage from the next file's truncated first frame.
620        let content = loop {
621            // Ensure we have enough data to check the expected position, but
622            // never buffer a declared count further than a frame can reach:
623            // past that, boundary scanning takes over.
624            let fill_to = (expected_end + 1).min(content_start + MAX_PARTIAL_FRAME);
625            while self.buf.len() <= fill_to && !self.eof {
626                if let Err(e) = self.fill_buf() {
627                    return Some(Err(ParseError::Io(e)));
628                }
629            }
630
631            // Check at expected position first (byte_count hint)
632            if expected_end < self.buf.len() && self.buf[expected_end] == 0x0B {
633                let has_newline =
634                    expected_end + 1 < self.buf.len() && self.buf[expected_end + 1] == b'\n';
635                let at_eof = expected_end + 1 >= self.buf.len() && self.eof;
636
637                if has_newline || at_eof {
638                    let content = self.buf[content_start..expected_end].to_vec();
639                    let drain_to = if has_newline {
640                        expected_end + 2
641                    } else {
642                        expected_end + 1
643                    };
644                    self.consume(drain_to);
645                    self.frame_count += 1;
646                    break content;
647                }
648            }
649
650            // Fall back to scanning for \x0B\n + valid header
651            if let Some(boundary_pos) = self.find_boundary(content_start) {
652                let content = self.buf[content_start..boundary_pos].to_vec();
653                let drain_to = boundary_pos + 2;
654                self.consume(drain_to);
655                self.frame_count += 1;
656
657                if content.len() != byte_count {
658                    debug!(
659                        frame = self.frame_count,
660                        expected = byte_count,
661                        actual = content.len(),
662                        "frame content size mismatch"
663                    );
664                }
665
666                break content;
667            }
668
669            if self.eof {
670                // Last frame — no trailing \x0B\n
671                let end = if self.buf.last() == Some(&0x0B) {
672                    self.buf.len() - 1
673                } else {
674                    self.buf.len()
675                };
676                let content = self.buf[content_start..end].to_vec();
677                let len = self.buf.len();
678                self.consume(len);
679                self.frame_count += 1;
680
681                if content.len() < byte_count {
682                    let missing = byte_count - content.len();
683                    debug!(
684                        frame = self.frame_count,
685                        expected = byte_count,
686                        actual = content.len(),
687                        missing,
688                        "incomplete frame at EOF"
689                    );
690                    self.stats.incomplete_frames += 1;
691                    self.stats.incomplete_frame_bytes += missing as u64;
692                    if self.skip_tracking != SkipTracking::CountOnly {
693                        self.stats.unparsed_regions.push(UnparsedRegion {
694                            offset: self.offset,
695                            length: missing as u64,
696                            reason: SkipReason::IncompleteFrame,
697                            data: None,
698                        });
699                    }
700                } else if content.len() != byte_count {
701                    debug!(
702                        frame = self.frame_count,
703                        expected = byte_count,
704                        actual = content.len(),
705                        "last frame content size mismatch"
706                    );
707                }
708
709                break content;
710            }
711
712            if let Err(e) = self.fill_buf() {
713                return Some(Err(ParseError::Io(e)));
714            }
715        };
716
717        Some(Ok(Frame {
718            direction,
719            byte_count,
720            transport,
721            address,
722            timestamp,
723            content,
724        }))
725    }
726}
727
728impl<R: Read> Iterator for FrameIterator<R> {
729    type Item = Result<Frame, ParseError>;
730
731    fn next(&mut self) -> Option<Self::Item> {
732        let header = loop {
733            if self.buf.is_empty() && !self.eof {
734                if let Err(e) = self.fill_buf() {
735                    return Some(Err(ParseError::Io(e)));
736                }
737            }
738            if self.buf.is_empty() {
739                return None;
740            }
741
742            if self.frame_count == 0 {
743                match self.sync_to_first_header() {
744                    SyncStep::Ready => {}
745                    SyncStep::End => return None,
746                    SyncStep::Failed(e) => return Some(Err(e)),
747                }
748            }
749
750            self.strip_padding();
751            if self.buf.is_empty() {
752                continue;
753            }
754
755            match self.read_header() {
756                HeaderStep::Got(header) => break header,
757                HeaderStep::Restart => continue,
758                HeaderStep::End => return None,
759                HeaderStep::Failed(e) => return Some(Err(e)),
760            }
761        };
762
763        self.read_content(header)
764    }
765}
766
767#[cfg(test)]
768mod tests {
769    use super::*;
770    use crate::types::SkipTracking;
771
772    #[test]
773    fn parse_recv_ipv4_tcp() {
774        let header = b"recv 100 bytes from tcp/192.168.1.1:5060 at 00:00:01.350874:\n";
775        let h = parse_frame_header(header).unwrap();
776        assert_eq!(h.direction, Direction::Recv);
777        assert_eq!(h.byte_count, 100);
778        assert_eq!(h.transport, Transport::Tcp);
779        assert_eq!(h.address, "192.168.1.1:5060");
780        assert_eq!(
781            h.timestamp,
782            Timestamp::TimeOnly {
783                hour: 0,
784                min: 0,
785                sec: 1,
786                usec: 350874
787            }
788        );
789        assert_eq!(h.header_len, header.len());
790    }
791
792    #[test]
793    fn parse_recv_ipv6_tcp() {
794        let header = b"recv 1440 bytes from tcp/[2001:4958:10:14::4]:30046 at 13:03:21.674883:\n";
795        let h = parse_frame_header(header).unwrap();
796        assert_eq!(h.direction, Direction::Recv);
797        assert_eq!(h.byte_count, 1440);
798        assert_eq!(h.transport, Transport::Tcp);
799        assert_eq!(h.address, "[2001:4958:10:14::4]:30046");
800        assert_eq!(
801            h.timestamp,
802            Timestamp::TimeOnly {
803                hour: 13,
804                min: 3,
805                sec: 21,
806                usec: 674883
807            }
808        );
809    }
810
811    #[test]
812    fn parse_sent_ipv6_tcp() {
813        let header = b"sent 681 bytes to tcp/[2001:4958:10:14::4]:30046 at 13:03:21.675500:\n";
814        let h = parse_frame_header(header).unwrap();
815        assert_eq!(h.direction, Direction::Sent);
816        assert_eq!(h.byte_count, 681);
817        assert_eq!(h.transport, Transport::Tcp);
818        assert_eq!(h.address, "[2001:4958:10:14::4]:30046");
819    }
820
821    #[test]
822    fn parse_recv_udp() {
823        let header = b"recv 457 bytes from udp/10.0.0.1:5060 at 00:19:47.123456:\n";
824        let h = parse_frame_header(header).unwrap();
825        assert_eq!(h.direction, Direction::Recv);
826        assert_eq!(h.transport, Transport::Udp);
827    }
828
829    #[test]
830    fn parse_sent_tls() {
831        let header = b"sent 500 bytes to tls/10.0.0.1:5061 at 12:00:00.000000:\n";
832        let h = parse_frame_header(header).unwrap();
833        assert_eq!(h.direction, Direction::Sent);
834        assert_eq!(h.byte_count, 500);
835        assert_eq!(h.transport, Transport::Tls);
836    }
837
838    #[test]
839    fn parse_full_datetime_timestamp() {
840        let header = b"recv 100 bytes from tcp/192.168.1.1:5060 at 2026-02-01 10:00:00.000000:\n";
841        let h = parse_frame_header(header).unwrap();
842        assert_eq!(
843            h.timestamp,
844            Timestamp::DateTime {
845                year: 2026,
846                month: 2,
847                day: 1,
848                hour: 10,
849                min: 0,
850                sec: 0,
851                usec: 0
852            }
853        );
854    }
855
856    #[test]
857    fn parse_invalid_header() {
858        assert!(parse_frame_header(b"invalid header\n").is_err());
859        assert!(
860            parse_frame_header(b"recv abc bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n")
861                .is_err()
862        );
863    }
864
865    #[test]
866    fn is_frame_header_valid() {
867        assert!(is_frame_header(
868            b"recv 100 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n"
869        ));
870        assert!(is_frame_header(
871            b"sent 681 bytes to tcp/[::1]:5060 at 00:00:00.000000:\n"
872        ));
873        assert!(!is_frame_header(b"not a header"));
874        assert!(!is_frame_header(b"recv abc bytes"));
875        assert!(!is_frame_header(b""));
876    }
877
878    #[test]
879    fn frame_iterator_single_frame() {
880        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
881        let frames: Vec<Frame> = FrameIterator::new(&data[..])
882            .collect::<Result<Vec<_>, _>>()
883            .unwrap();
884        assert_eq!(frames.len(), 1);
885        assert_eq!(frames[0].content, b"hello");
886        assert_eq!(frames[0].byte_count, 5);
887    }
888
889    #[test]
890    fn frame_iterator_multiple_frames() {
891        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\nsent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n";
892        let frames: Vec<Frame> = FrameIterator::new(&data[..])
893            .collect::<Result<Vec<_>, _>>()
894            .unwrap();
895        assert_eq!(frames.len(), 2);
896        assert_eq!(frames[0].content, b"hello");
897        assert_eq!(frames[0].direction, Direction::Recv);
898        assert_eq!(frames[1].content, b"world");
899        assert_eq!(frames[1].direction, Direction::Sent);
900    }
901
902    #[test]
903    fn frame_iterator_vt_in_content() {
904        // \x0B in content but not followed by valid header — should NOT split
905        let mut data = Vec::new();
906        data.extend_from_slice(b"recv 15 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n");
907        data.extend_from_slice(b"he\x0B\nllo world!!");
908        data.extend_from_slice(b"\x0B\n");
909        let frames: Vec<Frame> = FrameIterator::new(&data[..])
910            .collect::<Result<Vec<_>, _>>()
911            .unwrap();
912        assert_eq!(frames.len(), 1);
913        assert_eq!(frames[0].content, b"he\x0B\nllo world!!");
914    }
915
916    #[test]
917    fn frame_iterator_eof_without_boundary() {
918        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello";
919        let frames: Vec<Frame> = FrameIterator::new(&data[..])
920            .collect::<Result<Vec<_>, _>>()
921            .unwrap();
922        assert_eq!(frames.len(), 1);
923        assert_eq!(frames[0].content, b"hello");
924    }
925
926    #[test]
927    fn frame_iterator_eof_with_lone_vt() {
928        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B";
929        let frames: Vec<Frame> = FrameIterator::new(&data[..])
930            .collect::<Result<Vec<_>, _>>()
931            .unwrap();
932        assert_eq!(frames.len(), 1);
933        assert_eq!(frames[0].content, b"hello");
934    }
935
936    #[test]
937    fn frame_iterator_partial_first_frame() {
938        // Data starts with garbage, then a valid boundary + frame
939        let mut data = Vec::new();
940        data.extend_from_slice(b"partial garbage data");
941        data.extend_from_slice(b"\x0B\n");
942        data.extend_from_slice(
943            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
944        );
945        let frames: Vec<Frame> = FrameIterator::new(&data[..])
946            .collect::<Result<Vec<_>, _>>()
947            .unwrap();
948        assert_eq!(frames.len(), 1);
949        assert_eq!(frames[0].content, b"hello");
950    }
951
952    #[test]
953    fn frame_iterator_truncated_last_frame() {
954        // Complete frame followed by truncated frame at EOF (no \x0B\n)
955        let mut data = Vec::new();
956        data.extend_from_slice(
957            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
958        );
959        data.extend_from_slice(b"sent 3 bytes to tcp/1.1.1.1:5060 at 00:00:01.000000:\nbye");
960        let frames: Vec<Frame> = FrameIterator::new(&data[..])
961            .collect::<Result<Vec<_>, _>>()
962            .unwrap();
963        assert_eq!(frames.len(), 2);
964        assert_eq!(frames[0].content, b"hello");
965        assert_eq!(frames[1].content, b"bye");
966    }
967
968    #[test]
969    fn frame_iterator_file_concatenation() {
970        // Simulates `cat dump.20 dump.21 | parser`
971        // File 1: truncated start + valid frame + complete end
972        // File 2: truncated start (no header) + valid frame
973        let mut data = Vec::new();
974
975        // File 1: starts with valid header
976        data.extend_from_slice(
977            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
978        );
979        data.extend_from_slice(
980            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
981        );
982
983        // File 2: starts with truncated frame data (no header), then boundary, then valid frame
984        data.extend_from_slice(b"some truncated SIP content from previous rotation\r\n\r\n");
985        data.extend_from_slice(b"\x0B\n");
986        data.extend_from_slice(
987            b"recv 3 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\nfoo\x0B\n",
988        );
989
990        let items: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
991        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
992        assert_eq!(frames.len(), 3);
993        assert_eq!(frames[0].content, b"hello");
994        assert_eq!(frames[1].content, b"world");
995        assert_eq!(frames[2].content, b"foo");
996        assert_eq!(frames[2].address, "2.2.2.2:5060");
997    }
998
999    #[test]
1000    fn frame_iterator_file_concatenation_mid_stream_garbage() {
1001        // The join point between files produces garbage that looks like:
1002        // ...last_content\x0B\ntruncated_first_of_next_file\x0B\nvalid_header...
1003        // The truncated part is NOT a valid header, so recovery should skip it
1004        let mut data = Vec::new();
1005
1006        // Last frame of file 1
1007        data.extend_from_slice(
1008            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1009        );
1010
1011        // Truncated first frame of file 2 (mid-SIP content, no frame header)
1012        data.extend_from_slice(b"Content-Type: application/sdp\r\n\r\nv=0\r\n");
1013        data.extend_from_slice(b"\x0B\n");
1014
1015        // Valid second frame of file 2
1016        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1017
1018        let items: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1019        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1020        assert_eq!(frames.len(), 2);
1021        assert_eq!(frames[0].content, b"hello");
1022        assert_eq!(frames[1].content, b"bar");
1023    }
1024
1025    #[test]
1026    fn frame_iterator_empty_input() {
1027        let data: &[u8] = b"";
1028        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(data).collect();
1029        assert!(frames.is_empty());
1030    }
1031
1032    #[test]
1033    fn frame_iterator_only_garbage() {
1034        let data = b"this is not a SIP trace dump at all, just garbage text";
1035        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1036        let frames: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1037        assert!(frames.is_empty());
1038        let stats = iter.stats();
1039        assert_eq!(stats.bytes_read, data.len() as u64);
1040        assert_eq!(stats.bytes_skipped, data.len() as u64);
1041        assert_eq!(stats.unparsed_regions.len(), 1);
1042        assert_eq!(stats.unparsed_regions[0].reason, SkipReason::InvalidHeader);
1043    }
1044
1045    #[test]
1046    fn frame_iterator_truncated_header_at_eof() {
1047        let data = b"recv 5 bytes from tcp/1.1.1.1:5060";
1048        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1049        let frames: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1050        assert!(frames.is_empty());
1051        let stats = iter.stats();
1052        assert_eq!(stats.bytes_read, data.len() as u64);
1053        assert_eq!(stats.bytes_skipped, data.len() as u64);
1054    }
1055
1056    #[test]
1057    fn frame_iterator_dump_marker_at_eof() {
1058        // A dump restart marker at the end of input (with trailing \n\n as in real dumps)
1059        // should be silently consumed, not returned as an error.
1060        let mut data = Vec::new();
1061        data.extend_from_slice(
1062            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1063        );
1064        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1065
1066        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1067        assert_eq!(frames.len(), 1);
1068        assert!(frames[0].is_ok());
1069        assert_eq!(frames[0].as_ref().unwrap().content, b"hello");
1070    }
1071
1072    #[test]
1073    fn frame_iterator_dump_marker_mid_stream() {
1074        // A dump restart marker between two valid frames (with trailing \n\n as in real dumps)
1075        // should be skipped, and both frames should parse successfully.
1076        let mut data = Vec::new();
1077        data.extend_from_slice(
1078            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1079        );
1080        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1081        data.extend_from_slice(b"sent 3 bytes to tcp/2.2.2.2:5060 at 00:00:01.000000:\nbye\x0B\n");
1082
1083        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1084        assert_eq!(frames.len(), 2);
1085        assert_eq!(frames[0].as_ref().unwrap().content, b"hello");
1086        assert_eq!(frames[1].as_ref().unwrap().content, b"bye");
1087    }
1088
1089    #[test]
1090    fn frame_iterator_multiple_newlines_after_boundary() {
1091        // Multiple \n and \r\n between frames should all be stripped.
1092        let mut data = Vec::new();
1093        data.extend_from_slice(
1094            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1095        );
1096        data.extend_from_slice(b"\n\r\n\n");
1097        data.extend_from_slice(
1098            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
1099        );
1100
1101        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1102            .collect::<Result<Vec<_>, _>>()
1103            .unwrap();
1104        assert_eq!(frames.len(), 2);
1105        assert_eq!(frames[0].content, b"hello");
1106        assert_eq!(frames[1].content, b"world");
1107    }
1108
1109    #[test]
1110    fn stats_clean_input() {
1111        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
1112        let mut iter = FrameIterator::new(&data[..]);
1113        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1114        assert_eq!(frames.len(), 1);
1115        let stats = iter.stats();
1116        assert_eq!(stats.bytes_read, data.len() as u64);
1117        assert_eq!(stats.bytes_skipped, 0);
1118        assert!(stats.unparsed_regions.is_empty());
1119    }
1120
1121    #[test]
1122    fn stats_multiple_frames() {
1123        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\nsent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n";
1124        let mut iter = FrameIterator::new(&data[..]);
1125        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1126        assert_eq!(frames.len(), 2);
1127        let stats = iter.stats();
1128        assert_eq!(stats.bytes_read, data.len() as u64);
1129        assert_eq!(stats.bytes_skipped, 0);
1130        assert!(stats.unparsed_regions.is_empty());
1131    }
1132
1133    #[test]
1134    fn stats_partial_first_frame() {
1135        let mut data = Vec::new();
1136        data.extend_from_slice(b"partial garbage data");
1137        data.extend_from_slice(b"\x0B\n");
1138        data.extend_from_slice(
1139            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1140        );
1141        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1142        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1143        assert_eq!(frames.len(), 1);
1144        let stats = iter.stats();
1145        assert_eq!(stats.bytes_read, data.len() as u64);
1146        // "partial garbage data" + "\x0B\n" = 21 bytes skipped
1147        let skipped = b"partial garbage data\x0B\n".len() as u64;
1148        assert_eq!(stats.bytes_skipped, skipped);
1149        assert_eq!(stats.unparsed_regions.len(), 1);
1150        assert_eq!(stats.unparsed_regions[0].offset, 0);
1151        assert_eq!(stats.unparsed_regions[0].length, skipped);
1152        assert_eq!(
1153            stats.unparsed_regions[0].reason,
1154            crate::types::SkipReason::PartialFirstFrame
1155        );
1156        assert!(stats.unparsed_regions[0].data.is_none());
1157    }
1158
1159    #[test]
1160    fn stats_partial_first_frame_capture() {
1161        let mut data = Vec::new();
1162        data.extend_from_slice(b"partial garbage data");
1163        data.extend_from_slice(b"\x0B\n");
1164        data.extend_from_slice(
1165            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1166        );
1167        let mut iter = FrameIterator::new(&data[..]).capture_skipped(true);
1168        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1169        assert_eq!(frames.len(), 1);
1170        let stats = iter.stats();
1171        assert_eq!(stats.unparsed_regions.len(), 1);
1172        let region = &stats.unparsed_regions[0];
1173        assert_eq!(
1174            region.data.as_deref(),
1175            Some(b"partial garbage data\x0B\n".as_slice())
1176        );
1177    }
1178
1179    #[test]
1180    fn stats_mid_stream_partial_frame() {
1181        // SIP content between valid frames (file concatenation scenario)
1182        let mut data = Vec::new();
1183        data.extend_from_slice(
1184            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1185        );
1186        data.extend_from_slice(b"Content-Type: application/sdp\r\n\r\nv=0\r\n");
1187        data.extend_from_slice(b"\x0B\n");
1188        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1189
1190        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1191        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1192        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1193        assert_eq!(frames.len(), 2);
1194        let stats = iter.stats();
1195        assert!(stats.bytes_skipped > 0);
1196        assert_eq!(stats.unparsed_regions.len(), 1);
1197        assert_eq!(
1198            stats.unparsed_regions[0].reason,
1199            crate::types::SkipReason::MidStreamSkip
1200        );
1201    }
1202
1203    #[test]
1204    fn stats_replayed_frame() {
1205        // Simulate logrotate: a frame's tail (SIP headers ending with \r\n\r\n\x0B\n)
1206        // appears between two valid frames at a file boundary.
1207        let frame1 = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
1208        let replay = b"Route: <sip:10.0.0.1:5060;lr>\r\nContent-Length: 0\r\n\r\n\x0B\n";
1209        let frame2 = b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n";
1210
1211        let mut data = Vec::new();
1212        data.extend_from_slice(frame1);
1213        data.extend_from_slice(replay);
1214        data.extend_from_slice(frame2);
1215
1216        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1217        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1218        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1219        assert_eq!(frames.len(), 2);
1220        let stats = iter.stats();
1221        assert_eq!(stats.unparsed_regions.len(), 1);
1222        assert_eq!(
1223            stats.unparsed_regions[0].reason,
1224            crate::types::SkipReason::ReplayedFrame
1225        );
1226    }
1227
1228    #[test]
1229    fn stats_incomplete_frame_at_eof() {
1230        // Frame header says 100 bytes but only 20 bytes available before EOF
1231        let mut data = Vec::new();
1232        data.extend_from_slice(
1233            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1234        );
1235        data.extend_from_slice(b"recv 100 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\n");
1236        data.extend_from_slice(b"partial content only");
1237        // No \x0B\n boundary — EOF truncation
1238
1239        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1240        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1241        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1242        assert_eq!(frames.len(), 2, "truncated frame should still be returned");
1243        assert_eq!(frames[1].content, b"partial content only");
1244        assert_eq!(frames[1].byte_count, 100);
1245        let stats = iter.stats();
1246        assert_eq!(stats.unparsed_regions.len(), 1);
1247        assert_eq!(
1248            stats.unparsed_regions[0].reason,
1249            crate::types::SkipReason::IncompleteFrame
1250        );
1251    }
1252
1253    #[test]
1254    fn stats_incomplete_frame_counted_under_count_only() {
1255        let mut data = Vec::new();
1256        data.extend_from_slice(
1257            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1258        );
1259        data.extend_from_slice(b"recv 100 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\n");
1260        data.extend_from_slice(b"partial content only");
1261
1262        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::CountOnly);
1263        let _: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1264        let stats = iter.stats();
1265        assert!(stats.unparsed_regions.is_empty());
1266        assert_eq!(stats.incomplete_frames, 1);
1267        assert_eq!(stats.incomplete_frame_bytes, 80);
1268    }
1269
1270    #[test]
1271    fn stats_invalid_header_skip() {
1272        // Malformed frame header (starts with recv/sent but unparseable)
1273        let mut data = Vec::new();
1274        data.extend_from_slice(
1275            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1276        );
1277        data.extend_from_slice(b"recv CORRUPT HEADER garbage\n");
1278        data.extend_from_slice(b"\x0B\n");
1279        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1280
1281        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1282        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1283        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1284        assert_eq!(frames.len(), 2);
1285        let stats = iter.stats();
1286        assert!(stats.bytes_skipped > 0);
1287        assert_eq!(stats.unparsed_regions.len(), 1);
1288        assert_eq!(
1289            stats.unparsed_regions[0].reason,
1290            crate::types::SkipReason::InvalidHeader
1291        );
1292    }
1293
1294    #[test]
1295    fn stats_oversized_frame_at_start() {
1296        let mut data = Vec::new();
1297        data.resize(MAX_PARTIAL_FRAME + 1, b'x');
1298        data.extend_from_slice(b"\x0B\n");
1299        data.extend_from_slice(
1300            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1301        );
1302        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1303        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1304        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1305        assert_eq!(frames.len(), 1);
1306        let stats = iter.stats();
1307        assert_eq!(stats.unparsed_regions.len(), 1);
1308        assert_eq!(
1309            stats.unparsed_regions[0].reason,
1310            crate::types::SkipReason::OversizedFrame
1311        );
1312    }
1313
1314    #[test]
1315    fn stats_oversized_frame_mid_stream() {
1316        let mut data = Vec::new();
1317        data.extend_from_slice(
1318            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1319        );
1320        let garbage_len = MAX_PARTIAL_FRAME + 1;
1321        data.resize(data.len() + garbage_len, b'x');
1322        data.extend_from_slice(b"\x0B\n");
1323        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1324        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1325        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1326        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1327        assert_eq!(frames.len(), 2);
1328        let stats = iter.stats();
1329        assert_eq!(stats.unparsed_regions.len(), 1);
1330        assert_eq!(
1331            stats.unparsed_regions[0].reason,
1332            crate::types::SkipReason::OversizedFrame
1333        );
1334    }
1335
1336    #[test]
1337    fn overlong_byte_count_does_not_buffer_ahead() {
1338        let mut data = Vec::new();
1339        data.extend_from_slice(
1340            format!(
1341                "recv {} bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n",
1342                MAX_PARTIAL_FRAME * 10
1343            )
1344            .as_bytes(),
1345        );
1346        data.extend_from_slice(b"hello\x0B\n");
1347        while data.len() < MAX_PARTIAL_FRAME * 10 {
1348            data.extend_from_slice(
1349                b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n",
1350            );
1351        }
1352
1353        let mut iter = FrameIterator::new(&data[..]);
1354        let first = iter.next().unwrap().unwrap();
1355        assert_eq!(first.content, b"hello");
1356        assert!(
1357            iter.buf.len() <= MAX_PARTIAL_FRAME + READ_BUF_SIZE,
1358            "buffered {} bytes on a byte_count ten times the frame",
1359            iter.buf.len()
1360        );
1361        let second = iter.next().unwrap().unwrap();
1362        assert_eq!(second.content, b"bar");
1363    }
1364
1365    #[test]
1366    fn stats_partial_first_frame_within_limit() {
1367        // Content + \x0B\n boundary = MAX_PARTIAL_FRAME, should still be PartialFirstFrame
1368        let mut data = Vec::new();
1369        data.resize(MAX_PARTIAL_FRAME - 2, b'x');
1370        data.extend_from_slice(b"\x0B\n");
1371        data.extend_from_slice(
1372            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1373        );
1374        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1375        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1376        assert_eq!(frames.len(), 1);
1377        let stats = iter.stats();
1378        assert_eq!(stats.unparsed_regions.len(), 1);
1379        assert_eq!(
1380            stats.unparsed_regions[0].reason,
1381            crate::types::SkipReason::PartialFirstFrame
1382        );
1383    }
1384
1385    #[test]
1386    fn stats_dump_restart_marker() {
1387        let mut data = Vec::new();
1388        data.extend_from_slice(
1389            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1390        );
1391        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1392        data.extend_from_slice(b"sent 3 bytes to tcp/2.2.2.2:5060 at 00:00:01.000000:\nbye\x0B\n");
1393
1394        let mut iter = FrameIterator::new(&data[..]);
1395        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1396        assert_eq!(frames.len(), 2);
1397        let stats = iter.stats();
1398        // Dump restart marker is structural, not skipped
1399        assert_eq!(stats.bytes_skipped, 0);
1400        assert!(stats.unparsed_regions.is_empty());
1401    }
1402
1403    #[test]
1404    fn stats_count_only_no_regions() {
1405        let mut data = Vec::new();
1406        data.extend_from_slice(b"partial garbage data");
1407        data.extend_from_slice(b"\x0B\n");
1408        data.extend_from_slice(
1409            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1410        );
1411        let mut iter = FrameIterator::new(&data[..]);
1412        let frames: Vec<_> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1413        assert_eq!(frames.len(), 1);
1414        let stats = iter.stats();
1415        let skipped = b"partial garbage data\x0B\n".len() as u64;
1416        assert_eq!(stats.bytes_skipped, skipped);
1417        assert!(
1418            stats.unparsed_regions.is_empty(),
1419            "CountOnly should not accumulate regions"
1420        );
1421    }
1422
1423    #[test]
1424    fn frame_iterator_trailing_newlines_at_eof() {
1425        // Trailing newlines after the last boundary at EOF should not cause errors.
1426        let mut data = Vec::new();
1427        data.extend_from_slice(
1428            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1429        );
1430        data.extend_from_slice(b"\n\n");
1431
1432        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1433            .collect::<Result<Vec<_>, _>>()
1434            .unwrap();
1435        assert_eq!(frames.len(), 1);
1436        assert_eq!(frames[0].content, b"hello");
1437    }
1438}