File
Blob: firmware/vendor/str0m/tests/contiguous.rs
| 1 | use std::net::Ipv4Addr; |
| 2 | use std::time::{Duration, Instant}; |
| 3 | use str0m::Rtc; |
| 4 | use str0m::format::Codec; |
| 5 | use str0m::media::{Direction, MediaData, MediaKind}; |
| 6 | use str0m::rtp::RtpWrite; |
| 7 | use str0m::{Candidate, Event, RtcError}; |
| 8 | use tracing::info_span; |
| 9 | |
| 10 | mod common; |
| 11 | use common::{Peer, TestRtc, init_crypto_default, init_log, progress}; |
| 12 | |
| 13 | #[test] |
| 14 | pub 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] |
| 36 | pub 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] |
| 64 | pub 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] |
| 86 | pub 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 | |
| 109 | struct Server { |
| 110 | codec: Codec, |
| 111 | input_data: common::PcapData, |
| 112 | skip_packet: Option<u16>, |
| 113 | timeout: Option<Duration>, |
| 114 | } |
| 115 | |
| 116 | impl 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 | } |