1 use bytes::Bytes;
2 use tokio::time::Instant;
3 
4 use super::*;
5 use crate::rtp_transceiver::create_stream_info;
6 use crate::stats::stats_collector::StatsCollector;
7 use crate::stats::{
8     InboundRTPStats, OutboundRTPStats, RTCStatsType, RemoteInboundRTPStats, RemoteOutboundRTPStats,
9     StatsReportType,
10 };
11 use crate::track::TrackStream;
12 use crate::{SDES_REPAIR_RTP_STREAM_ID_URI, SDP_ATTRIBUTE_RID};
13 use arc_swap::ArcSwapOption;
14 use std::collections::VecDeque;
15 use std::sync::atomic::AtomicIsize;
16 use std::sync::Weak;
17 
18 pub(crate) struct PeerConnectionInternal {
19     /// a value containing the last known greater mid value
20     /// we internally generate mids as numbers. Needed since JSEP
21     /// requires that when reusing a media section a new unique mid
22     /// should be defined (see JSEP 3.4.1).
23     pub(super) greater_mid: AtomicIsize,
24     pub(super) sdp_origin: Mutex<::sdp::description::session::Origin>,
25     pub(super) last_offer: Mutex<String>,
26     pub(super) last_answer: Mutex<String>,
27 
28     pub(super) on_negotiation_needed_handler: Arc<ArcSwapOption<Mutex<OnNegotiationNeededHdlrFn>>>,
29     pub(super) is_closed: Arc<AtomicBool>,
30 
31     /// ops is an operations queue which will ensure the enqueued actions are
32     /// executed in order. It is used for asynchronously, but serially processing
33     /// remote and local descriptions
34     pub(crate) ops: Arc<Operations>,
35     pub(super) negotiation_needed_state: Arc<AtomicU8>,
36     pub(super) is_negotiation_needed: Arc<AtomicBool>,
37     pub(super) signaling_state: Arc<AtomicU8>,
38 
39     pub(super) ice_transport: Arc<RTCIceTransport>,
40     pub(super) dtls_transport: Arc<RTCDtlsTransport>,
41     pub(super) on_peer_connection_state_change_handler:
42         Arc<ArcSwapOption<Mutex<OnPeerConnectionStateChangeHdlrFn>>>,
43     pub(super) peer_connection_state: Arc<AtomicU8>,
44     pub(super) ice_connection_state: Arc<AtomicU8>,
45 
46     pub(super) sctp_transport: Arc<RTCSctpTransport>,
47     pub(super) rtp_transceivers: Arc<Mutex<Vec<Arc<RTCRtpTransceiver>>>>,
48 
49     pub(super) on_track_handler: Arc<ArcSwapOption<Mutex<OnTrackHdlrFn>>>,
50     pub(super) on_signaling_state_change_handler:
51         ArcSwapOption<Mutex<OnSignalingStateChangeHdlrFn>>,
52     pub(super) on_ice_connection_state_change_handler:
53         Arc<ArcSwapOption<Mutex<OnICEConnectionStateChangeHdlrFn>>>,
54     pub(super) on_data_channel_handler: Arc<ArcSwapOption<Mutex<OnDataChannelHdlrFn>>>,
55 
56     pub(super) ice_gatherer: Arc<RTCIceGatherer>,
57 
58     pub(super) current_local_description: Arc<Mutex<Option<RTCSessionDescription>>>,
59     pub(super) current_remote_description: Arc<Mutex<Option<RTCSessionDescription>>>,
60     pub(super) pending_local_description: Arc<Mutex<Option<RTCSessionDescription>>>,
61     pub(super) pending_remote_description: Arc<Mutex<Option<RTCSessionDescription>>>,
62 
63     // A reference to the associated API state used by this connection
64     pub(super) setting_engine: Arc<SettingEngine>,
65     pub(crate) media_engine: Arc<MediaEngine>,
66     pub(super) interceptor: Weak<dyn Interceptor + Send + Sync>,
67     stats_interceptor: Arc<stats::StatsInterceptor>,
68 }
69 
70 impl PeerConnectionInternal {
new( api: &API, interceptor: Weak<dyn Interceptor + Send + Sync>, stats_interceptor: Arc<stats::StatsInterceptor>, mut configuration: RTCConfiguration, ) -> Result<(Arc<Self>, RTCConfiguration)>71     pub(super) async fn new(
72         api: &API,
73         interceptor: Weak<dyn Interceptor + Send + Sync>,
74         stats_interceptor: Arc<stats::StatsInterceptor>,
75         mut configuration: RTCConfiguration,
76     ) -> Result<(Arc<Self>, RTCConfiguration)> {
77         let mut pc = PeerConnectionInternal {
78             greater_mid: AtomicIsize::new(-1),
79             sdp_origin: Mutex::new(Default::default()),
80             last_offer: Mutex::new("".to_owned()),
81             last_answer: Mutex::new("".to_owned()),
82 
83             on_negotiation_needed_handler: Arc::new(ArcSwapOption::empty()),
84             ops: Arc::new(Operations::new()),
85             is_closed: Arc::new(AtomicBool::new(false)),
86             is_negotiation_needed: Arc::new(AtomicBool::new(false)),
87             negotiation_needed_state: Arc::new(AtomicU8::new(NegotiationNeededState::Empty as u8)),
88             signaling_state: Arc::new(AtomicU8::new(RTCSignalingState::Stable as u8)),
89             ice_transport: Arc::new(Default::default()),
90             dtls_transport: Arc::new(Default::default()),
91             ice_connection_state: Arc::new(AtomicU8::new(RTCIceConnectionState::New as u8)),
92             sctp_transport: Arc::new(Default::default()),
93             rtp_transceivers: Arc::new(Default::default()),
94             on_track_handler: Arc::new(ArcSwapOption::empty()),
95             on_signaling_state_change_handler: ArcSwapOption::empty(),
96             on_ice_connection_state_change_handler: Arc::new(ArcSwapOption::empty()),
97             on_data_channel_handler: Arc::new(Default::default()),
98             ice_gatherer: Arc::new(Default::default()),
99             current_local_description: Arc::new(Default::default()),
100             current_remote_description: Arc::new(Default::default()),
101             pending_local_description: Arc::new(Default::default()),
102             peer_connection_state: Arc::new(AtomicU8::new(RTCPeerConnectionState::New as u8)),
103 
104             setting_engine: Arc::clone(&api.setting_engine),
105             media_engine: if !api.setting_engine.disable_media_engine_copy {
106                 Arc::new(api.media_engine.clone_to())
107             } else {
108                 Arc::clone(&api.media_engine)
109             },
110             interceptor,
111             stats_interceptor,
112             on_peer_connection_state_change_handler: Arc::new(ArcSwapOption::empty()),
113             pending_remote_description: Arc::new(Default::default()),
114         };
115 
116         // Create the ice gatherer
117         pc.ice_gatherer = Arc::new(api.new_ice_gatherer(RTCIceGatherOptions {
118             ice_servers: configuration.get_ice_servers(),
119             ice_gather_policy: configuration.ice_transport_policy,
120         })?);
121 
122         // Create the ice transport
123         pc.ice_transport = pc.create_ice_transport(api).await;
124 
125         // Create the DTLS transport
126         let certificates = configuration.certificates.drain(..).collect();
127         pc.dtls_transport =
128             Arc::new(api.new_dtls_transport(Arc::clone(&pc.ice_transport), certificates)?);
129 
130         // Create the SCTP transport
131         pc.sctp_transport = Arc::new(api.new_sctp_transport(Arc::clone(&pc.dtls_transport))?);
132 
133         // Wire up the on datachannel handler
134         let on_data_channel_handler = Arc::clone(&pc.on_data_channel_handler);
135         pc.sctp_transport
136             .on_data_channel(Box::new(move |d: Arc<RTCDataChannel>| {
137                 let on_data_channel_handler2 = Arc::clone(&on_data_channel_handler);
138                 Box::pin(async move {
139                     if let Some(handler) = &*on_data_channel_handler2.load() {
140                         let mut f = handler.lock().await;
141                         f(d).await;
142                     }
143                 })
144             }));
145 
146         Ok((Arc::new(pc), configuration))
147     }
148 
start_rtp( self: &Arc<Self>, is_renegotiation: bool, remote_desc: Arc<RTCSessionDescription>, ) -> Result<()>149     pub(super) async fn start_rtp(
150         self: &Arc<Self>,
151         is_renegotiation: bool,
152         remote_desc: Arc<RTCSessionDescription>,
153     ) -> Result<()> {
154         let mut track_details = if let Some(parsed) = &remote_desc.parsed {
155             track_details_from_sdp(parsed, false)
156         } else {
157             vec![]
158         };
159 
160         let current_transceivers = {
161             let current_transceivers = self.rtp_transceivers.lock().await;
162             current_transceivers.clone()
163         };
164 
165         if !is_renegotiation {
166             self.undeclared_media_processor();
167         } else {
168             for t in &current_transceivers {
169                 let receiver = t.receiver();
170                 let tracks = receiver.tracks().await;
171                 if tracks.is_empty() {
172                     continue;
173                 }
174 
175                 let mut receiver_needs_stopped = false;
176 
177                 for t in tracks {
178                     if !t.rid().is_empty() {
179                         if let Some(details) =
180                             track_details_for_rid(&track_details, t.rid().to_owned())
181                         {
182                             t.set_id(details.id.clone());
183                             t.set_stream_id(details.stream_id.clone());
184                             continue;
185                         }
186                     } else if t.ssrc() != 0 {
187                         if let Some(details) = track_details_for_ssrc(&track_details, t.ssrc()) {
188                             t.set_id(details.id.clone());
189                             t.set_stream_id(details.stream_id.clone());
190                             continue;
191                         }
192                     }
193 
194                     receiver_needs_stopped = true;
195                 }
196 
197                 if !receiver_needs_stopped {
198                     continue;
199                 }
200 
201                 log::info!("Stopping receiver {:?}", receiver);
202                 if let Err(err) = receiver.stop().await {
203                     log::warn!("Failed to stop RtpReceiver: {}", err);
204                     continue;
205                 }
206 
207                 let interceptor = self
208                     .interceptor
209                     .upgrade()
210                     .ok_or(Error::ErrInterceptorNotBind)?;
211 
212                 let receiver = Arc::new(RTCRtpReceiver::new(
213                     self.setting_engine.get_receive_mtu(),
214                     receiver.kind(),
215                     Arc::clone(&self.dtls_transport),
216                     Arc::clone(&self.media_engine),
217                     interceptor,
218                 ));
219                 t.set_receiver(receiver);
220             }
221         }
222 
223         self.start_rtp_receivers(&mut track_details, &current_transceivers)
224             .await?;
225         if let Some(parsed) = &remote_desc.parsed {
226             if have_application_media_section(parsed) {
227                 self.start_sctp().await;
228             }
229         }
230 
231         Ok(())
232     }
233 
234     /// undeclared_media_processor handles RTP/RTCP packets that don't match any a:ssrc lines
undeclared_media_processor(self: &Arc<Self>)235     fn undeclared_media_processor(self: &Arc<Self>) {
236         let dtls_transport = Arc::clone(&self.dtls_transport);
237         let is_closed = Arc::clone(&self.is_closed);
238         let pci = Arc::clone(self);
239 
240         // SRTP acceptor
241         tokio::spawn(async move {
242             let simulcast_routine_count = Arc::new(AtomicU64::new(0));
243             loop {
244                 let srtp_session = match dtls_transport.get_srtp_session().await {
245                     Some(s) => s,
246                     None => {
247                         log::warn!("undeclared_media_processor failed to open SrtpSession");
248                         return;
249                     }
250                 };
251 
252                 let stream = match srtp_session.accept().await {
253                     Ok(stream) => stream,
254                     Err(err) => {
255                         log::warn!("Failed to accept RTP {}", err);
256                         return;
257                     }
258                 };
259 
260                 if is_closed.load(Ordering::SeqCst) {
261                     if let Err(err) = stream.close().await {
262                         log::warn!("Failed to close RTP stream {}", err);
263                     }
264                     continue;
265                 }
266 
267                 if simulcast_routine_count.fetch_add(1, Ordering::SeqCst) + 1
268                     >= SIMULCAST_MAX_PROBE_ROUTINES
269                 {
270                     simulcast_routine_count.fetch_sub(1, Ordering::SeqCst);
271                     log::warn!("{:?}", Error::ErrSimulcastProbeOverflow);
272                     continue;
273                 }
274 
275                 {
276                     let dtls_transport = Arc::clone(&dtls_transport);
277                     let simulcast_routine_count = Arc::clone(&simulcast_routine_count);
278                     let pci = Arc::clone(&pci);
279                     tokio::spawn(async move {
280                         let ssrc = stream.get_ssrc();
281 
282                         dtls_transport
283                             .store_simulcast_stream(ssrc, Arc::clone(&stream))
284                             .await;
285 
286                         if let Err(err) = pci.handle_incoming_ssrc(stream, ssrc).await {
287                             log::error!(
288                                 "Incoming unhandled RTP ssrc({}), on_track will not be fired. {}",
289                                 ssrc,
290                                 err
291                             );
292                         }
293 
294                         simulcast_routine_count.fetch_sub(1, Ordering::SeqCst);
295                     });
296                 }
297             }
298         });
299 
300         // SRTCP acceptor
301         {
302             let dtls_transport = Arc::clone(&self.dtls_transport);
303             tokio::spawn(async move {
304                 loop {
305                     let srtcp_session = match dtls_transport.get_srtcp_session().await {
306                         Some(s) => s,
307                         None => {
308                             log::warn!("undeclared_media_processor failed to open SrtcpSession");
309                             return;
310                         }
311                     };
312 
313                     let stream = match srtcp_session.accept().await {
314                         Ok(stream) => stream,
315                         Err(err) => {
316                             log::warn!("Failed to accept RTCP {}", err);
317                             return;
318                         }
319                     };
320                     log::warn!(
321                         "Incoming unhandled RTCP ssrc({}), on_track will not be fired",
322                         stream.get_ssrc()
323                     );
324                 }
325             });
326         }
327     }
328 
329     /// start_rtp_receivers opens knows inbound SRTP streams from the remote_description
start_rtp_receivers( self: &Arc<Self>, incoming_tracks: &mut Vec<TrackDetails>, local_transceivers: &[Arc<RTCRtpTransceiver>], ) -> Result<()>330     async fn start_rtp_receivers(
331         self: &Arc<Self>,
332         incoming_tracks: &mut Vec<TrackDetails>,
333         local_transceivers: &[Arc<RTCRtpTransceiver>],
334     ) -> Result<()> {
335         // Ensure we haven't already started a transceiver for this ssrc
336         let mut filtered_tracks = incoming_tracks.clone();
337         for incoming_track in incoming_tracks {
338             // If we already have a TrackRemote for a given SSRC don't handle it again
339             for t in local_transceivers {
340                 let receiver = t.receiver();
341                 for track in receiver.tracks().await {
342                     for ssrc in &incoming_track.ssrcs {
343                         if *ssrc == track.ssrc() {
344                             filter_track_with_ssrc(&mut filtered_tracks, track.ssrc());
345                         }
346                     }
347                 }
348             }
349         }
350 
351         let mut unhandled_tracks = vec![]; // filtered_tracks[:0]
352         for incoming_track in filtered_tracks.iter() {
353             let mut track_handled = false;
354             for t in local_transceivers {
355                 if t.mid().as_ref() != Some(&incoming_track.mid) {
356                     continue;
357                 }
358 
359                 if (incoming_track.kind != t.kind())
360                     || (t.direction() != RTCRtpTransceiverDirection::Recvonly
361                         && t.direction() != RTCRtpTransceiverDirection::Sendrecv)
362                 {
363                     continue;
364                 }
365 
366                 let receiver = t.receiver();
367                 if receiver.have_received().await {
368                     continue;
369                 }
370                 PeerConnectionInternal::start_receiver(
371                     self.setting_engine.get_receive_mtu(),
372                     incoming_track,
373                     receiver,
374                     Arc::clone(t),
375                     Arc::clone(&self.on_track_handler),
376                 )
377                 .await;
378                 track_handled = true;
379             }
380 
381             if !track_handled {
382                 unhandled_tracks.push(incoming_track);
383             }
384         }
385 
386         Ok(())
387     }
388 
389     /// Start SCTP subsystem
start_sctp(&self)390     async fn start_sctp(&self) {
391         // Start sctp
392         if let Err(err) = self
393             .sctp_transport
394             .start(SCTPTransportCapabilities {
395                 max_message_size: 0,
396             })
397             .await
398         {
399             log::warn!("Failed to start SCTP: {}", err);
400             if let Err(err) = self.sctp_transport.stop().await {
401                 log::warn!("Failed to stop SCTPTransport: {}", err);
402             }
403 
404             return;
405         }
406 
407         // DataChannels that need to be opened now that SCTP is available
408         // make a copy we may have incoming DataChannels mutating this while we open
409         let data_channels = {
410             let data_channels = self.sctp_transport.data_channels.lock().await;
411             data_channels.clone()
412         };
413 
414         let mut opened_dc_count = 0;
415         for d in data_channels {
416             if d.ready_state() == RTCDataChannelState::Connecting {
417                 if let Err(err) = d.open(Arc::clone(&self.sctp_transport)).await {
418                     log::warn!("failed to open data channel: {}", err);
419                     continue;
420                 }
421                 opened_dc_count += 1;
422             }
423         }
424 
425         self.sctp_transport
426             .data_channels_opened
427             .fetch_add(opened_dc_count, Ordering::SeqCst);
428     }
429 
add_transceiver_from_kind( &self, kind: RTPCodecType, init: Option<RTCRtpTransceiverInit>, ) -> Result<Arc<RTCRtpTransceiver>>430     pub(super) async fn add_transceiver_from_kind(
431         &self,
432         kind: RTPCodecType,
433         init: Option<RTCRtpTransceiverInit>,
434     ) -> Result<Arc<RTCRtpTransceiver>> {
435         if self.is_closed.load(Ordering::SeqCst) {
436             return Err(Error::ErrConnectionClosed);
437         }
438 
439         let direction = init
440             .map(|value| value.direction)
441             .unwrap_or(RTCRtpTransceiverDirection::Sendrecv);
442 
443         if direction == RTCRtpTransceiverDirection::Unspecified {
444             return Err(Error::ErrPeerConnAddTransceiverFromKindSupport);
445         }
446 
447         let interceptor = self
448             .interceptor
449             .upgrade()
450             .ok_or(Error::ErrInterceptorNotBind)?;
451         let receiver = Arc::new(RTCRtpReceiver::new(
452             self.setting_engine.get_receive_mtu(),
453             kind,
454             Arc::clone(&self.dtls_transport),
455             Arc::clone(&self.media_engine),
456             Arc::clone(&interceptor),
457         ));
458 
459         let sender = Arc::new(
460             RTCRtpSender::new(
461                 self.setting_engine.get_receive_mtu(),
462                 None,
463                 Arc::clone(&self.dtls_transport),
464                 Arc::clone(&self.media_engine),
465                 interceptor,
466                 false,
467             )
468             .await,
469         );
470 
471         let t = RTCRtpTransceiver::new(
472             receiver,
473             sender,
474             direction,
475             kind,
476             vec![],
477             Arc::clone(&self.media_engine),
478             Some(Box::new(self.make_negotiation_needed_trigger())),
479         )
480         .await;
481 
482         self.add_rtp_transceiver(Arc::clone(&t)).await;
483 
484         Ok(t)
485     }
486 
new_transceiver_from_track( &self, direction: RTCRtpTransceiverDirection, track: Arc<dyn TrackLocal + Send + Sync>, ) -> Result<Arc<RTCRtpTransceiver>>487     pub(super) async fn new_transceiver_from_track(
488         &self,
489         direction: RTCRtpTransceiverDirection,
490         track: Arc<dyn TrackLocal + Send + Sync>,
491     ) -> Result<Arc<RTCRtpTransceiver>> {
492         let interceptor = self
493             .interceptor
494             .upgrade()
495             .ok_or(Error::ErrInterceptorNotBind)?;
496 
497         if direction == RTCRtpTransceiverDirection::Unspecified {
498             return Err(Error::ErrPeerConnAddTransceiverFromTrackSupport);
499         }
500 
501         let r = Arc::new(RTCRtpReceiver::new(
502             self.setting_engine.get_receive_mtu(),
503             track.kind(),
504             Arc::clone(&self.dtls_transport),
505             Arc::clone(&self.media_engine),
506             Arc::clone(&interceptor),
507         ));
508 
509         let s = Arc::new(
510             RTCRtpSender::new(
511                 self.setting_engine.get_receive_mtu(),
512                 Some(Arc::clone(&track)),
513                 Arc::clone(&self.dtls_transport),
514                 Arc::clone(&self.media_engine),
515                 Arc::clone(&interceptor),
516                 false,
517             )
518             .await,
519         );
520 
521         Ok(RTCRtpTransceiver::new(
522             r,
523             s,
524             direction,
525             track.kind(),
526             vec![],
527             Arc::clone(&self.media_engine),
528             Some(Box::new(self.make_negotiation_needed_trigger())),
529         )
530         .await)
531     }
532 
533     /// add_rtp_transceiver appends t into rtp_transceivers
534     /// and fires onNegotiationNeeded;
535     /// caller of this method should hold `self.mu` lock
add_rtp_transceiver(&self, t: Arc<RTCRtpTransceiver>)536     pub(super) async fn add_rtp_transceiver(&self, t: Arc<RTCRtpTransceiver>) {
537         {
538             let mut rtp_transceivers = self.rtp_transceivers.lock().await;
539             rtp_transceivers.push(t);
540         }
541         self.trigger_negotiation_needed().await;
542     }
543 
544     /// Helper to trigger a negotiation needed.
trigger_negotiation_needed(&self)545     pub(crate) async fn trigger_negotiation_needed(&self) {
546         RTCPeerConnection::do_negotiation_needed(self.create_negotiation_needed_params()).await;
547     }
548 
549     /// Creates the parameters needed to trigger a negotiation needed.
create_negotiation_needed_params(&self) -> NegotiationNeededParams550     fn create_negotiation_needed_params(&self) -> NegotiationNeededParams {
551         NegotiationNeededParams {
552             on_negotiation_needed_handler: Arc::clone(&self.on_negotiation_needed_handler),
553             is_closed: Arc::clone(&self.is_closed),
554             ops: Arc::clone(&self.ops),
555             negotiation_needed_state: Arc::clone(&self.negotiation_needed_state),
556             is_negotiation_needed: Arc::clone(&self.is_negotiation_needed),
557             signaling_state: Arc::clone(&self.signaling_state),
558             check_negotiation_needed_params: CheckNegotiationNeededParams {
559                 sctp_transport: Arc::clone(&self.sctp_transport),
560                 rtp_transceivers: Arc::clone(&self.rtp_transceivers),
561                 current_local_description: Arc::clone(&self.current_local_description),
562                 current_remote_description: Arc::clone(&self.current_remote_description),
563             },
564         }
565     }
566 
make_negotiation_needed_trigger( &self, ) -> impl Fn() -> Pin<Box<dyn Future<Output = ()> + Send + Sync>> + Send + Sync567     pub(crate) fn make_negotiation_needed_trigger(
568         &self,
569     ) -> impl Fn() -> Pin<Box<dyn Future<Output = ()> + Send + Sync>> + Send + Sync {
570         let params = self.create_negotiation_needed_params();
571         move || {
572             let params = params.clone();
573             Box::pin(async move {
574                 let params = params.clone();
575                 RTCPeerConnection::do_negotiation_needed(params).await;
576             })
577         }
578     }
579 
remote_description(&self) -> Option<RTCSessionDescription>580     pub(super) async fn remote_description(&self) -> Option<RTCSessionDescription> {
581         let pending_remote_description = self.pending_remote_description.lock().await;
582         if pending_remote_description.is_some() {
583             pending_remote_description.clone()
584         } else {
585             let current_remote_description = self.current_remote_description.lock().await;
586             current_remote_description.clone()
587         }
588     }
589 
set_gather_complete_handler(&self, f: OnGatheringCompleteHdlrFn)590     pub(super) fn set_gather_complete_handler(&self, f: OnGatheringCompleteHdlrFn) {
591         self.ice_gatherer.on_gathering_complete(f);
592     }
593 
594     /// Start all transports. PeerConnection now has enough state
start_transports( self: &Arc<Self>, ice_role: RTCIceRole, dtls_role: DTLSRole, remote_ufrag: String, remote_pwd: String, fingerprint: String, fingerprint_hash: String, )595     pub(super) async fn start_transports(
596         self: &Arc<Self>,
597         ice_role: RTCIceRole,
598         dtls_role: DTLSRole,
599         remote_ufrag: String,
600         remote_pwd: String,
601         fingerprint: String,
602         fingerprint_hash: String,
603     ) {
604         // Start the ice transport
605         if let Err(err) = self
606             .ice_transport
607             .start(
608                 &RTCIceParameters {
609                     username_fragment: remote_ufrag,
610                     password: remote_pwd,
611                     ice_lite: false,
612                 },
613                 Some(ice_role),
614             )
615             .await
616         {
617             log::warn!("Failed to start manager ice: {}", err);
618             return;
619         }
620 
621         // Start the dtls_transport transport
622         let result = self
623             .dtls_transport
624             .start(DTLSParameters {
625                 role: dtls_role,
626                 fingerprints: vec![RTCDtlsFingerprint {
627                     algorithm: fingerprint_hash,
628                     value: fingerprint,
629                 }],
630             })
631             .await;
632         RTCPeerConnection::update_connection_state(
633             &self.on_peer_connection_state_change_handler,
634             &self.is_closed,
635             &self.peer_connection_state,
636             self.ice_connection_state.load(Ordering::SeqCst).into(),
637             self.dtls_transport.state(),
638         )
639         .await;
640         if let Err(err) = result {
641             log::warn!("Failed to start manager dtls: {}", err);
642         }
643     }
644 
645     /// generate_unmatched_sdp generates an SDP that doesn't take remote state into account
646     /// This is used for the initial call for CreateOffer
generate_unmatched_sdp( &self, local_transceivers: Vec<Arc<RTCRtpTransceiver>>, use_identity: bool, ) -> Result<SessionDescription>647     pub(super) async fn generate_unmatched_sdp(
648         &self,
649         local_transceivers: Vec<Arc<RTCRtpTransceiver>>,
650         use_identity: bool,
651     ) -> Result<SessionDescription> {
652         let d = SessionDescription::new_jsep_session_description(use_identity);
653 
654         let ice_params = self.ice_gatherer.get_local_parameters().await?;
655 
656         let candidates = self.ice_gatherer.get_local_candidates().await?;
657 
658         let mut media_sections = vec![];
659 
660         for t in &local_transceivers {
661             if t.stopped.load(Ordering::SeqCst) {
662                 // An "m=" section is generated for each
663                 // RtpTransceiver that has been added to the PeerConnection, excluding
664                 // any stopped RtpTransceivers;
665                 continue;
666             }
667 
668             // TODO: This is dubious because of rollbacks.
669             t.sender().set_negotiated();
670             media_sections.push(MediaSection {
671                 id: t.mid().unwrap(),
672                 transceivers: vec![Arc::clone(t)],
673                 ..Default::default()
674             });
675         }
676 
677         if self
678             .sctp_transport
679             .data_channels_requested
680             .load(Ordering::SeqCst)
681             != 0
682         {
683             media_sections.push(MediaSection {
684                 id: format!("{}", media_sections.len()),
685                 data: true,
686                 ..Default::default()
687             });
688         }
689 
690         let dtls_fingerprints = if let Some(cert) = self.dtls_transport.certificates.first() {
691             cert.get_fingerprints()
692         } else {
693             return Err(Error::ErrNonCertificate);
694         };
695 
696         let params = PopulateSdpParams {
697             media_description_fingerprint: self.setting_engine.sdp_media_level_fingerprints,
698             is_icelite: self.setting_engine.candidates.ice_lite,
699             connection_role: DEFAULT_DTLS_ROLE_OFFER.to_connection_role(),
700             ice_gathering_state: self.ice_gathering_state(),
701         };
702         populate_sdp(
703             d,
704             &dtls_fingerprints,
705             &self.media_engine,
706             &candidates,
707             &ice_params,
708             &media_sections,
709             params,
710         )
711         .await
712     }
713 
714     /// generate_matched_sdp generates a SDP and takes the remote state into account
715     /// this is used everytime we have a remote_description
generate_matched_sdp( &self, mut local_transceivers: Vec<Arc<RTCRtpTransceiver>>, use_identity: bool, include_unmatched: bool, connection_role: ConnectionRole, ) -> Result<SessionDescription>716     pub(super) async fn generate_matched_sdp(
717         &self,
718         mut local_transceivers: Vec<Arc<RTCRtpTransceiver>>,
719         use_identity: bool,
720         include_unmatched: bool,
721         connection_role: ConnectionRole,
722     ) -> Result<SessionDescription> {
723         let d = SessionDescription::new_jsep_session_description(use_identity);
724 
725         let ice_params = self.ice_gatherer.get_local_parameters().await?;
726         let candidates = self.ice_gatherer.get_local_candidates().await?;
727 
728         let remote_description = self.remote_description().await;
729         let mut media_sections = vec![];
730         let mut already_have_application_media_section = false;
731         if let Some(remote_description) = remote_description.as_ref() {
732             if let Some(parsed) = &remote_description.parsed {
733                 for media in &parsed.media_descriptions {
734                     if let Some(mid_value) = get_mid_value(media) {
735                         if mid_value.is_empty() {
736                             return Err(Error::ErrPeerConnRemoteDescriptionWithoutMidValue);
737                         }
738 
739                         if media.media_name.media == MEDIA_SECTION_APPLICATION {
740                             media_sections.push(MediaSection {
741                                 id: mid_value.to_owned(),
742                                 data: true,
743                                 ..Default::default()
744                             });
745                             already_have_application_media_section = true;
746                             continue;
747                         }
748 
749                         let kind = RTPCodecType::from(media.media_name.media.as_str());
750                         let direction = get_peer_direction(media);
751                         if kind == RTPCodecType::Unspecified
752                             || direction == RTCRtpTransceiverDirection::Unspecified
753                         {
754                             continue;
755                         }
756 
757                         if let Some(t) = find_by_mid(mid_value, &mut local_transceivers).await {
758                             t.sender().set_negotiated();
759                             let media_transceivers = vec![t];
760 
761                             // NB: The below could use `then_some`, but with our current MSRV
762                             // it's not possible to actually do this. The clippy version that
763                             // ships with 1.64.0 complains about this so we disable it for now.
764                             #[allow(clippy::unnecessary_lazy_evaluations)]
765                             media_sections.push(MediaSection {
766                                 id: mid_value.to_owned(),
767                                 transceivers: media_transceivers,
768                                 rid_map: get_rids(media),
769                                 offered_direction: (!include_unmatched).then(|| direction),
770                                 ..Default::default()
771                             });
772                         } else {
773                             return Err(Error::ErrPeerConnTranscieverMidNil);
774                         }
775                     }
776                 }
777             }
778         }
779 
780         // If we are offering also include unmatched local transceivers
781         if include_unmatched {
782             for t in &local_transceivers {
783                 t.sender().set_negotiated();
784                 media_sections.push(MediaSection {
785                     id: t.mid().unwrap(),
786                     transceivers: vec![Arc::clone(t)],
787                     ..Default::default()
788                 });
789             }
790 
791             if self
792                 .sctp_transport
793                 .data_channels_requested
794                 .load(Ordering::SeqCst)
795                 != 0
796                 && !already_have_application_media_section
797             {
798                 media_sections.push(MediaSection {
799                     id: format!("{}", media_sections.len()),
800                     data: true,
801                     ..Default::default()
802                 });
803             }
804         }
805 
806         let dtls_fingerprints = if let Some(cert) = self.dtls_transport.certificates.first() {
807             cert.get_fingerprints()
808         } else {
809             return Err(Error::ErrNonCertificate);
810         };
811 
812         let params = PopulateSdpParams {
813             media_description_fingerprint: self.setting_engine.sdp_media_level_fingerprints,
814             is_icelite: self.setting_engine.candidates.ice_lite,
815             connection_role,
816             ice_gathering_state: self.ice_gathering_state(),
817         };
818         populate_sdp(
819             d,
820             &dtls_fingerprints,
821             &self.media_engine,
822             &candidates,
823             &ice_params,
824             &media_sections,
825             params,
826         )
827         .await
828     }
829 
ice_gathering_state(&self) -> RTCIceGatheringState830     pub(super) fn ice_gathering_state(&self) -> RTCIceGatheringState {
831         match self.ice_gatherer.state() {
832             RTCIceGathererState::New => RTCIceGatheringState::New,
833             RTCIceGathererState::Gathering => RTCIceGatheringState::Gathering,
834             _ => RTCIceGatheringState::Complete,
835         }
836     }
837 
handle_undeclared_ssrc( self: &Arc<Self>, ssrc: SSRC, remote_description: &SessionDescription, ) -> Result<bool>838     async fn handle_undeclared_ssrc(
839         self: &Arc<Self>,
840         ssrc: SSRC,
841         remote_description: &SessionDescription,
842     ) -> Result<bool> {
843         if remote_description.media_descriptions.len() != 1 {
844             return Ok(false);
845         }
846 
847         let only_media_section = &remote_description.media_descriptions[0];
848         let mut stream_id = "";
849         let mut id = "";
850 
851         for a in &only_media_section.attributes {
852             match a.key.as_str() {
853                 ATTR_KEY_MSID => {
854                     if let Some(value) = &a.value {
855                         let split: Vec<&str> = value.split(' ').collect();
856                         if split.len() == 2 {
857                             stream_id = split[0];
858                             id = split[1];
859                         }
860                     }
861                 }
862                 ATTR_KEY_SSRC => return Err(Error::ErrPeerConnSingleMediaSectionHasExplicitSSRC),
863                 SDP_ATTRIBUTE_RID => return Ok(false),
864                 _ => {}
865             };
866         }
867 
868         let mut incoming = TrackDetails {
869             ssrcs: vec![ssrc],
870             kind: RTPCodecType::Video,
871             stream_id: stream_id.to_owned(),
872             id: id.to_owned(),
873             ..Default::default()
874         };
875         if only_media_section.media_name.media == RTPCodecType::Audio.to_string() {
876             incoming.kind = RTPCodecType::Audio;
877         }
878 
879         let t = self
880             .add_transceiver_from_kind(
881                 incoming.kind,
882                 Some(RTCRtpTransceiverInit {
883                     direction: RTCRtpTransceiverDirection::Sendrecv,
884                     send_encodings: vec![],
885                 }),
886             )
887             .await?;
888 
889         let receiver = t.receiver();
890         PeerConnectionInternal::start_receiver(
891             self.setting_engine.get_receive_mtu(),
892             &incoming,
893             receiver,
894             t,
895             Arc::clone(&self.on_track_handler),
896         )
897         .await;
898         Ok(true)
899     }
900 
handle_incoming_ssrc( self: &Arc<Self>, rtp_stream: Arc<Stream>, ssrc: SSRC, ) -> Result<()>901     async fn handle_incoming_ssrc(
902         self: &Arc<Self>,
903         rtp_stream: Arc<Stream>,
904         ssrc: SSRC,
905     ) -> Result<()> {
906         let parsed = match self.remote_description().await.and_then(|rd| rd.parsed) {
907             Some(r) => r,
908             None => return Err(Error::ErrPeerConnRemoteDescriptionNil),
909         };
910         // If the remote SDP was only one media section the ssrc doesn't have to be explicitly declared
911         let handled = self.handle_undeclared_ssrc(ssrc, &parsed).await?;
912         if handled {
913             return Ok(());
914         }
915 
916         // Get MID extension ID
917         let (mid_extension_id, audio_supported, video_supported) = self
918             .media_engine
919             .get_header_extension_id(RTCRtpHeaderExtensionCapability {
920                 uri: ::sdp::extmap::SDES_MID_URI.to_owned(),
921             })
922             .await;
923         if !audio_supported && !video_supported {
924             return Err(Error::ErrPeerConnSimulcastMidRTPExtensionRequired);
925         }
926 
927         // Get RID extension ID
928         let (sid_extension_id, audio_supported, video_supported) = self
929             .media_engine
930             .get_header_extension_id(RTCRtpHeaderExtensionCapability {
931                 uri: ::sdp::extmap::SDES_RTP_STREAM_ID_URI.to_owned(),
932             })
933             .await;
934         if !audio_supported && !video_supported {
935             return Err(Error::ErrPeerConnSimulcastStreamIDRTPExtensionRequired);
936         }
937 
938         let (rsid_extension_id, _, _) = self
939             .media_engine
940             .get_header_extension_id(RTCRtpHeaderExtensionCapability {
941                 uri: SDES_REPAIR_RTP_STREAM_ID_URI.to_owned(),
942             })
943             .await;
944 
945         let mut buf = vec![0u8; self.setting_engine.get_receive_mtu()];
946         // Packets that we read as part of simulcast probing that we need to make available
947         // if we do find a track later.
948         let mut buffered_packets: VecDeque<(Bytes, Attributes)> = VecDeque::default();
949 
950         let n = rtp_stream.read(&mut buf).await?;
951 
952         let (mut mid, mut rid, mut rsid, payload_type) = handle_unknown_rtp_packet(
953             &buf[..n],
954             mid_extension_id as u8,
955             sid_extension_id as u8,
956             rsid_extension_id as u8,
957         )?;
958         // TODO: Can we have attributes on the first packets?
959         buffered_packets.push_back((Bytes::copy_from_slice(&buf[..n]), Attributes::new()));
960 
961         let params = self
962             .media_engine
963             .get_rtp_parameters_by_payload_type(payload_type)
964             .await?;
965 
966         let icpr = match self.interceptor.upgrade() {
967             Some(i) => i,
968             None => return Err(Error::ErrInterceptorNotBind),
969         };
970 
971         let stream_info = create_stream_info(
972             "".to_owned(),
973             ssrc,
974             params.codecs[0].payload_type,
975             params.codecs[0].capability.clone(),
976             &params.header_extensions,
977         );
978         let (rtp_read_stream, rtp_interceptor, rtcp_read_stream, rtcp_interceptor) = self
979             .dtls_transport
980             .streams_for_ssrc(ssrc, &stream_info, &icpr)
981             .await?;
982 
983         let a = Attributes::new();
984         for _ in 0..=SIMULCAST_PROBE_COUNT {
985             if mid.is_empty() || (rid.is_empty() && rsid.is_empty()) {
986                 let (n, _) = rtp_interceptor.read(&mut buf, &a).await?;
987                 let (m, r, rs, _) = handle_unknown_rtp_packet(
988                     &buf[..n],
989                     mid_extension_id as u8,
990                     sid_extension_id as u8,
991                     rsid_extension_id as u8,
992                 )?;
993                 mid = m;
994                 rid = r;
995                 rsid = rs;
996 
997                 buffered_packets.push_back((Bytes::copy_from_slice(&buf[..n]), a.clone()));
998                 continue;
999             }
1000 
1001             let transceivers = self.rtp_transceivers.lock().await;
1002             for t in &*transceivers {
1003                 if t.mid().as_ref() != Some(&mid) {
1004                     continue;
1005                 }
1006 
1007                 let receiver = t.receiver();
1008 
1009                 if !rsid.is_empty() {
1010                     return receiver
1011                         .receive_for_rtx(
1012                             0,
1013                             rsid,
1014                             TrackStream {
1015                                 stream_info: Some(stream_info.clone()),
1016                                 rtp_read_stream: Some(rtp_read_stream),
1017                                 rtp_interceptor: Some(rtp_interceptor),
1018                                 rtcp_read_stream: Some(rtcp_read_stream),
1019                                 rtcp_interceptor: Some(rtcp_interceptor),
1020                             },
1021                         )
1022                         .await;
1023                 }
1024 
1025                 let track = receiver
1026                     .receive_for_rid(
1027                         rid,
1028                         params,
1029                         TrackStream {
1030                             stream_info: Some(stream_info.clone()),
1031                             rtp_read_stream: Some(rtp_read_stream),
1032                             rtp_interceptor: Some(rtp_interceptor),
1033                             rtcp_read_stream: Some(rtcp_read_stream),
1034                             rtcp_interceptor: Some(rtcp_interceptor),
1035                         },
1036                     )
1037                     .await?;
1038                 track.prepopulate_peeked_data(buffered_packets).await;
1039 
1040                 RTCPeerConnection::do_track(
1041                     Arc::clone(&self.on_track_handler),
1042                     track,
1043                     receiver,
1044                     Arc::clone(t),
1045                 );
1046                 return Ok(());
1047             }
1048         }
1049 
1050         let _ = rtp_read_stream.close().await;
1051         let _ = rtcp_read_stream.close().await;
1052         icpr.unbind_remote_stream(&stream_info).await;
1053         self.dtls_transport.remove_simulcast_stream(ssrc).await;
1054 
1055         Err(Error::ErrPeerConnSimulcastIncomingSSRCFailed)
1056     }
1057 
start_receiver( receive_mtu: usize, incoming: &TrackDetails, receiver: Arc<RTCRtpReceiver>, transceiver: Arc<RTCRtpTransceiver>, on_track_handler: Arc<ArcSwapOption<Mutex<OnTrackHdlrFn>>>, )1058     async fn start_receiver(
1059         receive_mtu: usize,
1060         incoming: &TrackDetails,
1061         receiver: Arc<RTCRtpReceiver>,
1062         transceiver: Arc<RTCRtpTransceiver>,
1063         on_track_handler: Arc<ArcSwapOption<Mutex<OnTrackHdlrFn>>>,
1064     ) {
1065         receiver.start(incoming).await;
1066         for t in receiver.tracks().await {
1067             if t.ssrc() == 0 {
1068                 return;
1069             }
1070 
1071             let receiver = Arc::clone(&receiver);
1072             let transceiver = Arc::clone(&transceiver);
1073             let on_track_handler = Arc::clone(&on_track_handler);
1074             tokio::spawn(async move {
1075                 if let Some(track) = receiver.track().await {
1076                     let mut b = vec![0u8; receive_mtu];
1077                     let n = match track.peek(&mut b).await {
1078                         Ok((n, _)) => n,
1079                         Err(err) => {
1080                             log::warn!(
1081                                 "Could not determine PayloadType for SSRC {} ({})",
1082                                 track.ssrc(),
1083                                 err
1084                             );
1085                             return;
1086                         }
1087                     };
1088 
1089                     if let Err(err) = track.check_and_update_track(&b[..n]).await {
1090                         log::warn!(
1091                             "Failed to set codec settings for track SSRC {} ({})",
1092                             track.ssrc(),
1093                             err
1094                         );
1095                         return;
1096                     }
1097 
1098                     RTCPeerConnection::do_track(on_track_handler, track, receiver, transceiver);
1099                 }
1100             });
1101         }
1102     }
1103 
create_ice_transport(&self, api: &API) -> Arc<RTCIceTransport>1104     pub(super) async fn create_ice_transport(&self, api: &API) -> Arc<RTCIceTransport> {
1105         let ice_transport = Arc::new(api.new_ice_transport(Arc::clone(&self.ice_gatherer)));
1106 
1107         let ice_connection_state = Arc::clone(&self.ice_connection_state);
1108         let peer_connection_state = Arc::clone(&self.peer_connection_state);
1109         let is_closed = Arc::clone(&self.is_closed);
1110         let dtls_transport = Arc::clone(&self.dtls_transport);
1111         let on_ice_connection_state_change_handler =
1112             Arc::clone(&self.on_ice_connection_state_change_handler);
1113         let on_peer_connection_state_change_handler =
1114             Arc::clone(&self.on_peer_connection_state_change_handler);
1115 
1116         ice_transport.on_connection_state_change(Box::new(move |state: RTCIceTransportState| {
1117             let cs = match state {
1118                 RTCIceTransportState::New => RTCIceConnectionState::New,
1119                 RTCIceTransportState::Checking => RTCIceConnectionState::Checking,
1120                 RTCIceTransportState::Connected => RTCIceConnectionState::Connected,
1121                 RTCIceTransportState::Completed => RTCIceConnectionState::Completed,
1122                 RTCIceTransportState::Failed => RTCIceConnectionState::Failed,
1123                 RTCIceTransportState::Disconnected => RTCIceConnectionState::Disconnected,
1124                 RTCIceTransportState::Closed => RTCIceConnectionState::Closed,
1125                 _ => {
1126                     log::warn!("on_connection_state_change: unhandled ICE state: {}", state);
1127                     return Box::pin(async {});
1128                 }
1129             };
1130 
1131             let ice_connection_state2 = Arc::clone(&ice_connection_state);
1132             let on_ice_connection_state_change_handler2 =
1133                 Arc::clone(&on_ice_connection_state_change_handler);
1134             let on_peer_connection_state_change_handler2 =
1135                 Arc::clone(&on_peer_connection_state_change_handler);
1136             let is_closed2 = Arc::clone(&is_closed);
1137             let dtls_transport_state = dtls_transport.state();
1138             let peer_connection_state2 = Arc::clone(&peer_connection_state);
1139             Box::pin(async move {
1140                 RTCPeerConnection::do_ice_connection_state_change(
1141                     &on_ice_connection_state_change_handler2,
1142                     &ice_connection_state2,
1143                     cs,
1144                 )
1145                 .await;
1146 
1147                 RTCPeerConnection::update_connection_state(
1148                     &on_peer_connection_state_change_handler2,
1149                     &is_closed2,
1150                     &peer_connection_state2,
1151                     cs,
1152                     dtls_transport_state,
1153                 )
1154                 .await;
1155             })
1156         }));
1157 
1158         ice_transport
1159     }
1160 
1161     /// has_local_description_changed returns whether local media (rtp_transceivers) has changed
1162     /// caller of this method should hold `pc.mu` lock
has_local_description_changed(&self, desc: &RTCSessionDescription) -> bool1163     pub(super) async fn has_local_description_changed(&self, desc: &RTCSessionDescription) -> bool {
1164         let rtp_transceivers = self.rtp_transceivers.lock().await;
1165         for t in &*rtp_transceivers {
1166             let m = match t.mid().and_then(|mid| get_by_mid(&mid, desc)) {
1167                 Some(m) => m,
1168                 None => return true,
1169             };
1170 
1171             if get_peer_direction(m) != t.direction() {
1172                 return true;
1173             }
1174         }
1175         false
1176     }
1177 
get_stats(&self, stats_id: String) -> StatsCollector1178     pub(super) async fn get_stats(&self, stats_id: String) -> StatsCollector {
1179         let collector = StatsCollector::new();
1180         let transceivers = { self.rtp_transceivers.lock().await.clone() };
1181 
1182         tokio::join!(
1183             self.ice_gatherer.collect_stats(&collector),
1184             self.ice_transport.collect_stats(&collector),
1185             self.sctp_transport.collect_stats(&collector, stats_id),
1186             self.dtls_transport.collect_stats(&collector),
1187             self.media_engine.collect_stats(&collector),
1188             self.collect_inbound_stats(&collector, transceivers.clone()),
1189             self.collect_outbound_stats(&collector, transceivers)
1190         );
1191 
1192         collector
1193     }
1194 
collect_inbound_stats( &self, collector: &StatsCollector, transceivers: Vec<Arc<RTCRtpTransceiver>>, )1195     async fn collect_inbound_stats(
1196         &self,
1197         collector: &StatsCollector,
1198         transceivers: Vec<Arc<RTCRtpTransceiver>>,
1199     ) {
1200         // TODO: There's a lot of await points here that could run concurrently with `futures::join_all`.
1201         struct TrackInfo {
1202             ssrc: SSRC,
1203             mid: String,
1204             track_id: String,
1205             kind: &'static str,
1206         }
1207         let mut track_infos = vec![];
1208         for transeiver in transceivers {
1209             let receiver = transeiver.receiver();
1210 
1211             if let Some(mid) = transeiver.mid() {
1212                 let tracks = receiver.tracks().await;
1213 
1214                 for track in tracks {
1215                     let track_id = track.id();
1216                     let kind = match track.kind() {
1217                         RTPCodecType::Unspecified => continue,
1218                         RTPCodecType::Audio => "audio",
1219                         RTPCodecType::Video => "video",
1220                     };
1221 
1222                     track_infos.push(TrackInfo {
1223                         ssrc: track.ssrc(),
1224                         mid: mid.clone(),
1225                         track_id,
1226                         kind,
1227                     });
1228                 }
1229             }
1230         }
1231 
1232         let stream_stats = self
1233             .stats_interceptor
1234             .fetch_inbound_stats(track_infos.iter().map(|t| t.ssrc).collect())
1235             .await;
1236 
1237         for (stats, info) in
1238             (stream_stats.into_iter().zip(track_infos)).filter_map(|(s, i)| s.map(|s| (s, i)))
1239         {
1240             let ssrc = info.ssrc;
1241             let kind = info.kind;
1242 
1243             let id = format!("RTCInboundRTP{}Stream_{}", capitalize(kind), ssrc);
1244             let (
1245                 packets_received,
1246                 header_bytes_received,
1247                 bytes_received,
1248                 last_packet_received_timestamp,
1249                 nack_count,
1250                 remote_packets_sent,
1251                 remote_bytes_sent,
1252                 remote_reports_sent,
1253                 remote_round_trip_time,
1254                 remote_total_round_trip_time,
1255                 remote_round_trip_time_measurements,
1256             ) = (
1257                 stats.packets_received(),
1258                 stats.header_bytes_received(),
1259                 stats.payload_bytes_received(),
1260                 stats.last_packet_received_timestamp(),
1261                 stats.nacks_sent(),
1262                 stats.remote_packets_sent(),
1263                 stats.remote_bytes_sent(),
1264                 stats.remote_reports_sent(),
1265                 stats.remote_round_trip_time(),
1266                 stats.remote_total_round_trip_time(),
1267                 stats.remote_round_trip_time_measurements(),
1268             );
1269 
1270             collector.insert(
1271                 id.clone(),
1272                 crate::stats::StatsReportType::InboundRTP(InboundRTPStats {
1273                     timestamp: Instant::now(),
1274                     stats_type: RTCStatsType::InboundRTP,
1275                     id: id.clone(),
1276                     ssrc,
1277                     kind,
1278                     packets_received,
1279                     track_identifier: info.track_id,
1280                     mid: info.mid,
1281                     last_packet_received_timestamp,
1282                     header_bytes_received,
1283                     bytes_received,
1284                     nack_count,
1285 
1286                     fir_count: (info.kind == "video").then(|| stats.firs_sent()),
1287                     pli_count: (info.kind == "video").then(|| stats.plis_sent()),
1288                 }),
1289             );
1290 
1291             let local_id = id;
1292             let id = format!(
1293                 "RTCRemoteOutboundRTP{}Stream_{}",
1294                 capitalize(info.kind),
1295                 info.ssrc
1296             );
1297             collector.insert(
1298                 id.clone(),
1299                 crate::stats::StatsReportType::RemoteOutboundRTP(RemoteOutboundRTPStats {
1300                     timestamp: Instant::now(),
1301                     stats_type: RTCStatsType::RemoteOutboundRTP,
1302                     id,
1303 
1304                     ssrc,
1305                     kind,
1306 
1307                     packets_sent: remote_packets_sent as u64,
1308                     bytes_sent: remote_bytes_sent as u64,
1309                     local_id,
1310                     reports_sent: remote_reports_sent,
1311                     round_trip_time: remote_round_trip_time,
1312                     total_round_trip_time: remote_total_round_trip_time,
1313                     round_trip_time_measurements: remote_round_trip_time_measurements,
1314                 }),
1315             );
1316         }
1317     }
1318 
collect_outbound_stats( &self, collector: &StatsCollector, transceivers: Vec<Arc<RTCRtpTransceiver>>, )1319     async fn collect_outbound_stats(
1320         &self,
1321         collector: &StatsCollector,
1322         transceivers: Vec<Arc<RTCRtpTransceiver>>,
1323     ) {
1324         // TODO: There's a lot of await points here that could run concurrently with `futures::join_all`.
1325         struct TrackInfo {
1326             track_id: String,
1327             ssrc: SSRC,
1328             mid: String,
1329             rid: Option<String>,
1330             kind: &'static str,
1331         }
1332         let mut track_infos = vec![];
1333         for transceiver in transceivers {
1334             let sender = transceiver.sender();
1335 
1336             let mid = match transceiver.mid() {
1337                 Some(mid) => mid,
1338                 None => continue,
1339             };
1340 
1341             let track = match sender.track().await {
1342                 Some(track) => track,
1343                 None => continue,
1344             };
1345 
1346             let track_id = track.id().to_string();
1347             let kind = match track.kind() {
1348                 RTPCodecType::Unspecified => continue,
1349                 RTPCodecType::Audio => "audio",
1350                 RTPCodecType::Video => "video",
1351             };
1352 
1353             track_infos.push(TrackInfo {
1354                 track_id,
1355                 ssrc: sender.ssrc,
1356                 mid: mid.clone(),
1357                 rid: None,
1358                 kind,
1359             });
1360         }
1361 
1362         let stream_stats = self
1363             .stats_interceptor
1364             .fetch_outbound_stats(track_infos.iter().map(|t| t.ssrc).collect())
1365             .await;
1366 
1367         for (stats, info) in stream_stats
1368             .into_iter()
1369             .zip(track_infos)
1370             .filter_map(|(s, i)| s.map(|s| (s, i)))
1371         {
1372             // RTCOutboundRtpStreamStats
1373             let id = format!(
1374                 "RTCOutboundRTP{}Stream_{}",
1375                 capitalize(info.kind),
1376                 info.ssrc
1377             );
1378             let (
1379                 packets_sent,
1380                 bytes_sent,
1381                 header_bytes_sent,
1382                 nack_count,
1383                 remote_inbound_packets_received,
1384                 remote_inbound_packets_lost,
1385                 remote_rtt_ms,
1386                 remote_total_rtt_ms,
1387                 remote_rtt_measurements,
1388                 remote_fraction_lost,
1389             ) = (
1390                 stats.packets_sent(),
1391                 stats.payload_bytes_sent(),
1392                 stats.header_bytes_sent(),
1393                 stats.nacks_received(),
1394                 stats.remote_packets_received(),
1395                 stats.remote_total_lost(),
1396                 stats.remote_round_trip_time(),
1397                 stats.remote_total_round_trip_time(),
1398                 stats.remote_round_trip_time_measurements(),
1399                 stats.remote_fraction_lost(),
1400             );
1401 
1402             let TrackInfo {
1403                 mid,
1404                 ssrc,
1405                 rid,
1406                 kind,
1407                 track_id: track_identifier,
1408             } = info;
1409 
1410             collector.insert(
1411                 id.clone(),
1412                 crate::stats::StatsReportType::OutboundRTP(OutboundRTPStats {
1413                     timestamp: Instant::now(),
1414                     stats_type: RTCStatsType::OutboundRTP,
1415                     track_identifier,
1416                     id: id.clone(),
1417                     ssrc,
1418                     kind,
1419                     packets_sent,
1420                     mid,
1421                     rid,
1422                     header_bytes_sent,
1423                     bytes_sent,
1424                     nack_count,
1425 
1426                     fir_count: (info.kind == "video").then(|| stats.firs_received()),
1427                     pli_count: (info.kind == "video").then(|| stats.plis_received()),
1428                 }),
1429             );
1430 
1431             let local_id = id;
1432             let id = format!(
1433                 "RTCRemoteInboundRTP{}Stream_{}",
1434                 capitalize(info.kind),
1435                 info.ssrc
1436             );
1437 
1438             collector.insert(
1439                 id.clone(),
1440                 StatsReportType::RemoteInboundRTP(RemoteInboundRTPStats {
1441                     timestamp: Instant::now(),
1442                     stats_type: RTCStatsType::RemoteInboundRTP,
1443                     id,
1444                     ssrc,
1445                     kind,
1446 
1447                     packets_received: remote_inbound_packets_received,
1448                     packets_lost: remote_inbound_packets_lost as i64,
1449 
1450                     local_id,
1451 
1452                     round_trip_time: remote_rtt_ms,
1453                     total_round_trip_time: remote_total_rtt_ms,
1454                     fraction_lost: remote_fraction_lost.unwrap_or(0.0),
1455                     round_trip_time_measurements: remote_rtt_measurements,
1456                 }),
1457             );
1458         }
1459     }
1460 }
1461 
1462 type IResult<T> = std::result::Result<T, interceptor::Error>;
1463 
1464 #[async_trait]
1465 impl RTCPWriter for PeerConnectionInternal {
write( &self, pkts: &[Box<dyn rtcp::packet::Packet + Send + Sync>], _a: &Attributes, ) -> IResult<usize>1466     async fn write(
1467         &self,
1468         pkts: &[Box<dyn rtcp::packet::Packet + Send + Sync>],
1469         _a: &Attributes,
1470     ) -> IResult<usize> {
1471         Ok(self.dtls_transport.write_rtcp(pkts).await?)
1472     }
1473 }
1474 
capitalize(s: &str) -> String1475 fn capitalize(s: &str) -> String {
1476     let first = s
1477         .chars()
1478         .next()
1479         .expect("Must have at least one character to uppercase")
1480         .to_uppercase();
1481     let mut result = String::new();
1482 
1483     result.extend(first);
1484     result.extend(s.chars().skip(1));
1485 
1486     result
1487 }
1488