Skip to content
File

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

rust247 lines
1use std::net::Ipv4Addr;
2use std::time::{Duration, Instant};
3use str0m::Rtc;
4use str0m::format::Codec;
5use str0m::media::{Direction, MediaData, MediaKind};
6use str0m::rtp::RtpWrite;
7use str0m::{Candidate, Event, RtcError};
8use tracing::info_span;
9 
10mod common;
11use common::{Peer, TestRtc, init_crypto_default, init_log, progress};
12 
13#[test]
14pub fn contiguous_all_the_way() -> Result<(), RtcError> {
15 init_log();
16 init_crypto_default();
17 
18 let output = Server::with_vp8_input()
19 .timeout(Duration::from_secs(10))
20 .get_output()?;
21 let mut count = 0;
22 
23 // Contiguous all the way through.
24 for data in output {
25 assert!(data.contiguous);
26 count += 1;
27 }
28 
29 // We have 3 continuations: 104 - 3 == 101
30 assert_eq!(count, 101);
31 
32 Ok(())
33}
34 
35#[test]
36pub fn not_contiguous() -> Result<(), RtcError> {
37 init_log();
38 init_crypto_default();
39 
40 let output = Server::with_vp8_input()
41 .skip_packet(14337)
42 .timeout(Duration::from_secs(5))
43 .get_output()?;
44 let mut count = 0;
45 
46 // Contiguous all the way through.
47 for data in output {
48 count += 1;
49 // We dropped packet 14337, which means its dependendant 14338 is not
50 // emitted, and 14339 is emitted and marked as discontinuous.
51 let assume_contiguous = !data.seq_range.contains(&14339.into());
52 assert_eq!(assume_contiguous, data.contiguous);
53 }
54 
55 // assert!(false);
56 // We have 3 continuations, 2 missing packet (14337 14338)
57 // 104 - 3 - 2 == 99
58 assert_eq!(count, 99);
59 
60 Ok(())
61}
62 
63#[test]
64pub fn vp9_contiguous_all_the_way() -> Result<(), RtcError> {
65 init_log();
66 init_crypto_default();
67 
68 let output = Server::with_vp9_input().get_output()?;
69 let mut count = 0;
70 
71 // Contiguous all the way through.
72 for data in output {
73 assert!(data.contiguous);
74 count += 1;
75 }
76 
77 // The last packet is never flushed out because the depacketizer wants
78 // to see the next packet before releasing.
79 // We have one last packet missing: 16 - 1 == 15
80 assert_eq!(count, 15);
81 
82 Ok(())
83}
84 
85#[test]
86pub fn vp9_not_contiguous() -> Result<(), RtcError> {
87 init_log();
88 init_crypto_default();
89 
90 let output = Server::with_vp9_input().skip_packet(30952).get_output()?;
91 let mut count = 0;
92 
93 // Contiguous all the way through.
94 for data in output {
95 count += 1;
96 // We dropped packet 19365 and next packet is emitted and marked as discontinuous.
97 let assume_contiguous = !data.seq_range.contains(&30953.into());
98 assert_eq!(assume_contiguous, data.contiguous);
99 }
100 
101 // assert!(false);
102 // We 1 missing packet 30952, and one last
103 // packet missing: 16 - 1 - 1 == 14
104 assert_eq!(count, 14);
105 
106 Ok(())
107}
108 
109struct Server {
110 codec: Codec,
111 input_data: common::PcapData,
112 skip_packet: Option<u16>,
113 timeout: Option<Duration>,
114}
115 
116impl Server {
117 fn with_vp8_input() -> Self {
118 Self::new(Codec::Vp8, common::vp8_data())
119 }
120 
121 fn with_vp9_input() -> Self {
122 Self::new(Codec::Vp9, common::vp9_contiguous_data())
123 }
124 
125 fn new(codec: Codec, input_data: common::PcapData) -> Self {
126 Self {
127 codec,
128 input_data,
129 skip_packet: None,
130 timeout: None,
131 }
132 }
133 
134 fn skip_packet(mut self, packet: u16) -> Self {
135 self.skip_packet = Some(packet);
136 self
137 }
138 
139 fn timeout(mut self, timeout: Duration) -> Self {
140 self.timeout = Some(timeout);
141 self
142 }
143 
144 fn get_output(self) -> Result<Vec<MediaData>, RtcError> {
145 let mut l = TestRtc::new(Peer::Left);
146 
147 // We need to lower the default reordering buffer size, or we won't make it
148 // past the dropped packet.
149 let rtc_r = Rtc::builder()
150 .set_reordering_size_video(5)
151 .build(Instant::now());
152 
153 let mut r = TestRtc::new_with_rtc(info_span!("R"), rtc_r);
154 
155 l.add_local_candidate(Candidate::host(
156 (Ipv4Addr::new(1, 1, 1, 1), 1000).into(),
157 "udp",
158 )?);
159 r.add_local_candidate(Candidate::host(
160 (Ipv4Addr::new(2, 2, 2, 2), 2000).into(),
161 "udp",
162 )?);
163 
164 // The change is on the L (sending side) with Direction::SendRecv.
165 let mut change = l.sdp_api();
166 let mid = change.add_media(MediaKind::Video, Direction::SendOnly, None, None, None);
167 let (offer, pending) = change.apply().unwrap();
168 
169 let answer = r.rtc.sdp_api().accept_offer(offer)?;
170 l.rtc.sdp_api().accept_answer(pending, answer)?;
171 
172 loop {
173 if l.is_connected() || r.is_connected() {
174 break;
175 }
176 progress(&mut l, &mut r)?;
177 }
178 
179 let max = l.last.max(r.last);
180 l.last = max;
181 r.last = max;
182 
183 let params = match self.codec {
184 Codec::Vp8 => l.params_vp8(),
185 Codec::Vp9 => l.params_vp9(),
186 _ => unimplemented!(),
187 };
188 assert_eq!(params.spec().codec, self.codec);
189 
190 let pt = params.pt();
191 
192 for (relative, header, payload) in self.input_data {
193 // Drop a random packet in the middle.
194 if Some(header.sequence_number) == self.skip_packet {
195 continue;
196 }
197 
198 // Keep RTC time progressed to be "in sync" with the test data.
199 while (l.last - max) < relative {
200 progress(&mut l, &mut r)?;
201 }
202 
203 let absolute = max + relative;
204 
205 let mut direct = l.direct_api();
206 let tx = direct.stream_tx_by_mid(mid, None).unwrap();
207 tx.write_rtp(
208 RtpWrite::new(
209 pt,
210 header.sequence_number(None),
211 header.timestamp,
212 absolute,
213 payload,
214 )
215 .marker(header.marker)
216 .ext_vals(header.ext_vals)
217 .nackable(true),
218 );
219 
220 progress(&mut l, &mut r)?;
221 
222 if let Some(duration) = self.timeout {
223 if l.duration() > duration {
224 break;
225 }
226 }
227 }
228 
229 // Drain any remaining packets from the pacer
230 progress(&mut l, &mut r)?;
231 
232 let events = r
233 .events
234 .into_iter()
235 .filter_map(|(_, e)| {
236 if let Event::MediaData(d) = e {
237 Some(d)
238 } else {
239 None
240 }
241 })
242 .collect();
243 
244 Ok(events)
245 }
246}