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    /// Level 2 records its own losses in the same accounting.
373    pub(crate) fn stats_mut(&mut self) -> &mut ParseStats {
374        &mut self.stats
375    }
376
377    /// Take all accumulated unparsed regions, leaving the list empty.
378    pub fn drain_unparsed(&mut self) -> Vec<UnparsedRegion> {
379        self.stats.drain_regions()
380    }
381
382    fn consume(&mut self, n: usize) {
383        self.buf.drain(..n);
384        self.offset += n as u64;
385    }
386
387    fn consume_skipped(&mut self, n: usize, reason: SkipReason) {
388        if self.skip_tracking != SkipTracking::CountOnly {
389            let data = if self.skip_tracking == SkipTracking::CaptureData {
390                Some(self.buf[..n].to_vec())
391            } else {
392                None
393            };
394            self.stats.unparsed_regions.push(UnparsedRegion {
395                offset: self.offset,
396                length: n as u64,
397                reason,
398                data,
399            });
400        }
401        self.stats.bytes_skipped += n as u64;
402        self.consume(n);
403    }
404
405    fn fill_buf(&mut self) -> Result<bool, std::io::Error> {
406        if self.eof {
407            return Ok(false);
408        }
409        let old_len = self.buf.len();
410        self.buf.resize(old_len + READ_BUF_SIZE, 0);
411        let n = self.reader.read(&mut self.buf[old_len..])?;
412        self.buf.truncate(old_len + n);
413        if n == 0 {
414            self.eof = true;
415            return Ok(false);
416        }
417        self.stats.bytes_read += n as u64;
418        Ok(true)
419    }
420
421    /// A SIP header terminator right before a frame boundary is the tail of a
422    /// frame logrotate copied into both files.
423    fn is_replay(&self, skipped: &[u8]) -> bool {
424        if self.frame_count == 0 {
425            return false;
426        }
427        skipped.ends_with(b"\r\n\r\n\x0B\n")
428    }
429
430    /// Find the next `\x0B\n` boundary that is followed by a valid frame header.
431    fn find_boundary(&self, start: usize) -> Option<usize> {
432        let mut search_from = start;
433        loop {
434            let pos = BOUNDARY.find(&self.buf[search_from..])?;
435            let abs_pos = search_from + pos;
436            let after = abs_pos + 2;
437            if after >= self.buf.len() {
438                // Boundary at very end — could be real, but we can't validate header yet
439                // If EOF, accept it as boundary (content ends at \x0B)
440                if self.eof {
441                    return Some(abs_pos);
442                }
443                return None; // Need more data
444            }
445            if is_frame_header(&self.buf[after..]) {
446                return Some(abs_pos);
447            }
448            // \x0B\n in content, not a boundary — skip past it
449            trace!(
450                offset = abs_pos,
451                "found \\x0B\\n in content (not a boundary), skipping"
452            );
453            search_from = abs_pos + 2;
454        }
455    }
456
457    /// Skip to the first valid frame header in the buffer (for partial first frames).
458    fn skip_to_first_header(&mut self) -> Option<usize> {
459        if is_frame_header(&self.buf) {
460            return Some(0);
461        }
462        // Look for \x0B\n followed by a valid header
463        let mut search_from = 0;
464        loop {
465            let pos = BOUNDARY.find(&self.buf[search_from..])?;
466            let abs_pos = search_from + pos;
467            let after = abs_pos + 2;
468            if after < self.buf.len() && is_frame_header(&self.buf[after..]) {
469                debug!(skipped_bytes = after, "skipped partial first frame");
470                return Some(after);
471            }
472            search_from = abs_pos + 2;
473        }
474    }
475
476    /// Drop a partial first frame so the buffer starts at a frame header.
477    fn sync_to_first_header(&mut self) -> SyncStep {
478        loop {
479            match self.skip_to_first_header() {
480                Some(offset) => {
481                    if offset > 0 {
482                        let reason = if offset <= MAX_PARTIAL_FRAME {
483                            SkipReason::PartialFirstFrame
484                        } else {
485                            SkipReason::OversizedFrame
486                        };
487                        self.consume_skipped(offset, reason);
488                    }
489                    return SyncStep::Ready;
490                }
491                None => {
492                    if self.eof {
493                        debug!("no valid frame header found in entire input");
494                        let remaining = self.buf.len();
495                        if remaining > 0 {
496                            self.consume_skipped(remaining, SkipReason::InvalidHeader);
497                        }
498                        return SyncStep::End;
499                    }
500                    if let Err(e) = self.fill_buf() {
501                        return SyncStep::Failed(ParseError::Io(e));
502                    }
503                }
504            }
505        }
506    }
507
508    /// Drop `\n` and `\r\n` padding between frames.
509    fn strip_padding(&mut self) {
510        let mut strip = 0;
511        while strip < self.buf.len() {
512            if self.buf[strip] == b'\n' {
513                strip += 1;
514            } else if strip + 1 < self.buf.len()
515                && self.buf[strip] == b'\r'
516                && self.buf[strip + 1] == b'\n'
517            {
518                strip += 2;
519            } else {
520                break;
521            }
522        }
523        if strip > 0 {
524            self.consume(strip);
525        }
526    }
527
528    /// Parse the header at the head of the buffer, reading more as it needs.
529    fn read_header(&mut self) -> HeaderStep {
530        loop {
531            match classify_header(&self.buf) {
532                HeaderParse::Ok(h) => return HeaderStep::Got(h),
533                HeaderParse::NeedMore => {
534                    if self.eof {
535                        debug!("truncated frame header at EOF");
536                        let remaining = self.buf.len();
537                        if remaining > 0 {
538                            self.consume_skipped(remaining, SkipReason::InvalidHeader);
539                        }
540                        return HeaderStep::End;
541                    }
542                    if let Err(e) = self.fill_buf() {
543                        return HeaderStep::Failed(ParseError::Io(e));
544                    }
545                }
546                HeaderParse::Invalid(e) => {
547                    let header_preview: String = self
548                        .buf
549                        .iter()
550                        .take(200)
551                        .take_while(|&&b| b != b'\n')
552                        .map(|&b| {
553                            if b.is_ascii_graphic() || b == b' ' {
554                                b as char
555                            } else {
556                                '.'
557                            }
558                        })
559                        .collect();
560                    if header_preview.starts_with("dump started at ") {
561                        let skip = memchr::memchr(b'\n', &self.buf)
562                            .map(|p| {
563                                let mut end = p + 1;
564                                while end < self.buf.len() && self.buf[end] == b'\n' {
565                                    end += 1;
566                                }
567                                end
568                            })
569                            .unwrap_or(self.buf.len());
570                        debug!(
571                            header = %header_preview,
572                            skipped_bytes = skip,
573                            "skipped dump restart marker",
574                        );
575                        self.consume(skip);
576                        return HeaderStep::Restart;
577                    }
578                    let skip = if let Some(b) = self.find_boundary(0) {
579                        b + 2
580                    } else {
581                        memchr::memchr(b'\n', &self.buf)
582                            .map(|p| p + 1)
583                            .unwrap_or(self.buf.len())
584                    };
585                    let reason = self.classify_skip(skip);
586                    self.consume_skipped(skip, reason);
587                    return HeaderStep::Failed(e);
588                }
589            }
590        }
591    }
592
593    /// Reason for dropping `skip` bytes that failed header parsing.
594    fn classify_skip(&self, skip: usize) -> SkipReason {
595        if self.buf.starts_with(RECV_PREFIX) || self.buf.starts_with(SENT_PREFIX) {
596            SkipReason::InvalidHeader
597        } else if skip > MAX_PARTIAL_FRAME {
598            SkipReason::OversizedFrame
599        } else if self.frame_count == 0 {
600            SkipReason::PartialFirstFrame
601        } else if self.is_replay(&self.buf[..skip]) {
602            SkipReason::ReplayedFrame
603        } else {
604            SkipReason::MidStreamSkip
605        }
606    }
607
608    /// Read one frame's content, from the header's end to the frame boundary.
609    fn read_content(&mut self, header: FrameHeader) -> Option<Result<Frame, ParseError>> {
610        let FrameHeader {
611            direction,
612            byte_count,
613            transport,
614            address,
615            timestamp,
616            header_len,
617        } = header;
618        let content_start = header_len;
619        let expected_end = content_start + byte_count;
620
621        // Find the boundary for this frame.
622        // Strategy: first check at the expected position (content_start + byte_count),
623        // then fall back to scanning. This handles file concatenation where \x0B\n
624        // is followed by garbage from the next file's truncated first frame.
625        let content = loop {
626            // Ensure we have enough data to check the expected position, but
627            // never buffer a declared count further than a frame can reach:
628            // past that, boundary scanning takes over.
629            let fill_to = (expected_end + 1).min(content_start + MAX_PARTIAL_FRAME);
630            while self.buf.len() <= fill_to && !self.eof {
631                if let Err(e) = self.fill_buf() {
632                    return Some(Err(ParseError::Io(e)));
633                }
634            }
635
636            // Check at expected position first (byte_count hint)
637            if expected_end < self.buf.len() && self.buf[expected_end] == 0x0B {
638                let has_newline =
639                    expected_end + 1 < self.buf.len() && self.buf[expected_end + 1] == b'\n';
640                let at_eof = expected_end + 1 >= self.buf.len() && self.eof;
641
642                if has_newline || at_eof {
643                    let content = self.buf[content_start..expected_end].to_vec();
644                    let drain_to = if has_newline {
645                        expected_end + 2
646                    } else {
647                        expected_end + 1
648                    };
649                    self.consume(drain_to);
650                    self.frame_count += 1;
651                    break content;
652                }
653            }
654
655            // Fall back to scanning for \x0B\n + valid header
656            if let Some(boundary_pos) = self.find_boundary(content_start) {
657                let content = self.buf[content_start..boundary_pos].to_vec();
658                let drain_to = boundary_pos + 2;
659                self.consume(drain_to);
660                self.frame_count += 1;
661
662                if content.len() != byte_count {
663                    debug!(
664                        frame = self.frame_count,
665                        expected = byte_count,
666                        actual = content.len(),
667                        "frame content size mismatch"
668                    );
669                }
670
671                break content;
672            }
673
674            if self.eof {
675                // Last frame — no trailing \x0B\n
676                let end = if self.buf.last() == Some(&0x0B) {
677                    self.buf.len() - 1
678                } else {
679                    self.buf.len()
680                };
681                let content = self.buf[content_start..end].to_vec();
682                let len = self.buf.len();
683                self.consume(len);
684                self.frame_count += 1;
685
686                if content.len() < byte_count {
687                    let missing = byte_count - content.len();
688                    debug!(
689                        frame = self.frame_count,
690                        expected = byte_count,
691                        actual = content.len(),
692                        missing,
693                        "incomplete frame at EOF"
694                    );
695                    self.stats.incomplete_frames += 1;
696                    self.stats.incomplete_frame_bytes += missing as u64;
697                    if self.skip_tracking != SkipTracking::CountOnly {
698                        self.stats.unparsed_regions.push(UnparsedRegion {
699                            offset: self.offset,
700                            length: missing as u64,
701                            reason: SkipReason::IncompleteFrame,
702                            data: None,
703                        });
704                    }
705                } else if content.len() != byte_count {
706                    debug!(
707                        frame = self.frame_count,
708                        expected = byte_count,
709                        actual = content.len(),
710                        "last frame content size mismatch"
711                    );
712                }
713
714                break content;
715            }
716
717            if let Err(e) = self.fill_buf() {
718                return Some(Err(ParseError::Io(e)));
719            }
720        };
721
722        Some(Ok(Frame {
723            direction,
724            byte_count,
725            transport,
726            address,
727            timestamp,
728            content,
729        }))
730    }
731}
732
733impl<R: Read> Iterator for FrameIterator<R> {
734    type Item = Result<Frame, ParseError>;
735
736    fn next(&mut self) -> Option<Self::Item> {
737        let header = loop {
738            if self.buf.is_empty() && !self.eof {
739                if let Err(e) = self.fill_buf() {
740                    return Some(Err(ParseError::Io(e)));
741                }
742            }
743            if self.buf.is_empty() {
744                return None;
745            }
746
747            if self.frame_count == 0 {
748                match self.sync_to_first_header() {
749                    SyncStep::Ready => {}
750                    SyncStep::End => return None,
751                    SyncStep::Failed(e) => return Some(Err(e)),
752                }
753            }
754
755            self.strip_padding();
756            if self.buf.is_empty() {
757                continue;
758            }
759
760            match self.read_header() {
761                HeaderStep::Got(header) => break header,
762                HeaderStep::Restart => continue,
763                HeaderStep::End => return None,
764                HeaderStep::Failed(e) => return Some(Err(e)),
765            }
766        };
767
768        self.read_content(header)
769    }
770}
771
772#[cfg(test)]
773mod tests {
774    use super::*;
775    use crate::types::SkipTracking;
776
777    #[test]
778    fn parse_recv_ipv4_tcp() {
779        let header = b"recv 100 bytes from tcp/192.168.1.1:5060 at 00:00:01.350874:\n";
780        let h = parse_frame_header(header).unwrap();
781        assert_eq!(h.direction, Direction::Recv);
782        assert_eq!(h.byte_count, 100);
783        assert_eq!(h.transport, Transport::Tcp);
784        assert_eq!(h.address, "192.168.1.1:5060");
785        assert_eq!(
786            h.timestamp,
787            Timestamp::TimeOnly {
788                hour: 0,
789                min: 0,
790                sec: 1,
791                usec: 350874
792            }
793        );
794        assert_eq!(h.header_len, header.len());
795    }
796
797    #[test]
798    fn parse_recv_ipv6_tcp() {
799        let header = b"recv 1440 bytes from tcp/[2001:4958:10:14::4]:30046 at 13:03:21.674883:\n";
800        let h = parse_frame_header(header).unwrap();
801        assert_eq!(h.direction, Direction::Recv);
802        assert_eq!(h.byte_count, 1440);
803        assert_eq!(h.transport, Transport::Tcp);
804        assert_eq!(h.address, "[2001:4958:10:14::4]:30046");
805        assert_eq!(
806            h.timestamp,
807            Timestamp::TimeOnly {
808                hour: 13,
809                min: 3,
810                sec: 21,
811                usec: 674883
812            }
813        );
814    }
815
816    #[test]
817    fn parse_sent_ipv6_tcp() {
818        let header = b"sent 681 bytes to tcp/[2001:4958:10:14::4]:30046 at 13:03:21.675500:\n";
819        let h = parse_frame_header(header).unwrap();
820        assert_eq!(h.direction, Direction::Sent);
821        assert_eq!(h.byte_count, 681);
822        assert_eq!(h.transport, Transport::Tcp);
823        assert_eq!(h.address, "[2001:4958:10:14::4]:30046");
824    }
825
826    #[test]
827    fn parse_recv_udp() {
828        let header = b"recv 457 bytes from udp/10.0.0.1:5060 at 00:19:47.123456:\n";
829        let h = parse_frame_header(header).unwrap();
830        assert_eq!(h.direction, Direction::Recv);
831        assert_eq!(h.transport, Transport::Udp);
832    }
833
834    #[test]
835    fn parse_sent_tls() {
836        let header = b"sent 500 bytes to tls/10.0.0.1:5061 at 12:00:00.000000:\n";
837        let h = parse_frame_header(header).unwrap();
838        assert_eq!(h.direction, Direction::Sent);
839        assert_eq!(h.byte_count, 500);
840        assert_eq!(h.transport, Transport::Tls);
841    }
842
843    #[test]
844    fn parse_full_datetime_timestamp() {
845        let header = b"recv 100 bytes from tcp/192.168.1.1:5060 at 2026-02-01 10:00:00.000000:\n";
846        let h = parse_frame_header(header).unwrap();
847        assert_eq!(
848            h.timestamp,
849            Timestamp::DateTime {
850                year: 2026,
851                month: 2,
852                day: 1,
853                hour: 10,
854                min: 0,
855                sec: 0,
856                usec: 0
857            }
858        );
859    }
860
861    #[test]
862    fn parse_invalid_header() {
863        assert!(parse_frame_header(b"invalid header\n").is_err());
864        assert!(
865            parse_frame_header(b"recv abc bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n")
866                .is_err()
867        );
868    }
869
870    #[test]
871    fn is_frame_header_valid() {
872        assert!(is_frame_header(
873            b"recv 100 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n"
874        ));
875        assert!(is_frame_header(
876            b"sent 681 bytes to tcp/[::1]:5060 at 00:00:00.000000:\n"
877        ));
878        assert!(!is_frame_header(b"not a header"));
879        assert!(!is_frame_header(b"recv abc bytes"));
880        assert!(!is_frame_header(b""));
881    }
882
883    #[test]
884    fn frame_iterator_single_frame() {
885        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
886        let frames: Vec<Frame> = FrameIterator::new(&data[..])
887            .collect::<Result<Vec<_>, _>>()
888            .unwrap();
889        assert_eq!(frames.len(), 1);
890        assert_eq!(frames[0].content, b"hello");
891        assert_eq!(frames[0].byte_count, 5);
892    }
893
894    #[test]
895    fn frame_iterator_multiple_frames() {
896        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";
897        let frames: Vec<Frame> = FrameIterator::new(&data[..])
898            .collect::<Result<Vec<_>, _>>()
899            .unwrap();
900        assert_eq!(frames.len(), 2);
901        assert_eq!(frames[0].content, b"hello");
902        assert_eq!(frames[0].direction, Direction::Recv);
903        assert_eq!(frames[1].content, b"world");
904        assert_eq!(frames[1].direction, Direction::Sent);
905    }
906
907    #[test]
908    fn frame_iterator_vt_in_content() {
909        // \x0B in content but not followed by valid header — should NOT split
910        let mut data = Vec::new();
911        data.extend_from_slice(b"recv 15 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n");
912        data.extend_from_slice(b"he\x0B\nllo world!!");
913        data.extend_from_slice(b"\x0B\n");
914        let frames: Vec<Frame> = FrameIterator::new(&data[..])
915            .collect::<Result<Vec<_>, _>>()
916            .unwrap();
917        assert_eq!(frames.len(), 1);
918        assert_eq!(frames[0].content, b"he\x0B\nllo world!!");
919    }
920
921    #[test]
922    fn frame_iterator_eof_without_boundary() {
923        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello";
924        let frames: Vec<Frame> = FrameIterator::new(&data[..])
925            .collect::<Result<Vec<_>, _>>()
926            .unwrap();
927        assert_eq!(frames.len(), 1);
928        assert_eq!(frames[0].content, b"hello");
929    }
930
931    #[test]
932    fn frame_iterator_eof_with_lone_vt() {
933        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B";
934        let frames: Vec<Frame> = FrameIterator::new(&data[..])
935            .collect::<Result<Vec<_>, _>>()
936            .unwrap();
937        assert_eq!(frames.len(), 1);
938        assert_eq!(frames[0].content, b"hello");
939    }
940
941    #[test]
942    fn frame_iterator_partial_first_frame() {
943        // Data starts with garbage, then a valid boundary + frame
944        let mut data = Vec::new();
945        data.extend_from_slice(b"partial garbage data");
946        data.extend_from_slice(b"\x0B\n");
947        data.extend_from_slice(
948            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
949        );
950        let frames: Vec<Frame> = FrameIterator::new(&data[..])
951            .collect::<Result<Vec<_>, _>>()
952            .unwrap();
953        assert_eq!(frames.len(), 1);
954        assert_eq!(frames[0].content, b"hello");
955    }
956
957    #[test]
958    fn frame_iterator_truncated_last_frame() {
959        // Complete frame followed by truncated frame at EOF (no \x0B\n)
960        let mut data = Vec::new();
961        data.extend_from_slice(
962            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
963        );
964        data.extend_from_slice(b"sent 3 bytes to tcp/1.1.1.1:5060 at 00:00:01.000000:\nbye");
965        let frames: Vec<Frame> = FrameIterator::new(&data[..])
966            .collect::<Result<Vec<_>, _>>()
967            .unwrap();
968        assert_eq!(frames.len(), 2);
969        assert_eq!(frames[0].content, b"hello");
970        assert_eq!(frames[1].content, b"bye");
971    }
972
973    #[test]
974    fn frame_iterator_file_concatenation() {
975        // Simulates `cat dump.20 dump.21 | parser`
976        // File 1: truncated start + valid frame + complete end
977        // File 2: truncated start (no header) + valid frame
978        let mut data = Vec::new();
979
980        // File 1: starts with valid header
981        data.extend_from_slice(
982            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
983        );
984        data.extend_from_slice(
985            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
986        );
987
988        // File 2: starts with truncated frame data (no header), then boundary, then valid frame
989        data.extend_from_slice(b"some truncated SIP content from previous rotation\r\n\r\n");
990        data.extend_from_slice(b"\x0B\n");
991        data.extend_from_slice(
992            b"recv 3 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\nfoo\x0B\n",
993        );
994
995        let items: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
996        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
997        assert_eq!(frames.len(), 3);
998        assert_eq!(frames[0].content, b"hello");
999        assert_eq!(frames[1].content, b"world");
1000        assert_eq!(frames[2].content, b"foo");
1001        assert_eq!(frames[2].address, "2.2.2.2:5060");
1002    }
1003
1004    #[test]
1005    fn frame_iterator_file_concatenation_mid_stream_garbage() {
1006        // The join point between files produces garbage that looks like:
1007        // ...last_content\x0B\ntruncated_first_of_next_file\x0B\nvalid_header...
1008        // The truncated part is NOT a valid header, so recovery should skip it
1009        let mut data = Vec::new();
1010
1011        // Last frame of file 1
1012        data.extend_from_slice(
1013            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1014        );
1015
1016        // Truncated first frame of file 2 (mid-SIP content, no frame header)
1017        data.extend_from_slice(b"Content-Type: application/sdp\r\n\r\nv=0\r\n");
1018        data.extend_from_slice(b"\x0B\n");
1019
1020        // Valid second frame of file 2
1021        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1022
1023        let items: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1024        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1025        assert_eq!(frames.len(), 2);
1026        assert_eq!(frames[0].content, b"hello");
1027        assert_eq!(frames[1].content, b"bar");
1028    }
1029
1030    #[test]
1031    fn frame_iterator_empty_input() {
1032        let data: &[u8] = b"";
1033        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(data).collect();
1034        assert!(frames.is_empty());
1035    }
1036
1037    #[test]
1038    fn frame_iterator_only_garbage() {
1039        let data = b"this is not a SIP trace dump at all, just garbage text";
1040        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1041        let frames: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1042        assert!(frames.is_empty());
1043        let stats = iter.stats();
1044        assert_eq!(stats.bytes_read, data.len() as u64);
1045        assert_eq!(stats.bytes_skipped, data.len() as u64);
1046        assert_eq!(stats.unparsed_regions.len(), 1);
1047        assert_eq!(stats.unparsed_regions[0].reason, SkipReason::InvalidHeader);
1048    }
1049
1050    #[test]
1051    fn frame_iterator_truncated_header_at_eof() {
1052        let data = b"recv 5 bytes from tcp/1.1.1.1:5060";
1053        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1054        let frames: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1055        assert!(frames.is_empty());
1056        let stats = iter.stats();
1057        assert_eq!(stats.bytes_read, data.len() as u64);
1058        assert_eq!(stats.bytes_skipped, data.len() as u64);
1059    }
1060
1061    #[test]
1062    fn frame_iterator_dump_marker_at_eof() {
1063        // A dump restart marker at the end of input (with trailing \n\n as in real dumps)
1064        // should be silently consumed, not returned as an error.
1065        let mut data = Vec::new();
1066        data.extend_from_slice(
1067            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1068        );
1069        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1070
1071        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1072        assert_eq!(frames.len(), 1);
1073        assert!(frames[0].is_ok());
1074        assert_eq!(frames[0].as_ref().unwrap().content, b"hello");
1075    }
1076
1077    #[test]
1078    fn frame_iterator_dump_marker_mid_stream() {
1079        // A dump restart marker between two valid frames (with trailing \n\n as in real dumps)
1080        // should be skipped, and both frames should parse successfully.
1081        let mut data = Vec::new();
1082        data.extend_from_slice(
1083            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1084        );
1085        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1086        data.extend_from_slice(b"sent 3 bytes to tcp/2.2.2.2:5060 at 00:00:01.000000:\nbye\x0B\n");
1087
1088        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1089        assert_eq!(frames.len(), 2);
1090        assert_eq!(frames[0].as_ref().unwrap().content, b"hello");
1091        assert_eq!(frames[1].as_ref().unwrap().content, b"bye");
1092    }
1093
1094    #[test]
1095    fn frame_iterator_multiple_newlines_after_boundary() {
1096        // Multiple \n and \r\n between frames should all be stripped.
1097        let mut data = Vec::new();
1098        data.extend_from_slice(
1099            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1100        );
1101        data.extend_from_slice(b"\n\r\n\n");
1102        data.extend_from_slice(
1103            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
1104        );
1105
1106        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1107            .collect::<Result<Vec<_>, _>>()
1108            .unwrap();
1109        assert_eq!(frames.len(), 2);
1110        assert_eq!(frames[0].content, b"hello");
1111        assert_eq!(frames[1].content, b"world");
1112    }
1113
1114    #[test]
1115    fn stats_clean_input() {
1116        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
1117        let mut iter = FrameIterator::new(&data[..]);
1118        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1119        assert_eq!(frames.len(), 1);
1120        let stats = iter.stats();
1121        assert_eq!(stats.bytes_read, data.len() as u64);
1122        assert_eq!(stats.bytes_skipped, 0);
1123        assert!(stats.unparsed_regions.is_empty());
1124    }
1125
1126    #[test]
1127    fn stats_multiple_frames() {
1128        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";
1129        let mut iter = FrameIterator::new(&data[..]);
1130        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1131        assert_eq!(frames.len(), 2);
1132        let stats = iter.stats();
1133        assert_eq!(stats.bytes_read, data.len() as u64);
1134        assert_eq!(stats.bytes_skipped, 0);
1135        assert!(stats.unparsed_regions.is_empty());
1136    }
1137
1138    #[test]
1139    fn stats_partial_first_frame() {
1140        let mut data = Vec::new();
1141        data.extend_from_slice(b"partial garbage data");
1142        data.extend_from_slice(b"\x0B\n");
1143        data.extend_from_slice(
1144            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1145        );
1146        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1147        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1148        assert_eq!(frames.len(), 1);
1149        let stats = iter.stats();
1150        assert_eq!(stats.bytes_read, data.len() as u64);
1151        // "partial garbage data" + "\x0B\n" = 21 bytes skipped
1152        let skipped = b"partial garbage data\x0B\n".len() as u64;
1153        assert_eq!(stats.bytes_skipped, skipped);
1154        assert_eq!(stats.unparsed_regions.len(), 1);
1155        assert_eq!(stats.unparsed_regions[0].offset, 0);
1156        assert_eq!(stats.unparsed_regions[0].length, skipped);
1157        assert_eq!(
1158            stats.unparsed_regions[0].reason,
1159            crate::types::SkipReason::PartialFirstFrame
1160        );
1161        assert!(stats.unparsed_regions[0].data.is_none());
1162    }
1163
1164    #[test]
1165    fn stats_partial_first_frame_capture() {
1166        let mut data = Vec::new();
1167        data.extend_from_slice(b"partial garbage data");
1168        data.extend_from_slice(b"\x0B\n");
1169        data.extend_from_slice(
1170            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1171        );
1172        let mut iter = FrameIterator::new(&data[..]).capture_skipped(true);
1173        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1174        assert_eq!(frames.len(), 1);
1175        let stats = iter.stats();
1176        assert_eq!(stats.unparsed_regions.len(), 1);
1177        let region = &stats.unparsed_regions[0];
1178        assert_eq!(
1179            region.data.as_deref(),
1180            Some(b"partial garbage data\x0B\n".as_slice())
1181        );
1182    }
1183
1184    #[test]
1185    fn stats_mid_stream_partial_frame() {
1186        // SIP content between valid frames (file concatenation scenario)
1187        let mut data = Vec::new();
1188        data.extend_from_slice(
1189            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1190        );
1191        data.extend_from_slice(b"Content-Type: application/sdp\r\n\r\nv=0\r\n");
1192        data.extend_from_slice(b"\x0B\n");
1193        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1194
1195        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1196        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1197        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1198        assert_eq!(frames.len(), 2);
1199        let stats = iter.stats();
1200        assert!(stats.bytes_skipped > 0);
1201        assert_eq!(stats.unparsed_regions.len(), 1);
1202        assert_eq!(
1203            stats.unparsed_regions[0].reason,
1204            crate::types::SkipReason::MidStreamSkip
1205        );
1206    }
1207
1208    #[test]
1209    fn stats_replayed_frame() {
1210        // Simulate logrotate: a frame's tail (SIP headers ending with \r\n\r\n\x0B\n)
1211        // appears between two valid frames at a file boundary.
1212        let frame1 = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
1213        let replay = b"Route: <sip:10.0.0.1:5060;lr>\r\nContent-Length: 0\r\n\r\n\x0B\n";
1214        let frame2 = b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n";
1215
1216        let mut data = Vec::new();
1217        data.extend_from_slice(frame1);
1218        data.extend_from_slice(replay);
1219        data.extend_from_slice(frame2);
1220
1221        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1222        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1223        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1224        assert_eq!(frames.len(), 2);
1225        let stats = iter.stats();
1226        assert_eq!(stats.unparsed_regions.len(), 1);
1227        assert_eq!(
1228            stats.unparsed_regions[0].reason,
1229            crate::types::SkipReason::ReplayedFrame
1230        );
1231    }
1232
1233    #[test]
1234    fn stats_incomplete_frame_at_eof() {
1235        // Frame header says 100 bytes but only 20 bytes available before EOF
1236        let mut data = Vec::new();
1237        data.extend_from_slice(
1238            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1239        );
1240        data.extend_from_slice(b"recv 100 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\n");
1241        data.extend_from_slice(b"partial content only");
1242        // No \x0B\n boundary — EOF truncation
1243
1244        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1245        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1246        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1247        assert_eq!(frames.len(), 2, "truncated frame should still be returned");
1248        assert_eq!(frames[1].content, b"partial content only");
1249        assert_eq!(frames[1].byte_count, 100);
1250        let stats = iter.stats();
1251        assert_eq!(stats.unparsed_regions.len(), 1);
1252        assert_eq!(
1253            stats.unparsed_regions[0].reason,
1254            crate::types::SkipReason::IncompleteFrame
1255        );
1256    }
1257
1258    #[test]
1259    fn stats_incomplete_frame_counted_under_count_only() {
1260        let mut data = Vec::new();
1261        data.extend_from_slice(
1262            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1263        );
1264        data.extend_from_slice(b"recv 100 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\n");
1265        data.extend_from_slice(b"partial content only");
1266
1267        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::CountOnly);
1268        let _: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1269        let stats = iter.stats();
1270        assert!(stats.unparsed_regions.is_empty());
1271        assert_eq!(stats.incomplete_frames, 1);
1272        assert_eq!(stats.incomplete_frame_bytes, 80);
1273    }
1274
1275    #[test]
1276    fn stats_invalid_header_skip() {
1277        // Malformed frame header (starts with recv/sent but unparseable)
1278        let mut data = Vec::new();
1279        data.extend_from_slice(
1280            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1281        );
1282        data.extend_from_slice(b"recv CORRUPT HEADER garbage\n");
1283        data.extend_from_slice(b"\x0B\n");
1284        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1285
1286        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1287        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1288        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1289        assert_eq!(frames.len(), 2);
1290        let stats = iter.stats();
1291        assert!(stats.bytes_skipped > 0);
1292        assert_eq!(stats.unparsed_regions.len(), 1);
1293        assert_eq!(
1294            stats.unparsed_regions[0].reason,
1295            crate::types::SkipReason::InvalidHeader
1296        );
1297    }
1298
1299    #[test]
1300    fn stats_oversized_frame_at_start() {
1301        let mut data = Vec::new();
1302        data.resize(MAX_PARTIAL_FRAME + 1, b'x');
1303        data.extend_from_slice(b"\x0B\n");
1304        data.extend_from_slice(
1305            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1306        );
1307        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1308        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1309        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1310        assert_eq!(frames.len(), 1);
1311        let stats = iter.stats();
1312        assert_eq!(stats.unparsed_regions.len(), 1);
1313        assert_eq!(
1314            stats.unparsed_regions[0].reason,
1315            crate::types::SkipReason::OversizedFrame
1316        );
1317    }
1318
1319    #[test]
1320    fn stats_oversized_frame_mid_stream() {
1321        let mut data = Vec::new();
1322        data.extend_from_slice(
1323            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1324        );
1325        let garbage_len = MAX_PARTIAL_FRAME + 1;
1326        data.resize(data.len() + garbage_len, b'x');
1327        data.extend_from_slice(b"\x0B\n");
1328        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1329        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1330        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1331        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1332        assert_eq!(frames.len(), 2);
1333        let stats = iter.stats();
1334        assert_eq!(stats.unparsed_regions.len(), 1);
1335        assert_eq!(
1336            stats.unparsed_regions[0].reason,
1337            crate::types::SkipReason::OversizedFrame
1338        );
1339    }
1340
1341    #[test]
1342    fn overlong_byte_count_does_not_buffer_ahead() {
1343        let mut data = Vec::new();
1344        data.extend_from_slice(
1345            format!(
1346                "recv {} bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n",
1347                MAX_PARTIAL_FRAME * 10
1348            )
1349            .as_bytes(),
1350        );
1351        data.extend_from_slice(b"hello\x0B\n");
1352        while data.len() < MAX_PARTIAL_FRAME * 10 {
1353            data.extend_from_slice(
1354                b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n",
1355            );
1356        }
1357
1358        let mut iter = FrameIterator::new(&data[..]);
1359        let first = iter.next().unwrap().unwrap();
1360        assert_eq!(first.content, b"hello");
1361        assert!(
1362            iter.buf.len() <= MAX_PARTIAL_FRAME + READ_BUF_SIZE,
1363            "buffered {} bytes on a byte_count ten times the frame",
1364            iter.buf.len()
1365        );
1366        let second = iter.next().unwrap().unwrap();
1367        assert_eq!(second.content, b"bar");
1368    }
1369
1370    #[test]
1371    fn stats_partial_first_frame_within_limit() {
1372        // Content + \x0B\n boundary = MAX_PARTIAL_FRAME, should still be PartialFirstFrame
1373        let mut data = Vec::new();
1374        data.resize(MAX_PARTIAL_FRAME - 2, b'x');
1375        data.extend_from_slice(b"\x0B\n");
1376        data.extend_from_slice(
1377            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1378        );
1379        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1380        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1381        assert_eq!(frames.len(), 1);
1382        let stats = iter.stats();
1383        assert_eq!(stats.unparsed_regions.len(), 1);
1384        assert_eq!(
1385            stats.unparsed_regions[0].reason,
1386            crate::types::SkipReason::PartialFirstFrame
1387        );
1388    }
1389
1390    #[test]
1391    fn stats_dump_restart_marker() {
1392        let mut data = Vec::new();
1393        data.extend_from_slice(
1394            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1395        );
1396        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1397        data.extend_from_slice(b"sent 3 bytes to tcp/2.2.2.2:5060 at 00:00:01.000000:\nbye\x0B\n");
1398
1399        let mut iter = FrameIterator::new(&data[..]);
1400        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1401        assert_eq!(frames.len(), 2);
1402        let stats = iter.stats();
1403        // Dump restart marker is structural, not skipped
1404        assert_eq!(stats.bytes_skipped, 0);
1405        assert!(stats.unparsed_regions.is_empty());
1406    }
1407
1408    #[test]
1409    fn stats_count_only_no_regions() {
1410        let mut data = Vec::new();
1411        data.extend_from_slice(b"partial garbage data");
1412        data.extend_from_slice(b"\x0B\n");
1413        data.extend_from_slice(
1414            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1415        );
1416        let mut iter = FrameIterator::new(&data[..]);
1417        let frames: Vec<_> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1418        assert_eq!(frames.len(), 1);
1419        let stats = iter.stats();
1420        let skipped = b"partial garbage data\x0B\n".len() as u64;
1421        assert_eq!(stats.bytes_skipped, skipped);
1422        assert!(
1423            stats.unparsed_regions.is_empty(),
1424            "CountOnly should not accumulate regions"
1425        );
1426    }
1427
1428    #[test]
1429    fn frame_iterator_trailing_newlines_at_eof() {
1430        // Trailing newlines after the last boundary at EOF should not cause errors.
1431        let mut data = Vec::new();
1432        data.extend_from_slice(
1433            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1434        );
1435        data.extend_from_slice(b"\n\n");
1436
1437        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1438            .collect::<Result<Vec<_>, _>>()
1439            .unwrap();
1440        assert_eq!(frames.len(), 1);
1441        assert_eq!(frames[0].content, b"hello");
1442    }
1443}