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 ";
14const MAX_PARTIAL_FRAME: usize = 65537;
17
18#[derive(Debug)]
23#[non_exhaustive]
24pub enum ParseError {
25 InvalidHeader(String),
27 InvalidMessage(String),
29 TransportNoise {
32 bytes: usize,
34 transport: Transport,
36 address: String,
38 },
39 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
83fn 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
96fn parse_timestamp(bytes: &[u8]) -> Option<Timestamp> {
98 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 let (hour, min, sec, usec) = parse_time_part(bytes)?;
116 Some(Timestamp::TimeOnly {
117 hour,
118 min,
119 sec,
120 usec,
121 })
122}
123
124fn 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#[derive(Debug)]
141pub struct FrameHeader {
142 pub direction: Direction,
144 pub byte_count: usize,
146 pub transport: Transport,
148 pub address: String,
150 pub timestamp: Timestamp,
152 pub header_len: usize,
154}
155
156pub 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 let line = line.strip_suffix(b"\r").unwrap_or(line);
166 let line = line
168 .strip_suffix(b":")
169 .ok_or_else(|| ParseError::InvalidHeader("header does not end with ':'".into()))?;
170
171 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 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 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 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 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 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
262fn 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
276pub 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 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
305pub 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 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 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 pub fn skip_tracking(mut self, tracking: SkipTracking) -> Self {
363 self.skip_tracking = tracking;
364 self
365 }
366
367 pub fn stats(&self) -> &ParseStats {
369 &self.stats
370 }
371
372 pub(crate) fn stats_mut(&mut self) -> &mut ParseStats {
374 &mut self.stats
375 }
376
377 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 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 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 if self.eof {
441 return Some(abs_pos);
442 }
443 return None; }
445 if is_frame_header(&self.buf[after..]) {
446 return Some(abs_pos);
447 }
448 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 fn skip_to_first_header(&mut self) -> Option<usize> {
459 if is_frame_header(&self.buf) {
460 return Some(0);
461 }
462 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 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 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 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 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 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 let content = loop {
630 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 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 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 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 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 #[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 #[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 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 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 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 let mut data = Vec::new();
1034
1035 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 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 let mut data = Vec::new();
1065
1066 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}