File
Blob: firmware/vendor/str0m/src/rtp/rtcp/mod.rs
| 1 | #![allow(clippy::unusual_byte_groupings)] |
| 2 | |
| 3 | mod header; |
| 4 | use std::collections::VecDeque; |
| 5 | |
| 6 | pub use header::{RtcpHeader, RtcpType}; |
| 7 | |
| 8 | mod list; |
| 9 | pub use list::ReportList; |
| 10 | use list::private::WordSized; |
| 11 | |
| 12 | mod fmt; |
| 13 | pub use fmt::{FeedbackMessageType, PayloadType, TransportType}; |
| 14 | |
| 15 | mod sr; |
| 16 | pub use sr::{SenderInfo, SenderReport}; |
| 17 | |
| 18 | mod rr; |
| 19 | pub use rr::{ReceiverReport, ReceptionReport}; |
| 20 | |
| 21 | mod xr; |
| 22 | pub use xr::{Dlrr, DlrrItem, ExtendedReport, ReportBlock, Rrtr}; |
| 23 | |
| 24 | mod sdes; |
| 25 | pub use sdes::{Descriptions, Sdes, SdesType}; |
| 26 | |
| 27 | mod bb; |
| 28 | pub use bb::Goodbye; |
| 29 | |
| 30 | mod nack; |
| 31 | pub use nack::{Nack, NackEntry}; |
| 32 | |
| 33 | mod pli; |
| 34 | pub use pli::Pli; |
| 35 | |
| 36 | mod fir; |
| 37 | pub use fir::{Fir, FirEntry}; |
| 38 | |
| 39 | mod twcc; |
| 40 | pub use twcc::{Twcc, TwccPacketId, TwccRecvRegister, TwccSendRecord, TwccSendRegister}; |
| 41 | |
| 42 | mod rtcpfb; |
| 43 | pub use rtcpfb::RtcpFb; |
| 44 | |
| 45 | mod psfbapp; |
| 46 | mod remb; |
| 47 | pub use psfbapp::AppSpecificFeedback; |
| 48 | pub use remb::Remb; |
| 49 | |
| 50 | use super::SeqNo; |
| 51 | use super::Ssrc; |
| 52 | use super::extend_u16; |
| 53 | |
| 54 | pub trait RtcpPacket { |
| 55 | /// The... |
| 56 | fn header(&self) -> RtcpHeader; |
| 57 | |
| 58 | /// Length of entire RTCP packet (including header) in words (4 bytes). |
| 59 | fn length_words(&self) -> usize; |
| 60 | |
| 61 | /// Write this packet to the buffer. |
| 62 | /// |
| 63 | /// Panics if the buffer doesn't have capacity to hold length_words * 4 bytes. |
| 64 | fn write_to(&self, buf: &mut [u8]) -> usize; |
| 65 | } |
| 66 | |
| 67 | /// RTCP reports handled by str0m. |
| 68 | #[allow(clippy::large_enum_variant)] |
| 69 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 70 | pub enum Rtcp { |
| 71 | /// Sender report. Also known as SR. |
| 72 | SenderReport(SenderReport), |
| 73 | /// Receiver report. Also known as RR. |
| 74 | ReceiverReport(ReceiverReport), |
| 75 | /// Extended receiver report. Sometimes called XR. |
| 76 | /// |
| 77 | /// Always sent together with a receiver report. |
| 78 | ExtendedReport(ExtendedReport), |
| 79 | /// Description of Synchronization Sources (senders). |
| 80 | SourceDescription(Descriptions), |
| 81 | /// BYE. When a stream is over. |
| 82 | Goodbye(Goodbye), |
| 83 | /// Reports missing packets. |
| 84 | Nack(Nack), |
| 85 | /// Picture Loss Indiciation. When decoding a picture is not possible. |
| 86 | Pli(Pli), |
| 87 | /// Full Intra Request. Complete restart of a video decoder. |
| 88 | Fir(Fir), |
| 89 | /// Transport Wide Congestion Control. Feedback for every received RTP packet. |
| 90 | Twcc(Twcc), |
| 91 | /// Receiver Estimated Maximum Bitrate. Feedback to the sender about the maximum bitrate. |
| 92 | Remb(Remb), |
| 93 | /// Application-specific Payload-Specific Feedback (FMT=15, PT=206) that is not REMB. |
| 94 | /// Used for application-specific PSFB feedback (RFC 4585 Section 6.4). |
| 95 | AppSpecificFeedback(AppSpecificFeedback), |
| 96 | } |
| 97 | |
| 98 | impl Rtcp { |
| 99 | pub(crate) fn read_packet(buf: &[u8], feedback: &mut VecDeque<Rtcp>) { |
| 100 | let mut buf = buf; |
| 101 | loop { |
| 102 | if buf.is_empty() { |
| 103 | break; |
| 104 | } |
| 105 | |
| 106 | let header: RtcpHeader = match buf.try_into() { |
| 107 | Ok(v) => v, |
| 108 | Err(e) => { |
| 109 | debug!("{}", e); |
| 110 | break; |
| 111 | } |
| 112 | }; |
| 113 | let has_padding = buf[0] & 0b00_1_00000 > 0; |
| 114 | let full_length = header.length_words() * 4; |
| 115 | |
| 116 | if full_length > buf.len() { |
| 117 | // this length is incorrect. |
| 118 | break; |
| 119 | } |
| 120 | |
| 121 | let unpadded_length = if has_padding { |
| 122 | let pad = buf[full_length - 1] as usize; |
| 123 | if full_length < pad { |
| 124 | debug!("buf.len() is less than padding: {} < {}", full_length, pad); |
| 125 | break; |
| 126 | } |
| 127 | full_length - pad |
| 128 | } else { |
| 129 | full_length |
| 130 | }; |
| 131 | |
| 132 | match (&buf[..unpadded_length]).try_into() { |
| 133 | Ok(v) => feedback.push_back(v), |
| 134 | Err(e) => debug!("{}", e), |
| 135 | } |
| 136 | |
| 137 | buf = &buf[full_length..]; |
| 138 | } |
| 139 | } |
| 140 | |
| 141 | pub(crate) fn write_packet( |
| 142 | feedback: &mut VecDeque<Rtcp>, |
| 143 | buf: &mut [u8], |
| 144 | mut output: impl FnMut(Rtcp), |
| 145 | ) -> usize { |
| 146 | if feedback.is_empty() { |
| 147 | return 0; |
| 148 | } |
| 149 | |
| 150 | // Total length, in bytes, shrunk to be on the pad_to boundary. |
| 151 | let total_len = buf.len(); |
| 152 | |
| 153 | // Capacity in words |
| 154 | let word_capacity = total_len / 4; |
| 155 | |
| 156 | // Pack RTCP feedback packets. Merge together ones of the same type. |
| 157 | Rtcp::pack(feedback, word_capacity); |
| 158 | |
| 159 | let mut offset = 0; |
| 160 | while let Some(fb) = feedback.front() { |
| 161 | // Length of next item. |
| 162 | let item_len = fb.length_words() * 4; |
| 163 | |
| 164 | // Capacity left in the buffer. |
| 165 | let capacity = total_len - offset; |
| 166 | if capacity < item_len { |
| 167 | break; |
| 168 | } |
| 169 | |
| 170 | // We definitely can fit the next RTCP item. |
| 171 | let fb = feedback.pop_front().unwrap(); |
| 172 | let written = fb.write_to(&mut buf[offset..]); |
| 173 | |
| 174 | assert_eq!( |
| 175 | written, item_len, |
| 176 | "length_words equals write_to length: {fb:?}" |
| 177 | ); |
| 178 | |
| 179 | // When debugging we can pass an output to get the serialized packets. |
| 180 | output(fb); |
| 181 | |
| 182 | // Move offsets for the amount written. |
| 183 | offset += item_len; |
| 184 | } |
| 185 | |
| 186 | offset |
| 187 | } |
| 188 | |
| 189 | fn merge(&mut self, other: &mut Rtcp, words_left: usize) -> bool { |
| 190 | match (self, other) { |
| 191 | // Stack receiver reports into sender reports. |
| 192 | (Rtcp::SenderReport(sr), Rtcp::ReceiverReport(rr)) => { |
| 193 | let n = sr.reports.append_all_possible(&mut rr.reports, words_left); |
| 194 | n > 0 |
| 195 | } |
| 196 | |
| 197 | // Stack receiver reports. |
| 198 | (Rtcp::ReceiverReport(r1), Rtcp::ReceiverReport(r2)) => { |
| 199 | let n = r1.reports.append_all_possible(&mut r2.reports, words_left); |
| 200 | n > 0 |
| 201 | } |
| 202 | |
| 203 | // Stack source descriptions. |
| 204 | (Rtcp::SourceDescription(s1), Rtcp::SourceDescription(s2)) => { |
| 205 | let n = s1.reports.append_all_possible(&mut s2.reports, words_left); |
| 206 | n > 0 |
| 207 | } |
| 208 | |
| 209 | // Stack source descriptions. |
| 210 | (Rtcp::Goodbye(g1), Rtcp::Goodbye(g2)) => { |
| 211 | let n = g1.reports.append_all_possible(&mut g2.reports, words_left); |
| 212 | n > 0 |
| 213 | } |
| 214 | |
| 215 | // Stack Nack |
| 216 | (Rtcp::Nack(n1), Rtcp::Nack(n2)) if n1.ssrc == n2.ssrc => { |
| 217 | let n = n1.reports.append_all_possible(&mut n2.reports, words_left); |
| 218 | n > 0 |
| 219 | } |
| 220 | |
| 221 | // Stack source descriptions. |
| 222 | (Rtcp::Fir(f1), Rtcp::Fir(f2)) => { |
| 223 | let n = f1.reports.append_all_possible(&mut f2.reports, words_left); |
| 224 | n > 0 |
| 225 | } |
| 226 | |
| 227 | // No merge possible |
| 228 | _ => false, |
| 229 | } |
| 230 | } |
| 231 | |
| 232 | fn is_full(&self) -> bool { |
| 233 | match self { |
| 234 | Rtcp::SenderReport(v) => v.reports.is_full(), |
| 235 | Rtcp::ReceiverReport(v) => v.reports.is_full(), |
| 236 | Rtcp::ExtendedReport(_) => true, |
| 237 | Rtcp::SourceDescription(v) => v.reports.is_full(), |
| 238 | Rtcp::Goodbye(v) => v.reports.is_full(), |
| 239 | Rtcp::Nack(v) => v.reports.is_full(), |
| 240 | Rtcp::Pli(_) => true, |
| 241 | Rtcp::Fir(v) => v.reports.is_full(), |
| 242 | Rtcp::Twcc(_) => true, |
| 243 | Rtcp::Remb(_) => true, |
| 244 | Rtcp::AppSpecificFeedback(_) => true, |
| 245 | } |
| 246 | } |
| 247 | |
| 248 | /// If this RtcpFb contains no reports (anymore). This can happen after |
| 249 | /// merging reports together. |
| 250 | fn is_empty(&self) -> bool { |
| 251 | match self { |
| 252 | // A SenderReport always has, at least, the SenderInfo part. |
| 253 | Rtcp::SenderReport(_) => false, |
| 254 | // ReceiverReport can become empty. |
| 255 | Rtcp::ReceiverReport(v) => v.reports.is_empty(), |
| 256 | // ExtendedReport can become empty. |
| 257 | Rtcp::ExtendedReport(v) => v.blocks.is_empty(), |
| 258 | // SourceDescription can become empty. |
| 259 | Rtcp::SourceDescription(v) => v.reports.is_empty(), |
| 260 | // Goodbye can become empty, |
| 261 | Rtcp::Goodbye(v) => v.reports.is_empty(), |
| 262 | // Nack can become empty |
| 263 | Rtcp::Nack(v) => v.reports.is_empty(), |
| 264 | // Nack is never empty |
| 265 | Rtcp::Pli(_) => false, |
| 266 | // Fir can be merged to empty. |
| 267 | Rtcp::Fir(v) => v.reports.is_empty(), |
| 268 | // A twcc report is never empty. |
| 269 | Rtcp::Twcc(_) => false, |
| 270 | // A REMB report is never empty. |
| 271 | Rtcp::Remb(_) => false, |
| 272 | // An AppSpecificFeedback report is never empty. |
| 273 | Rtcp::AppSpecificFeedback(_) => false, |
| 274 | } |
| 275 | } |
| 276 | |
| 277 | fn pack(feedback: &mut VecDeque<Self>, mut word_capacity: usize) { |
| 278 | // Index into feedback of item we are to pack into. |
| 279 | let mut i = 0; |
| 280 | let len = feedback.len(); |
| 281 | |
| 282 | // Need at least on feedback to pack into, and one to take from. |
| 283 | if len < 2 { |
| 284 | return; |
| 285 | } |
| 286 | |
| 287 | // SenderReport/ReceiveReport first for SRTCP. |
| 288 | feedback.make_contiguous().sort_by_key(Self::order_no); |
| 289 | |
| 290 | 'outer: loop { |
| 291 | // If we reach last element, there is no more packing to do. |
| 292 | if i == len - 1 { |
| 293 | break; |
| 294 | } |
| 295 | |
| 296 | let (pack_into, pack_from) = feedback.make_contiguous().split_at_mut(i + 1); |
| 297 | let fb_a = pack_into.last_mut().unwrap(); |
| 298 | |
| 299 | // abort if fb_a won't fit in the spare capacity. |
| 300 | if word_capacity < fb_a.length_words() { |
| 301 | break 'outer; |
| 302 | } |
| 303 | |
| 304 | // if we manage to merge anything into fb_a. |
| 305 | let mut any_change = false; |
| 306 | |
| 307 | // fb_b goes from the item _after_ i |
| 308 | for fb_b in pack_from { |
| 309 | // if fb_a is full (or empty), we don't want to move any more elements into fb_a. |
| 310 | if fb_a.is_full() || fb_a.is_empty() { |
| 311 | break; |
| 312 | } |
| 313 | |
| 314 | // abort if fb_a won't fit in the spare capacity. |
| 315 | if word_capacity < fb_a.length_words() { |
| 316 | break 'outer; |
| 317 | } |
| 318 | |
| 319 | // amount of capacity (in words) left to fill. |
| 320 | let capacity = word_capacity - fb_a.length_words(); |
| 321 | |
| 322 | // attempt to merge some elements into fb_a from fb_b. |
| 323 | let did_merge = fb_a.merge(fb_b, capacity); |
| 324 | any_change |= did_merge; |
| 325 | } |
| 326 | |
| 327 | if !any_change { |
| 328 | word_capacity -= fb_a.length_words(); |
| 329 | i += 1; |
| 330 | } |
| 331 | } |
| 332 | |
| 333 | // Prune empty. |
| 334 | feedback.retain(|f| !f.is_empty()); |
| 335 | } |
| 336 | |
| 337 | fn order_no(&self) -> u8 { |
| 338 | use Rtcp::*; |
| 339 | match self { |
| 340 | // SenderReport/ReceiverReport first since they possibly contain |
| 341 | // the SSRC for the SRTCP encryption. |
| 342 | SenderReport(_) => 0, |
| 343 | ReceiverReport(_) => 1, |
| 344 | |
| 345 | SourceDescription(_) => 2, |
| 346 | Nack(_) => 3, |
| 347 | Pli(_) => 4, |
| 348 | Fir(_) => 5, |
| 349 | Twcc(_) => 6, |
| 350 | Remb(_) => 7, |
| 351 | AppSpecificFeedback(_) => 8, |
| 352 | ExtendedReport(_) => 10, |
| 353 | |
| 354 | // Goodbye last since they remove stuff. |
| 355 | Goodbye(_) => 11, |
| 356 | } |
| 357 | } |
| 358 | } |
| 359 | |
| 360 | impl RtcpPacket for Rtcp { |
| 361 | fn header(&self) -> RtcpHeader { |
| 362 | match self { |
| 363 | Rtcp::SenderReport(v) => v.header(), |
| 364 | Rtcp::ReceiverReport(v) => v.header(), |
| 365 | Rtcp::ExtendedReport(v) => v.header(), |
| 366 | Rtcp::SourceDescription(v) => v.header(), |
| 367 | Rtcp::Goodbye(v) => v.header(), |
| 368 | Rtcp::Nack(v) => v.header(), |
| 369 | Rtcp::Pli(v) => v.header(), |
| 370 | Rtcp::Fir(v) => v.header(), |
| 371 | Rtcp::Twcc(v) => v.header(), |
| 372 | Rtcp::Remb(v) => v.header(), |
| 373 | Rtcp::AppSpecificFeedback(v) => v.header(), |
| 374 | } |
| 375 | } |
| 376 | |
| 377 | fn length_words(&self) -> usize { |
| 378 | match self { |
| 379 | Rtcp::SenderReport(v) => v.length_words(), |
| 380 | Rtcp::ReceiverReport(v) => v.length_words(), |
| 381 | Rtcp::ExtendedReport(v) => v.length_words(), |
| 382 | Rtcp::SourceDescription(v) => v.length_words(), |
| 383 | Rtcp::Goodbye(v) => v.length_words(), |
| 384 | Rtcp::Nack(v) => v.length_words(), |
| 385 | Rtcp::Pli(v) => v.length_words(), |
| 386 | Rtcp::Fir(v) => v.length_words(), |
| 387 | Rtcp::Twcc(v) => v.length_words(), |
| 388 | Rtcp::Remb(v) => v.length_words(), |
| 389 | Rtcp::AppSpecificFeedback(v) => v.length_words(), |
| 390 | } |
| 391 | } |
| 392 | |
| 393 | fn write_to(&self, buf: &mut [u8]) -> usize { |
| 394 | match self { |
| 395 | Rtcp::SenderReport(v) => v.write_to(buf), |
| 396 | Rtcp::ReceiverReport(v) => v.write_to(buf), |
| 397 | Rtcp::ExtendedReport(v) => v.write_to(buf), |
| 398 | Rtcp::SourceDescription(v) => v.write_to(buf), |
| 399 | Rtcp::Goodbye(v) => v.write_to(buf), |
| 400 | Rtcp::Nack(v) => v.write_to(buf), |
| 401 | Rtcp::Pli(v) => v.write_to(buf), |
| 402 | Rtcp::Fir(v) => v.write_to(buf), |
| 403 | Rtcp::Twcc(v) => v.write_to(buf), |
| 404 | Rtcp::Remb(v) => v.write_to(buf), |
| 405 | Rtcp::AppSpecificFeedback(v) => v.write_to(buf), |
| 406 | } |
| 407 | } |
| 408 | } |
| 409 | |
| 410 | impl<'a> TryFrom<&'a [u8]> for Rtcp { |
| 411 | type Error = &'static str; |
| 412 | |
| 413 | fn try_from(buf: &'a [u8]) -> Result<Self, Self::Error> { |
| 414 | let header: RtcpHeader = buf.try_into()?; |
| 415 | |
| 416 | // By constraining the length, all subparsing can go |
| 417 | // until they exhaust the buffer length. This presupposes |
| 418 | // padding is removed from the input. |
| 419 | let buf = &buf[4..]; |
| 420 | |
| 421 | Ok(match header.rtcp_type() { |
| 422 | RtcpType::SenderReport => Rtcp::SenderReport(buf.try_into()?), |
| 423 | RtcpType::ReceiverReport => Rtcp::ReceiverReport(buf.try_into()?), |
| 424 | RtcpType::SourceDescription => Rtcp::SourceDescription(buf.try_into()?), |
| 425 | RtcpType::Goodbye => Rtcp::Goodbye((header.count(), buf).try_into()?), |
| 426 | RtcpType::ApplicationDefined => return Err("Ignore RTCP type: ApplicationDefined"), |
| 427 | RtcpType::TransportLayerFeedback => { |
| 428 | let tlfb = match header.feedback_message_type() { |
| 429 | FeedbackMessageType::TransportFeedback(v) => v, |
| 430 | _ => return Err("Expected TransportFeedback in FeedbackMessageType"), |
| 431 | }; |
| 432 | |
| 433 | match tlfb { |
| 434 | TransportType::Nack => Rtcp::Nack(buf.try_into()?), |
| 435 | TransportType::TransportWide => Rtcp::Twcc(buf.try_into()?), |
| 436 | } |
| 437 | } |
| 438 | RtcpType::PayloadSpecificFeedback => { |
| 439 | let plfb = match header.feedback_message_type() { |
| 440 | FeedbackMessageType::PayloadFeedback(v) => v, |
| 441 | _ => return Err("Expected PayloadFeedback in FeedbackMessageType"), |
| 442 | }; |
| 443 | |
| 444 | match plfb { |
| 445 | PayloadType::PictureLossIndication => Rtcp::Pli(buf.try_into()?), |
| 446 | PayloadType::SliceLossIndication => return Err("Ignore PayloadType type: SLI"), |
| 447 | PayloadType::ReferencePictureSelectionIndication => { |
| 448 | return Err("Ignore PayloadType type: RPSI"); |
| 449 | } |
| 450 | PayloadType::FullIntraRequest => Rtcp::Fir(buf.try_into()?), |
| 451 | PayloadType::ApplicationLayer => { |
| 452 | if header.rtcp_type() == RtcpType::PayloadSpecificFeedback { |
| 453 | if let Ok(remb) = Remb::try_from(buf) { |
| 454 | return Ok(Rtcp::Remb(remb)); |
| 455 | } |
| 456 | // Not REMB โ parse as generic application-specific feedback |
| 457 | if let Ok(fb) = AppSpecificFeedback::try_from(buf) { |
| 458 | return Ok(Rtcp::AppSpecificFeedback(fb)); |
| 459 | } |
| 460 | } |
| 461 | return Err("Ignore PayloadType: ApplicationLayer"); |
| 462 | } |
| 463 | } |
| 464 | } |
| 465 | RtcpType::ExtendedReport => Rtcp::ExtendedReport(buf.try_into()?), |
| 466 | }) |
| 467 | } |
| 468 | } |
| 469 | |
| 470 | impl WordSized for Ssrc { |
| 471 | fn word_size(&self) -> usize { |
| 472 | 1 |
| 473 | } |
| 474 | } |
| 475 | |
| 476 | /// Pad up to the next word (4 byte) boundary. |
| 477 | fn pad_bytes_to_word(n: usize) -> usize { |
| 478 | let pad = 4 - n % 4; |
| 479 | if pad == 4 { n } else { n + pad } |
| 480 | } |
| 481 | |
| 482 | #[cfg(test)] |
| 483 | mod test { |
| 484 | use std::time::{Duration, SystemTime}; |
| 485 | |
| 486 | use crate::rtp_::MediaTime; |
| 487 | |
| 488 | use super::twcc::{Delta, PacketChunk, PacketStatus}; |
| 489 | use super::*; |
| 490 | |
| 491 | #[test] |
| 492 | fn padding_of_rtcp() { |
| 493 | let mut queue = VecDeque::new(); |
| 494 | let mut twcc = Twcc { |
| 495 | sender_ssrc: 1.into(), |
| 496 | ssrc: 0.into(), |
| 497 | base_seq: 82, |
| 498 | status_count: 3, |
| 499 | reference_time: 25, |
| 500 | feedback_count: 17, |
| 501 | chunks: VecDeque::new(), |
| 502 | delta: VecDeque::new(), |
| 503 | }; |
| 504 | twcc.chunks |
| 505 | .push_back(PacketChunk::Run(PacketStatus::ReceivedSmallDelta, 3)); |
| 506 | twcc.delta.push_back(Delta::Small(0x7c)); |
| 507 | twcc.delta.push_back(Delta::Small(0x93)); |
| 508 | twcc.delta.push_back(Delta::Small(0x84)); |
| 509 | queue.push_back(Rtcp::Twcc(twcc)); |
| 510 | let mut buf = vec![0; 1500]; |
| 511 | let n = Rtcp::write_packet(&mut queue, &mut buf, |_| {}); |
| 512 | buf.truncate(n); |
| 513 | println!("{buf:02x?}"); |
| 514 | assert_eq!( |
| 515 | &buf, |
| 516 | &[ |
| 517 | // TWCC 0xaf got padding bit set |
| 518 | 0xaf, 0xcd, 0x00, 0x06, // |
| 519 | 0x00, 0x00, 0x00, 0x01, // sender SSRC |
| 520 | 0x00, 0x00, 0x00, 0x00, // media SSRC |
| 521 | 0x00, 0x52, // base seq |
| 522 | 0x00, 0x03, // status count |
| 523 | 0x00, 0x00, 0x19, // reference time |
| 524 | 0x11, // feedback count |
| 525 | 0x20, 0x03, // run of 3 |
| 526 | 0x7c, 0x93, 0x84, // three small delta |
| 527 | 0x00, 0x00, 0x03 // padding |
| 528 | ] |
| 529 | ); |
| 530 | } |
| 531 | |
| 532 | #[test] |
| 533 | fn pack_sr_4_rr() { |
| 534 | let now = SystemTime::now(); |
| 535 | let mut queue = VecDeque::new(); |
| 536 | queue.push_back(rr(3)); |
| 537 | queue.push_back(rr(4)); |
| 538 | queue.push_back(rr(5)); |
| 539 | queue.push_back(sr(1, now)); // should be sorted to front |
| 540 | |
| 541 | Rtcp::pack(&mut queue, 350); |
| 542 | |
| 543 | assert_eq!(queue.len(), 1); |
| 544 | |
| 545 | let sr = match queue.pop_front().unwrap() { |
| 546 | Rtcp::SenderReport(v) => v, |
| 547 | _ => unreachable!(), |
| 548 | }; |
| 549 | |
| 550 | assert_eq!(sr.reports.len(), 4); |
| 551 | let mut iter = sr.reports.iter(); |
| 552 | assert_eq!(iter.next().unwrap(), &report(2)); |
| 553 | assert_eq!(iter.next().unwrap(), &report(3)); |
| 554 | assert_eq!(iter.next().unwrap(), &report(4)); |
| 555 | assert_eq!(iter.next().unwrap(), &report(5)); |
| 556 | } |
| 557 | |
| 558 | #[test] |
| 559 | fn pack_4_rr() { |
| 560 | let mut queue = VecDeque::new(); |
| 561 | queue.push_back(rr(1)); |
| 562 | queue.push_back(rr(2)); |
| 563 | queue.push_back(rr(3)); |
| 564 | queue.push_back(rr(4)); |
| 565 | |
| 566 | Rtcp::pack(&mut queue, 350); |
| 567 | |
| 568 | assert_eq!(queue.len(), 1); |
| 569 | |
| 570 | let sr = match queue.pop_front().unwrap() { |
| 571 | Rtcp::ReceiverReport(v) => v, |
| 572 | _ => unreachable!(), |
| 573 | }; |
| 574 | |
| 575 | assert_eq!(sr.reports.len(), 4); |
| 576 | let mut iter = sr.reports.iter(); |
| 577 | assert_eq!(iter.next().unwrap(), &report(1)); |
| 578 | assert_eq!(iter.next().unwrap(), &report(2)); |
| 579 | assert_eq!(iter.next().unwrap(), &report(3)); |
| 580 | assert_eq!(iter.next().unwrap(), &report(4)); |
| 581 | } |
| 582 | |
| 583 | #[test] |
| 584 | fn roundtrip_sr_rr() { |
| 585 | let now = SystemTime::now(); |
| 586 | let mut feedback = VecDeque::new(); |
| 587 | feedback.push_back(sr(1, now)); |
| 588 | feedback.push_back(rr(3)); |
| 589 | feedback.push_back(rr(4)); |
| 590 | feedback.push_back(rr(5)); |
| 591 | |
| 592 | let mut buf = vec![0_u8; 1360]; |
| 593 | let n = Rtcp::write_packet(&mut feedback, &mut buf, |_| {}); |
| 594 | buf.truncate(n); |
| 595 | |
| 596 | let mut parsed = VecDeque::new(); |
| 597 | Rtcp::read_packet(&buf, &mut parsed); |
| 598 | |
| 599 | let Rtcp::SenderReport(s) = parsed.get(0).unwrap() else { |
| 600 | panic!("Not a SenderReport in Rtcp"); |
| 601 | }; |
| 602 | let now2 = s.sender_info.ntp_time; |
| 603 | |
| 604 | let mut compare = VecDeque::new(); |
| 605 | compare.push_back(sr(1, now2)); |
| 606 | compare.push_back(rr(3)); |
| 607 | compare.push_back(rr(4)); |
| 608 | compare.push_back(rr(5)); |
| 609 | Rtcp::pack(&mut compare, 1400); |
| 610 | |
| 611 | assert_eq!(parsed, compare); |
| 612 | |
| 613 | // Ensure ntp_time is not too far off. |
| 614 | let abs = abs_time_delta(now, now2); |
| 615 | assert!(abs < Duration::from_millis(1)); |
| 616 | } |
| 617 | |
| 618 | fn abs_time_delta(st1: SystemTime, st2: SystemTime) -> Duration { |
| 619 | let delta = if st1 > st2 { |
| 620 | st1.duration_since(st2) |
| 621 | } else { |
| 622 | st2.duration_since(st1) |
| 623 | }; |
| 624 | delta.expect("delta should be absolute") |
| 625 | } |
| 626 | |
| 627 | fn sr(ssrc: u32, ntp_time: SystemTime) -> Rtcp { |
| 628 | Rtcp::SenderReport(SenderReport { |
| 629 | sender_info: SenderInfo { |
| 630 | ssrc: ssrc.into(), |
| 631 | ntp_time, |
| 632 | rtp_time: MediaTime::from_secs(4), |
| 633 | sender_packet_count: 5, |
| 634 | sender_octet_count: 6, |
| 635 | }, |
| 636 | reports: report(2).into(), |
| 637 | }) |
| 638 | } |
| 639 | |
| 640 | fn rr(ssrc: u32) -> Rtcp { |
| 641 | Rtcp::ReceiverReport(ReceiverReport { |
| 642 | sender_ssrc: 42.into(), |
| 643 | reports: report(ssrc).into(), |
| 644 | }) |
| 645 | } |
| 646 | |
| 647 | fn report(ssrc: u32) -> ReceptionReport { |
| 648 | ReceptionReport { |
| 649 | ssrc: ssrc.into(), |
| 650 | fraction_lost: 3, |
| 651 | packets_lost: 1234, |
| 652 | max_seq: 4000, |
| 653 | jitter: 5, |
| 654 | last_sr_time: 12, |
| 655 | last_sr_delay: 1, |
| 656 | } |
| 657 | } |
| 658 | |
| 659 | // fn sdes(ssrc: u32) -> RtcpFb { |
| 660 | // RtcpFb::Sdes(Sdes { |
| 661 | // ssrc: ssrc.into(), |
| 662 | // values: vec![ |
| 663 | // (SdesType::NAME, "Martin".into()), |
| 664 | // (SdesType::TOOL, "str0m".into()), |
| 665 | // (SdesType::NOTE, "Writing things right here".into()), |
| 666 | // ], |
| 667 | // }) |
| 668 | // } |
| 669 | |
| 670 | // fn nack(ssrc: u32, pid: u16) -> RtcpFb { |
| 671 | // RtcpFb::Nack(Nack { |
| 672 | // ssrc: ssrc.into(), |
| 673 | // pid, |
| 674 | // blp: 0b1010_0101, |
| 675 | // }) |
| 676 | // } |
| 677 | |
| 678 | // fn gb(ssrc: u32) -> RtcpFb { |
| 679 | // RtcpFb::Goodbye(ssrc.into()) |
| 680 | // } |
| 681 | |
| 682 | #[test] |
| 683 | fn fuzz_failures() { |
| 684 | const TESTS: &[&[u8]] = &[ |
| 685 | // |
| 686 | &[133, 201, 0, 0], |
| 687 | &[191, 202, 54, 74], |
| 688 | &[166, 202, 0, 2, 218, 54, 214, 222, 160, 2, 146, 0, 251], |
| 689 | &[ |
| 690 | 151, 203, 0, 40, 88, 236, 217, 19, 82, 62, 73, 84, 112, 252, 69, 78, 38, 72, 43, 4, |
| 691 | 21, 136, 90, 29, 89, 70, 90, 196, 149, 168, 54, 1, 57, 16, 128, 8, 53, 172, 192, |
| 692 | 248, 175, 7, 92, 54, 82, 153, 179, 204, 181, 64, 94, 211, 67, 77, 110, 252, 181, |
| 693 | 18, 53, 48, 180, 179, 205, 234, 139, 61, 179, 54, 19, 120, 79, 119, 232, 208, 210, |
| 694 | 73, 78, 28, 242, 156, 242, 239, 19, 246, 183, 10, 49, 114, 216, 64, 105, 161, 50, |
| 695 | 99, 156, 113, 153, 90, 207, 53, 145, 96, 158, 198, 224, 114, 9, 20, 30, 156, 220, |
| 696 | 56, 151, 216, 164, 129, 156, 40, 85, 70, 189, 210, 146, 242, 242, 55, 70, 144, 113, |
| 697 | 9, 44, 74, 22, 123, 180, 153, 18, 88, 1, 185, 85, 227, 200, 62, 53, 142, 89, 28, |
| 698 | 37, 128, 223, 36, 248, 117, 26, 182, 173, 112, 42, 1, 2, 117, 203, 114, 179, |
| 699 | ], |
| 700 | &[ |
| 701 | 150, 202, 0, 54, 0, 149, 201, 0, 0, 138, 201, 0, 0, 152, 201, 0, 0, 151, 201, 0, 0, |
| 702 | 150, 201, 0, 0, 141, 201, 0, 0, 159, 201, 0, 0, 150, 201, 0, 0, 159, 201, 0, 0, |
| 703 | 134, 201, 0, 0, 143, 201, 0, 0, 162, 201, 0, 0, 166, 201, 0, 0, 177, 201, 0, 0, |
| 704 | 182, 201, 0, 0, 131, 201, 0, 0, 164, 201, 0, 0, 133, 201, 0, 0, 143, 201, 0, 0, |
| 705 | 174, 201, 0, 0, 186, 201, 0, 0, 165, 201, 0, 0, 173, 201, 0, 0, 186, 201, 0, 0, |
| 706 | 166, 201, 0, 0, 159, 201, 0, 0, 158, 201, 0, 0, 190, 201, 0, 0, 156, 201, 0, 0, |
| 707 | 147, 201, 0, 0, 169, 201, 0, 0, 135, 201, 0, 0, 148, 201, 0, 0, 132, 201, 0, 0, |
| 708 | 138, 201, 0, 0, 162, 201, 0, 0, 185, 201, 0, 0, 157, 201, 0, 0, 183, 201, 0, 0, |
| 709 | 145, 201, 0, 0, 130, 201, 0, 0, 183, 201, 0, 0, 152, 201, 0, 0, 153, 201, 0, 0, |
| 710 | 154, 201, 0, 0, 138, 201, 0, 0, 148, 201, 0, 0, 158, 201, 0, 0, 156, 201, 0, 0, |
| 711 | 181, 201, 0, 0, 173, 201, 0, 0, 171, 201, 0, 0, 169, 201, 0, 0, 167, 201, 41, 216, |
| 712 | ], |
| 713 | &[ |
| 714 | 143, 205, 0, 8, 143, 93, 208, 93, 201, 4, 131, 131, 131, 3, 0, 143, 1, 143, 0, 143, |
| 715 | 0, 80, 143, 231, 231, 0, 143, 181, 202, 0, 143, 236, 242, 0, 238, 21, |
| 716 | ], |
| 717 | ]; |
| 718 | |
| 719 | let mut parsed = VecDeque::new(); |
| 720 | |
| 721 | for t in TESTS { |
| 722 | parsed.clear(); |
| 723 | Rtcp::read_packet(t, &mut parsed); |
| 724 | } |
| 725 | } |
| 726 | } |