1 use super::*; 2 use crate::mock::mock_stream::MockStream; 3 use crate::mock::mock_time::MockTime; 4 //use bytes::Bytes; 5 use chrono::prelude::*; 6 use rtp::extension::abs_send_time_extension::unix2ntp; 7 8 #[tokio::test] 9 async fn test_receiver_interceptor_before_any_packet() -> Result<()> { 10 let mt = Arc::new(MockTime::default()); 11 let time_gen = { 12 let mt = Arc::clone(&mt); 13 Arc::new(move || mt.now()) 14 }; 15 16 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 17 .with_interval(Duration::from_millis(50)) 18 .with_now_fn(time_gen) 19 .build("")?; 20 21 let stream = MockStream::new( 22 &StreamInfo { 23 ssrc: 123456, 24 clock_rate: 90000, 25 ..Default::default() 26 }, 27 icpr, 28 ) 29 .await; 30 31 let pkts = stream.written_rtcp().await.unwrap(); 32 assert_eq!(pkts.len(), 1); 33 34 if let Some(rr) = pkts[0] 35 .as_any() 36 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 37 { 38 assert_eq!(1, rr.reports.len()); 39 assert_eq!( 40 rtcp::reception_report::ReceptionReport { 41 ssrc: 123456, 42 last_sequence_number: 0, 43 last_sender_report: 0, 44 fraction_lost: 0, 45 total_lost: 0, 46 delay: 0, 47 jitter: 0, 48 }, 49 rr.reports[0] 50 ) 51 } else { 52 assert!(false); 53 } 54 55 stream.close().await?; 56 57 Ok(()) 58 } 59 60 #[tokio::test] 61 async fn test_receiver_interceptor_after_rtp_packets() -> Result<()> { 62 let mt = Arc::new(MockTime::default()); 63 let time_gen = { 64 let mt = Arc::clone(&mt); 65 Arc::new(move || mt.now()) 66 }; 67 68 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 69 .with_interval(Duration::from_millis(50)) 70 .with_now_fn(time_gen) 71 .build("")?; 72 73 let stream = MockStream::new( 74 &StreamInfo { 75 ssrc: 123456, 76 clock_rate: 90000, 77 ..Default::default() 78 }, 79 icpr, 80 ) 81 .await; 82 83 for i in 0..10u16 { 84 stream 85 .receive_rtp(rtp::packet::Packet { 86 header: rtp::header::Header { 87 sequence_number: i, 88 ..Default::default() 89 }, 90 ..Default::default() 91 }) 92 .await; 93 } 94 95 let pkts = stream.written_rtcp().await.unwrap(); 96 assert_eq!(pkts.len(), 1); 97 if let Some(rr) = pkts[0] 98 .as_any() 99 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 100 { 101 assert_eq!(1, rr.reports.len()); 102 assert_eq!( 103 rtcp::reception_report::ReceptionReport { 104 ssrc: 123456, 105 last_sequence_number: 9, 106 last_sender_report: 0, 107 fraction_lost: 0, 108 total_lost: 0, 109 delay: 0, 110 jitter: 0, 111 }, 112 rr.reports[0] 113 ) 114 } else { 115 assert!(false); 116 } 117 118 stream.close().await?; 119 120 Ok(()) 121 } 122 123 #[tokio::test] 124 async fn test_receiver_interceptor_after_rtp_and_rtcp_packets() -> Result<()> { 125 let rtp_time: SystemTime = Utc.with_ymd_and_hms(2009, 10, 23, 0, 0, 0).unwrap().into(); 126 127 let mt = Arc::new(MockTime::default()); 128 let time_gen = { 129 let mt = Arc::clone(&mt); 130 Arc::new(move || mt.now()) 131 }; 132 133 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 134 .with_interval(Duration::from_millis(50)) 135 .with_now_fn(time_gen) 136 .build("")?; 137 138 let stream = MockStream::new( 139 &StreamInfo { 140 ssrc: 123456, 141 clock_rate: 90000, 142 ..Default::default() 143 }, 144 icpr, 145 ) 146 .await; 147 148 for i in 0..10u16 { 149 stream 150 .receive_rtp(rtp::packet::Packet { 151 header: rtp::header::Header { 152 sequence_number: i, 153 ..Default::default() 154 }, 155 ..Default::default() 156 }) 157 .await; 158 } 159 160 let now: SystemTime = Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 1).unwrap().into(); 161 let rt = 987654321u32.wrapping_add( 162 (now.duration_since(rtp_time) 163 .unwrap_or(Duration::from_secs(0)) 164 .as_secs_f64() 165 * 90000.0) as u32, 166 ); 167 stream 168 .receive_rtcp(vec![Box::new(rtcp::sender_report::SenderReport { 169 ssrc: 123456, 170 ntp_time: unix2ntp(now), 171 rtp_time: rt, 172 packet_count: 10, 173 octet_count: 0, 174 ..Default::default() 175 })]) 176 .await; 177 178 let pkts = stream.written_rtcp().await.unwrap(); 179 assert_eq!(pkts.len(), 1); 180 if let Some(rr) = pkts[0] 181 .as_any() 182 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 183 { 184 assert_eq!(1, rr.reports.len()); 185 assert_eq!( 186 rtcp::reception_report::ReceptionReport { 187 ssrc: 123456, 188 last_sequence_number: 9, 189 last_sender_report: 1861287936, 190 fraction_lost: 0, 191 total_lost: 0, 192 delay: rr.reports[0].delay, 193 jitter: 0, 194 }, 195 rr.reports[0] 196 ) 197 } else { 198 assert!(false); 199 } 200 201 stream.close().await?; 202 203 Ok(()) 204 } 205 206 #[tokio::test] 207 async fn test_receiver_interceptor_overflow() -> Result<()> { 208 let mt = Arc::new(MockTime::default()); 209 let _mt2 = Arc::clone(&mt); 210 let time_gen = { 211 let mt = Arc::clone(&mt); 212 Arc::new(move || mt.now()) 213 }; 214 215 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 216 .with_interval(Duration::from_millis(50)) 217 .with_now_fn(time_gen) 218 .build("")?; 219 220 let stream = MockStream::new( 221 &StreamInfo { 222 ssrc: 123456, 223 clock_rate: 90000, 224 ..Default::default() 225 }, 226 icpr, 227 ) 228 .await; 229 230 stream 231 .receive_rtp(rtp::packet::Packet { 232 header: rtp::header::Header { 233 sequence_number: 0xffff, 234 ..Default::default() 235 }, 236 ..Default::default() 237 }) 238 .await; 239 240 stream 241 .receive_rtp(rtp::packet::Packet { 242 header: rtp::header::Header { 243 sequence_number: 0, 244 ..Default::default() 245 }, 246 ..Default::default() 247 }) 248 .await; 249 250 let pkts = stream.written_rtcp().await.unwrap(); 251 assert_eq!(pkts.len(), 1); 252 if let Some(rr) = pkts[0] 253 .as_any() 254 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 255 { 256 assert_eq!(1, rr.reports.len()); 257 assert_eq!( 258 rtcp::reception_report::ReceptionReport { 259 ssrc: 123456, 260 last_sequence_number: (1 << 16) | 0x0000, 261 last_sender_report: 0, 262 fraction_lost: 0, 263 total_lost: 0, 264 delay: rr.reports[0].delay, 265 jitter: 0, 266 }, 267 rr.reports[0] 268 ) 269 } else { 270 assert!(false); 271 } 272 273 stream.close().await?; 274 Ok(()) 275 } 276 277 #[tokio::test] 278 async fn test_receiver_interceptor_overflow_five_pkts() -> Result<()> { 279 let mt = Arc::new(MockTime::default()); 280 let time_gen = { 281 let mt = Arc::clone(&mt); 282 Arc::new(move || mt.now()) 283 }; 284 285 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 286 .with_interval(Duration::from_millis(50)) 287 .with_now_fn(time_gen) 288 .build("")?; 289 290 let stream = MockStream::new( 291 &StreamInfo { 292 ssrc: 123456, 293 clock_rate: 90000, 294 ..Default::default() 295 }, 296 icpr, 297 ) 298 .await; 299 300 stream 301 .receive_rtp(rtp::packet::Packet { 302 header: rtp::header::Header { 303 sequence_number: 0xfffd, 304 ..Default::default() 305 }, 306 ..Default::default() 307 }) 308 .await; 309 310 stream 311 .receive_rtp(rtp::packet::Packet { 312 header: rtp::header::Header { 313 sequence_number: 0xfffe, 314 ..Default::default() 315 }, 316 ..Default::default() 317 }) 318 .await; 319 320 stream 321 .receive_rtp(rtp::packet::Packet { 322 header: rtp::header::Header { 323 sequence_number: 0xffff, 324 ..Default::default() 325 }, 326 ..Default::default() 327 }) 328 .await; 329 330 stream 331 .receive_rtp(rtp::packet::Packet { 332 header: rtp::header::Header { 333 sequence_number: 0, 334 ..Default::default() 335 }, 336 ..Default::default() 337 }) 338 .await; 339 340 stream 341 .receive_rtp(rtp::packet::Packet { 342 header: rtp::header::Header { 343 sequence_number: 1, 344 ..Default::default() 345 }, 346 ..Default::default() 347 }) 348 .await; 349 350 let pkts = stream.written_rtcp().await.unwrap(); 351 assert_eq!(pkts.len(), 1); 352 if let Some(rr) = pkts[0] 353 .as_any() 354 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 355 { 356 assert_eq!(1, rr.reports.len()); 357 assert_eq!( 358 rtcp::reception_report::ReceptionReport { 359 ssrc: 123456, 360 last_sequence_number: (1 << 16) | 0x0001, 361 last_sender_report: 0, 362 fraction_lost: 0, 363 total_lost: 0, 364 delay: rr.reports[0].delay, 365 jitter: 0, 366 }, 367 rr.reports[0] 368 ) 369 } else { 370 assert!(false); 371 } 372 373 stream.close().await?; 374 Ok(()) 375 } 376 377 #[tokio::test] 378 async fn test_receiver_interceptor_packet_loss() -> Result<()> { 379 let rtp_time: SystemTime = Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 0).unwrap().into(); 380 381 let mt = Arc::new(MockTime::default()); 382 let time_gen = { 383 let mt = Arc::clone(&mt); 384 Arc::new(move || mt.now()) 385 }; 386 387 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 388 .with_interval(Duration::from_millis(50)) 389 .with_now_fn(time_gen) 390 .build("")?; 391 392 let stream = MockStream::new( 393 &StreamInfo { 394 ssrc: 123456, 395 clock_rate: 90000, 396 ..Default::default() 397 }, 398 icpr, 399 ) 400 .await; 401 402 stream 403 .receive_rtp(rtp::packet::Packet { 404 header: rtp::header::Header { 405 sequence_number: 0x01, 406 ..Default::default() 407 }, 408 ..Default::default() 409 }) 410 .await; 411 412 stream 413 .receive_rtp(rtp::packet::Packet { 414 header: rtp::header::Header { 415 sequence_number: 0x03, 416 ..Default::default() 417 }, 418 ..Default::default() 419 }) 420 .await; 421 422 let pkts = stream.written_rtcp().await.unwrap(); 423 assert_eq!(pkts.len(), 1); 424 if let Some(rr) = pkts[0] 425 .as_any() 426 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 427 { 428 assert_eq!(1, rr.reports.len()); 429 assert_eq!( 430 rtcp::reception_report::ReceptionReport { 431 ssrc: 123456, 432 last_sequence_number: 0x03, 433 last_sender_report: 0, 434 fraction_lost: (256u16 * 1 / 3) as u8, 435 total_lost: 1, 436 delay: 0, 437 jitter: 0, 438 }, 439 rr.reports[0] 440 ) 441 } else { 442 assert!(false); 443 } 444 445 let now: SystemTime = Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 1).unwrap().into(); 446 let rt = 987654321u32.wrapping_add( 447 (now.duration_since(rtp_time) 448 .unwrap_or(Duration::from_secs(0)) 449 .as_secs_f64() 450 * 90000.0) as u32, 451 ); 452 stream 453 .receive_rtcp(vec![Box::new(rtcp::sender_report::SenderReport { 454 ssrc: 123456, 455 ntp_time: unix2ntp(now), 456 rtp_time: rt, 457 packet_count: 10, 458 octet_count: 0, 459 ..Default::default() 460 })]) 461 .await; 462 463 let pkts = stream.written_rtcp().await.unwrap(); 464 assert_eq!(pkts.len(), 1); 465 if let Some(rr) = pkts[0] 466 .as_any() 467 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 468 { 469 assert_eq!(1, rr.reports.len()); 470 assert_eq!( 471 rtcp::reception_report::ReceptionReport { 472 ssrc: 123456, 473 last_sequence_number: 0x03, 474 last_sender_report: 1861287936, 475 fraction_lost: 0, 476 total_lost: 1, 477 delay: rr.reports[0].delay, 478 jitter: 0, 479 }, 480 rr.reports[0] 481 ) 482 } else { 483 assert!(false); 484 } 485 486 stream.close().await?; 487 Ok(()) 488 } 489 490 #[tokio::test] 491 async fn test_receiver_interceptor_overflow_and_packet_loss() -> Result<()> { 492 let mt = Arc::new(MockTime::default()); 493 let time_gen = { 494 let mt = Arc::clone(&mt); 495 Arc::new(move || mt.now()) 496 }; 497 498 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 499 .with_interval(Duration::from_millis(50)) 500 .with_now_fn(time_gen) 501 .build("")?; 502 503 let stream = MockStream::new( 504 &StreamInfo { 505 ssrc: 123456, 506 clock_rate: 90000, 507 ..Default::default() 508 }, 509 icpr, 510 ) 511 .await; 512 513 stream 514 .receive_rtp(rtp::packet::Packet { 515 header: rtp::header::Header { 516 sequence_number: 0xffff, 517 ..Default::default() 518 }, 519 ..Default::default() 520 }) 521 .await; 522 523 stream 524 .receive_rtp(rtp::packet::Packet { 525 header: rtp::header::Header { 526 sequence_number: 0x01, 527 ..Default::default() 528 }, 529 ..Default::default() 530 }) 531 .await; 532 533 let pkts = stream.written_rtcp().await.unwrap(); 534 assert_eq!(pkts.len(), 1); 535 if let Some(rr) = pkts[0] 536 .as_any() 537 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 538 { 539 assert_eq!(1, rr.reports.len()); 540 assert_eq!( 541 rtcp::reception_report::ReceptionReport { 542 ssrc: 123456, 543 last_sequence_number: 1 << 16 | 0x01, 544 last_sender_report: 0, 545 fraction_lost: (256u16 * 1 / 3) as u8, 546 total_lost: 1, 547 delay: 0, 548 jitter: 0, 549 }, 550 rr.reports[0] 551 ) 552 } else { 553 assert!(false); 554 } 555 556 stream.close().await?; 557 Ok(()) 558 } 559 560 #[tokio::test] 561 async fn test_receiver_interceptor_reordered_packets() -> Result<()> { 562 let mt = Arc::new(MockTime::default()); 563 let time_gen = { 564 let mt = Arc::clone(&mt); 565 Arc::new(move || mt.now()) 566 }; 567 568 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 569 .with_interval(Duration::from_millis(50)) 570 .with_now_fn(time_gen) 571 .build("")?; 572 573 let stream = MockStream::new( 574 &StreamInfo { 575 ssrc: 123456, 576 clock_rate: 90000, 577 ..Default::default() 578 }, 579 icpr, 580 ) 581 .await; 582 583 for sequence_number in [0x01, 0x03, 0x02, 0x04] { 584 stream 585 .receive_rtp(rtp::packet::Packet { 586 header: rtp::header::Header { 587 sequence_number, 588 ..Default::default() 589 }, 590 ..Default::default() 591 }) 592 .await; 593 } 594 595 let pkts = stream.written_rtcp().await.unwrap(); 596 assert_eq!(pkts.len(), 1); 597 if let Some(rr) = pkts[0] 598 .as_any() 599 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 600 { 601 assert_eq!(1, rr.reports.len()); 602 assert_eq!( 603 rtcp::reception_report::ReceptionReport { 604 ssrc: 123456, 605 last_sequence_number: 0x04, 606 last_sender_report: 0, 607 fraction_lost: 0, 608 total_lost: 0, 609 delay: 0, 610 jitter: 0, 611 }, 612 rr.reports[0] 613 ) 614 } else { 615 assert!(false); 616 } 617 618 stream.close().await?; 619 Ok(()) 620 } 621 622 #[tokio::test(start_paused = true)] 623 async fn test_receiver_interceptor_jitter() -> Result<()> { 624 let mt = Arc::new(MockTime::default()); 625 let time_gen = { 626 let mt = Arc::clone(&mt); 627 Arc::new(move || mt.now()) 628 }; 629 630 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 631 .with_interval(Duration::from_millis(50)) 632 .with_now_fn(time_gen) 633 .build("")?; 634 635 let stream = MockStream::new( 636 &StreamInfo { 637 ssrc: 123456, 638 clock_rate: 90000, 639 ..Default::default() 640 }, 641 icpr, 642 ) 643 .await; 644 645 mt.set_now(Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 0).unwrap().into()); 646 stream 647 .receive_rtp(rtp::packet::Packet { 648 header: rtp::header::Header { 649 sequence_number: 0x01, 650 timestamp: 42378934, 651 ..Default::default() 652 }, 653 ..Default::default() 654 }) 655 .await; 656 stream.read_rtp().await; 657 658 mt.set_now(Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 1).unwrap().into()); 659 stream 660 .receive_rtp(rtp::packet::Packet { 661 header: rtp::header::Header { 662 sequence_number: 0x02, 663 timestamp: 42378934 + 60000, 664 ..Default::default() 665 }, 666 ..Default::default() 667 }) 668 .await; 669 670 // Advance the time to generate a report 671 tokio::time::advance(Duration::from_millis(60)).await; 672 // Yield to let the reporting task run 673 tokio::task::yield_now().await; 674 675 let pkts = stream.last_written_rtcp().await.unwrap(); 676 assert_eq!(pkts.len(), 1); 677 if let Some(rr) = pkts[0] 678 .as_any() 679 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 680 { 681 assert_eq!(1, rr.reports.len()); 682 assert_eq!( 683 rtcp::reception_report::ReceptionReport { 684 ssrc: 123456, 685 last_sequence_number: 0x02, 686 last_sender_report: 0, 687 fraction_lost: 0, 688 total_lost: 0, 689 delay: 0, 690 jitter: 30000 / 16, 691 }, 692 rr.reports[0] 693 ) 694 } else { 695 assert!(false); 696 } 697 698 stream.close().await?; 699 Ok(()) 700 } 701 702 #[tokio::test] 703 async fn test_receiver_interceptor_delay() -> Result<()> { 704 let mt = Arc::new(MockTime::default()); 705 let time_gen = { 706 let mt = Arc::clone(&mt); 707 Arc::new(move || mt.now()) 708 }; 709 710 let icpr: Arc<dyn Interceptor + Send + Sync> = ReceiverReport::builder() 711 .with_interval(Duration::from_millis(50)) 712 .with_now_fn(time_gen) 713 .build("")?; 714 715 let stream = MockStream::new( 716 &StreamInfo { 717 ssrc: 123456, 718 clock_rate: 90000, 719 ..Default::default() 720 }, 721 icpr, 722 ) 723 .await; 724 725 mt.set_now(Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 0).unwrap().into()); 726 stream 727 .receive_rtcp(vec![Box::new(rtcp::sender_report::SenderReport { 728 ssrc: 123456, 729 ntp_time: unix2ntp(Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 0).unwrap().into()), 730 rtp_time: 987654321, 731 packet_count: 0, 732 octet_count: 0, 733 ..Default::default() 734 })]) 735 .await; 736 stream.read_rtcp().await; 737 738 mt.set_now(Utc.with_ymd_and_hms(2009, 11, 10, 23, 0, 1).unwrap().into()); 739 740 let pkts = stream.written_rtcp().await.unwrap(); 741 assert_eq!(pkts.len(), 1); 742 if let Some(rr) = pkts[0] 743 .as_any() 744 .downcast_ref::<rtcp::receiver_report::ReceiverReport>() 745 { 746 assert_eq!(1, rr.reports.len()); 747 assert_eq!( 748 rtcp::reception_report::ReceptionReport { 749 ssrc: 123456, 750 last_sequence_number: 0, 751 last_sender_report: 1861222400, 752 fraction_lost: 0, 753 total_lost: 0, 754 delay: 65536, 755 jitter: 0, 756 }, 757 rr.reports[0] 758 ) 759 } else { 760 assert!(false); 761 } 762 763 stream.close().await?; 764 Ok(()) 765 } 766