Skip to main content

freeswitch_sofia_trace_parser/
message.rs

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
13/// Level 2 streaming parser: reassembles TCP segments into complete SIP messages.
14///
15/// Wraps a [`FrameIterator`] and groups TCP frames by `(Direction, Address)`.
16/// Messages are emitted when headers and Content-Length body bytes are fully
17/// available. UDP frames pass through as-is (1:1 mapping).
18///
19/// Stale TCP connection buffers are evicted after 2 hours of inactivity to
20/// maintain constant memory on multi-day dump streams.
21///
22/// # Example
23///
24/// ```no_run
25/// use std::fs::File;
26/// use freeswitch_sofia_trace_parser::MessageIterator;
27///
28/// let file = File::open("profile.dump").unwrap();
29/// for msg in MessageIterator::new(file) {
30///     let msg = msg.unwrap();
31///     println!("{} {} {} ({} frames)",
32///         msg.timestamp, msg.direction, msg.address, msg.frame_count);
33/// }
34/// ```
35pub 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    /// Create a new message iterator reading from the given source.
59    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    /// Enable capturing of skipped bytes in the underlying [`FrameIterator`];
70    /// `false` selects [`SkipTracking::CountOnly`]. Whichever of this and
71    /// [`skip_tracking`](Self::skip_tracking) is called last wins.
72    pub fn capture_skipped(mut self, enable: bool) -> Self {
73        self.frames = self.frames.capture_skipped(enable);
74        self
75    }
76
77    /// Set the level of detail for unparsed region tracking.
78    pub fn skip_tracking(mut self, tracking: SkipTracking) -> Self {
79        self.frames = self.frames.skip_tracking(tracking);
80        self
81    }
82
83    /// Borrow the accumulated parse statistics from the underlying frame parser.
84    pub fn parse_stats(&self) -> &ParseStats {
85        self.frames.stats()
86    }
87
88    /// Take all accumulated unparsed regions, leaving the list empty.
89    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
220/// Where the buffer stands relative to the next SIP message start.
221enum Resync {
222    Ready,
223    Retry,
224    Wait,
225}
226
227/// Move every message the buffer already holds in full into `ready`.
228fn 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
248/// Drop whatever precedes the next SIP message start in the buffer.
249fn 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, // Start line incomplete, wait for more data
253        SipStart::No => {}
254    }
255
256    // Drain leading whitespace (CRLF padding, bare LF keep-alives, etc.)
257    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, // No SIP start found, wait for more data
300    }
301}
302
303/// Split one complete message off the front of the buffer. False means the
304/// message is still arriving and the buffer keeps its bytes.
305fn 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, // Headers incomplete, wait for more data
313    };
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; // Body incomplete, wait for more data
321            }
322            end
323        }
324        None => body_start, // No CL = no body (RFC 3261 Section 18.3)
325    };
326
327    let remaining = buf.content.split_off(msg_end);
328    let msg_content = std::mem::replace(&mut buf.content, remaining);
329
330    // Skip trailing CRLF between messages
331    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
366/// Find Content-Length header value in the header block ending at `header_end`.
367fn 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; // skip \r\n
383    }
384    None
385}
386
387/// RFC 3261 HCOLON allows SP/HTAB between the header name and the colon, and
388/// `sip_header` accepts it, so Level 2 framing must see the same header.
389fn 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
431/// Scan for the first SIP message start at a CRLF boundary within data. A
432/// candidate whose start line is still arriving stops the scan there, so the
433/// bytes before it are dropped and the line itself is waited for.
434fn 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        // Frame 1: recv from A (partial INVITE)
533        // Frame 2: sent to A (response on same connection — interrupts)
534        // Frame 3: recv from A (rest of INVITE)
535        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        // 100 Trying completes first (single frame)
560        let trying = &msgs[0];
561        assert_eq!(trying.direction, Direction::Sent);
562        assert_eq!(trying.content, response);
563
564        // INVITE completes when frame 3 arrives (reassembled from frames 1+3)
565        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        // Two different addresses both sending multi-frame messages,
576        // frames arriving interleaved:
577        //   Frame 1: recv from A (partial INVITE)
578        //   Frame 2: recv from B (partial NOTIFY)
579        //   Frame 3: recv from A (rest of INVITE — completes A)
580        //   Frame 4: recv from B (rest of NOTIFY — completes B)
581        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        // INVITE from [::1] completes first (frame 3 arrives before frame 4)
613        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        // NOTIFY from [::2] completes second
621        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        // Headers complete but body is missing
895        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        // Headers not complete (no \r\n\r\n)
906        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        // Simulate REQUEST.md scenario: NOTIFY with Content-Length: 6424
916        // body split across 5 TCP frames (like NG9-1-1 abandoned call JSON)
917        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        // Split into 5 frames like real TCP segments
936        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        // Same scenario but verified through Level 3 (ParsedSipMessage)
998        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        // Split into 3 frames
1015        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        // Partial frame from stale addr at 10:00:00
1175        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        // Complete msg from another addr at 12:00:01 (>2h later, triggers sweep)
1183        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        // Partial frame at 23:59:00
1270        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        // Complete msg at 02:00:01 (next day — rollover detected, >2h from 23:59)
1278        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}