File
Blob: firmware/vendor/str0m/tests/rtp-direct-with-roc.rs
| 1 | use std::collections::VecDeque; |
| 2 | use std::time::Duration; |
| 3 | |
| 4 | use str0m::format::Codec; |
| 5 | use str0m::media::MediaKind; |
| 6 | use str0m::rtp::SeqNo; |
| 7 | use str0m::rtp::{ExtensionValues, RtpWrite, Ssrc}; |
| 8 | use str0m::{Event, RtcError}; |
| 9 | |
| 10 | mod common; |
| 11 | use common::{connect_l_r, init_crypto_default, init_log, progress}; |
| 12 | |
| 13 | #[test] |
| 14 | pub fn rtp_direct_with_roc() -> Result<(), RtcError> { |
| 15 | init_log(); |
| 16 | init_crypto_default(); |
| 17 | |
| 18 | let (mut l, mut r) = connect_l_r(); |
| 19 | |
| 20 | let mid = "aud".into(); |
| 21 | |
| 22 | // In this example we are using MID only (no RID) to identify the incoming media. |
| 23 | let ssrc_tx: Ssrc = 42.into(); |
| 24 | |
| 25 | l.direct_api().declare_media(mid, MediaKind::Audio); |
| 26 | |
| 27 | l.direct_api().declare_stream_tx(ssrc_tx, None, mid, None); |
| 28 | |
| 29 | r.direct_api().declare_media(mid, MediaKind::Audio); |
| 30 | |
| 31 | let mut d = r.direct_api(); |
| 32 | let rx = d.expect_stream_rx(ssrc_tx, None, mid, None); |
| 33 | |
| 34 | // Above 2^16, which means we have ROC:ed. |
| 35 | let seq_no_offset: SeqNo = 100_000.into(); |
| 36 | |
| 37 | // By telling the receiver side to start at a specific ROC, we can send first ever |
| 38 | // packet from a high sequence number. |
| 39 | rx.reset_roc(seq_no_offset.roc()); |
| 40 | |
| 41 | let max = l.last.max(r.last); |
| 42 | l.last = max; |
| 43 | r.last = max; |
| 44 | |
| 45 | let params = l.params_opus(); |
| 46 | let ssrc = l.direct_api().stream_tx_by_mid(mid, None).unwrap().ssrc(); |
| 47 | assert_eq!(params.spec().codec, Codec::Opus); |
| 48 | let pt = params.pt(); |
| 49 | |
| 50 | let to_write: Vec<&[u8]> = vec![ |
| 51 | // 1 |
| 52 | &[0x1, 0x2, 0x3, 0x4], |
| 53 | // 3 |
| 54 | &[0x9, 0xa, 0xb, 0xc], |
| 55 | // 2 |
| 56 | &[0x5, 0x6, 0x7, 0x8], |
| 57 | ]; |
| 58 | |
| 59 | let mut to_write: VecDeque<_> = to_write.into(); |
| 60 | |
| 61 | let mut write_at = l.last + Duration::from_millis(300); |
| 62 | |
| 63 | let mut counts: Vec<u64> = vec![0, 3, 1]; |
| 64 | |
| 65 | loop { |
| 66 | if l.start + l.duration() > write_at { |
| 67 | write_at = l.last + Duration::from_millis(300); |
| 68 | if let Some(packet) = to_write.pop_front() { |
| 69 | let wallclock = l.start + l.duration(); |
| 70 | |
| 71 | let mut direct = l.direct_api(); |
| 72 | let stream = direct.stream_tx(&ssrc).unwrap(); |
| 73 | |
| 74 | let count = counts.remove(0); |
| 75 | let time = (count * 1000 + 47_000_000) as u32; |
| 76 | |
| 77 | // this seqno is already past first ROC |
| 78 | let seq_no = (*seq_no_offset + count).into(); |
| 79 | |
| 80 | let exts = ExtensionValues { |
| 81 | audio_level: Some(-42 - count as i8), |
| 82 | voice_activity: Some(false), |
| 83 | ..Default::default() |
| 84 | }; |
| 85 | |
| 86 | stream.write_rtp(RtpWrite::new(pt, seq_no, time, wallclock, packet).ext_vals(exts)); |
| 87 | } |
| 88 | } |
| 89 | |
| 90 | progress(&mut l, &mut r)?; |
| 91 | |
| 92 | if l.duration() > Duration::from_secs(10) { |
| 93 | break; |
| 94 | } |
| 95 | } |
| 96 | |
| 97 | let media: Vec<_> = r |
| 98 | .events |
| 99 | .iter() |
| 100 | .filter_map(|(_, e)| { |
| 101 | if let Event::RtpPacket(v) = e { |
| 102 | Some(v) |
| 103 | } else { |
| 104 | None |
| 105 | } |
| 106 | }) |
| 107 | .collect(); |
| 108 | |
| 109 | assert_eq!(media.len(), 3); |
| 110 | |
| 111 | assert!(l.media(mid).is_some()); |
| 112 | assert!(l.direct_api().stream_tx_by_mid(mid, None).is_some()); |
| 113 | l.direct_api().remove_media(mid); |
| 114 | assert!(l.media(mid).is_none()); |
| 115 | assert!(l.direct_api().stream_tx_by_mid(mid, None).is_none()); |
| 116 | |
| 117 | assert!(r.media(mid).is_some()); |
| 118 | assert!(r.direct_api().stream_rx_by_mid(mid, None).is_some()); |
| 119 | r.direct_api().remove_media(mid); |
| 120 | assert!(r.media(mid).is_none()); |
| 121 | assert!(r.direct_api().stream_rx_by_mid(mid, None).is_none()); |
| 122 | |
| 123 | Ok(()) |
| 124 | } |