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