1 use super::*;
2 
3 use crate::api::media_engine::MIME_TYPE_VP8;
4 use crate::api::APIBuilder;
5 use crate::ice_transport::ice_candidate_pair::RTCIceCandidatePair;
6 use crate::rtp_transceiver::rtp_codec::RTCRtpCodecCapability;
7 use crate::stats::StatsReportType;
8 use crate::track::track_local::track_local_static_sample::TrackLocalStaticSample;
9 use bytes::Bytes;
10 use media::Sample;
11 use std::sync::atomic::AtomicU32;
12 use tokio::time::Duration;
13 use util::vnet::net::{Net, NetConfig};
14 use util::vnet::router::{Router, RouterConfig};
15 use waitgroup::WaitGroup;
16 
create_vnet_pair( ) -> Result<(RTCPeerConnection, RTCPeerConnection, Arc<Mutex<Router>>)>17 pub(crate) async fn create_vnet_pair(
18 ) -> Result<(RTCPeerConnection, RTCPeerConnection, Arc<Mutex<Router>>)> {
19     // Create a root router
20     let wan = Arc::new(Mutex::new(Router::new(RouterConfig {
21         cidr: "1.2.3.0/24".to_owned(),
22         ..Default::default()
23     })?));
24 
25     // Create a network interface for offerer
26     let offer_vnet = Arc::new(Net::new(Some(NetConfig {
27         static_ips: vec!["1.2.3.4".to_owned()],
28         ..Default::default()
29     })));
30 
31     // Add the network interface to the router
32     let nic = offer_vnet.get_nic()?;
33     {
34         let mut w = wan.lock().await;
35         w.add_net(Arc::clone(&nic)).await?;
36     }
37     {
38         let n = nic.lock().await;
39         n.set_router(Arc::clone(&wan)).await?;
40     }
41 
42     let mut offer_setting_engine = SettingEngine::default();
43     offer_setting_engine.set_vnet(Some(offer_vnet));
44     offer_setting_engine.set_ice_timeouts(
45         Some(Duration::from_secs(1)),
46         Some(Duration::from_secs(1)),
47         Some(Duration::from_millis(200)),
48     );
49 
50     // Create a network interface for answerer
51     let answer_vnet = Arc::new(Net::new(Some(NetConfig {
52         static_ips: vec!["1.2.3.5".to_owned()],
53         ..Default::default()
54     })));
55 
56     // Add the network interface to the router
57     let nic = answer_vnet.get_nic()?;
58     {
59         let mut w = wan.lock().await;
60         w.add_net(Arc::clone(&nic)).await?;
61     }
62     {
63         let n = nic.lock().await;
64         n.set_router(Arc::clone(&wan)).await?;
65     }
66 
67     let mut answer_setting_engine = SettingEngine::default();
68     answer_setting_engine.set_vnet(Some(answer_vnet));
69     answer_setting_engine.set_ice_timeouts(
70         Some(Duration::from_secs(1)),
71         Some(Duration::from_secs(1)),
72         Some(Duration::from_millis(200)),
73     );
74 
75     // Start the virtual network by calling Start() on the root router
76     {
77         let mut w = wan.lock().await;
78         w.start().await?;
79     }
80 
81     let mut offer_media_engine = MediaEngine::default();
82     offer_media_engine.register_default_codecs()?;
83     let offer_peer_connection = APIBuilder::new()
84         .with_setting_engine(offer_setting_engine)
85         .with_media_engine(offer_media_engine)
86         .build()
87         .new_peer_connection(RTCConfiguration::default())
88         .await?;
89 
90     let mut answer_media_engine = MediaEngine::default();
91     answer_media_engine.register_default_codecs()?;
92     let answer_peer_connection = APIBuilder::new()
93         .with_setting_engine(answer_setting_engine)
94         .with_media_engine(answer_media_engine)
95         .build()
96         .new_peer_connection(RTCConfiguration::default())
97         .await?;
98 
99     Ok((offer_peer_connection, answer_peer_connection, wan))
100 }
101 
102 /// new_pair creates two new peer connections (an offerer and an answerer)
103 /// *without* using an api (i.e. using the default settings).
new_pair(api: &API) -> Result<(RTCPeerConnection, RTCPeerConnection)>104 pub(crate) async fn new_pair(api: &API) -> Result<(RTCPeerConnection, RTCPeerConnection)> {
105     let pca = api.new_peer_connection(RTCConfiguration::default()).await?;
106     let pcb = api.new_peer_connection(RTCConfiguration::default()).await?;
107 
108     Ok((pca, pcb))
109 }
110 
signal_pair( pc_offer: &mut RTCPeerConnection, pc_answer: &mut RTCPeerConnection, ) -> Result<()>111 pub(crate) async fn signal_pair(
112     pc_offer: &mut RTCPeerConnection,
113     pc_answer: &mut RTCPeerConnection,
114 ) -> Result<()> {
115     // Note(albrow): We need to create a data channel in order to trigger ICE
116     // candidate gathering in the background for the JavaScript/Wasm bindings. If
117     // we don't do this, the complete offer including ICE candidates will never be
118     // generated.
119     pc_offer
120         .create_data_channel("initial_data_channel", None)
121         .await?;
122 
123     let offer = pc_offer.create_offer(None).await?;
124 
125     let mut offer_gathering_complete = pc_offer.gathering_complete_promise().await;
126     pc_offer.set_local_description(offer).await?;
127 
128     let _ = offer_gathering_complete.recv().await;
129 
130     pc_answer
131         .set_remote_description(
132             pc_offer
133                 .local_description()
134                 .await
135                 .ok_or(Error::new("non local description".to_owned()))?,
136         )
137         .await?;
138 
139     let answer = pc_answer.create_answer(None).await?;
140 
141     let mut answer_gathering_complete = pc_answer.gathering_complete_promise().await;
142     pc_answer.set_local_description(answer).await?;
143 
144     let _ = answer_gathering_complete.recv().await;
145 
146     pc_offer
147         .set_remote_description(
148             pc_answer
149                 .local_description()
150                 .await
151                 .ok_or(Error::new("non local description".to_owned()))?,
152         )
153         .await
154 }
155 
close_pair_now(pc1: &RTCPeerConnection, pc2: &RTCPeerConnection)156 pub(crate) async fn close_pair_now(pc1: &RTCPeerConnection, pc2: &RTCPeerConnection) {
157     let mut fail = false;
158     if let Err(err) = pc1.close().await {
159         log::error!("Failed to close PeerConnection: {}", err);
160         fail = true;
161     }
162     if let Err(err) = pc2.close().await {
163         log::error!("Failed to close PeerConnection: {}", err);
164         fail = true;
165     }
166 
167     assert!(!fail);
168 }
169 
close_pair( pc1: &RTCPeerConnection, pc2: &RTCPeerConnection, mut done_rx: mpsc::Receiver<()>, )170 pub(crate) async fn close_pair(
171     pc1: &RTCPeerConnection,
172     pc2: &RTCPeerConnection,
173     mut done_rx: mpsc::Receiver<()>,
174 ) {
175     let timeout = tokio::time::sleep(Duration::from_secs(10));
176     tokio::pin!(timeout);
177 
178     tokio::select! {
179         _ = timeout.as_mut() =>{
180             panic!("close_pair timed out waiting for done signal");
181         }
182         _ = done_rx.recv() =>{
183             close_pair_now(pc1, pc2).await;
184         }
185     }
186 }
187 
188 /*
189 func offerMediaHasDirection(offer SessionDescription, kind RTPCodecType, direction RTPTransceiverDirection) bool {
190     parsed := &sdp.SessionDescription{}
191     if err := parsed.Unmarshal([]byte(offer.SDP)); err != nil {
192         return false
193     }
194 
195     for _, media := range parsed.MediaDescriptions {
196         if media.MediaName.Media == kind.String() {
197             _, exists := media.Attribute(direction.String())
198             return exists
199         }
200     }
201     return false
202 }*/
203 
send_video_until_done( mut done_rx: mpsc::Receiver<()>, tracks: Vec<Arc<TrackLocalStaticSample>>, data: Bytes, max_sends: Option<usize>, ) -> bool204 pub(crate) async fn send_video_until_done(
205     mut done_rx: mpsc::Receiver<()>,
206     tracks: Vec<Arc<TrackLocalStaticSample>>,
207     data: Bytes,
208     max_sends: Option<usize>,
209 ) -> bool {
210     let mut sends = 0;
211 
212     loop {
213         let timeout = tokio::time::sleep(Duration::from_millis(20));
214         tokio::pin!(timeout);
215 
216         tokio::select! {
217             biased;
218 
219             _ = done_rx.recv() =>{
220                 log::debug!("sendVideoUntilDone received done");
221                 return false;
222             }
223 
224             _ = timeout.as_mut() =>{
225                 if max_sends.map(|s| sends >= s).unwrap_or(false) {
226                     continue;
227                 }
228 
229                 log::debug!("sendVideoUntilDone timeout");
230                 for track in &tracks {
231                     log::debug!("sendVideoUntilDone track.WriteSample");
232                     let result = track.write_sample(&Sample{
233                         data: data.clone(),
234                         duration: Duration::from_secs(1),
235                         ..Default::default()
236                     }).await;
237                     assert!(result.is_ok());
238                     sends += 1;
239                 }
240             }
241         }
242     }
243 }
244 
until_connection_state( pc: &mut RTCPeerConnection, wg: &WaitGroup, state: RTCPeerConnectionState, )245 pub(crate) async fn until_connection_state(
246     pc: &mut RTCPeerConnection,
247     wg: &WaitGroup,
248     state: RTCPeerConnectionState,
249 ) {
250     let w = Arc::new(Mutex::new(Some(wg.worker())));
251     pc.on_peer_connection_state_change(Box::new(move |pcs: RTCPeerConnectionState| {
252         let w2 = Arc::clone(&w);
253         Box::pin(async move {
254             if pcs == state {
255                 let mut worker = w2.lock().await;
256                 worker.take();
257             }
258         })
259     }));
260 }
261 
262 #[tokio::test]
test_get_stats() -> Result<()>263 async fn test_get_stats() -> Result<()> {
264     let mut m = MediaEngine::default();
265     m.register_default_codecs()?;
266     let api = APIBuilder::new().with_media_engine(m).build();
267 
268     let (mut pc_offer, mut pc_answer) = new_pair(&api).await?;
269 
270     let (ice_complete_tx, mut ice_complete_rx) = mpsc::channel::<()>(1);
271     let ice_complete_tx = Arc::new(Mutex::new(Some(ice_complete_tx)));
272     pc_answer.on_ice_connection_state_change(Box::new(move |ice_state: RTCIceConnectionState| {
273         let ice_complete_tx2 = Arc::clone(&ice_complete_tx);
274         Box::pin(async move {
275             if ice_state == RTCIceConnectionState::Connected {
276                 tokio::time::sleep(Duration::from_secs(1)).await;
277                 let mut done = ice_complete_tx2.lock().await;
278                 done.take();
279             }
280         })
281     }));
282 
283     let sender_called_candidate_change = Arc::new(AtomicU32::new(0));
284     let sender_called_candidate_change2 = Arc::clone(&sender_called_candidate_change);
285     pc_offer
286         .sctp()
287         .transport()
288         .ice_transport()
289         .on_selected_candidate_pair_change(Box::new(move |_: RTCIceCandidatePair| {
290             sender_called_candidate_change2.store(1, Ordering::SeqCst);
291             Box::pin(async {})
292         }));
293     let track = Arc::new(TrackLocalStaticSample::new(
294         RTCRtpCodecCapability {
295             mime_type: MIME_TYPE_VP8.to_owned(),
296             ..Default::default()
297         },
298         "video".to_owned(),
299         "webrtc-rs".to_owned(),
300     ));
301     pc_offer
302         .add_track(track.clone())
303         .await
304         .expect("Failed to add track");
305     let (packet_tx, packet_rx) = mpsc::channel(1);
306 
307     pc_answer.on_track(Box::new(move |track, _, _| {
308         let packet_tx = packet_tx.clone();
309         tokio::spawn(async move {
310             while let Ok((pkt, _)) = track.read_rtp().await {
311                 dbg!(&pkt);
312                 let last = pkt.payload[pkt.payload.len() - 1];
313 
314                 if last == 0xAA {
315                     let _ = packet_tx.send(()).await;
316                     break;
317                 }
318             }
319         });
320 
321         Box::pin(async move {})
322     }));
323 
324     signal_pair(&mut pc_offer, &mut pc_answer).await?;
325 
326     let _ = ice_complete_rx.recv().await;
327     send_video_until_done(
328         packet_rx,
329         vec![track],
330         Bytes::from_static(b"\xDE\xAD\xBE\xEF\xAA"),
331         Some(1),
332     )
333     .await;
334 
335     let offer_stats = pc_offer.get_stats().await;
336     assert!(!offer_stats.reports.is_empty());
337 
338     match offer_stats.reports.get("ice_transport") {
339         Some(StatsReportType::Transport(ice_transport_stats)) => {
340             assert!(ice_transport_stats.bytes_received > 0);
341             assert!(ice_transport_stats.bytes_sent > 0);
342         }
343         Some(_other) => panic!("found the wrong type"),
344         None => panic!("missed it"),
345     }
346     let outbound_stats = offer_stats
347         .reports
348         .values()
349         .find_map(|v| match v {
350             StatsReportType::OutboundRTP(d) => Some(d),
351             _ => None,
352         })
353         .expect("Should have produced an RTP Outbound stat");
354     assert_eq!(outbound_stats.packets_sent, 1);
355     assert_eq!(outbound_stats.kind, "video");
356     assert_eq!(outbound_stats.bytes_sent, 8);
357     assert_eq!(outbound_stats.header_bytes_sent, 12);
358 
359     let answer_stats = pc_answer.get_stats().await;
360     let inbound_stats = answer_stats
361         .reports
362         .values()
363         .find_map(|v| match v {
364             StatsReportType::InboundRTP(d) => Some(d),
365             _ => None,
366         })
367         .expect("Should have produced an RTP inbound stat");
368     assert_eq!(inbound_stats.packets_received, 1);
369     assert_eq!(inbound_stats.kind, "video");
370     assert_eq!(inbound_stats.bytes_received, 8);
371     assert_eq!(inbound_stats.header_bytes_received, 12);
372 
373     close_pair_now(&pc_offer, &pc_answer).await;
374 
375     Ok(())
376 }
377