Skip to content
File

Blob: firmware/vendor/str0m/src/streams/rtx_cache.rs

rust193 lines
1use std::time::Duration;
2use std::time::Instant;
3 
4use crate::io::DATAGRAM_MAX_PACKET_SIZE;
5use crate::rtp_::SeqNo;
6 
7use super::RtpPacket;
8use super::rtx_cache_buf::EvictingBuffer;
9 
10const RTX_CACHE_SIZE_QUANTIZER: usize = 25;
11const RTX_CACHE_QUANTIZE_SLOTS: usize = DATAGRAM_MAX_PACKET_SIZE / RTX_CACHE_SIZE_QUANTIZER;
12 
13#[derive(Debug)]
14pub(crate) struct RtxCache {
15 // Data, new additions here probably need to be cleared in [`clear`].
16 packet_by_seq_no: EvictingBuffer<RtpPacket>,
17 
18 // Technically we want [Option<SeqNo>; X] to indicate the absence of
19 // a SeqNo. However We can half the storage space by using the sentinel
20 // values SeqNo::MAX to indicate None
21 seq_no_by_quantized_size: [SeqNo; RTX_CACHE_QUANTIZE_SLOTS],
22}
23 
24impl RtxCache {
25 pub fn new(max_packet_count: usize, max_packet_age: Duration) -> Self {
26 Self {
27 packet_by_seq_no: EvictingBuffer::new(10, max_packet_age, max_packet_count),
28 seq_no_by_quantized_size: [SeqNo::MAX; RTX_CACHE_QUANTIZE_SLOTS],
29 }
30 }
31 
32 pub fn cache_sent_packet(&mut self, packet: RtpPacket, now: Instant) {
33 assert!(packet.nackable);
34 let seq_no = packet.seq_no;
35 let quantized_size = packet.payload.len() / RTX_CACHE_SIZE_QUANTIZER;
36 self.packet_by_seq_no.push(*seq_no, now, packet);
37 self.seq_no_by_quantized_size[quantized_size] = seq_no;
38 self.remove_old_packets(now);
39 }
40 
41 pub fn last_cached_seq_no(&self) -> Option<SeqNo> {
42 Some(self.packet_by_seq_no.last_position()?.into())
43 }
44 
45 pub fn get_cached_packet_by_seq_no(&mut self, seq_no: SeqNo) -> Option<&mut RtpPacket> {
46 self.packet_by_seq_no.get_mut(*seq_no)
47 }
48 
49 pub fn get_cached_packet_smaller_than(&mut self, max_size: usize) -> Option<&mut RtpPacket> {
50 let quantized_size = max_size / RTX_CACHE_SIZE_QUANTIZER;
51 
52 let seq_no = self.seq_no_by_quantized_size[..quantized_size]
53 .iter()
54 .rev()
55 .filter(|seq_no| !seq_no.is_max())
56 .find(|seq_no| self.packet_by_seq_no.contains(***seq_no))?;
57 
58 self.get_cached_packet_by_seq_no(*seq_no)
59 }
60 
61 fn remove_old_packets(&mut self, now: Instant) {
62 self.packet_by_seq_no.maybe_evict(now);
63 }
64 
65 pub(crate) fn last_packet(&self) -> Option<&[u8]> {
66 let packet = self.packet_by_seq_no.last()?;
67 Some(packet.payload.as_ref())
68 }
69 
70 pub(crate) fn clear(&mut self) {
71 self.packet_by_seq_no.clear();
72 self.seq_no_by_quantized_size = [SeqNo::MAX; RTX_CACHE_QUANTIZE_SLOTS];
73 }
74}
75 
76#[cfg(test)]
77mod test {
78 use crate::rtp_::MediaTime;
79 use crate::rtp_::RtpHeader;
80 
81 use super::*;
82 
83 fn after(now: Instant, millis: u64) -> Instant {
84 now + Duration::from_millis(millis)
85 }
86 
87 fn packet(now: Instant, seq_no: u64, millis: u64) -> RtpPacket {
88 RtpPacket {
89 header: RtpHeader::default(),
90 seq_no: seq_no.into(),
91 time: MediaTime::from_90khz(0),
92 payload: millis.to_be_bytes().into(),
93 vp8_patch: None,
94 timestamp: after(now, millis),
95 last_sender_info: None,
96 nackable: true,
97 }
98 }
99 
100 #[test]
101 fn rtx_cache_0_sized() {
102 let now = Instant::now();
103 let mut rtx_cache = RtxCache::new(0, Duration::from_secs(3));
104 rtx_cache.cache_sent_packet(packet(now, 1, 10), after(now, 10));
105 assert_eq!(
106 Some(&mut packet(now, 1, 10)),
107 rtx_cache.get_cached_packet_by_seq_no(1.into())
108 );
109 assert_eq!(
110 Some(&mut packet(now, 1, 10)),
111 rtx_cache.get_cached_packet_smaller_than(1000)
112 );
113 }
114 
115 #[test]
116 fn rtx_cache_0_duration() {
117 let now = Instant::now();
118 let mut rtx_cache = RtxCache::new(10, Duration::from_secs(0));
119 rtx_cache.cache_sent_packet(packet(now, 1, 10), after(now, 10));
120 assert_eq!(None, rtx_cache.get_cached_packet_by_seq_no(1.into()));
121 assert_eq!(None, rtx_cache.get_cached_packet_smaller_than(1000));
122 }
123 
124 #[test]
125 fn rtx_cache_1_sized() {
126 let now = Instant::now();
127 let mut rtx_cache = RtxCache::new(1, Duration::from_secs(3));
128 rtx_cache.cache_sent_packet(packet(now, 1, 10), after(now, 10));
129 assert_eq!(
130 Some(&mut packet(now, 1, 10)),
131 rtx_cache.get_cached_packet_by_seq_no(1.into())
132 );
133 assert_eq!(
134 Some(&mut packet(now, 1, 10)),
135 rtx_cache.get_cached_packet_smaller_than(1000)
136 );
137 assert_eq!(
138 Some(&mut packet(now, 1, 10)),
139 rtx_cache.get_cached_packet_smaller_than(25)
140 );
141 assert_eq!(None, rtx_cache.get_cached_packet_smaller_than(24));
142 rtx_cache.cache_sent_packet(packet(now, 2, 20), after(now, 20));
143 assert_eq!(
144 Some(&mut packet(now, 2, 20)),
145 rtx_cache.get_cached_packet_by_seq_no(2.into())
146 );
147 }
148 
149 #[test]
150 fn rtx_cache_100_sized_backwards() {
151 let now = Instant::now();
152 let mut rtx_cache = RtxCache::new(100, Duration::from_secs(3));
153 
154 for i in 1..=200 {
155 let seq_no = 201 - i;
156 let pkt = packet(now, seq_no, (201 - i) * 10);
157 let now = after(now, i * 10);
158 rtx_cache.cache_sent_packet(pkt, now);
159 }
160 
161 assert_eq!(
162 Some(&mut packet(now, 200, 2000)),
163 rtx_cache.get_cached_packet_by_seq_no(200.into())
164 );
165 assert_eq!(None, rtx_cache.get_cached_packet_by_seq_no(100.into()));
166 assert_eq!(None, rtx_cache.get_cached_packet_by_seq_no(101.into()));
167 }
168 
169 #[test]
170 fn rtx_cache_eviction() {
171 let now = Instant::now();
172 // Cache can take all entries
173 let mut rtx_cache = RtxCache::new(400, Duration::from_secs(1));
174 for i in 1..=200 {
175 let seq_no = i;
176 let pkt = packet(now, seq_no, i * 10);
177 let now = after(now, i * 10);
178 rtx_cache.cache_sent_packet(pkt, now);
179 }
180 
181 assert_eq!(None, rtx_cache.get_cached_packet_by_seq_no(99.into()));
182 
183 assert_eq!(
184 Some(&mut packet(now, 100, 1000)),
185 rtx_cache.get_cached_packet_by_seq_no(100.into())
186 );
187 assert_eq!(
188 Some(&mut packet(now, 200, 2000)),
189 rtx_cache.get_cached_packet_by_seq_no(200.into())
190 );
191 }
192}