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 ¤t_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, ¤t_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 ¶ms.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