Skip to content
File

Blob: firmware/vendor/str0m/tests/nack.rs

rust314 lines
1use std::collections::VecDeque;
2use std::time::{Duration, Instant};
3 
4use netem::{NetemConfig, Probability, RandomLoss};
5use str0m::RtcError;
6use str0m::format::Codec;
7use str0m::media::MediaKind;
8use str0m::rtp::rtcp::Rtcp;
9use str0m::rtp::{ExtensionValues, RawPacket, RtpWrite, SeqNo, Ssrc};
10 
11mod common;
12use common::{connect_l_r, init_crypto_default, init_log, progress};
13 
14#[test]
15pub fn loss_recovery() -> Result<(), RtcError> {
16 init_log();
17 init_crypto_default();
18 
19 let (mut l, mut r) = connect_l_r();
20 
21 // Configure 5% random loss on R's incoming queue (L -> R has loss)
22 let loss_config = NetemConfig::new()
23 .loss(RandomLoss::new(Probability::new(0.05)))
24 .seed(42);
25 r.set_netem(loss_config);
26 
27 let mid = "vid".into();
28 
29 // In this example we are using MID only (no RID) to identify the incoming media.
30 let ssrc_tx: Ssrc = 42.into();
31 let ssrc_rtx: Ssrc = 44.into();
32 
33 l.direct_api().declare_media(mid, MediaKind::Video);
34 
35 l.direct_api()
36 .declare_stream_tx(ssrc_tx, Some(ssrc_rtx), mid, None);
37 
38 r.direct_api().declare_media(mid, MediaKind::Video);
39 
40 r.direct_api()
41 .expect_stream_rx(ssrc_tx, Some(ssrc_rtx), mid, None);
42 
43 let max = l.last.max(r.last);
44 l.last = max;
45 r.last = max;
46 
47 let params = l.params_vp8();
48 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
49 assert_eq!(params.spec().codec, Codec::Vp8);
50 let pt = params.pt();
51 
52 let to_write = [0x1, 0x2, 0x3, 0x4];
53 let num_packets: usize = 1000;
54 
55 // write all packets num_packets
56 for index in 0..num_packets {
57 let wallclock = l.start + l.duration();
58 
59 let mut direct = l.direct_api();
60 let stream = direct.stream_tx(&ssrc).unwrap();
61 
62 let time = (index * 1000 + 47_000_000) as u32;
63 let seq_no = (47_000 + index as u64).into();
64 
65 stream.write_rtp(RtpWrite::new(pt, seq_no, time, wallclock, to_write).nackable(true));
66 
67 // Disable loss near start and end to let retransmission algo stabilize
68 // (see MISORDER_DELAY in register.rs)
69 if !(10..=990).contains(&index) {
70 r.set_netem(NetemConfig::new()); // No loss
71 }
72 
73 progress(&mut l, &mut r)?;
74 
75 // Re-enable loss for middle packets
76 if index == 9 {
77 let loss_config = NetemConfig::new()
78 .loss(RandomLoss::new(Probability::new(0.05)))
79 .seed(42);
80 r.set_netem(loss_config);
81 }
82 }
83 
84 // let some time pass for retransmission to happen
85 let settle_time = l.duration() + Duration::from_secs(10);
86 loop {
87 progress(&mut l, &mut r)?;
88 
89 if l.duration() > settle_time {
90 break;
91 }
92 }
93 
94 // some nacks have been transmitted
95 let nacks_tx = r
96 .events
97 .iter()
98 .filter_map(|(_, e)| match e.as_raw_packet() {
99 Some(RawPacket::RtcpTx(Rtcp::Nack(p))) => Some(p),
100 _ => None,
101 })
102 .collect::<Vec<_>>();
103 
104 assert!(!nacks_tx.is_empty());
105 
106 // some nacks have been received
107 let nacks_rx = l
108 .events
109 .iter()
110 .filter_map(|(_, e)| match e.as_raw_packet() {
111 Some(RawPacket::RtcpRx(Rtcp::Nack(p))) => Some(p),
112 _ => None,
113 })
114 .collect::<Vec<_>>();
115 
116 assert!(!nacks_rx.is_empty());
117 
118 // all packets were received in the end
119 let mut packets_rx = r
120 .events
121 .iter()
122 .filter_map(|(_, e)| match e.as_raw_packet() {
123 Some(RawPacket::RtpRx(p, b)) => {
124 //
125 if p.payload_type == params.resend().unwrap() {
126 // read original seq no
127 let seq_no = u16::from_be_bytes(b.get(0..2)?.try_into().ok()?);
128 Some(seq_no)
129 } else {
130 Some(p.sequence_number)
131 }
132 }
133 
134 _ => None,
135 })
136 .collect::<Vec<_>>();
137 
138 packets_rx.sort();
139 
140 let discontinuities = packets_rx
141 .windows(2)
142 .filter_map(|slice| {
143 let a = slice.first()?;
144 let b = slice.get(1)?;
145 if a + 1 != *b { Some((*a, *b)) } else { None }
146 })
147 .collect::<Vec<_>>();
148 
149 let min = packets_rx.first().unwrap();
150 let max = packets_rx.last().unwrap();
151 
152 // useful for debugging
153 println!(
154 "min: {}, max: {}, total_rx: {}, discontinuities: {:?}",
155 min,
156 max,
157 packets_rx.len(),
158 discontinuities
159 );
160 
161 assert_eq!(*min, 47_000);
162 assert_eq!(*max, 47_999);
163 
164 assert_eq!(discontinuities.len(), 0);
165 assert_eq!(packets_rx.len(), num_packets);
166 
167 Ok(())
168}
169 
170#[test]
171pub fn nack_delay() -> Result<(), RtcError> {
172 init_log();
173 init_crypto_default();
174 
175 let (mut l, mut r) = connect_l_r();
176 
177 let mid = "vid".into();
178 
179 // In this example we are using MID only (no RID) to identify the incoming media.
180 let ssrc_tx: Ssrc = 42.into();
181 let ssrc_rtx: Ssrc = 44.into();
182 
183 l.direct_api().declare_media(mid, MediaKind::Video);
184 
185 l.direct_api()
186 .declare_stream_tx(ssrc_tx, Some(ssrc_rtx), mid, None);
187 
188 r.direct_api().declare_media(mid, MediaKind::Video);
189 
190 r.direct_api()
191 .expect_stream_rx(ssrc_tx, Some(ssrc_rtx), mid, None);
192 
193 let max = l.last.max(r.last);
194 l.last = max;
195 r.last = max;
196 
197 let params = l.params_vp8();
198 let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc();
199 assert_eq!(params.spec().codec, Codec::Vp8);
200 let pt = params.pt();
201 
202 let to_write: Vec<&[u8]> = vec![
203 &[0x1, 0x2, 0x3, 0x4],
204 &[0x9, 0xa, 0xb, 0xc],
205 &[0x5, 0x6, 0x7, 0x8],
206 &[0x1, 0x2, 0x3, 0x4],
207 &[0x9, 0xa, 0xb, 0xc],
208 &[0x5, 0x6, 0x7, 0x8],
209 &[0x1, 0x2, 0x3, 0x4],
210 &[0x9, 0xa, 0xb, 0xc],
211 &[0x5, 0x6, 0x7, 0x8],
212 &[0x1, 0x2, 0x3, 0x4],
213 &[0x9, 0xa, 0xb, 0xc],
214 ];
215 
216 let mut to_write: VecDeque<_> = to_write.into();
217 
218 let mut write_at = l.last + Duration::from_millis(5);
219 
220 let mut counts: Vec<u64> = vec![0, 1, 2, 4, 3, 5, 6, 7, 8, 9, 10];
221 
222 let mut dropped = (Instant::now(), 0.into());
223 
224 loop {
225 if l.start + l.duration() > write_at {
226 write_at = l.last + Duration::from_millis(5);
227 if let Some(packet) = to_write.pop_front() {
228 let wallclock = l.start + l.duration();
229 
230 let mut direct = l.direct_api();
231 let stream = direct.stream_tx(&ssrc).unwrap();
232 
233 let count = counts.remove(0);
234 let time = (count * 1000 + 47_000_000) as u32;
235 let seq_no = (47_000 + count).into();
236 
237 if count == 5 {
238 // Drop a packet
239 dropped = (wallclock, seq_no);
240 continue;
241 }
242 
243 let exts = ExtensionValues {
244 audio_level: Some(-42 - count as i8),
245 voice_activity: Some(false),
246 ..Default::default()
247 };
248 
249 stream.write_rtp(
250 RtpWrite::new(pt, seq_no, time, wallclock, packet)
251 .ext_vals(exts)
252 .nackable(true),
253 );
254 }
255 }
256 
257 progress(&mut l, &mut r)?;
258 
259 if l.duration() > Duration::from_secs(10) {
260 break;
261 }
262 }
263 
264 let nacks_tx = r
265 .events
266 .iter()
267 .filter_map(|(t, e)| match e.as_raw_packet() {
268 Some(RawPacket::RtcpTx(Rtcp::Nack(p))) => {
269 if p.reports
270 .iter()
271 .any(|r| SeqNo::from(r.pid as u64) == dropped.1)
272 {
273 Some(*t - dropped.0)
274 } else {
275 None
276 }
277 }
278 _ => None,
279 })
280 .collect::<Vec<_>>();
281 
282 let first_nack_tx = nacks_tx.first().expect("nack");
283 
284 assert!(first_nack_tx < &Duration::from_millis(100));
285 assert!(nacks_tx.iter().all(|f| f < &Duration::from_millis(200)));
286 
287 let nacks_rx = l
288 .events
289 .iter()
290 .filter_map(|(t, e)| match e.as_raw_packet() {
291 Some(RawPacket::RtcpRx(Rtcp::Nack(p))) => {
292 if p.reports
293 .iter()
294 .any(|r| SeqNo::from(r.pid as u64) == dropped.1)
295 {
296 Some(*t - dropped.0)
297 } else {
298 None
299 }
300 }
301 _ => None,
302 })
303 .collect::<Vec<_>>();
304 
305 let first_nack_rx = nacks_rx.first().expect("nack");
306 
307 assert!(first_nack_rx < &Duration::from_millis(100));
308 assert!(nacks_rx.iter().all(|f| f < &Duration::from_millis(200)));
309 
310 assert_eq!(nacks_rx.len(), nacks_tx.len());
311 
312 Ok(())
313}