1use std::collections::{HashMap, VecDeque};
2
3use tracing::{debug, trace, warn};
4
5use crate::finders::{CRLF, CRLFCRLF};
6use crate::frame::{FrameIterator, ParseError};
7use crate::startline::{sip_start, SipStart};
8use crate::types::{
9 Direction, ParseStats, SipMessage, SkipTracking, StaleClock, Timestamp, Transport,
10 UnparsedRegion,
11};
12
13const EVICTION_SAMPLE: usize = 4;
15
16pub struct MessageIterator<R> {
39 frames: FrameIterator<R>,
40 buffers: HashMap<ConnectionKey, ConnectionBuffer>,
41 ready: VecDeque<SipMessage>,
42 exhausted: bool,
43 clock: StaleClock,
44}
45
46#[derive(Clone, PartialEq, Eq, Hash)]
47struct ConnectionKey {
48 direction: Direction,
49 address: String,
50}
51
52struct ConnectionBuffer {
53 transport: Transport,
54 timestamp: Timestamp,
55 content: Vec<u8>,
56 frame_count: usize,
57 last_seen: u64,
58}
59
60impl<R: std::io::Read> MessageIterator<R> {
61 pub fn new(reader: R) -> Self {
63 MessageIterator {
64 frames: FrameIterator::new(reader),
65 buffers: HashMap::new(),
66 ready: VecDeque::new(),
67 exhausted: false,
68 clock: StaleClock::new(),
69 }
70 }
71
72 pub fn capture_skipped(mut self, enable: bool) -> Self {
76 self.frames = self.frames.capture_skipped(enable);
77 self
78 }
79
80 pub fn skip_tracking(mut self, tracking: SkipTracking) -> Self {
82 self.frames = self.frames.skip_tracking(tracking);
83 self
84 }
85
86 pub fn parse_stats(&self) -> &ParseStats {
88 self.frames.stats()
89 }
90
91 pub fn drain_unparsed(&mut self) -> Vec<UnparsedRegion> {
93 self.frames.drain_unparsed()
94 }
95
96 fn sweep_stale_buffers(&mut self) {
97 let clock = &self.clock;
98 let mut incomplete = 0usize;
99 let mut pending_bytes = 0usize;
100 let mut dropped: Vec<String> = Vec::new();
101 self.buffers.retain(|key, buf| {
102 if !clock.is_stale(buf.last_seen) {
103 return true;
104 }
105 if buf.content.is_empty() {
106 trace!(
107 address = %key.address,
108 direction = %key.direction,
109 elapsed_secs = clock.now().saturating_sub(buf.last_seen),
110 "evicted empty stale connection buffer"
111 );
112 } else {
113 incomplete += 1;
114 pending_bytes += buf.content.len();
115 if dropped.len() < EVICTION_SAMPLE {
116 dropped.push(format!(
117 "{}/{} {}",
118 buf.transport, key.direction, key.address
119 ));
120 }
121 }
122 false
123 });
124 if incomplete > 0 {
125 let stats = self.frames.stats_mut();
126 stats.stale_evictions += incomplete as u64;
127 stats.stale_evicted_bytes += pending_bytes as u64;
128 warn!(
129 buffers = incomplete,
130 pending_bytes,
131 connections = %dropped.join(", "),
132 undisplayed = incomplete - dropped.len(),
133 "evicted stale connection buffers with incomplete data"
134 );
135 }
136 }
137
138 fn flush_all(&mut self) {
139 let mut loss = ResyncLoss::default();
140 for (key, mut buf) in std::mem::take(&mut self.buffers) {
141 extract_complete(&mut buf, &key, &mut self.ready, &mut loss);
142
143 if !buf.content.is_empty() {
144 let content = std::mem::take(&mut buf.content);
145 let frame_count = buf.frame_count;
146 self.ready
147 .push_back(message(&buf, &key, content, frame_count));
148 }
149 }
150 self.record_resync_loss(loss);
151 }
152
153 fn record_resync_loss(&mut self, loss: ResyncLoss) {
154 if loss.count == 0 {
155 return;
156 }
157 let stats = self.frames.stats_mut();
158 stats.non_sip_prefixes += loss.count;
159 stats.non_sip_prefix_bytes += loss.bytes;
160 }
161}
162
163impl<R: std::io::Read> Iterator for MessageIterator<R> {
164 type Item = Result<SipMessage, ParseError>;
165
166 fn next(&mut self) -> Option<Self::Item> {
167 if let Some(msg) = self.ready.pop_front() {
168 return Some(Ok(msg));
169 }
170
171 if self.exhausted {
172 return None;
173 }
174
175 loop {
176 match self.frames.next() {
177 Some(Ok(frame)) => {
178 if frame.transport == Transport::Udp {
179 return Some(Ok(SipMessage {
180 direction: frame.direction,
181 transport: frame.transport,
182 address: frame.address,
183 timestamp: frame.timestamp,
184 content: frame.content,
185 frame_count: 1,
186 }));
187 }
188
189 let now = self.clock.observe(frame.timestamp);
190 if self.clock.sweep_due() {
191 self.sweep_stale_buffers();
192 }
193
194 let key = ConnectionKey {
195 direction: frame.direction,
196 address: frame.address,
197 };
198
199 let buf = match self.buffers.get_mut(&key) {
200 Some(buf) => buf,
201 None => self.buffers.entry(key.clone()).or_insert(ConnectionBuffer {
202 transport: frame.transport,
203 timestamp: frame.timestamp,
204 content: Vec::new(),
205 frame_count: 0,
206 last_seen: now,
207 }),
208 };
209
210 buf.last_seen = now;
211
212 if buf.content.is_empty() {
213 buf.timestamp = frame.timestamp;
214 }
215
216 trace!(
217 frame = buf.frame_count + 1,
218 bytes = frame.content.len(),
219 address = %key.address,
220 "buffering TCP frame"
221 );
222
223 buf.content.extend_from_slice(&frame.content);
224 buf.frame_count += 1;
225
226 let mut loss = ResyncLoss::default();
227 extract_complete(buf, &key, &mut self.ready, &mut loss);
228
229 if buf.content.is_empty() {
230 self.buffers.remove(&key);
231 }
232 self.record_resync_loss(loss);
233
234 if let Some(msg) = self.ready.pop_front() {
235 return Some(Ok(msg));
236 }
237 }
238 Some(Err(e)) => return Some(Err(e)),
239 None => {
240 self.exhausted = true;
241 self.flush_all();
242 return self.ready.pop_front().map(Ok);
243 }
244 }
245 }
246 }
247}
248
249enum Resync {
251 Ready,
252 Retry,
253 Wait,
254}
255
256#[derive(Default)]
259struct ResyncLoss {
260 count: u64,
261 bytes: u64,
262}
263
264fn extract_complete(
266 buf: &mut ConnectionBuffer,
267 key: &ConnectionKey,
268 ready: &mut VecDeque<SipMessage>,
269 loss: &mut ResyncLoss,
270) {
271 loop {
272 if buf.content.is_empty() {
273 return;
274 }
275 match resync_to_sip_start(buf, key, loss) {
276 Resync::Ready => {}
277 Resync::Retry => continue,
278 Resync::Wait => return,
279 }
280 if !split_one_message(buf, key, ready) {
281 return;
282 }
283 }
284}
285
286fn resync_to_sip_start(
288 buf: &mut ConnectionBuffer,
289 key: &ConnectionKey,
290 loss: &mut ResyncLoss,
291) -> Resync {
292 match sip_start(&buf.content) {
293 SipStart::Yes => return Resync::Ready,
294 SipStart::NeedMore => return Resync::Wait, SipStart::No => {}
296 }
297
298 let ws_len = buf
300 .content
301 .iter()
302 .position(|&b| !matches!(b, b'\r' | b'\n' | b' ' | b'\t'))
303 .unwrap_or(buf.content.len());
304
305 if ws_len > 0 {
306 if ws_len == buf.content.len() {
307 trace!(
308 bytes = ws_len,
309 address = %key.address,
310 "drained transport whitespace"
311 );
312 buf.content.clear();
313 buf.frame_count = 0;
314 return Resync::Wait;
315 }
316 match sip_start(&buf.content[ws_len..]) {
317 SipStart::Yes => {
318 trace!(bytes = ws_len, "drained inter-message whitespace padding");
319 buf.content.drain(..ws_len);
320 return Resync::Retry;
321 }
322 SipStart::NeedMore => {
323 trace!(bytes = ws_len, "drained inter-message whitespace padding");
324 buf.content.drain(..ws_len);
325 return Resync::Wait;
326 }
327 SipStart::No => {}
328 }
329 }
330
331 match find_sip_start(&buf.content) {
332 Some(offset) if offset > 0 => {
333 debug!(
334 skipped_bytes = offset,
335 address = %key.address,
336 "skipped non-SIP prefix in TCP buffer"
337 );
338 loss.count += 1;
339 loss.bytes += offset as u64;
340 buf.content.drain(..offset);
341 Resync::Retry
342 }
343 _ => Resync::Wait, }
345}
346
347fn split_one_message(
350 buf: &mut ConnectionBuffer,
351 key: &ConnectionKey,
352 ready: &mut VecDeque<SipMessage>,
353) -> bool {
354 let header_end = match CRLFCRLF.find(&buf.content) {
355 Some(offset) => offset,
356 None => return false, };
358 let body_start = header_end + 4;
359
360 let msg_end = match find_content_length(&buf.content, header_end) {
361 Some(cl) => {
362 let end = body_start + cl;
363 if end > buf.content.len() {
364 return false; }
366 end
367 }
368 None => body_start, };
370
371 let remaining = buf.content.split_off(msg_end);
372 let msg_content = std::mem::replace(&mut buf.content, remaining);
373
374 while buf.content.len() >= 2 && buf.content[0] == b'\r' && buf.content[1] == b'\n' {
376 buf.content.drain(..2);
377 }
378
379 let frame_count = buf.frame_count;
380 if frame_count > 1 {
381 debug!(
382 frame_count,
383 bytes = msg_content.len(),
384 address = %key.address,
385 "extracted reassembled TCP message"
386 );
387 }
388
389 ready.push_back(message(buf, key, msg_content, frame_count));
390 buf.frame_count = 0;
391 true
392}
393
394fn message(
395 buf: &ConnectionBuffer,
396 key: &ConnectionKey,
397 content: Vec<u8>,
398 frame_count: usize,
399) -> SipMessage {
400 SipMessage {
401 direction: key.direction,
402 transport: buf.transport,
403 address: key.address.clone(),
404 timestamp: buf.timestamp,
405 content,
406 frame_count,
407 }
408}
409
410fn find_content_length(data: &[u8], header_end: usize) -> Option<usize> {
412 let headers = &data[..header_end];
413
414 let mut pos = 0;
415 while pos < headers.len() {
416 let line_end = CRLF.find(&headers[pos..]).unwrap_or(headers.len() - pos);
417 let line = &headers[pos..pos + line_end];
418
419 if let Some(value) = extract_header_value(line, b"Content-Length") {
420 return parse_content_length(value);
421 }
422 if let Some(value) = extract_compact_header_value(line, b'l') {
423 return parse_content_length(value);
424 }
425
426 pos += line_end + 2; }
428 None
429}
430
431fn extract_header_value<'a>(line: &'a [u8], name: &[u8]) -> Option<&'a [u8]> {
434 let rest = line.get(name.len()..)?;
435 if !line[..name.len()].eq_ignore_ascii_case(name) {
436 return None;
437 }
438 let pad = rest
439 .iter()
440 .position(|&c| c != b' ' && c != b'\t')
441 .unwrap_or(rest.len());
442 match rest.get(pad..)?.split_first() {
443 Some((b':', value)) => Some(trim_bytes(value)),
444 _ => None,
445 }
446}
447
448fn extract_compact_header_value(line: &[u8], compact: u8) -> Option<&[u8]> {
449 if line.len() < 2 {
450 return None;
451 }
452 if line[0] != compact || line[1] != b':' {
453 return None;
454 }
455 Some(trim_bytes(&line[2..]))
456}
457
458fn trim_bytes(b: &[u8]) -> &[u8] {
459 let start = b
460 .iter()
461 .position(|&c| c != b' ' && c != b'\t')
462 .unwrap_or(b.len());
463 let end = b
464 .iter()
465 .rposition(|&c| c != b' ' && c != b'\t')
466 .map_or(start, |p| p + 1);
467 &b[start..end]
468}
469
470fn parse_content_length(value: &[u8]) -> Option<usize> {
471 let s = std::str::from_utf8(value).ok()?;
472 s.parse().ok()
473}
474
475fn find_sip_start(data: &[u8]) -> Option<usize> {
479 if !matches!(sip_start(data), SipStart::No) {
480 return Some(0);
481 }
482 let mut pos = 0;
483 while let Some(offset) = CRLF.find(&data[pos..]) {
484 let candidate = pos + offset + 2;
485 if candidate >= data.len() {
486 break;
487 }
488 if !matches!(sip_start(&data[candidate..]), SipStart::No) {
489 return Some(candidate);
490 }
491 pos = candidate;
492 }
493 None
494}
495
496#[cfg(test)]
497mod tests {
498 use super::*;
499 use crate::types::Direction;
500
501 fn buffer_with(content: Vec<u8>) -> ConnectionBuffer {
502 ConnectionBuffer {
503 transport: Transport::Tcp,
504 timestamp: Timestamp::TimeOnly {
505 hour: 0,
506 min: 0,
507 sec: 0,
508 usec: 0,
509 },
510 content,
511 frame_count: 1,
512 last_seen: 0,
513 }
514 }
515
516 fn extracted(buf: &mut ConnectionBuffer) -> Vec<SipMessage> {
517 let key = ConnectionKey {
518 direction: Direction::Recv,
519 address: "[::1]:5060".to_string(),
520 };
521 let mut ready = VecDeque::new();
522 extract_complete(buf, &key, &mut ready, &mut ResyncLoss::default());
523 ready.into()
524 }
525
526 fn content_length(data: &[u8]) -> Option<usize> {
527 find_content_length(data, CRLFCRLF.find(data)?)
528 }
529
530 fn make_frame(
531 direction: Direction,
532 transport: Transport,
533 addr: &str,
534 content: &[u8],
535 ) -> Vec<u8> {
536 make_frame_at(direction, transport, addr, content, "00:00:00.000000")
537 }
538
539 #[test]
540 fn single_udp_message() {
541 let content = b"OPTIONS sip:user@host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
542 let data = make_frame(Direction::Recv, Transport::Udp, "1.1.1.1:5060", content);
543 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
544 .collect::<Result<Vec<_>, _>>()
545 .unwrap();
546 assert_eq!(msgs.len(), 1);
547 assert_eq!(msgs[0].content, content);
548 assert_eq!(msgs[0].frame_count, 1);
549 assert_eq!(msgs[0].transport, Transport::Udp);
550 }
551
552 #[test]
553 fn tcp_reassembly_two_frames() {
554 let part1 = b"NOTIFY sip:user@host SIP/2.0\r\n";
555 let part2 = b"Content-Length: 0\r\n\r\n";
556 let mut data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", part1);
557 data.extend_from_slice(&make_frame(
558 Direction::Recv,
559 Transport::Tcp,
560 "[::1]:5060",
561 part2,
562 ));
563 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
564 .collect::<Result<Vec<_>, _>>()
565 .unwrap();
566 assert_eq!(msgs.len(), 1);
567 assert_eq!(msgs[0].frame_count, 2);
568 let mut expected = Vec::new();
569 expected.extend_from_slice(part1);
570 expected.extend_from_slice(part2);
571 assert_eq!(msgs[0].content, expected);
572 }
573
574 #[test]
575 fn tcp_reassembly_across_interleaved_frames() {
576 let part1 = b"INVITE sip:user@host SIP/2.0\r\n";
580 let part2 = b"Content-Length: 3\r\n\r\nSDP";
581 let response = b"SIP/2.0 100 Trying\r\nContent-Length: 0\r\n\r\n";
582
583 let mut data = make_frame(Direction::Recv, Transport::Tcp, "10.0.0.1:5060", part1);
584 data.extend_from_slice(&make_frame(
585 Direction::Sent,
586 Transport::Tcp,
587 "10.0.0.1:5060",
588 response,
589 ));
590 data.extend_from_slice(&make_frame(
591 Direction::Recv,
592 Transport::Tcp,
593 "10.0.0.1:5060",
594 part2,
595 ));
596
597 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
598 .collect::<Result<Vec<_>, _>>()
599 .unwrap();
600
601 assert_eq!(msgs.len(), 2);
602
603 let trying = &msgs[0];
605 assert_eq!(trying.direction, Direction::Sent);
606 assert_eq!(trying.content, response);
607
608 let invite = &msgs[1];
610 assert_eq!(invite.direction, Direction::Recv);
611 let mut expected_invite = Vec::new();
612 expected_invite.extend_from_slice(part1);
613 expected_invite.extend_from_slice(part2);
614 assert_eq!(invite.content, expected_invite);
615 }
616
617 #[test]
618 fn tcp_reassembly_interleaved_different_addresses() {
619 let a_part1 = b"INVITE sip:user@host SIP/2.0\r\n";
626 let a_part2 = b"Content-Length: 3\r\n\r\nSDP";
627 let b_part1 = b"NOTIFY sip:user@host SIP/2.0\r\n";
628 let b_part2 = b"Content-Length: 4\r\n\r\nBODY";
629
630 let mut data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", a_part1);
631 data.extend_from_slice(&make_frame(
632 Direction::Recv,
633 Transport::Tcp,
634 "[::2]:5060",
635 b_part1,
636 ));
637 data.extend_from_slice(&make_frame(
638 Direction::Recv,
639 Transport::Tcp,
640 "[::1]:5060",
641 a_part2,
642 ));
643 data.extend_from_slice(&make_frame(
644 Direction::Recv,
645 Transport::Tcp,
646 "[::2]:5060",
647 b_part2,
648 ));
649
650 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
651 .collect::<Result<Vec<_>, _>>()
652 .unwrap();
653
654 assert_eq!(msgs.len(), 2);
655
656 assert_eq!(msgs[0].address, "[::1]:5060");
658 assert_eq!(msgs[0].frame_count, 2);
659 let mut expected_a = Vec::new();
660 expected_a.extend_from_slice(a_part1);
661 expected_a.extend_from_slice(a_part2);
662 assert_eq!(msgs[0].content, expected_a);
663
664 assert_eq!(msgs[1].address, "[::2]:5060");
666 assert_eq!(msgs[1].frame_count, 2);
667 let mut expected_b = Vec::new();
668 expected_b.extend_from_slice(b_part1);
669 expected_b.extend_from_slice(b_part2);
670 assert_eq!(msgs[1].content, expected_b);
671 }
672
673 #[test]
674 fn direction_change_splits_messages() {
675 let recv_content = b"OPTIONS sip:user@host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
676 let sent_content = b"SIP/2.0 200 OK\r\nContent-Length: 0\r\n\r\n";
677 let mut data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", recv_content);
678 data.extend_from_slice(&make_frame(
679 Direction::Sent,
680 Transport::Tcp,
681 "[::1]:5060",
682 sent_content,
683 ));
684 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
685 .collect::<Result<Vec<_>, _>>()
686 .unwrap();
687 assert_eq!(msgs.len(), 2);
688 assert_eq!(msgs[0].direction, Direction::Recv);
689 assert_eq!(msgs[1].direction, Direction::Sent);
690 }
691
692 #[test]
693 fn address_change_splits_messages() {
694 let content = b"OPTIONS sip:user@host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
695 let mut data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", content);
696 data.extend_from_slice(&make_frame(
697 Direction::Recv,
698 Transport::Tcp,
699 "[::2]:5060",
700 content,
701 ));
702 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
703 .collect::<Result<Vec<_>, _>>()
704 .unwrap();
705 assert_eq!(msgs.len(), 2);
706 assert_eq!(msgs[0].address, "[::1]:5060");
707 assert_eq!(msgs[1].address, "[::2]:5060");
708 }
709
710 #[test]
711 fn udp_no_reassembly() {
712 let content1 = b"OPTIONS sip:a SIP/2.0\r\nContent-Length: 0\r\n\r\n";
713 let content2 = b"OPTIONS sip:b SIP/2.0\r\nContent-Length: 0\r\n\r\n";
714 let mut data = make_frame(Direction::Recv, Transport::Udp, "1.1.1.1:5060", content1);
715 data.extend_from_slice(&make_frame(
716 Direction::Recv,
717 Transport::Udp,
718 "1.1.1.1:5060",
719 content2,
720 ));
721 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
722 .collect::<Result<Vec<_>, _>>()
723 .unwrap();
724 assert_eq!(msgs.len(), 2, "UDP frames should not be reassembled");
725 assert_eq!(msgs[0].frame_count, 1);
726 assert_eq!(msgs[1].frame_count, 1);
727 }
728
729 #[test]
730 fn aggregated_messages_split_by_content_length() {
731 let msg1 = b"NOTIFY sip:a SIP/2.0\r\nContent-Length: 5\r\n\r\nhello";
732 let msg2 = b"SIP/2.0 200 OK\r\nContent-Length: 0\r\n\r\n";
733 let mut combined = Vec::new();
734 combined.extend_from_slice(msg1);
735 combined.extend_from_slice(msg2);
736 let data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", &combined);
737 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
738 .collect::<Result<Vec<_>, _>>()
739 .unwrap();
740 assert_eq!(msgs.len(), 2);
741 assert_eq!(msgs[0].content, msg1);
742 assert_eq!(msgs[1].content, msg2);
743 }
744
745 #[test]
746 fn find_content_length_standard() {
747 let data = b"NOTIFY sip:a SIP/2.0\r\nContent-Length: 42\r\n\r\n";
748 assert_eq!(content_length(data), Some(42));
749 }
750
751 #[test]
752 fn find_content_length_compact() {
753 let data = b"NOTIFY sip:a SIP/2.0\r\nl: 42\r\n\r\n";
754 assert_eq!(content_length(data), Some(42));
755 }
756
757 #[test]
758 fn find_content_length_padded_before_colon() {
759 let data = b"NOTIFY sip:a SIP/2.0\r\nContent-Length \t: 5\r\n\r\nhello";
760 assert_eq!(content_length(data), Some(5));
761 }
762
763 #[test]
764 fn padded_colon_splits_aggregated_messages() {
765 let msg1 = b"NOTIFY sip:a SIP/2.0\r\nContent-Length : 5\r\n\r\nhello";
766 let msg2 = b"SIP/2.0 200 OK\r\nContent-Length: 0\r\n\r\n";
767 let mut combined = Vec::new();
768 combined.extend_from_slice(msg1);
769 combined.extend_from_slice(msg2);
770 let data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", &combined);
771 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
772 .collect::<Result<Vec<_>, _>>()
773 .unwrap();
774 assert_eq!(msgs.len(), 2);
775 assert_eq!(msgs[0].content, msg1);
776 assert_eq!(msgs[1].content, msg2);
777 }
778
779 #[test]
780 fn find_content_length_missing() {
781 let data = b"NOTIFY sip:a SIP/2.0\r\nCSeq: 1 NOTIFY\r\n\r\n";
782 assert_eq!(content_length(data), None);
783 }
784
785 #[test]
786 fn sip_start_request() {
787 assert!(matches!(
788 sip_start(b"INVITE sip:user@host SIP/2.0\r\n"),
789 SipStart::Yes
790 ));
791 assert!(matches!(
792 sip_start(b"XYZZY sip:user@host SIP/2.0\r\n"),
793 SipStart::Yes
794 ));
795 assert!(matches!(
796 sip_start(b"ACK sip:user@host SIP/2.0\r\n"),
797 SipStart::Yes
798 ));
799 }
800
801 #[test]
802 fn sip_start_response() {
803 assert!(matches!(sip_start(b"SIP/2.0 200 OK\r\n"), SipStart::Yes));
804 assert!(matches!(
805 sip_start(b"SIP/2.0 100 Trying\r\n"),
806 SipStart::Yes
807 ));
808 }
809
810 #[test]
811 fn sip_start_not_sip() {
812 assert!(matches!(sip_start(b"some random data\r\n"), SipStart::No));
813 assert!(matches!(sip_start(b"HTTP/1.1 200 OK\r\n"), SipStart::No));
814 assert!(matches!(
815 sip_start(b"INVITE sip:user@host HTTP/1.1\r\n"),
816 SipStart::No
817 ));
818 }
819
820 #[test]
821 fn sip_start_needs_more() {
822 assert!(matches!(sip_start(b"INVI"), SipStart::NeedMore));
823 assert!(matches!(sip_start(b"SIP/2."), SipStart::NeedMore));
824 assert!(matches!(
825 sip_start(b"INVITE sip:user@host SIP/2.0\r"),
826 SipStart::NeedMore
827 ));
828 }
829
830 #[test]
831 fn find_sip_start_at_beginning() {
832 let data = b"INVITE sip:user@host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
833 assert_eq!(find_sip_start(data), Some(0));
834 }
835
836 #[test]
837 fn find_sip_start_after_prefix() {
838 let data = b"</xml>\r\nNOTIFY sip:user@host SIP/2.0\r\n";
839 assert_eq!(find_sip_start(data), Some(8));
840 }
841
842 #[test]
843 fn find_sip_start_none() {
844 let data = b"no SIP here\r\nat all\r\n";
845 assert_eq!(find_sip_start(data), None);
846 }
847
848 #[test]
849 fn message_preserves_metadata() {
850 let content = b"OPTIONS sip:user@host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
851 let data = make_frame(
852 Direction::Sent,
853 Transport::Tls,
854 "[2001:db8::1]:5061",
855 content,
856 );
857 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
858 .collect::<Result<Vec<_>, _>>()
859 .unwrap();
860 assert_eq!(msgs.len(), 1);
861 assert_eq!(msgs[0].direction, Direction::Sent);
862 assert_eq!(msgs[0].transport, Transport::Tls);
863 assert_eq!(msgs[0].address, "[2001:db8::1]:5061");
864 assert_eq!(
865 msgs[0].timestamp,
866 Timestamp::TimeOnly {
867 hour: 0,
868 min: 0,
869 sec: 0,
870 usec: 0
871 }
872 );
873 }
874
875 #[test]
876 fn extract_handles_crlf_between_messages() {
877 let msg1 = b"NOTIFY sip:a SIP/2.0\r\nContent-Length: 5\r\n\r\nhello";
878 let msg2 = b"SIP/2.0 200 OK\r\nContent-Length: 0\r\n\r\n";
879 let mut content = Vec::new();
880 content.extend_from_slice(msg1);
881 content.extend_from_slice(b"\r\n");
882 content.extend_from_slice(msg2);
883
884 let mut buf = buffer_with(content);
885 let msgs = extracted(&mut buf);
886 assert_eq!(msgs.len(), 2);
887 assert_eq!(msgs[0].content, msg1);
888 assert_eq!(msgs[1].content, msg2);
889 }
890
891 #[test]
892 fn extract_skips_non_sip_prefix() {
893 let prefix = b"</conference-info>\r\n";
894 let msg = b"NOTIFY sip:a SIP/2.0\r\nContent-Length: 0\r\n\r\n";
895 let mut content = Vec::new();
896 content.extend_from_slice(prefix);
897 content.extend_from_slice(msg);
898
899 let mut buf = buffer_with(content);
900 let msgs = extracted(&mut buf);
901 assert_eq!(msgs.len(), 1);
902 assert_eq!(msgs[0].content, msg);
903 }
904
905 #[test]
906 fn non_sip_prefix_counted_in_stats() {
907 let prefix = b"</conference-info>\r\n";
908 let msg = b"NOTIFY sip:a SIP/2.0\r\nContent-Length: 0\r\n\r\n";
909 let mut content = Vec::new();
910 content.extend_from_slice(prefix);
911 content.extend_from_slice(msg);
912 let data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", &content);
913
914 let mut iter = MessageIterator::new(&data[..]);
915 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
916 assert_eq!(msgs.len(), 1);
917 let stats = iter.parse_stats();
918 assert_eq!(stats.non_sip_prefixes(), 1);
919 assert_eq!(stats.non_sip_prefix_bytes(), prefix.len() as u64);
920 }
921
922 #[test]
923 fn extract_resyncs_on_extension_method() {
924 let prefix = b"</conference-info>\r\n";
925 let msg = b"XYZZY sip:a SIP/2.0\r\nContent-Length: 0\r\n\r\n";
926 let mut content = Vec::new();
927 content.extend_from_slice(prefix);
928 content.extend_from_slice(msg);
929 let data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", &content);
930
931 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
932 .collect::<Result<Vec<_>, _>>()
933 .unwrap();
934 assert_eq!(msgs.len(), 1);
935 assert_eq!(msgs[0].content, msg);
936 }
937
938 #[test]
939 fn extract_waits_for_half_received_request_line() {
940 let mut content = Vec::new();
941 content.extend_from_slice(b"</conference-info>\r\n");
942 content.extend_from_slice(b"XYZZY sip:a SI");
943
944 let mut buf = buffer_with(content);
945 let msgs = extracted(&mut buf);
946 assert!(msgs.is_empty());
947 assert_eq!(
948 buf.content, b"XYZZY sip:a SI",
949 "an unterminated request line waits instead of being scanned past"
950 );
951 }
952
953 #[test]
954 fn extract_waits_for_incomplete_body() {
955 let content = b"INVITE sip:a SIP/2.0\r\nContent-Length: 100\r\n\r\npartial".to_vec();
957
958 let mut buf = buffer_with(content);
959 let msgs = extracted(&mut buf);
960 assert!(msgs.is_empty(), "should wait for body to complete");
961 assert!(!buf.content.is_empty(), "buffer should retain data");
962 }
963
964 #[test]
965 fn extract_waits_for_incomplete_headers() {
966 let content = b"INVITE sip:a SIP/2.0\r\nContent-Length: 0\r\n".to_vec();
968
969 let mut buf = buffer_with(content);
970 let msgs = extracted(&mut buf);
971 assert!(msgs.is_empty(), "should wait for headers to complete");
972 }
973
974 #[test]
975 fn tcp_body_split_across_five_frames() {
976 let body_len: usize = 6424;
979 let body: Vec<u8> = (0..body_len).map(|i| b'A' + (i % 26) as u8).collect();
980
981 let mut headers = Vec::new();
982 headers.extend_from_slice(b"NOTIFY sip:user@host SIP/2.0\r\n");
983 headers
984 .extend_from_slice(b"Via: SIP/2.0/TCP [2001:4958:10:11::6]:45538;branch=z9hG4bK-1\r\n");
985 headers.extend_from_slice(b"Call-ID: fragmented-notify@host\r\n");
986 headers.extend_from_slice(b"CSeq: 1 NOTIFY\r\n");
987 headers.extend_from_slice(
988 b"Content-Type: application/emergencyCallData.AbandonedCall+json\r\n",
989 );
990 headers.extend_from_slice(format!("Content-Length: {body_len}\r\n").as_bytes());
991 headers.extend_from_slice(b"\r\n");
992
993 let mut full_content = headers.clone();
994 full_content.extend_from_slice(&body);
995
996 let frame1_len = 1500.min(full_content.len());
998 let remaining = &full_content[frame1_len..];
999 let frame2_len = 1428.min(remaining.len());
1000 let remaining = &remaining[frame2_len..];
1001 let frame3_len = 1428.min(remaining.len());
1002 let remaining = &remaining[frame3_len..];
1003 let frame4_len = 1428.min(remaining.len());
1004 let remaining = &remaining[frame4_len..];
1005 let frame5_len = remaining.len();
1006
1007 let addr = "[2001:4958:10:11::6]:45538";
1008 let mut data = make_frame(
1009 Direction::Recv,
1010 Transport::Tcp,
1011 addr,
1012 &full_content[..frame1_len],
1013 );
1014 let mut offset = frame1_len;
1015 for len in [frame2_len, frame3_len, frame4_len, frame5_len] {
1016 data.extend_from_slice(&make_frame(
1017 Direction::Recv,
1018 Transport::Tcp,
1019 addr,
1020 &full_content[offset..offset + len],
1021 ));
1022 offset += len;
1023 }
1024 assert_eq!(offset, full_content.len());
1025
1026 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
1027 .collect::<Result<Vec<_>, _>>()
1028 .unwrap();
1029
1030 assert_eq!(
1031 msgs.len(),
1032 1,
1033 "should produce exactly one reassembled message"
1034 );
1035 assert_eq!(msgs[0].frame_count, 5, "should track all 5 frames");
1036 assert_eq!(
1037 msgs[0].content, full_content,
1038 "content should be fully reassembled"
1039 );
1040 assert_eq!(msgs[0].direction, Direction::Recv);
1041 assert_eq!(msgs[0].address, addr);
1042 }
1043
1044 #[test]
1045 fn parse_stats_delegates() {
1046 let content = b"OPTIONS sip:user@host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1047 let data = make_frame(Direction::Recv, Transport::Udp, "1.1.1.1:5060", content);
1048 let mut iter = MessageIterator::new(&data[..]);
1049 let msgs: Vec<_> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1050 assert_eq!(msgs.len(), 1);
1051 let stats = iter.parse_stats();
1052 assert_eq!(stats.bytes_read, data.len() as u64);
1053 assert_eq!(stats.bytes_skipped, 0);
1054 }
1055
1056 #[test]
1057 fn tcp_body_split_parsed_message() {
1058 let body_len: usize = 6424;
1060 let body: Vec<u8> = (0..body_len).map(|i| b'A' + (i % 26) as u8).collect();
1061
1062 let mut headers = Vec::new();
1063 headers.extend_from_slice(b"NOTIFY sip:user@host SIP/2.0\r\n");
1064 headers.extend_from_slice(b"Call-ID: fragmented-parsed@host\r\n");
1065 headers.extend_from_slice(b"CSeq: 1 NOTIFY\r\n");
1066 headers.extend_from_slice(
1067 b"Content-Type: application/emergencyCallData.AbandonedCall+json\r\n",
1068 );
1069 headers.extend_from_slice(format!("Content-Length: {body_len}\r\n").as_bytes());
1070 headers.extend_from_slice(b"\r\n");
1071
1072 let mut full_content = headers.clone();
1073 full_content.extend_from_slice(&body);
1074
1075 let split1 = 1500.min(full_content.len());
1077 let split2 = (split1 + 3000).min(full_content.len());
1078
1079 let addr = "[2001:db8::1]:5060";
1080 let mut data = make_frame(
1081 Direction::Recv,
1082 Transport::Tcp,
1083 addr,
1084 &full_content[..split1],
1085 );
1086 data.extend_from_slice(&make_frame(
1087 Direction::Recv,
1088 Transport::Tcp,
1089 addr,
1090 &full_content[split1..split2],
1091 ));
1092 data.extend_from_slice(&make_frame(
1093 Direction::Recv,
1094 Transport::Tcp,
1095 addr,
1096 &full_content[split2..],
1097 ));
1098
1099 let parsed: Vec<crate::types::ParsedSipMessage> =
1100 crate::sip::ParsedMessageIterator::new(&data[..])
1101 .collect::<Result<Vec<_>, _>>()
1102 .unwrap();
1103
1104 assert_eq!(parsed.len(), 1, "should produce one parsed message");
1105 assert_eq!(parsed[0].content_length(), Some(body_len));
1106 assert_eq!(parsed[0].body.len(), body_len, "body should be complete");
1107 assert_eq!(parsed[0].body, body, "body content should match");
1108 assert_eq!(parsed[0].frame_count, 3);
1109 assert_eq!(parsed[0].method(), Some("NOTIFY"));
1110 }
1111
1112 #[test]
1113 fn tls_keepalive_single_lf_drained() {
1114 let data = make_frame(Direction::Recv, Transport::Tls, "[10.0.0.1]:5061", b"\n");
1115 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
1116 .collect::<Result<Vec<_>, _>>()
1117 .unwrap();
1118 assert_eq!(msgs.len(), 0, "keep-alive \\n should produce no messages");
1119 }
1120
1121 #[test]
1122 fn tls_keepalive_multiple_lf_drained() {
1123 let addr = "[10.0.0.1]:5061";
1124 let mut data = make_frame(Direction::Recv, Transport::Tls, addr, b"\n");
1125 data.extend_from_slice(&make_frame(Direction::Recv, Transport::Tls, addr, b"\n"));
1126 data.extend_from_slice(&make_frame(Direction::Recv, Transport::Tls, addr, b"\n"));
1127 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
1128 .collect::<Result<Vec<_>, _>>()
1129 .unwrap();
1130 assert_eq!(
1131 msgs.len(),
1132 0,
1133 "multiple keep-alive \\n should produce no messages"
1134 );
1135 }
1136
1137 #[test]
1138 fn tls_keepalive_interleaved_with_sip() {
1139 let addr = "[10.0.0.1]:5061";
1140 let sip = b"OPTIONS sip:host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1141 let mut data = make_frame(Direction::Recv, Transport::Tls, addr, b"\n");
1142 data.extend_from_slice(&make_frame(Direction::Recv, Transport::Tls, addr, sip));
1143 data.extend_from_slice(&make_frame(Direction::Recv, Transport::Tls, addr, b"\n"));
1144 let mut iter = MessageIterator::new(&data[..]);
1145 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1146 assert_eq!(msgs.len(), 1, "only the SIP message should be emitted");
1147 assert_eq!(msgs[0].content, sip);
1148 assert_eq!(
1149 msgs[0].frame_count, 1,
1150 "the drained keep-alive frame must not count toward the message"
1151 );
1152 assert!(
1153 iter.buffers.is_empty(),
1154 "a buffer drained to nothing must not be retained"
1155 );
1156 }
1157
1158 #[test]
1159 fn tls_bare_lf_before_sip_start() {
1160 let addr = "[10.0.0.1]:5061";
1161 let sip_part1 = b"\nOPTIONS sip:host SIP/2.0\r\n";
1162 let sip_part2 = b"Content-Length: 0\r\n\r\n";
1163 let mut data = make_frame(Direction::Recv, Transport::Tls, addr, sip_part1);
1164 data.extend_from_slice(&make_frame(
1165 Direction::Recv,
1166 Transport::Tls,
1167 addr,
1168 sip_part2,
1169 ));
1170 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
1171 .collect::<Result<Vec<_>, _>>()
1172 .unwrap();
1173 assert_eq!(msgs.len(), 1, "SIP after bare LF should be extracted");
1174 assert!(
1175 msgs[0].content.starts_with(b"OPTIONS"),
1176 "message should start with SIP method, not \\n"
1177 );
1178 }
1179
1180 fn make_frame_at(
1181 direction: Direction,
1182 transport: Transport,
1183 addr: &str,
1184 content: &[u8],
1185 timestamp: &str,
1186 ) -> Vec<u8> {
1187 let dir_str = match direction {
1188 Direction::Recv => "recv",
1189 Direction::Sent => "sent",
1190 };
1191 let prep = match direction {
1192 Direction::Recv => "from",
1193 Direction::Sent => "to",
1194 };
1195 let transport_str = match transport {
1196 Transport::Tcp => "tcp",
1197 Transport::Udp => "udp",
1198 Transport::Tls => "tls",
1199 Transport::Wss => "wss",
1200 };
1201 let header = format!(
1202 "{dir_str} {} bytes {prep} {transport_str}/{addr} at {timestamp}:\n",
1203 content.len()
1204 );
1205 let mut data = header.into_bytes();
1206 data.extend_from_slice(content);
1207 data.extend_from_slice(b"\x0B\n");
1208 data
1209 }
1210
1211 #[test]
1212 fn empty_buffer_removed_after_complete_message() {
1213 let sip = b"OPTIONS sip:host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1214 let addr_a = "[::1]:5060";
1215 let addr_b = "[::2]:5060";
1216
1217 let mut data = make_frame(Direction::Recv, Transport::Tcp, addr_a, sip);
1218 data.extend_from_slice(&make_frame(Direction::Recv, Transport::Tcp, addr_b, sip));
1219
1220 let mut iter = MessageIterator::new(&data[..]);
1221 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1222 assert_eq!(msgs.len(), 2);
1223 assert!(
1224 iter.buffers.is_empty(),
1225 "all buffers should be removed after complete messages are extracted"
1226 );
1227 }
1228
1229 #[test]
1230 fn stale_buffer_evicted_after_timeout() {
1231 let sip = b"OPTIONS sip:host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1232 let partial = b"INVITE sip:host SIP/2.0\r\n";
1233
1234 let mut data = Vec::new();
1235 data.extend_from_slice(&make_frame_at(
1237 Direction::Recv,
1238 Transport::Tls,
1239 "[::99]:44444",
1240 partial,
1241 "2026-02-16 10:00:00.000000",
1242 ));
1243 data.extend_from_slice(&make_frame_at(
1245 Direction::Recv,
1246 Transport::Tls,
1247 "[::1]:5060",
1248 sip,
1249 "2026-02-16 12:00:01.000000",
1250 ));
1251
1252 let mut iter = MessageIterator::new(&data[..]);
1253 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1254
1255 assert_eq!(msgs.len(), 1, "should produce the complete message");
1256 assert_eq!(msgs[0].address, "[::1]:5060");
1257 assert!(
1258 iter.buffers.is_empty(),
1259 "stale buffer for [::99]:44444 should have been evicted"
1260 );
1261 let stats = iter.parse_stats();
1262 assert_eq!(stats.stale_evictions(), 1);
1263 assert_eq!(stats.stale_evicted_bytes(), partial.len() as u64);
1264 }
1265
1266 #[test]
1267 fn stale_buffer_evicted_across_dates() {
1268 let sip = b"OPTIONS sip:host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1269 let partial = b"INVITE sip:host SIP/2.0\r\n";
1270
1271 let mut data = make_frame_at(
1272 Direction::Recv,
1273 Transport::Tls,
1274 "[::99]:44444",
1275 partial,
1276 "2026-02-16 10:00:00.000000",
1277 );
1278 data.extend_from_slice(&make_frame_at(
1279 Direction::Recv,
1280 Transport::Tls,
1281 "[::1]:5060",
1282 sip,
1283 "2026-02-19 10:00:01.000000",
1284 ));
1285
1286 let mut iter = MessageIterator::new(&data[..]);
1287 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1288
1289 assert_eq!(msgs.len(), 1, "the three-day-old buffer must be evicted");
1290 assert_eq!(msgs[0].address, "[::1]:5060");
1291 assert!(iter.buffers.is_empty());
1292 }
1293
1294 #[test]
1295 fn timestamp_format_change_evicts_nothing() {
1296 let sip = b"OPTIONS sip:host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1297 let partial = b"INVITE sip:host SIP/2.0\r\n";
1298
1299 let mut data = make_frame_at(
1300 Direction::Recv,
1301 Transport::Tls,
1302 "[::99]:44444",
1303 partial,
1304 "10:00:00.000000",
1305 );
1306 data.extend_from_slice(&make_frame_at(
1307 Direction::Recv,
1308 Transport::Tls,
1309 "[::1]:5060",
1310 sip,
1311 "2026-02-16 12:00:01.000000",
1312 ));
1313
1314 let msgs: Vec<SipMessage> = MessageIterator::new(&data[..])
1315 .collect::<Result<Vec<_>, _>>()
1316 .unwrap();
1317
1318 assert_eq!(
1319 msgs.len(),
1320 2,
1321 "the two timestamp formats share no epoch, so nothing is stale"
1322 );
1323 assert_eq!(msgs[1].address, "[::99]:44444");
1324 assert_eq!(msgs[1].content, partial);
1325 }
1326
1327 #[test]
1328 fn day_rollover_detection_with_time_only() {
1329 let sip = b"OPTIONS sip:host SIP/2.0\r\nContent-Length: 0\r\n\r\n";
1330 let partial = b"INVITE sip:host SIP/2.0\r\n";
1331
1332 let mut data = Vec::new();
1333 data.extend_from_slice(&make_frame_at(
1335 Direction::Recv,
1336 Transport::Tcp,
1337 "[::99]:44444",
1338 partial,
1339 "23:59:00.000000",
1340 ));
1341 data.extend_from_slice(&make_frame_at(
1343 Direction::Recv,
1344 Transport::Tcp,
1345 "[::1]:5060",
1346 sip,
1347 "02:00:01.000000",
1348 ));
1349
1350 let mut iter = MessageIterator::new(&data[..]);
1351 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1352
1353 assert_eq!(msgs.len(), 1);
1354 assert_eq!(msgs[0].address, "[::1]:5060");
1355 assert_eq!(
1356 iter.clock.now(),
1357 86400 + 2 * 3600 + 1,
1358 "should have detected one day rollover"
1359 );
1360 assert!(
1361 iter.buffers.is_empty(),
1362 "stale buffer should have been evicted after day rollover"
1363 );
1364 }
1365
1366 #[test]
1367 fn flush_all_clears_buffers() {
1368 let partial = b"INVITE sip:host SIP/2.0\r\n";
1369 let data = make_frame(Direction::Recv, Transport::Tcp, "[::1]:5060", partial);
1370
1371 let mut iter = MessageIterator::new(&data[..]);
1372 let msgs: Vec<SipMessage> = iter.by_ref().collect::<Result<Vec<_>, _>>().unwrap();
1373 assert_eq!(msgs.len(), 1, "partial should be flushed at EOF");
1374 assert!(
1375 iter.buffers.is_empty(),
1376 "flush_all should clear the HashMap"
1377 );
1378 }
1379}