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