Skip to main content

rtc_interceptor/jitterbuffer/
buffer.rs

1//! The per-stream packet store: ordered, deduplicating and bounded.
2//!
3//! This is the data structure only. Deciding *when* a packet is due — the playout policy — is the
4//! interceptor's job, and lives above this.
5
6use super::sequence::SequenceExtender;
7use crate::{Packet, TaggedPacket};
8use std::collections::BTreeMap;
9
10/// Counters describing what a buffer has had to cope with.
11///
12/// Upstream reports these through a listener callback. There is no natural sans-I/O analogue for
13/// a callback, and these are diagnostics rather than control signals, so they are plain counters
14/// on the buffer instead.
15#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
16pub struct JitterBufferStats {
17    /// Packets that arrived after a higher sequence number had already been seen.
18    pub out_of_order: u64,
19    /// Packets discarded because the same sequence number was already held.
20    pub duplicates: u64,
21    /// Packets dropped because the buffer was full.
22    pub overflow: u64,
23    /// Packets discarded because they arrived after the buffer had already released that
24    /// position — too late to be played out in order.
25    pub late: u64,
26    /// Packets rejected because they carried a different SSRC than this buffer's stream.
27    pub foreign_ssrc: u64,
28    /// Times playout asked for a packet and the buffer had none to give.
29    pub underflow: u64,
30}
31
32/// Whether the buffer is still filling or is handing packets out.
33///
34/// The *trigger* for the transition is not here: upstream moves to `Emitting` once
35/// `minStartCount` packets have accumulated, whereas the policy this port is built for is
36/// time-based (a depth in milliseconds). So the buffer owns the state and the playout policy above
37/// it decides when to call [`JitterBuffer::begin_emitting`].
38#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
39pub enum State {
40    /// Filling; playout has not started, so nothing is handed out yet.
41    #[default]
42    Buffering,
43    /// Handing packets out as they come due.
44    Emitting,
45}
46
47/// Why a packet was not stored.
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub enum Rejected {
50    /// The same sequence number is already buffered.
51    Duplicate,
52    /// The buffer is full and this packet did not displace anything.
53    Overflow,
54    /// The position has already been released; emitting it now would be out of order.
55    Late,
56    /// The packet belongs to a different stream.
57    ForeignSsrc,
58}
59
60/// A single stream's packets, ordered by extended sequence number.
61///
62/// **One buffer per SSRC.** The buffer records the SSRC of the first packet it accepts and
63/// rejects any other, so two streams cannot interleave into one sequence-number ordering. That is
64/// a deliberate correction: pion's `ReceiverInterceptor` holds a single buffer for every remote
65/// stream and its `BindRemoteStream` ignores `info.SSRC`, so two streams' sequence numbers sort
66/// against each other.
67///
68/// Ordering is by *extended* sequence number — the 16-bit value plus its wrap count — so a packet
69/// that arrives
70/// after a wrap still sorts into its true position; and because the map is keyed by that value,
71/// duplicates collapse rather than being inserted twice as they are upstream.
72pub struct JitterBuffer {
73    /// The stream this buffer belongs to; `None` until the first packet arrives.
74    ssrc: Option<u32>,
75    /// Packets keyed by extended sequence number — ordering and deduplication in one structure.
76    packets: BTreeMap<u64, TaggedPacket>,
77    extender: SequenceExtender,
78    /// Highest extended sequence number already released; nothing at or below it may be stored.
79    released_through: Option<u64>,
80    /// Maximum packets held before the oldest is dropped.
81    capacity: usize,
82    state: State,
83    stats: JitterBufferStats,
84}
85
86/// Summarises the buffer rather than dumping its contents: `TaggedPacket` is not `Debug`, and a
87/// list of buffered packets is not what anyone wants from a debug print anyway.
88impl std::fmt::Debug for JitterBuffer {
89    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90        f.debug_struct("JitterBuffer")
91            .field("ssrc", &self.ssrc)
92            .field("held", &self.packets.len())
93            .field("capacity", &self.capacity)
94            .field("next_due", &self.front_sequence())
95            .field("released_through", &self.released_through)
96            .field("state", &self.state)
97            .field("stats", &self.stats)
98            .finish()
99    }
100}
101
102impl JitterBuffer {
103    /// Create an empty buffer holding at most `capacity` packets.
104    pub fn new(capacity: usize) -> Self {
105        Self {
106            ssrc: None,
107            packets: BTreeMap::new(),
108            extender: SequenceExtender::new(),
109            released_through: None,
110            capacity: capacity.max(1),
111            state: State::Buffering,
112            stats: JitterBufferStats::default(),
113        }
114    }
115
116    /// The SSRC this buffer is bound to, once a packet has established it.
117    pub fn ssrc(&self) -> Option<u32> {
118        self.ssrc
119    }
120
121    /// How many packets are currently held.
122    pub fn len(&self) -> usize {
123        self.packets.len()
124    }
125
126    /// Whether the buffer holds nothing.
127    pub fn is_empty(&self) -> bool {
128        self.packets.is_empty()
129    }
130
131    /// Counters describing what this buffer has coped with.
132    pub fn stats(&self) -> JitterBufferStats {
133        self.stats
134    }
135
136    /// Store a packet, returning its extended sequence number, or why it was not stored.
137    ///
138    /// A non-RTP packet is rejected as foreign: this buffer orders by sequence number, and RTCP
139    /// has none.
140    pub fn push(&mut self, packet: TaggedPacket) -> Result<u64, Rejected> {
141        let Packet::Rtp(rtp) = &packet.message else {
142            self.stats.foreign_ssrc += 1;
143            return Err(Rejected::ForeignSsrc);
144        };
145
146        let ssrc = rtp.header.ssrc;
147        match self.ssrc {
148            None => self.ssrc = Some(ssrc),
149            Some(bound) if bound != ssrc => {
150                self.stats.foreign_ssrc += 1;
151                return Err(Rejected::ForeignSsrc);
152            }
153            Some(_) => {}
154        }
155
156        let sequence_number = rtp.header.sequence_number;
157        let previous_highest = self.extender.highest();
158        let extended = self.extender.extend(sequence_number);
159
160        if previous_highest.is_some_and(|highest| extended < highest) {
161            self.stats.out_of_order += 1;
162        }
163
164        // Already played out: storing it would either emit out of order or emit twice.
165        if self
166            .released_through
167            .is_some_and(|released| extended <= released)
168        {
169            self.stats.late += 1;
170            return Err(Rejected::Late);
171        }
172
173        if self.packets.contains_key(&extended) {
174            self.stats.duplicates += 1;
175            return Err(Rejected::Duplicate);
176        }
177
178        if self.packets.len() >= self.capacity {
179            // Full: the oldest packet is the one closest to being played out, so dropping the
180            // *new* packet when it is even older keeps the buffer's contents contiguous.
181            let oldest = self
182                .packets
183                .keys()
184                .next()
185                .copied()
186                .expect("capacity is at least 1, so a full buffer is non-empty");
187            if extended < oldest {
188                self.stats.overflow += 1;
189                return Err(Rejected::Overflow);
190            }
191            self.packets.remove(&oldest);
192            self.stats.overflow += 1;
193        }
194
195        self.packets.insert(extended, packet);
196        Ok(extended)
197    }
198
199    /// The extended sequence number of the packet nearest playout, if any.
200    pub fn front_sequence(&self) -> Option<u64> {
201        self.packets.keys().next().copied()
202    }
203
204    /// Look at the packet nearest playout without removing it.
205    pub fn peek(&self) -> Option<&TaggedPacket> {
206        self.packets.values().next()
207    }
208
209    /// Whether the buffer is filling or emitting.
210    pub fn state(&self) -> State {
211        self.state
212    }
213
214    /// Start handing packets out.
215    ///
216    /// Called by the playout policy when its start condition is met — a buffered depth here,
217    /// where upstream uses a packet count.
218    pub fn begin_emitting(&mut self) {
219        self.state = State::Emitting;
220    }
221
222    /// Stop handing packets out and fill again.
223    ///
224    /// A stream that has run dry, or restarted across a discontinuity, has to re-accumulate
225    /// before playout is smooth again.
226    pub fn begin_buffering(&mut self) {
227        self.state = State::Buffering;
228    }
229
230    /// Remove and return the packet nearest playout.
231    ///
232    /// Yields nothing while [`State::Buffering`] — upstream returns `ErrPopWhileBuffering` here,
233    /// a sentinel its synchronous reader needs; in sans-I/O "nothing yet" is just `None`.
234    ///
235    /// Releasing a packet also marks its position played out, so a later-arriving copy or an
236    /// even older straggler is rejected rather than emitted behind it.
237    ///
238    /// **Gaps are skipped, not waited for.** Upstream pops strictly at its playout head and
239    /// counts an underflow when that exact sequence number is missing, which stalls a stream on
240    /// any un-recovered loss. Whether a gap is worth waiting for is a question about deadlines,
241    /// so it belongs to the playout policy above this: it decides when a packet is due, and this
242    /// hands over whatever is due now.
243    pub fn pop(&mut self) -> Option<TaggedPacket> {
244        if self.state == State::Buffering {
245            return None;
246        }
247        let Some((&extended, _)) = self.packets.iter().next() else {
248            self.stats.underflow += 1;
249            return None;
250        };
251        let packet = self.packets.remove(&extended)?;
252        self.mark_released(extended);
253        Some(packet)
254    }
255
256    /// Remove and return the packet with this wire sequence number, if it is held.
257    ///
258    /// Takes the 16-bit number off the wire, as a caller naturally holds; the extension to the
259    /// internal ordering key happens here.
260    pub fn pop_at_sequence(&mut self, sequence_number: u16) -> Option<TaggedPacket> {
261        self.pop_at(self.extended_of(sequence_number)?)
262    }
263
264    /// Borrow the packet with this wire sequence number without removing it.
265    pub fn peek_at_sequence(&self, sequence_number: u16) -> Option<&TaggedPacket> {
266        self.find(self.extended_of(sequence_number)?)
267    }
268
269    /// Remove and return the first held packet carrying this RTP timestamp.
270    ///
271    /// One packet, not a whole frame: a video frame spans several packets sharing a timestamp, so
272    /// releasing the frame means calling this until it yields `None`. That matches upstream's
273    /// `PopAtTimestamp`, and keeps the "how much of a frame is releasable" decision in the
274    /// playout policy where the deadline lives.
275    pub fn pop_at_timestamp(&mut self, timestamp: u32) -> Option<TaggedPacket> {
276        let extended =
277            self.packets
278                .iter()
279                .find_map(|(&extended, packet)| match &packet.message {
280                    Packet::Rtp(rtp) if rtp.header.timestamp == timestamp => Some(extended),
281                    _ => None,
282                })?;
283        self.pop_at(extended)
284    }
285
286    /// Remove and return the packet at `extended`, if it is held.
287    ///
288    /// The extended key is the buffer's own ordering space; [`pop_at_sequence`](Self::pop_at_sequence)
289    /// is the one to reach for from outside.
290    pub fn pop_at(&mut self, extended: u64) -> Option<TaggedPacket> {
291        let packet = self.packets.remove(&extended)?;
292        self.mark_released(extended);
293        Some(packet)
294    }
295
296    /// Borrow the packet at `extended` without removing it.
297    pub fn find(&self, extended: u64) -> Option<&TaggedPacket> {
298        self.packets.get(&extended)
299    }
300
301    /// The extended key of a held packet with this wire sequence number.
302    ///
303    /// Searches rather than extending arithmetically: extending would need to mutate the
304    /// extender's anchor, and a lookup must not move the ordering origin. At most one held packet
305    /// can carry a given wire number in a buffer this shallow, so the first match is the packet.
306    fn extended_of(&self, sequence_number: u16) -> Option<u64> {
307        self.packets
308            .iter()
309            .find_map(|(&extended, packet)| match &packet.message {
310                Packet::Rtp(rtp) if rtp.header.sequence_number == sequence_number => Some(extended),
311                _ => None,
312            })
313    }
314
315    /// The RTP timestamp of the packet nearest playout.
316    pub fn front_timestamp(&self) -> Option<u32> {
317        match &self.peek()?.message {
318            Packet::Rtp(rtp) => Some(rtp.header.timestamp),
319            _ => None,
320        }
321    }
322
323    /// Drop everything, keeping the SSRC binding and stats.
324    ///
325    /// Used when a stream restarts: the ordering anchors are meaningless across a discontinuity,
326    /// but what the buffer has coped with so far is still worth reporting.
327    pub fn reset(&mut self) {
328        self.packets.clear();
329        self.extender = SequenceExtender::new();
330        self.released_through = None;
331        self.state = State::Buffering;
332    }
333
334    fn mark_released(&mut self, extended: u64) {
335        self.released_through = Some(match self.released_through {
336            Some(previous) => previous.max(extended),
337            None => extended,
338        });
339    }
340}
341
342#[cfg(test)]
343mod tests {
344    use super::*;
345    use shared::TransportContext;
346    use std::time::Instant;
347
348    fn packet(ssrc: u32, sequence_number: u16, timestamp: u32) -> TaggedPacket {
349        TaggedPacket {
350            now: Instant::now(),
351            transport: TransportContext::default(),
352            message: Packet::Rtp(rtp::Packet {
353                header: rtp::header::Header {
354                    ssrc,
355                    sequence_number,
356                    timestamp,
357                    ..Default::default()
358                },
359                ..Default::default()
360            }),
361        }
362    }
363
364    fn sequence_of(packet: &TaggedPacket) -> u16 {
365        match &packet.message {
366            Packet::Rtp(rtp) => rtp.header.sequence_number,
367            _ => panic!("not RTP"),
368        }
369    }
370
371    /// Drains everything held. Starts playout first: `pop` yields nothing while buffering, which
372    /// is the state a fresh buffer is in.
373    fn drain(buffer: &mut JitterBuffer) -> Vec<u16> {
374        buffer.begin_emitting();
375        let mut out = Vec::new();
376        while let Some(packet) = buffer.pop() {
377            out.push(sequence_of(&packet));
378        }
379        out
380    }
381
382    #[test]
383    fn packets_come_out_in_sequence_order_however_they_went_in() {
384        let mut buffer = JitterBuffer::new(64);
385        for sequence_number in [3u16, 1, 4, 2, 0] {
386            buffer.push(packet(1, sequence_number, 0)).expect("push");
387        }
388        assert_eq!(vec![0, 1, 2, 3, 4], drain(&mut buffer));
389    }
390
391    #[test]
392    fn reordering_is_counted() {
393        let mut buffer = JitterBuffer::new(64);
394        buffer.push(packet(1, 0, 0)).expect("push");
395        buffer.push(packet(1, 2, 0)).expect("push");
396        assert_eq!(0, buffer.stats().out_of_order);
397        buffer.push(packet(1, 1, 0)).expect("push");
398        assert_eq!(1, buffer.stats().out_of_order);
399    }
400
401    /// Upstream's queue inserts equal priorities, so a retransmission that races its original
402    /// leaves two copies of the same packet in the ordering.
403    #[test]
404    fn a_duplicate_is_rejected_rather_than_stored_twice() {
405        let mut buffer = JitterBuffer::new(64);
406        buffer.push(packet(1, 7, 0)).expect("push");
407        assert_eq!(Err(Rejected::Duplicate), buffer.push(packet(1, 7, 0)));
408        assert_eq!(1, buffer.len());
409        assert_eq!(1, buffer.stats().duplicates);
410        assert_eq!(vec![7], drain(&mut buffer));
411    }
412
413    /// The wrap case, which raw `u16` ordering gets wrong by a whole cycle.
414    #[test]
415    fn ordering_survives_a_wrap_around() {
416        let mut buffer = JitterBuffer::new(64);
417        for sequence_number in [65534u16, 65535, 0, 1] {
418            buffer.push(packet(1, sequence_number, 0)).expect("push");
419        }
420        assert_eq!(vec![65534, 65535, 0, 1], drain(&mut buffer));
421    }
422
423    #[test]
424    fn ordering_survives_two_wrap_arounds_with_reordering() {
425        let mut buffer = JitterBuffer::new(4096);
426
427        // Two full cycles, each pair swapped on the way in.
428        let mut expected = Vec::new();
429        let mut sequence_number = 65000u16;
430        let mut pushed = Vec::new();
431        for _ in 0..(2 * 65536 / 2) {
432            let a = sequence_number;
433            let b = sequence_number.wrapping_add(1);
434            pushed.push((b, a)); // swapped
435            expected.push(a);
436            expected.push(b);
437            sequence_number = sequence_number.wrapping_add(2);
438        }
439
440        // Push and drain incrementally so the buffer stays within capacity.
441        let mut emitted = Vec::new();
442        buffer.begin_emitting();
443        for (b, a) in pushed {
444            let _ = buffer.push(packet(1, b, 0));
445            let _ = buffer.push(packet(1, a, 0));
446            while buffer.len() > 2 {
447                emitted.push(sequence_of(&buffer.pop().expect("non-empty")));
448            }
449        }
450        emitted.extend(drain(&mut buffer));
451
452        assert_eq!(
453            expected.len(),
454            emitted.len(),
455            "every packet came out exactly once across two wraps"
456        );
457        assert_eq!(expected, emitted, "and in sequence order throughout");
458    }
459
460    /// The test pion's shared-buffer design cannot pass: one buffer belongs to one stream.
461    #[test]
462    fn a_second_ssrc_cannot_interleave_into_this_buffer() {
463        let mut buffer = JitterBuffer::new(64);
464        buffer.push(packet(1, 10, 0)).expect("push");
465
466        assert_eq!(
467            Err(Rejected::ForeignSsrc),
468            buffer.push(packet(2, 5, 0)),
469            "a different stream's packet must not sort against this stream's sequence numbers"
470        );
471        assert_eq!(1, buffer.stats().foreign_ssrc);
472        assert_eq!(Some(1), buffer.ssrc());
473        assert_eq!(vec![10], drain(&mut buffer));
474    }
475
476    #[test]
477    fn rtcp_is_not_stored() {
478        let mut buffer = JitterBuffer::new(64);
479        let rtcp = TaggedPacket {
480            now: Instant::now(),
481            transport: TransportContext::default(),
482            message: Packet::Rtcp(vec![]),
483        };
484        assert_eq!(Err(Rejected::ForeignSsrc), buffer.push(rtcp));
485        assert!(buffer.is_empty());
486    }
487
488    #[test]
489    fn a_packet_arriving_after_its_position_was_released_is_rejected() {
490        let mut buffer = JitterBuffer::new(64);
491        buffer.push(packet(1, 5, 0)).expect("push");
492        buffer.begin_emitting();
493        buffer.pop().expect("release 5");
494
495        assert_eq!(
496            Err(Rejected::Late),
497            buffer.push(packet(1, 5, 0)),
498            "the same packet again would be emitted twice"
499        );
500        assert_eq!(
501            Err(Rejected::Late),
502            buffer.push(packet(1, 4, 0)),
503            "an older straggler would be emitted out of order"
504        );
505        assert_eq!(2, buffer.stats().late);
506        assert!(buffer.is_empty());
507    }
508
509    #[test]
510    fn a_full_buffer_drops_its_oldest_packet() {
511        let mut buffer = JitterBuffer::new(3);
512        for sequence_number in 0..3u16 {
513            buffer.push(packet(1, sequence_number, 0)).expect("push");
514        }
515        buffer.push(packet(1, 3, 0)).expect("push displaces");
516
517        assert_eq!(3, buffer.len(), "capacity is respected");
518        assert_eq!(1, buffer.stats().overflow);
519        assert_eq!(vec![1, 2, 3], drain(&mut buffer), "the oldest gave way");
520    }
521
522    /// When the buffer is full and the arrival is older than everything in it, dropping the
523    /// arrival keeps the held run contiguous — evicting the oldest to make room for something
524    /// even older would leave a hole at both ends.
525    #[test]
526    fn a_full_buffer_drops_an_arrival_older_than_everything_held() {
527        let mut buffer = JitterBuffer::new(3);
528        for sequence_number in [10u16, 11, 12] {
529            buffer.push(packet(1, sequence_number, 0)).expect("push");
530        }
531
532        assert_eq!(Err(Rejected::Overflow), buffer.push(packet(1, 9, 0)));
533        assert_eq!(vec![10, 11, 12], drain(&mut buffer));
534    }
535
536    #[test]
537    fn pop_at_and_find_address_packets_by_extended_sequence_number() {
538        let mut buffer = JitterBuffer::new(64);
539        let first = buffer.push(packet(1, 100, 0)).expect("push");
540        let second = buffer.push(packet(1, 101, 0)).expect("push");
541
542        assert!(buffer.find(second).is_some());
543        assert!(buffer.find(second + 10).is_none());
544
545        assert_eq!(101, sequence_of(&buffer.pop_at(second).expect("pop_at")));
546        assert!(buffer.find(second).is_none());
547        assert_eq!(100, sequence_of(&buffer.pop_at(first).expect("pop_at")));
548        assert!(buffer.is_empty());
549    }
550
551    #[test]
552    fn the_front_reports_the_next_packet_due() {
553        let mut buffer = JitterBuffer::new(64);
554        assert_eq!(None, buffer.front_sequence());
555        assert_eq!(None, buffer.front_timestamp());
556
557        buffer.push(packet(1, 8, 9000)).expect("push");
558        buffer.push(packet(1, 7, 3000)).expect("push");
559
560        assert_eq!(Some(7), buffer.front_sequence());
561        assert_eq!(Some(3000), buffer.front_timestamp());
562        assert_eq!(7, sequence_of(buffer.peek().expect("peek")));
563        assert_eq!(2, buffer.len(), "peeking does not remove");
564    }
565
566    /// Upstream returns `ErrPopWhileBuffering` from `Pop` — a "try again later" sentinel its
567    /// synchronous reader needs. Sans-I/O has no use for it: not being ready yet is `None`.
568    #[test]
569    fn nothing_is_handed_out_while_buffering() {
570        let mut buffer = JitterBuffer::new(64);
571        buffer.push(packet(1, 1, 0)).expect("push");
572
573        assert_eq!(
574            State::Buffering,
575            buffer.state(),
576            "a fresh buffer fills first"
577        );
578        assert!(buffer.pop().is_none(), "no packets while buffering");
579        assert_eq!(1, buffer.len(), "and the packet is still held, not dropped");
580        assert_eq!(
581            0,
582            buffer.stats().underflow,
583            "declining to emit while buffering is not an underflow"
584        );
585
586        buffer.begin_emitting();
587        assert_eq!(1, sequence_of(&buffer.pop().expect("now emitting")));
588    }
589
590    #[test]
591    fn asking_an_empty_emitting_buffer_for_a_packet_is_an_underflow() {
592        let mut buffer = JitterBuffer::new(64);
593        buffer.begin_emitting();
594
595        assert!(buffer.pop().is_none());
596        assert_eq!(1, buffer.stats().underflow, "playout asked and got nothing");
597
598        buffer.begin_buffering();
599        assert!(buffer.pop().is_none());
600        assert_eq!(
601            1,
602            buffer.stats().underflow,
603            "but a buffering stream is not underflowing — it has not started"
604        );
605    }
606
607    /// A gap does not stall the stream. Upstream pops strictly at its playout head and underflows
608    /// on any missing sequence number; here the deadline above decides whether a gap is worth
609    /// waiting for, and this hands over what is due.
610    #[test]
611    fn a_gap_does_not_stall_playout() {
612        let mut buffer = JitterBuffer::new(64);
613        buffer.push(packet(1, 1, 0)).expect("push");
614        buffer.push(packet(1, 3, 0)).expect("push"); // 2 never arrives
615
616        assert_eq!(vec![1, 3], drain(&mut buffer));
617    }
618
619    #[test]
620    fn packets_are_addressable_by_their_wire_sequence_number() {
621        let mut buffer = JitterBuffer::new(64);
622        buffer.push(packet(1, 100, 0)).expect("push");
623        buffer.push(packet(1, 101, 0)).expect("push");
624
625        assert_eq!(
626            101,
627            sequence_of(buffer.peek_at_sequence(101).expect("peek"))
628        );
629        assert!(buffer.peek_at_sequence(999).is_none());
630        assert_eq!(2, buffer.len(), "peeking does not remove");
631
632        assert_eq!(101, sequence_of(&buffer.pop_at_sequence(101).expect("pop")));
633        assert!(buffer.peek_at_sequence(101).is_none());
634        assert!(buffer.pop_at_sequence(101).is_none(), "already taken");
635        assert_eq!(1, buffer.len());
636    }
637
638    /// Addressing by wire number has to keep working across a wrap, where two held packets differ
639    /// by a cycle in the ordering but are ordinary neighbours on the wire.
640    #[test]
641    fn addressing_by_wire_sequence_number_works_across_a_wrap() {
642        let mut buffer = JitterBuffer::new(64);
643        for sequence_number in [65534u16, 65535, 0, 1] {
644            buffer.push(packet(1, sequence_number, 0)).expect("push");
645        }
646
647        assert_eq!(
648            0,
649            sequence_of(buffer.peek_at_sequence(0).expect("post-wrap"))
650        );
651        assert_eq!(
652            65535,
653            sequence_of(buffer.peek_at_sequence(65535).expect("pre-wrap"))
654        );
655        assert_eq!(0, sequence_of(&buffer.pop_at_sequence(0).expect("pop 0")));
656        assert_eq!(vec![65534, 65535, 1], drain(&mut buffer));
657    }
658
659    /// One packet per call, matching upstream's `PopAtTimestamp`: a video frame spans several
660    /// packets sharing a timestamp, and releasing the frame means calling until it yields `None`.
661    #[test]
662    fn packets_are_addressable_by_rtp_timestamp() {
663        let mut buffer = JitterBuffer::new(64);
664        buffer.push(packet(1, 10, 9000)).expect("push");
665        buffer.push(packet(1, 11, 9000)).expect("push"); // same frame
666        buffer.push(packet(1, 12, 12000)).expect("push"); // next frame
667
668        assert_eq!(
669            10,
670            sequence_of(&buffer.pop_at_timestamp(9000).expect("first of the frame")),
671            "the lowest sequence number of that timestamp comes out first"
672        );
673        assert_eq!(
674            11,
675            sequence_of(&buffer.pop_at_timestamp(9000).expect("second of the frame"))
676        );
677        assert!(
678            buffer.pop_at_timestamp(9000).is_none(),
679            "the frame is fully released"
680        );
681        assert!(buffer.pop_at_timestamp(4242).is_none(), "no such timestamp");
682
683        assert_eq!(1, buffer.len(), "the next frame is untouched");
684        assert_eq!(vec![12], drain(&mut buffer));
685    }
686
687    /// Taking a packet out by sequence number or timestamp still marks that position played out,
688    /// or a straggler could be emitted behind it.
689    #[test]
690    fn addressed_removal_also_marks_the_position_released() {
691        let mut buffer = JitterBuffer::new(64);
692        buffer.push(packet(1, 5, 0)).expect("push");
693        buffer.push(packet(1, 6, 0)).expect("push");
694
695        buffer.pop_at_sequence(6).expect("pop 6");
696        assert_eq!(
697            Err(Rejected::Late),
698            buffer.push(packet(1, 6, 0)),
699            "6 has been played out"
700        );
701        assert_eq!(
702            Err(Rejected::Late),
703            buffer.push(packet(1, 5, 0)),
704            "and so has everything before it"
705        );
706    }
707
708    #[test]
709    fn reset_clears_the_contents_but_keeps_the_stream_and_counters() {
710        let mut buffer = JitterBuffer::new(64);
711        buffer.push(packet(1, 5, 0)).expect("push");
712        buffer.begin_emitting();
713        buffer.pop().expect("release");
714        buffer.push(packet(1, 5, 0)).ok(); // counted late
715
716        buffer.reset();
717
718        assert!(buffer.is_empty());
719        assert_eq!(Some(1), buffer.ssrc(), "still this stream's buffer");
720        assert_eq!(1, buffer.stats().late, "history is not erased");
721        assert_eq!(
722            State::Buffering,
723            buffer.state(),
724            "a restarted stream re-accumulates before playout resumes"
725        );
726        // The release anchor is gone, so the stream may restart at any sequence number.
727        assert!(buffer.push(packet(1, 5, 0)).is_ok());
728    }
729}