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(
610        &mut self,
611        header: FrameHeader,
612        offset: u64,
613    ) -> Option<Result<Frame, ParseError>> {
614        let FrameHeader {
615            direction,
616            byte_count,
617            transport,
618            address,
619            timestamp,
620            header_len,
621        } = header;
622        let content_start = header_len;
623        let expected_end = content_start + byte_count;
624
625        // Find the boundary for this frame.
626        // Strategy: first check at the expected position (content_start + byte_count),
627        // then fall back to scanning. This handles file concatenation where \x0B\n
628        // is followed by garbage from the next file's truncated first frame.
629        let content = loop {
630            // Ensure we have enough data to check the expected position, but
631            // never buffer a declared count further than a frame can reach:
632            // past that, boundary scanning takes over.
633            let fill_to = (expected_end + 1).min(content_start + MAX_PARTIAL_FRAME);
634            while self.buf.len() <= fill_to && !self.eof {
635                if let Err(e) = self.fill_buf() {
636                    return Some(Err(ParseError::Io(e)));
637                }
638            }
639
640            // Check at expected position first (byte_count hint)
641            if expected_end < self.buf.len() && self.buf[expected_end] == 0x0B {
642                let has_newline =
643                    expected_end + 1 < self.buf.len() && self.buf[expected_end + 1] == b'\n';
644                let at_eof = expected_end + 1 >= self.buf.len() && self.eof;
645
646                if has_newline || at_eof {
647                    let content = self.buf[content_start..expected_end].to_vec();
648                    let drain_to = if has_newline {
649                        expected_end + 2
650                    } else {
651                        expected_end + 1
652                    };
653                    self.consume(drain_to);
654                    self.frame_count += 1;
655                    break content;
656                }
657            }
658
659            // Fall back to scanning for \x0B\n + valid header
660            if let Some(boundary_pos) = self.find_boundary(content_start) {
661                let content = self.buf[content_start..boundary_pos].to_vec();
662                let drain_to = boundary_pos + 2;
663                self.consume(drain_to);
664                self.frame_count += 1;
665
666                if content.len() != byte_count {
667                    debug!(
668                        frame = self.frame_count,
669                        expected = byte_count,
670                        actual = content.len(),
671                        "frame content size mismatch"
672                    );
673                }
674
675                break content;
676            }
677
678            if self.eof {
679                // Last frame — no trailing \x0B\n
680                let end = if self.buf.last() == Some(&0x0B) {
681                    self.buf.len() - 1
682                } else {
683                    self.buf.len()
684                };
685                let content = self.buf[content_start..end].to_vec();
686                let len = self.buf.len();
687                self.consume(len);
688                self.frame_count += 1;
689
690                if content.len() < byte_count {
691                    let missing = byte_count - content.len();
692                    debug!(
693                        frame = self.frame_count,
694                        expected = byte_count,
695                        actual = content.len(),
696                        missing,
697                        "incomplete frame at EOF"
698                    );
699                    self.stats.incomplete_frames += 1;
700                    self.stats.incomplete_frame_bytes += missing as u64;
701                    if self.skip_tracking != SkipTracking::CountOnly {
702                        self.stats.unparsed_regions.push(UnparsedRegion {
703                            offset: self.offset,
704                            length: missing as u64,
705                            reason: SkipReason::IncompleteFrame,
706                            data: None,
707                        });
708                    }
709                } else if content.len() != byte_count {
710                    debug!(
711                        frame = self.frame_count,
712                        expected = byte_count,
713                        actual = content.len(),
714                        "last frame content size mismatch"
715                    );
716                }
717
718                break content;
719            }
720
721            if let Err(e) = self.fill_buf() {
722                return Some(Err(ParseError::Io(e)));
723            }
724        };
725
726        Some(Ok(Frame {
727            direction,
728            byte_count,
729            transport,
730            address,
731            timestamp,
732            content,
733            offset,
734        }))
735    }
736}
737
738impl<R: Read> Iterator for FrameIterator<R> {
739    type Item = Result<Frame, ParseError>;
740
741    fn next(&mut self) -> Option<Self::Item> {
742        let (header, start) = loop {
743            if self.buf.is_empty() && !self.eof {
744                if let Err(e) = self.fill_buf() {
745                    return Some(Err(ParseError::Io(e)));
746                }
747            }
748            if self.buf.is_empty() {
749                return None;
750            }
751
752            if self.frame_count == 0 {
753                match self.sync_to_first_header() {
754                    SyncStep::Ready => {}
755                    SyncStep::End => return None,
756                    SyncStep::Failed(e) => return Some(Err(e)),
757                }
758            }
759
760            self.strip_padding();
761            if self.buf.is_empty() {
762                continue;
763            }
764
765            // Before the header is consumed: buf[0] sits at this offset, so a
766            // frame's position is its header's, not its content's.
767            let start = self.offset;
768            match self.read_header() {
769                HeaderStep::Got(header) => break (header, start),
770                HeaderStep::Restart => continue,
771                HeaderStep::End => return None,
772                HeaderStep::Failed(e) => return Some(Err(e)),
773            }
774        };
775
776        self.read_content(header, start)
777    }
778}
779
780#[cfg(test)]
781mod tests {
782    use super::*;
783    use crate::types::SkipTracking;
784
785    #[test]
786    fn parse_recv_ipv4_tcp() {
787        let header = b"recv 100 bytes from tcp/192.168.1.1:5060 at 00:00:01.350874:\n";
788        let h = parse_frame_header(header).unwrap();
789        assert_eq!(h.direction, Direction::Recv);
790        assert_eq!(h.byte_count, 100);
791        assert_eq!(h.transport, Transport::Tcp);
792        assert_eq!(h.address, "192.168.1.1:5060");
793        assert_eq!(
794            h.timestamp,
795            Timestamp::TimeOnly {
796                hour: 0,
797                min: 0,
798                sec: 1,
799                usec: 350874
800            }
801        );
802        assert_eq!(h.header_len, header.len());
803    }
804
805    #[test]
806    fn parse_recv_ipv6_tcp() {
807        let header = b"recv 1440 bytes from tcp/[2001:4958:10:14::4]:30046 at 13:03:21.674883:\n";
808        let h = parse_frame_header(header).unwrap();
809        assert_eq!(h.direction, Direction::Recv);
810        assert_eq!(h.byte_count, 1440);
811        assert_eq!(h.transport, Transport::Tcp);
812        assert_eq!(h.address, "[2001:4958:10:14::4]:30046");
813        assert_eq!(
814            h.timestamp,
815            Timestamp::TimeOnly {
816                hour: 13,
817                min: 3,
818                sec: 21,
819                usec: 674883
820            }
821        );
822    }
823
824    #[test]
825    fn parse_sent_ipv6_tcp() {
826        let header = b"sent 681 bytes to tcp/[2001:4958:10:14::4]:30046 at 13:03:21.675500:\n";
827        let h = parse_frame_header(header).unwrap();
828        assert_eq!(h.direction, Direction::Sent);
829        assert_eq!(h.byte_count, 681);
830        assert_eq!(h.transport, Transport::Tcp);
831        assert_eq!(h.address, "[2001:4958:10:14::4]:30046");
832    }
833
834    #[test]
835    fn parse_recv_udp() {
836        let header = b"recv 457 bytes from udp/10.0.0.1:5060 at 00:19:47.123456:\n";
837        let h = parse_frame_header(header).unwrap();
838        assert_eq!(h.direction, Direction::Recv);
839        assert_eq!(h.transport, Transport::Udp);
840    }
841
842    #[test]
843    fn parse_sent_tls() {
844        let header = b"sent 500 bytes to tls/10.0.0.1:5061 at 12:00:00.000000:\n";
845        let h = parse_frame_header(header).unwrap();
846        assert_eq!(h.direction, Direction::Sent);
847        assert_eq!(h.byte_count, 500);
848        assert_eq!(h.transport, Transport::Tls);
849    }
850
851    #[test]
852    fn parse_full_datetime_timestamp() {
853        let header = b"recv 100 bytes from tcp/192.168.1.1:5060 at 2026-02-01 10:00:00.000000:\n";
854        let h = parse_frame_header(header).unwrap();
855        assert_eq!(
856            h.timestamp,
857            Timestamp::DateTime {
858                year: 2026,
859                month: 2,
860                day: 1,
861                hour: 10,
862                min: 0,
863                sec: 0,
864                usec: 0
865            }
866        );
867    }
868
869    #[test]
870    fn parse_invalid_header() {
871        assert!(parse_frame_header(b"invalid header\n").is_err());
872        assert!(
873            parse_frame_header(b"recv abc bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n")
874                .is_err()
875        );
876    }
877
878    #[test]
879    fn is_frame_header_valid() {
880        assert!(is_frame_header(
881            b"recv 100 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n"
882        ));
883        assert!(is_frame_header(
884            b"sent 681 bytes to tcp/[::1]:5060 at 00:00:00.000000:\n"
885        ));
886        assert!(!is_frame_header(b"not a header"));
887        assert!(!is_frame_header(b"recv abc bytes"));
888        assert!(!is_frame_header(b""));
889    }
890
891    #[test]
892    fn frame_iterator_single_frame() {
893        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
894        let frames: Vec<Frame> = FrameIterator::new(&data[..])
895            .collect::<Result<Vec<_>, _>>()
896            .unwrap();
897        assert_eq!(frames.len(), 1);
898        assert_eq!(frames[0].content, b"hello");
899        assert_eq!(frames[0].byte_count, 5);
900    }
901
902    #[test]
903    fn frame_iterator_multiple_frames() {
904        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";
905        let frames: Vec<Frame> = FrameIterator::new(&data[..])
906            .collect::<Result<Vec<_>, _>>()
907            .unwrap();
908        assert_eq!(frames.len(), 2);
909        assert_eq!(frames[0].content, b"hello");
910        assert_eq!(frames[0].direction, Direction::Recv);
911        assert_eq!(frames[1].content, b"world");
912        assert_eq!(frames[1].direction, Direction::Sent);
913    }
914
915    /// A frame's offset is where its header begins, so a caller chaining two
916    /// sources can say which one a frame came out of.
917    #[test]
918    fn frame_iterator_states_where_each_frame_began() {
919        let first = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
920        let second = b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n";
921        let mut data = first.to_vec();
922        data.extend_from_slice(second);
923
924        let frames: Vec<Frame> = FrameIterator::new(&data[..])
925            .collect::<Result<Vec<_>, _>>()
926            .unwrap();
927        assert_eq!(frames[0].offset, 0);
928        assert_eq!(frames[1].offset, first.len() as u64);
929    }
930
931    /// The reader's chunking is not the frame's: a source handing back one byte
932    /// per call must not move where a frame is reported to start.
933    #[test]
934    fn a_short_reader_does_not_move_a_frame_offset() {
935        struct OneByteAtATime<'a>(&'a [u8]);
936        impl std::io::Read for OneByteAtATime<'_> {
937            fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
938                match (self.0.first(), buf.is_empty()) {
939                    (Some(&b), false) => {
940                        buf[0] = b;
941                        self.0 = &self.0[1..];
942                        Ok(1)
943                    }
944                    _ => Ok(0),
945                }
946            }
947        }
948
949        let first = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
950        let mut data = first.to_vec();
951        data.extend_from_slice(
952            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
953        );
954
955        let frames: Vec<Frame> = FrameIterator::new(OneByteAtATime(&data))
956            .collect::<Result<Vec<_>, _>>()
957            .unwrap();
958        assert_eq!(frames[0].offset, 0);
959        assert_eq!(frames[1].offset, first.len() as u64);
960    }
961
962    #[test]
963    fn frame_iterator_vt_in_content() {
964        // \x0B in content but not followed by valid header — should NOT split
965        let mut data = Vec::new();
966        data.extend_from_slice(b"recv 15 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n");
967        data.extend_from_slice(b"he\x0B\nllo world!!");
968        data.extend_from_slice(b"\x0B\n");
969        let frames: Vec<Frame> = FrameIterator::new(&data[..])
970            .collect::<Result<Vec<_>, _>>()
971            .unwrap();
972        assert_eq!(frames.len(), 1);
973        assert_eq!(frames[0].content, b"he\x0B\nllo world!!");
974    }
975
976    #[test]
977    fn frame_iterator_eof_without_boundary() {
978        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello";
979        let frames: Vec<Frame> = FrameIterator::new(&data[..])
980            .collect::<Result<Vec<_>, _>>()
981            .unwrap();
982        assert_eq!(frames.len(), 1);
983        assert_eq!(frames[0].content, b"hello");
984    }
985
986    #[test]
987    fn frame_iterator_eof_with_lone_vt() {
988        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B";
989        let frames: Vec<Frame> = FrameIterator::new(&data[..])
990            .collect::<Result<Vec<_>, _>>()
991            .unwrap();
992        assert_eq!(frames.len(), 1);
993        assert_eq!(frames[0].content, b"hello");
994    }
995
996    #[test]
997    fn frame_iterator_partial_first_frame() {
998        // Data starts with garbage, then a valid boundary + frame
999        let mut data = Vec::new();
1000        data.extend_from_slice(b"partial garbage data");
1001        data.extend_from_slice(b"\x0B\n");
1002        data.extend_from_slice(
1003            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1004        );
1005        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1006            .collect::<Result<Vec<_>, _>>()
1007            .unwrap();
1008        assert_eq!(frames.len(), 1);
1009        assert_eq!(frames[0].content, b"hello");
1010    }
1011
1012    #[test]
1013    fn frame_iterator_truncated_last_frame() {
1014        // Complete frame followed by truncated frame at EOF (no \x0B\n)
1015        let mut data = Vec::new();
1016        data.extend_from_slice(
1017            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1018        );
1019        data.extend_from_slice(b"sent 3 bytes to tcp/1.1.1.1:5060 at 00:00:01.000000:\nbye");
1020        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1021            .collect::<Result<Vec<_>, _>>()
1022            .unwrap();
1023        assert_eq!(frames.len(), 2);
1024        assert_eq!(frames[0].content, b"hello");
1025        assert_eq!(frames[1].content, b"bye");
1026    }
1027
1028    #[test]
1029    fn frame_iterator_file_concatenation() {
1030        // Simulates `cat dump.20 dump.21 | parser`
1031        // File 1: truncated start + valid frame + complete end
1032        // File 2: truncated start (no header) + valid frame
1033        let mut data = Vec::new();
1034
1035        // File 1: starts with valid header
1036        data.extend_from_slice(
1037            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1038        );
1039        data.extend_from_slice(
1040            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
1041        );
1042
1043        // File 2: starts with truncated frame data (no header), then boundary, then valid frame
1044        data.extend_from_slice(b"some truncated SIP content from previous rotation\r\n\r\n");
1045        data.extend_from_slice(b"\x0B\n");
1046        data.extend_from_slice(
1047            b"recv 3 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\nfoo\x0B\n",
1048        );
1049
1050        let items: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1051        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1052        assert_eq!(frames.len(), 3);
1053        assert_eq!(frames[0].content, b"hello");
1054        assert_eq!(frames[1].content, b"world");
1055        assert_eq!(frames[2].content, b"foo");
1056        assert_eq!(frames[2].address, "2.2.2.2:5060");
1057    }
1058
1059    #[test]
1060    fn frame_iterator_file_concatenation_mid_stream_garbage() {
1061        // The join point between files produces garbage that looks like:
1062        // ...last_content\x0B\ntruncated_first_of_next_file\x0B\nvalid_header...
1063        // The truncated part is NOT a valid header, so recovery should skip it
1064        let mut data = Vec::new();
1065
1066        // Last frame of file 1
1067        data.extend_from_slice(
1068            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1069        );
1070
1071        // Truncated first frame of file 2 (mid-SIP content, no frame header)
1072        data.extend_from_slice(b"Content-Type: application/sdp\r\n\r\nv=0\r\n");
1073        data.extend_from_slice(b"\x0B\n");
1074
1075        // Valid second frame of file 2
1076        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1077
1078        let items: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1079        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1080        assert_eq!(frames.len(), 2);
1081        assert_eq!(frames[0].content, b"hello");
1082        assert_eq!(frames[1].content, b"bar");
1083    }
1084
1085    #[test]
1086    fn frame_iterator_empty_input() {
1087        let data: &[u8] = b"";
1088        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(data).collect();
1089        assert!(frames.is_empty());
1090    }
1091
1092    #[test]
1093    fn frame_iterator_only_garbage() {
1094        let data = b"this is not a SIP trace dump at all, just garbage text";
1095        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1096        let frames: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1097        assert!(frames.is_empty());
1098        let stats = iter.stats();
1099        assert_eq!(stats.bytes_read, data.len() as u64);
1100        assert_eq!(stats.bytes_skipped, data.len() as u64);
1101        assert_eq!(stats.unparsed_regions.len(), 1);
1102        assert_eq!(stats.unparsed_regions[0].reason, SkipReason::InvalidHeader);
1103    }
1104
1105    #[test]
1106    fn frame_iterator_truncated_header_at_eof() {
1107        let data = b"recv 5 bytes from tcp/1.1.1.1:5060";
1108        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1109        let frames: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1110        assert!(frames.is_empty());
1111        let stats = iter.stats();
1112        assert_eq!(stats.bytes_read, data.len() as u64);
1113        assert_eq!(stats.bytes_skipped, data.len() as u64);
1114    }
1115
1116    #[test]
1117    fn frame_iterator_dump_marker_at_eof() {
1118        // A dump restart marker at the end of input (with trailing \n\n as in real dumps)
1119        // should be silently consumed, not returned as an error.
1120        let mut data = Vec::new();
1121        data.extend_from_slice(
1122            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1123        );
1124        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1125
1126        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1127        assert_eq!(frames.len(), 1);
1128        assert!(frames[0].is_ok());
1129        assert_eq!(frames[0].as_ref().unwrap().content, b"hello");
1130    }
1131
1132    #[test]
1133    fn frame_iterator_dump_marker_mid_stream() {
1134        // A dump restart marker between two valid frames (with trailing \n\n as in real dumps)
1135        // should be skipped, and both frames should parse successfully.
1136        let mut data = Vec::new();
1137        data.extend_from_slice(
1138            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1139        );
1140        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1141        data.extend_from_slice(b"sent 3 bytes to tcp/2.2.2.2:5060 at 00:00:01.000000:\nbye\x0B\n");
1142
1143        let frames: Vec<Result<Frame, ParseError>> = FrameIterator::new(&data[..]).collect();
1144        assert_eq!(frames.len(), 2);
1145        assert_eq!(frames[0].as_ref().unwrap().content, b"hello");
1146        assert_eq!(frames[1].as_ref().unwrap().content, b"bye");
1147    }
1148
1149    #[test]
1150    fn frame_iterator_multiple_newlines_after_boundary() {
1151        // Multiple \n and \r\n between frames should all be stripped.
1152        let mut data = Vec::new();
1153        data.extend_from_slice(
1154            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1155        );
1156        data.extend_from_slice(b"\n\r\n\n");
1157        data.extend_from_slice(
1158            b"sent 5 bytes to tcp/1.1.1.1:5060 at 00:00:00.000001:\nworld\x0B\n",
1159        );
1160
1161        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1162            .collect::<Result<Vec<_>, _>>()
1163            .unwrap();
1164        assert_eq!(frames.len(), 2);
1165        assert_eq!(frames[0].content, b"hello");
1166        assert_eq!(frames[1].content, b"world");
1167    }
1168
1169    #[test]
1170    fn stats_clean_input() {
1171        let data = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
1172        let mut iter = FrameIterator::new(&data[..]);
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.bytes_read, data.len() as u64);
1177        assert_eq!(stats.bytes_skipped, 0);
1178        assert!(stats.unparsed_regions.is_empty());
1179    }
1180
1181    #[test]
1182    fn stats_multiple_frames() {
1183        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";
1184        let mut iter = FrameIterator::new(&data[..]);
1185        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1186        assert_eq!(frames.len(), 2);
1187        let stats = iter.stats();
1188        assert_eq!(stats.bytes_read, data.len() as u64);
1189        assert_eq!(stats.bytes_skipped, 0);
1190        assert!(stats.unparsed_regions.is_empty());
1191    }
1192
1193    #[test]
1194    fn stats_partial_first_frame() {
1195        let mut data = Vec::new();
1196        data.extend_from_slice(b"partial garbage data");
1197        data.extend_from_slice(b"\x0B\n");
1198        data.extend_from_slice(
1199            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1200        );
1201        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1202        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1203        assert_eq!(frames.len(), 1);
1204        let stats = iter.stats();
1205        assert_eq!(stats.bytes_read, data.len() as u64);
1206        // "partial garbage data" + "\x0B\n" = 21 bytes skipped
1207        let skipped = b"partial garbage data\x0B\n".len() as u64;
1208        assert_eq!(stats.bytes_skipped, skipped);
1209        assert_eq!(stats.unparsed_regions.len(), 1);
1210        assert_eq!(stats.unparsed_regions[0].offset, 0);
1211        assert_eq!(stats.unparsed_regions[0].length, skipped);
1212        assert_eq!(
1213            stats.unparsed_regions[0].reason,
1214            crate::types::SkipReason::PartialFirstFrame
1215        );
1216        assert!(stats.unparsed_regions[0].data.is_none());
1217    }
1218
1219    #[test]
1220    fn stats_partial_first_frame_capture() {
1221        let mut data = Vec::new();
1222        data.extend_from_slice(b"partial garbage data");
1223        data.extend_from_slice(b"\x0B\n");
1224        data.extend_from_slice(
1225            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1226        );
1227        let mut iter = FrameIterator::new(&data[..]).capture_skipped(true);
1228        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1229        assert_eq!(frames.len(), 1);
1230        let stats = iter.stats();
1231        assert_eq!(stats.unparsed_regions.len(), 1);
1232        let region = &stats.unparsed_regions[0];
1233        assert_eq!(
1234            region.data.as_deref(),
1235            Some(b"partial garbage data\x0B\n".as_slice())
1236        );
1237    }
1238
1239    #[test]
1240    fn stats_mid_stream_partial_frame() {
1241        // SIP content between valid frames (file concatenation scenario)
1242        let mut data = Vec::new();
1243        data.extend_from_slice(
1244            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1245        );
1246        data.extend_from_slice(b"Content-Type: application/sdp\r\n\r\nv=0\r\n");
1247        data.extend_from_slice(b"\x0B\n");
1248        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1249
1250        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1251        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1252        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1253        assert_eq!(frames.len(), 2);
1254        let stats = iter.stats();
1255        assert!(stats.bytes_skipped > 0);
1256        assert_eq!(stats.unparsed_regions.len(), 1);
1257        assert_eq!(
1258            stats.unparsed_regions[0].reason,
1259            crate::types::SkipReason::MidStreamSkip
1260        );
1261    }
1262
1263    #[test]
1264    fn stats_replayed_frame() {
1265        // Simulate logrotate: a frame's tail (SIP headers ending with \r\n\r\n\x0B\n)
1266        // appears between two valid frames at a file boundary.
1267        let frame1 = b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n";
1268        let replay = b"Route: <sip:10.0.0.1:5060;lr>\r\nContent-Length: 0\r\n\r\n\x0B\n";
1269        let frame2 = b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n";
1270
1271        let mut data = Vec::new();
1272        data.extend_from_slice(frame1);
1273        data.extend_from_slice(replay);
1274        data.extend_from_slice(frame2);
1275
1276        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1277        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1278        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1279        assert_eq!(frames.len(), 2);
1280        let stats = iter.stats();
1281        assert_eq!(stats.unparsed_regions.len(), 1);
1282        assert_eq!(
1283            stats.unparsed_regions[0].reason,
1284            crate::types::SkipReason::ReplayedFrame
1285        );
1286    }
1287
1288    #[test]
1289    fn stats_incomplete_frame_at_eof() {
1290        // Frame header says 100 bytes but only 20 bytes available before EOF
1291        let mut data = Vec::new();
1292        data.extend_from_slice(
1293            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1294        );
1295        data.extend_from_slice(b"recv 100 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\n");
1296        data.extend_from_slice(b"partial content only");
1297        // No \x0B\n boundary — EOF truncation
1298
1299        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1300        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1301        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1302        assert_eq!(frames.len(), 2, "truncated frame should still be returned");
1303        assert_eq!(frames[1].content, b"partial content only");
1304        assert_eq!(frames[1].byte_count, 100);
1305        let stats = iter.stats();
1306        assert_eq!(stats.unparsed_regions.len(), 1);
1307        assert_eq!(
1308            stats.unparsed_regions[0].reason,
1309            crate::types::SkipReason::IncompleteFrame
1310        );
1311    }
1312
1313    #[test]
1314    fn stats_incomplete_frame_counted_under_count_only() {
1315        let mut data = Vec::new();
1316        data.extend_from_slice(
1317            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1318        );
1319        data.extend_from_slice(b"recv 100 bytes from tcp/2.2.2.2:5060 at 01:00:00.000000:\n");
1320        data.extend_from_slice(b"partial content only");
1321
1322        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::CountOnly);
1323        let _: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1324        let stats = iter.stats();
1325        assert!(stats.unparsed_regions.is_empty());
1326        assert_eq!(stats.incomplete_frames, 1);
1327        assert_eq!(stats.incomplete_frame_bytes, 80);
1328    }
1329
1330    #[test]
1331    fn stats_invalid_header_skip() {
1332        // Malformed frame header (starts with recv/sent but unparseable)
1333        let mut data = Vec::new();
1334        data.extend_from_slice(
1335            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1336        );
1337        data.extend_from_slice(b"recv CORRUPT HEADER garbage\n");
1338        data.extend_from_slice(b"\x0B\n");
1339        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1340
1341        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1342        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1343        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1344        assert_eq!(frames.len(), 2);
1345        let stats = iter.stats();
1346        assert!(stats.bytes_skipped > 0);
1347        assert_eq!(stats.unparsed_regions.len(), 1);
1348        assert_eq!(
1349            stats.unparsed_regions[0].reason,
1350            crate::types::SkipReason::InvalidHeader
1351        );
1352    }
1353
1354    #[test]
1355    fn stats_oversized_frame_at_start() {
1356        let mut data = Vec::new();
1357        data.resize(MAX_PARTIAL_FRAME + 1, b'x');
1358        data.extend_from_slice(b"\x0B\n");
1359        data.extend_from_slice(
1360            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1361        );
1362        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1363        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1364        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1365        assert_eq!(frames.len(), 1);
1366        let stats = iter.stats();
1367        assert_eq!(stats.unparsed_regions.len(), 1);
1368        assert_eq!(
1369            stats.unparsed_regions[0].reason,
1370            crate::types::SkipReason::OversizedFrame
1371        );
1372    }
1373
1374    #[test]
1375    fn stats_oversized_frame_mid_stream() {
1376        let mut data = Vec::new();
1377        data.extend_from_slice(
1378            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1379        );
1380        let garbage_len = MAX_PARTIAL_FRAME + 1;
1381        data.resize(data.len() + garbage_len, b'x');
1382        data.extend_from_slice(b"\x0B\n");
1383        data.extend_from_slice(b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n");
1384        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1385        let items: Vec<Result<Frame, ParseError>> = iter.by_ref().collect();
1386        let frames: Vec<Frame> = items.into_iter().filter_map(Result::ok).collect();
1387        assert_eq!(frames.len(), 2);
1388        let stats = iter.stats();
1389        assert_eq!(stats.unparsed_regions.len(), 1);
1390        assert_eq!(
1391            stats.unparsed_regions[0].reason,
1392            crate::types::SkipReason::OversizedFrame
1393        );
1394    }
1395
1396    #[test]
1397    fn overlong_byte_count_does_not_buffer_ahead() {
1398        let mut data = Vec::new();
1399        data.extend_from_slice(
1400            format!(
1401                "recv {} bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\n",
1402                MAX_PARTIAL_FRAME * 10
1403            )
1404            .as_bytes(),
1405        );
1406        data.extend_from_slice(b"hello\x0B\n");
1407        while data.len() < MAX_PARTIAL_FRAME * 10 {
1408            data.extend_from_slice(
1409                b"sent 3 bytes to tcp/3.3.3.3:5060 at 02:00:00.000000:\nbar\x0B\n",
1410            );
1411        }
1412
1413        let mut iter = FrameIterator::new(&data[..]);
1414        let first = iter.next().unwrap().unwrap();
1415        assert_eq!(first.content, b"hello");
1416        assert!(
1417            iter.buf.len() <= MAX_PARTIAL_FRAME + READ_BUF_SIZE,
1418            "buffered {} bytes on a byte_count ten times the frame",
1419            iter.buf.len()
1420        );
1421        let second = iter.next().unwrap().unwrap();
1422        assert_eq!(second.content, b"bar");
1423    }
1424
1425    #[test]
1426    fn stats_partial_first_frame_within_limit() {
1427        // Content + \x0B\n boundary = MAX_PARTIAL_FRAME, should still be PartialFirstFrame
1428        let mut data = Vec::new();
1429        data.resize(MAX_PARTIAL_FRAME - 2, b'x');
1430        data.extend_from_slice(b"\x0B\n");
1431        data.extend_from_slice(
1432            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1433        );
1434        let mut iter = FrameIterator::new(&data[..]).skip_tracking(SkipTracking::TrackRegions);
1435        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1436        assert_eq!(frames.len(), 1);
1437        let stats = iter.stats();
1438        assert_eq!(stats.unparsed_regions.len(), 1);
1439        assert_eq!(
1440            stats.unparsed_regions[0].reason,
1441            crate::types::SkipReason::PartialFirstFrame
1442        );
1443    }
1444
1445    #[test]
1446    fn stats_dump_restart_marker() {
1447        let mut data = Vec::new();
1448        data.extend_from_slice(
1449            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1450        );
1451        data.extend_from_slice(b"dump started at Thu Aug 22 11:38:11 2024\n\n\n");
1452        data.extend_from_slice(b"sent 3 bytes to tcp/2.2.2.2:5060 at 00:00:01.000000:\nbye\x0B\n");
1453
1454        let mut iter = FrameIterator::new(&data[..]);
1455        let frames: Vec<Frame> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1456        assert_eq!(frames.len(), 2);
1457        let stats = iter.stats();
1458        // Dump restart marker is structural, not skipped
1459        assert_eq!(stats.bytes_skipped, 0);
1460        assert!(stats.unparsed_regions.is_empty());
1461    }
1462
1463    #[test]
1464    fn stats_count_only_no_regions() {
1465        let mut data = Vec::new();
1466        data.extend_from_slice(b"partial garbage data");
1467        data.extend_from_slice(b"\x0B\n");
1468        data.extend_from_slice(
1469            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1470        );
1471        let mut iter = FrameIterator::new(&data[..]);
1472        let frames: Vec<_> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1473        assert_eq!(frames.len(), 1);
1474        let stats = iter.stats();
1475        let skipped = b"partial garbage data\x0B\n".len() as u64;
1476        assert_eq!(stats.bytes_skipped, skipped);
1477        assert!(
1478            stats.unparsed_regions.is_empty(),
1479            "CountOnly should not accumulate regions"
1480        );
1481    }
1482
1483    #[test]
1484    fn frame_iterator_trailing_newlines_at_eof() {
1485        // Trailing newlines after the last boundary at EOF should not cause errors.
1486        let mut data = Vec::new();
1487        data.extend_from_slice(
1488            b"recv 5 bytes from tcp/1.1.1.1:5060 at 00:00:00.000000:\nhello\x0B\n",
1489        );
1490        data.extend_from_slice(b"\n\n");
1491
1492        let frames: Vec<Frame> = FrameIterator::new(&data[..])
1493            .collect::<Result<Vec<_>, _>>()
1494            .unwrap();
1495        assert_eq!(frames.len(), 1);
1496        assert_eq!(frames[0].content, b"hello");
1497    }
1498}