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