File
Blob: firmware/vendor/str0m/src/session.rs
| 1 | use std::collections::{HashMap, VecDeque}; |
| 2 | use std::ops::RangeInclusive; |
| 3 | use std::sync::Arc; |
| 4 | use std::time::{Duration, Instant}; |
| 5 | |
| 6 | use crate::Event; |
| 7 | use crate::bwe::BweKind; |
| 8 | use crate::bwe_::Bwe; |
| 9 | use crate::config::KeyingMaterial; |
| 10 | use crate::config_mod::RtcpReportIntervals; |
| 11 | use crate::crypto::CryptoProvider; |
| 12 | use crate::crypto::dtls::SrtpProfile; |
| 13 | use crate::format::CodecConfig; |
| 14 | use crate::format::PayloadParams; |
| 15 | use crate::format::Vp9PacketizerMode; |
| 16 | use crate::io::DatagramSend; |
| 17 | use crate::media::AppSpecificFeedback; |
| 18 | use crate::media::Media; |
| 19 | use crate::media::{KeyframeRequestKind, MID_PROBE}; |
| 20 | use crate::media::{MediaAdded, MediaChanged}; |
| 21 | use crate::pacer::PacerControl; |
| 22 | use crate::pacer::{Pacer, PacerImpl}; |
| 23 | use crate::rtp::{Extension, RawPacket}; |
| 24 | use crate::rtp_::Direction; |
| 25 | use crate::rtp_::MidRid; |
| 26 | use crate::rtp_::Pli; |
| 27 | use crate::rtp_::Pt; |
| 28 | use crate::rtp_::SRTCP_OVERHEAD; |
| 29 | use crate::rtp_::SeqNo; |
| 30 | use crate::rtp_::{Bitrate, ExtensionMap, Goodbye, Mid, ReportList, Rtcp, RtcpFb}; |
| 31 | use crate::rtp_::{Dlrr, DlrrItem, ExtendedReport, ReportBlock}; |
| 32 | use crate::rtp_::{RtpHeader, SessionId, TwccPacketId, extend_u16}; |
| 33 | use crate::rtp_::{SrtpContext, Ssrc}; |
| 34 | use crate::rtp_::{TwccRecvRegister, TwccSendRegister}; |
| 35 | use crate::stats::StatsSnapshot; |
| 36 | use crate::streams::{RtpPacket, StreamTimeoutConfig, Streams}; |
| 37 | use crate::util::{Soonest, SystemTimeExt, already_happened, not_happening}; |
| 38 | use crate::{Reason, net}; |
| 39 | use crate::{RtcConfig, RtcError}; |
| 40 | |
| 41 | /// Minimum time we delay between sending nacks. This should be |
| 42 | /// set high enough to not cause additional problems in very bad |
| 43 | /// network conditions. |
| 44 | const NACK_MIN_INTERVAL: Duration = Duration::from_millis(33); |
| 45 | |
| 46 | /// Delay between reports of TWCC. This is deliberately very low. |
| 47 | const TWCC_INTERVAL: Duration = Duration::from_millis(50); |
| 48 | |
| 49 | /// Maximum number of pending RRTR entries before oldest are dropped. |
| 50 | const MAX_PENDING_RRTRS: usize = 300; |
| 51 | |
| 52 | /// Maximum number of DLRR items to include in a single ExtendedReport. |
| 53 | const MAX_DLRR_PER_REPORT: usize = 50; |
| 54 | |
| 55 | pub(crate) struct Session { |
| 56 | id: SessionId, |
| 57 | |
| 58 | // These fields are pub to allow session_sdp.rs modify them. |
| 59 | // Notice the fields are maybe not in m-line index order since the app |
| 60 | // might be spliced in somewhere. |
| 61 | pub medias: Vec<Media>, |
| 62 | |
| 63 | // The actual RTP encoded streams. |
| 64 | pub streams: Streams, |
| 65 | |
| 66 | /// The app m-line. Spliced into medias above. |
| 67 | app: Option<(Mid, usize)>, |
| 68 | |
| 69 | reordering_size_audio: usize, |
| 70 | reordering_size_video: usize, |
| 71 | intervals: RtcpReportIntervals, |
| 72 | pub send_buffer_audio: usize, |
| 73 | pub send_buffer_video: usize, |
| 74 | |
| 75 | /// Extension mappings are _per BUNDLE_, but we can only have one a=group BUNDLE |
| 76 | /// in WebRTC (one ice connection), so they are effectively per session. |
| 77 | pub exts: ExtensionMap, |
| 78 | |
| 79 | // Configuration of how we are sending/receiving media. |
| 80 | pub codec_config: CodecConfig, |
| 81 | |
| 82 | // Payload params for SSRC 0 BWE probes. |
| 83 | // Lazy init when we see the PT. |
| 84 | probe_payload_params: Option<PayloadParams>, |
| 85 | |
| 86 | srtp_rx: Option<SrtpContext>, |
| 87 | srtp_tx: Option<SrtpContext>, |
| 88 | last_nack: Instant, |
| 89 | last_twcc: Instant, |
| 90 | twcc: u64, |
| 91 | twcc_rx_register: TwccRecvRegister, |
| 92 | twcc_tx_register: TwccSendRegister, |
| 93 | max_rx_seq_lookup: HashMap<Ssrc, SeqNo>, |
| 94 | |
| 95 | bwe: Option<Bwe>, |
| 96 | |
| 97 | enable_twcc_feedback: bool, |
| 98 | |
| 99 | /// A pacer for sending RTP at specific rate. |
| 100 | pacer: PacerImpl, |
| 101 | pacer_control: PacerControl, |
| 102 | |
| 103 | // temporary buffer when getting the next (unencrypted) RTP packet from Media line. |
| 104 | poll_packet_buf: Vec<u8>, |
| 105 | |
| 106 | // Next packet for RtpPacket event. |
| 107 | pending_packet: Option<RtpPacket>, |
| 108 | |
| 109 | // Whether we sent a single outgoing RTP packet. |
| 110 | packet_first_sent: bool, |
| 111 | |
| 112 | // Total RTP payload bytes for media, including retransmissions. |
| 113 | media_bytes_rx: u64, |
| 114 | media_bytes_tx: u64, |
| 115 | |
| 116 | pub ice_lite: bool, |
| 117 | |
| 118 | /// Whether we are running in RTP-mode. |
| 119 | pub rtp_mode: bool, |
| 120 | |
| 121 | /// VP9 packetizer mode. |
| 122 | vp9_packetizer_mode: Vp9PacketizerMode, |
| 123 | |
| 124 | feedback_tx: VecDeque<Rtcp>, |
| 125 | feedback_rx: VecDeque<Rtcp>, |
| 126 | rtp_closing: bool, |
| 127 | |
| 128 | /// Pending RRTRs from remote SSRCs in FIFO order, deduplicated by SSRC. |
| 129 | /// A new RRTR for an already-pending SSRC updates its value in place. |
| 130 | /// Capped at 300 entries; at most 50 are consumed per DLRR report. |
| 131 | pending_rrtrs: VecDeque<(Ssrc, LastRrtr)>, |
| 132 | |
| 133 | raw_packets: Option<VecDeque<Box<RawPacket>>>, |
| 134 | |
| 135 | // Pending application-specific feedback (PSFB FMT=15) messages to emit as events. |
| 136 | pending_app_feedback: Option<AppSpecificFeedback>, |
| 137 | |
| 138 | /// Target MTU (start) and warn threshold (end). Buffer sizing uses the |
| 139 | /// target; oversized outgoing datagrams above the warn threshold log a warning. |
| 140 | mtu: RangeInclusive<usize>, |
| 141 | |
| 142 | #[cfg(feature = "_internal_test_exports")] |
| 143 | pending_probe: Option<crate::bwe_::ProbeClusterConfig>, |
| 144 | } |
| 145 | |
| 146 | /// Last RRTR (Receiver Reference Time Report, RFC 3611 §4.4) received from a remote SSRC. |
| 147 | /// Stored at session level so we can respond with a DLRR regardless of which local stream |
| 148 | /// the remote is observing. |
| 149 | #[derive(Debug, Clone, Copy)] |
| 150 | struct LastRrtr { |
| 151 | /// Middle 32 bits of the RRTR NTP timestamp (compact NTP / LRR). |
| 152 | lrr: u32, |
| 153 | /// When the RRTR arrived locally. |
| 154 | received_at: Instant, |
| 155 | } |
| 156 | |
| 157 | impl Session { |
| 158 | pub fn new(config: &RtcConfig) -> Self { |
| 159 | let mut id = SessionId::new(); |
| 160 | // Max 2^62 - 1: https://bugzilla.mozilla.org/show_bug.cgi?id=861895 |
| 161 | const MAX_ID: u64 = 2_u64.pow(62) - 1; |
| 162 | while *id > MAX_ID { |
| 163 | id = (*id >> 1).into(); |
| 164 | } |
| 165 | let (pacer, bwe) = if let Some(config) = &config.bwe_config { |
| 166 | let rate = config.initial_bitrate; |
| 167 | let pacer = PacerImpl::leaky_bucket(rate * 2.0); |
| 168 | let bwe = Bwe::new(rate); |
| 169 | (pacer, Some(bwe)) |
| 170 | } else { |
| 171 | (PacerImpl::null(), None) |
| 172 | }; |
| 173 | |
| 174 | let enable_stats = config.stats_interval.is_some(); |
| 175 | |
| 176 | Session { |
| 177 | id, |
| 178 | medias: vec![], |
| 179 | streams: Streams::new(enable_stats, *config.mtu.end()), |
| 180 | app: None, |
| 181 | reordering_size_audio: config.reordering_size_audio, |
| 182 | reordering_size_video: config.reordering_size_video, |
| 183 | intervals: config.intervals, |
| 184 | send_buffer_audio: config.send_buffer_audio, |
| 185 | send_buffer_video: config.send_buffer_video, |
| 186 | exts: config.exts.clone(), |
| 187 | |
| 188 | // Both sending and receiving starts from the configured codecs. |
| 189 | // These can then be changed in the SDP OFFER/ANSWER dance. |
| 190 | codec_config: config.codec_config.clone(), |
| 191 | probe_payload_params: None, |
| 192 | |
| 193 | srtp_rx: None, |
| 194 | srtp_tx: None, |
| 195 | last_nack: already_happened(), |
| 196 | last_twcc: already_happened(), |
| 197 | twcc: 0, |
| 198 | twcc_rx_register: TwccRecvRegister::new(100), |
| 199 | twcc_tx_register: TwccSendRegister::new(1000), |
| 200 | max_rx_seq_lookup: HashMap::new(), |
| 201 | bwe, |
| 202 | enable_twcc_feedback: false, |
| 203 | pacer, |
| 204 | pacer_control: PacerControl::new(), |
| 205 | poll_packet_buf: vec![0; 2000], |
| 206 | pending_packet: None, |
| 207 | packet_first_sent: false, |
| 208 | media_bytes_rx: 0, |
| 209 | media_bytes_tx: 0, |
| 210 | ice_lite: config.ice_lite, |
| 211 | rtp_mode: config.rtp_mode, |
| 212 | vp9_packetizer_mode: config.vp9_packetizer_mode, |
| 213 | feedback_tx: VecDeque::new(), |
| 214 | feedback_rx: VecDeque::new(), |
| 215 | rtp_closing: false, |
| 216 | pending_rrtrs: VecDeque::new(), |
| 217 | raw_packets: if config.enable_raw_packets { |
| 218 | Some(VecDeque::new()) |
| 219 | } else { |
| 220 | None |
| 221 | }, |
| 222 | pending_app_feedback: None, |
| 223 | mtu: config.mtu.clone(), |
| 224 | #[cfg(feature = "_internal_test_exports")] |
| 225 | pending_probe: None, |
| 226 | } |
| 227 | } |
| 228 | |
| 229 | fn mtu(&self) -> usize { |
| 230 | *self.mtu.start() |
| 231 | } |
| 232 | |
| 233 | fn mtu_warn(&self) -> usize { |
| 234 | *self.mtu.end() |
| 235 | } |
| 236 | |
| 237 | pub fn id(&self) -> SessionId { |
| 238 | self.id |
| 239 | } |
| 240 | |
| 241 | pub fn set_app(&mut self, mid: Mid, index: usize) -> Result<(), String> { |
| 242 | if let Some((mid_existing, index_existing)) = self.app { |
| 243 | if mid_existing != mid { |
| 244 | return Err(format!("App mid changed {} != {}", mid, mid_existing,)); |
| 245 | } |
| 246 | if index_existing != index { |
| 247 | return Err(format!("App index changed {} != {}", index, index_existing,)); |
| 248 | } |
| 249 | } else { |
| 250 | self.app = Some((mid, index)); |
| 251 | } |
| 252 | Ok(()) |
| 253 | } |
| 254 | |
| 255 | pub fn app(&self) -> &Option<(Mid, usize)> { |
| 256 | &self.app |
| 257 | } |
| 258 | |
| 259 | pub fn set_keying_material( |
| 260 | &mut self, |
| 261 | mat: KeyingMaterial, |
| 262 | crypto: &CryptoProvider, |
| 263 | srtp_profile: SrtpProfile, |
| 264 | active: bool, |
| 265 | ) { |
| 266 | // TODO: rename this to `initialise_srtp_context`? |
| 267 | // Whether we're active or passive determines if we use the left or right |
| 268 | // hand side of the key material to derive input/output. |
| 269 | let left = active; |
| 270 | |
| 271 | self.srtp_rx = Some(SrtpContext::new(crypto, srtp_profile, &mat, !left)); |
| 272 | self.srtp_tx = Some(SrtpContext::new(crypto, srtp_profile, &mat, left)); |
| 273 | } |
| 274 | |
| 275 | pub fn handle_timeout(&mut self, now: Instant) -> Result<(), RtcError> { |
| 276 | if self.rtp_closing { |
| 277 | return Ok(()); |
| 278 | } |
| 279 | |
| 280 | // Payload any waiting frames |
| 281 | self.do_payload()?; |
| 282 | |
| 283 | let sender_ssrc = self.streams.first_ssrc_local(); |
| 284 | |
| 285 | let do_nack = now >= self.nack_at().unwrap_or(not_happening()); |
| 286 | |
| 287 | self.streams.handle_timeout( |
| 288 | now, |
| 289 | sender_ssrc, |
| 290 | do_nack, |
| 291 | &self.medias, |
| 292 | StreamTimeoutConfig { |
| 293 | codecs: &self.codec_config, |
| 294 | intervals: self.intervals, |
| 295 | }, |
| 296 | &mut self.feedback_tx, |
| 297 | ); |
| 298 | |
| 299 | // Flush any pending RRTRs as a DLRR alongside the Sender Report compound packet. |
| 300 | // We piggyback on the SR cadence so the DLRR doesn't generate extra standalone |
| 301 | // RTCP traffic that could perturb bandwidth estimation. |
| 302 | // Consume at most 50 entries per report; the rest remain for subsequent SRs. |
| 303 | if !self.pending_rrtrs.is_empty() |
| 304 | && self |
| 305 | .feedback_tx |
| 306 | .iter() |
| 307 | .any(|r| matches!(r, Rtcp::SenderReport(_))) |
| 308 | { |
| 309 | let n = self.pending_rrtrs.len().min(MAX_DLRR_PER_REPORT); |
| 310 | let items: Vec<DlrrItem> = self |
| 311 | .pending_rrtrs |
| 312 | .drain(..n) |
| 313 | .map(|(ssrc, rrtr)| { |
| 314 | let delay = now.saturating_duration_since(rrtr.received_at); |
| 315 | let last_rr_delay = ((delay.as_micros() * 65_536) / 1_000_000) as u32; |
| 316 | DlrrItem { |
| 317 | ssrc, |
| 318 | last_rr_time: rrtr.lrr, |
| 319 | last_rr_delay, |
| 320 | } |
| 321 | }) |
| 322 | .collect(); |
| 323 | |
| 324 | self.feedback_tx |
| 325 | .push_back(Rtcp::ExtendedReport(ExtendedReport { |
| 326 | ssrc: sender_ssrc, |
| 327 | blocks: vec![ReportBlock::Dlrr(Dlrr { items })], |
| 328 | })); |
| 329 | } |
| 330 | |
| 331 | if do_nack { |
| 332 | self.last_nack = now; |
| 333 | } |
| 334 | |
| 335 | self.update_queue_state(now); |
| 336 | |
| 337 | if let Some(twcc_at) = self.twcc_at() { |
| 338 | if now >= twcc_at { |
| 339 | self.create_twcc_feedback(sender_ssrc, now); |
| 340 | } |
| 341 | } |
| 342 | |
| 343 | self.handle_timeout_bwe(now); |
| 344 | |
| 345 | Ok(()) |
| 346 | } |
| 347 | |
| 348 | fn handle_timeout_bwe(&mut self, now: Instant) { |
| 349 | let Some(bwe) = self.bwe.as_mut() else { |
| 350 | return; |
| 351 | }; |
| 352 | |
| 353 | // We can only run probes after first packet is sent and there |
| 354 | // are any queues that can handle padding requests. |
| 355 | let do_probe = self.packet_first_sent && self.pacer.has_padding_queue(); |
| 356 | |
| 357 | if let Some(probe_config) = bwe.handle_timeout(now, do_probe) { |
| 358 | // Only start the probe in the pacer if the estimator accepted it. |
| 359 | if bwe.start_probe(probe_config, now) { |
| 360 | #[cfg(feature = "_internal_test_exports")] |
| 361 | { |
| 362 | self.pending_probe = Some(probe_config); |
| 363 | } |
| 364 | self.pacer.start_probe(probe_config); |
| 365 | } |
| 366 | } |
| 367 | |
| 368 | // Check if active probe just completed |
| 369 | if let Some(cluster_id) = self.pacer.check_probe_complete(now) { |
| 370 | bwe.end_probe(now, cluster_id); |
| 371 | } |
| 372 | } |
| 373 | |
| 374 | fn update_queue_state(&mut self, now: Instant) { |
| 375 | let iter = self.streams.streams_tx().map(|m| m.queue_state(now)); |
| 376 | |
| 377 | let Some(padding_request) = self.pacer.handle_timeout(now, iter) else { |
| 378 | return; |
| 379 | }; |
| 380 | |
| 381 | let stream = self |
| 382 | .streams |
| 383 | .stream_tx_by_midrid(padding_request.midrid) |
| 384 | .expect("pacer to use an existing stream"); |
| 385 | |
| 386 | stream.generate_padding(padding_request.padding); |
| 387 | } |
| 388 | |
| 389 | fn create_twcc_feedback(&mut self, sender_ssrc: Ssrc, now: Instant) -> Option<()> { |
| 390 | self.last_twcc = now; |
| 391 | let mut twcc = self.twcc_rx_register.build_report(self.mtu() - 100)?; |
| 392 | |
| 393 | // These SSRC are on media level, but twcc is on session level, |
| 394 | // we fill in the first discovered media SSRC in each direction. |
| 395 | twcc.sender_ssrc = sender_ssrc; |
| 396 | twcc.ssrc = self.streams.first_ssrc_remote(); |
| 397 | |
| 398 | trace!("Created feedback TWCC: {:?}", twcc); |
| 399 | self.feedback_tx.push_front(Rtcp::Twcc(twcc)); |
| 400 | Some(()) |
| 401 | } |
| 402 | |
| 403 | /// Enqueue an application-specific feedback message (PSFB FMT=15, PT=206) for |
| 404 | /// transmission as part of the regular RTCP compound packet. |
| 405 | pub(crate) fn send_app_specific_feedback( |
| 406 | &mut self, |
| 407 | sender_ssrc: Ssrc, |
| 408 | media_ssrc: Ssrc, |
| 409 | payload: impl Into<Arc<[u8]>>, |
| 410 | ) { |
| 411 | use crate::rtp_::AppSpecificFeedback as RtcpAppFeedback; |
| 412 | let feedback = RtcpAppFeedback { |
| 413 | sender_ssrc, |
| 414 | media_ssrc, |
| 415 | payload: payload.into(), |
| 416 | }; |
| 417 | self.feedback_tx |
| 418 | .push_back(Rtcp::AppSpecificFeedback(feedback)); |
| 419 | } |
| 420 | |
| 421 | /// Enqueue a Picture Loss Indication (PLI, PT=206 FMT=1) for `media_ssrc`, to be |
| 422 | /// sent as part of the regular RTCP compound packet. The feedback packet's sender |
| 423 | /// SSRC is set to `sender_ssrc`. |
| 424 | pub(crate) fn send_pli_feedback(&mut self, sender_ssrc: Ssrc, media_ssrc: Ssrc) { |
| 425 | self.feedback_tx.push_back(Rtcp::Pli(Pli { |
| 426 | sender_ssrc, |
| 427 | ssrc: media_ssrc, |
| 428 | })); |
| 429 | } |
| 430 | |
| 431 | pub fn handle_rtp_receive(&mut self, now: Instant, message: &[u8]) { |
| 432 | let Some(header) = RtpHeader::parse(message, &self.exts) else { |
| 433 | trace!("Failed to parse RTP header"); |
| 434 | return; |
| 435 | }; |
| 436 | |
| 437 | self.handle_rtp(now, header, message); |
| 438 | } |
| 439 | |
| 440 | pub fn handle_rtcp_receive(&mut self, now: Instant, message: &[u8]) { |
| 441 | // According to spec, the outer enclosing SRTCP packet should always be a SR or RR, |
| 442 | // even if it's irrelevant and empty. |
| 443 | // In practice I'm not sure that is happening, because libWebRTC hates empty packets. |
| 444 | self.handle_rtcp(now, message); |
| 445 | } |
| 446 | |
| 447 | fn mid_and_ssrc_for_header(&mut self, now: Instant, header: &RtpHeader) -> Option<(Mid, Ssrc)> { |
| 448 | let ssrc_header = header.ssrc; |
| 449 | |
| 450 | if let Some(r) = self.streams.mid_ssrc_rx_by_ssrc_or_rtx(now, ssrc_header) { |
| 451 | return Some(r); |
| 452 | } |
| 453 | |
| 454 | // Attempt to dynamically map this header to some Media/ReceiveStream. |
| 455 | self.map_dynamic(header); |
| 456 | |
| 457 | // The dynamic mapping might have added an entry by now. |
| 458 | if let Some(r) = self.streams.mid_ssrc_rx_by_ssrc_or_rtx(now, ssrc_header) { |
| 459 | return Some(r); |
| 460 | } |
| 461 | |
| 462 | // SSRC 0 is used for non-media BWE probes from libwebrtc. |
| 463 | // These probes need TWCC feedback but don't carry actual media. |
| 464 | if ssrc_header.is_probe() { |
| 465 | return self.ensure_probe_stream(header.payload_type); |
| 466 | } |
| 467 | |
| 468 | None |
| 469 | } |
| 470 | |
| 471 | /// Creates the probe stream on-demand for handling SSRC 0 BWE probes. |
| 472 | /// |
| 473 | /// No Media is created since probes don't carry real media - they only need |
| 474 | /// SRTP decryption and TWCC feedback. |
| 475 | fn ensure_probe_stream(&mut self, pt: Pt) -> Option<(Mid, Ssrc)> { |
| 476 | let ssrc: Ssrc = 0.into(); |
| 477 | let midrid = MidRid(MID_PROBE, None); |
| 478 | |
| 479 | // Add PayloadParams for this PT if not already configured. |
| 480 | // Probes may use different PTs, so we add them as we see them. |
| 481 | self.probe_payload_params = Some(PayloadParams::new_probe(pt)); |
| 482 | |
| 483 | // Create the stream with NACK suppressed (probes don't need retransmission) |
| 484 | self.streams.expect_stream_rx(ssrc, None, midrid, true); |
| 485 | |
| 486 | Some((MID_PROBE, ssrc)) |
| 487 | } |
| 488 | |
| 489 | fn map_dynamic(&mut self, header: &RtpHeader) { |
| 490 | // There are two strategies for dynamically mapping SSRC. Both use the RTP "mid" |
| 491 | // header extension. |
| 492 | // A) Mid+Rid - used when doing simulcast. Rid points out which |
| 493 | // simulcast layer is in use. There is a separate header |
| 494 | // to indicate repair (RTX) stream. |
| 495 | // B) Mid+PT - when not doing simulcast, the PT identifies whether |
| 496 | // this is a repair stream. |
| 497 | |
| 498 | let Some(mid) = header.ext_vals.mid else { |
| 499 | return; |
| 500 | }; |
| 501 | let rid = header.ext_vals.rid.or(header.ext_vals.rid_repair); |
| 502 | |
| 503 | // The media the mid points out. Bail if the mid points to something |
| 504 | // we don't know about. |
| 505 | let Some(media) = self.medias.iter_mut().find(|m| m.mid() == mid) else { |
| 506 | return; |
| 507 | }; |
| 508 | |
| 509 | // Figure out which payload the PT maps to. Either main or RTX. |
| 510 | let maybe_payload = self |
| 511 | .codec_config |
| 512 | .iter() |
| 513 | .find(|p| p.pt() == header.payload_type || p.resend() == Some(header.payload_type)); |
| 514 | |
| 515 | // If we don't find it, bail out. |
| 516 | let Some(payload) = maybe_payload else { |
| 517 | return; |
| 518 | }; |
| 519 | |
| 520 | if let Some(rid) = rid { |
| 521 | // Case A - use the rid_repair header to identify RTX. |
| 522 | let is_main = header.ext_vals.rid.is_some(); |
| 523 | |
| 524 | let midrid = MidRid(mid, Some(rid)); |
| 525 | |
| 526 | self.streams |
| 527 | .map_dynamic_by_rid(header.ssrc, midrid, media, *payload, is_main); |
| 528 | } else { |
| 529 | // Case B - the payload type identifies RTX. |
| 530 | let is_main = payload.pt() == header.payload_type; |
| 531 | |
| 532 | let midrid = MidRid(mid, None); |
| 533 | |
| 534 | self.streams |
| 535 | .map_dynamic_by_pt(header.ssrc, midrid, media, *payload, is_main); |
| 536 | } |
| 537 | } |
| 538 | |
| 539 | pub(crate) fn handle_rtp(&mut self, now: Instant, mut header: RtpHeader, buf: &[u8]) { |
| 540 | // Rewrite absolute-send-time (if present) to be relative to now. |
| 541 | header.ext_vals.update_absolute_send_time(now); |
| 542 | |
| 543 | trace!("Handle RTP: {:?}", header); |
| 544 | |
| 545 | // The ssrc is the _main_ ssrc (no the rtx, that might be in the header). |
| 546 | let Some((mid, ssrc)) = self.mid_and_ssrc_for_header(now, &header) else { |
| 547 | debug!("No mid/SSRC for header: {:?}", header); |
| 548 | return; |
| 549 | }; |
| 550 | |
| 551 | let srtp = match self.srtp_rx.as_mut() { |
| 552 | Some(v) => v, |
| 553 | None => { |
| 554 | trace!("Rejecting SRTP while missing SrtpContext"); |
| 555 | return; |
| 556 | } |
| 557 | }; |
| 558 | |
| 559 | // mid_and_ssrc_for_header guarantees stream exists. |
| 560 | // Media may not exist for internal probe streams (SSRC 0). |
| 561 | let stream = self.streams.stream_rx(&ssrc).unwrap(); |
| 562 | |
| 563 | let maybe_params = if ssrc.is_probe() { |
| 564 | self.probe_payload_params.as_ref() |
| 565 | } else { |
| 566 | main_payload_params(&self.codec_config, header.payload_type) |
| 567 | }; |
| 568 | |
| 569 | let params = match maybe_params { |
| 570 | Some(p) => p, |
| 571 | None => { |
| 572 | trace!( |
| 573 | "No payload params could be found (main or RTX) for {:?}", |
| 574 | header.payload_type |
| 575 | ); |
| 576 | return; |
| 577 | } |
| 578 | }; |
| 579 | let clock_rate = params.spec().rtp_clock_rate(); |
| 580 | let pt = params.pt(); |
| 581 | let is_repair = pt != header.payload_type; |
| 582 | |
| 583 | let max_seq_lookup = make_max_seq_lookup(&self.max_rx_seq_lookup); |
| 584 | |
| 585 | // is_repair controls whether update is updating the main register or the RTX register. |
| 586 | // Either way we get a seq_no_outer which is used to decrypt the SRTP. |
| 587 | let mut seq_no = stream.extend_seq(&header, is_repair, max_seq_lookup); |
| 588 | |
| 589 | if !stream.is_new_packet(is_repair, seq_no) { |
| 590 | // Dupe packet. This could be a potential SRTP replay attack, which means |
| 591 | // we should not spend any CPU cycles towards decrypting it. |
| 592 | trace!( |
| 593 | "Ignoring dupe packet mid: {} seq_no: {} is_repair: {}", |
| 594 | mid, seq_no, is_repair |
| 595 | ); |
| 596 | return; |
| 597 | } |
| 598 | |
| 599 | // The decrypted plaintext borrows from the SRTP scratch buffer. We |
| 600 | // narrow it down to the actual payload via slicing only, then build |
| 601 | // the final `Arc<[u8]>` in a single allocation at the end. |
| 602 | let data = match srtp.unprotect_rtp(buf, &header, *seq_no) { |
| 603 | Some(d) => d, |
| 604 | None => { |
| 605 | trace!( |
| 606 | "Failed to unprotect SRTP for SSRC: {} pt: {} mid: {} \ |
| 607 | rid: {:?} seq_no: {} is_repair: {}", |
| 608 | header.ssrc, |
| 609 | pt, |
| 610 | stream.mid(), |
| 611 | stream.rid(), |
| 612 | seq_no, |
| 613 | is_repair |
| 614 | ); |
| 615 | return; |
| 616 | } |
| 617 | }; |
| 618 | |
| 619 | let data = if header.has_padding { |
| 620 | match RtpHeader::unpad_payload(data) { |
| 621 | Some(d) => d, |
| 622 | None => { |
| 623 | trace!("unpadding of unprotected payload failed"); |
| 624 | return; |
| 625 | } |
| 626 | } |
| 627 | } else { |
| 628 | data |
| 629 | }; |
| 630 | |
| 631 | if let Some(raw_packets) = &mut self.raw_packets { |
| 632 | raw_packets.push_back(Box::new(RawPacket::RtpRx(header.clone(), data.to_vec()))); |
| 633 | } |
| 634 | |
| 635 | // Mark as received for TWCC purposes |
| 636 | if let Some(transport_cc) = header.ext_vals.transport_cc { |
| 637 | let prev = self.twcc_rx_register.max_seq(); |
| 638 | let extended = extend_u16(prev.map(|s| *s), transport_cc); |
| 639 | self.twcc_rx_register.update_seq(extended.into(), now); |
| 640 | } |
| 641 | |
| 642 | // Store largest seen seq_no for the SSRC. This is used in case we get SSRC changes |
| 643 | // like A -> B -> A. When we go back to A, we must keep the ROC. |
| 644 | update_max_seq(&mut self.max_rx_seq_lookup, header.ssrc, seq_no); |
| 645 | |
| 646 | // Register reception in nack registers. |
| 647 | let receipt_outer = stream.update_register(now, &header, clock_rate, is_repair, seq_no); |
| 648 | |
| 649 | // RTX packets must be rewritten to be a normal packet. This only changes the |
| 650 | // the seq_no, however MediaTime might be different when interpreted against the |
| 651 | // the "main" register. |
| 652 | let (receipt, data) = if is_repair { |
| 653 | // Drop RTX packets that are just empty padding. The payload here |
| 654 | // is empty because we would have done RtpHeader::unpad_payload above. |
| 655 | // For unpausing, it's enough with the stream.update() already done above. |
| 656 | if data.is_empty() { |
| 657 | return; |
| 658 | } |
| 659 | |
| 660 | // Rewrite the header, and strip the resent seq_no prefix from the body. |
| 661 | let data = stream.un_rtx(&mut header, data, pt); |
| 662 | |
| 663 | let max_seq_lookup = make_max_seq_lookup(&self.max_rx_seq_lookup); |
| 664 | |
| 665 | // Header has changed, which means we extend a new seq_no. This time |
| 666 | // without is_repair since this is the wrapped resend. This is the |
| 667 | // extended number of the main stream. |
| 668 | seq_no = stream.extend_seq(&header, false, max_seq_lookup); |
| 669 | |
| 670 | // header.ssrc is changed by un_rtx() above to be main SSRC. |
| 671 | update_max_seq(&mut self.max_rx_seq_lookup, header.ssrc, seq_no); |
| 672 | |
| 673 | // Now update the "main" register with the repaired packet info. |
| 674 | let receipt = stream.update_register(now, &header, clock_rate, false, seq_no); |
| 675 | (receipt, data) |
| 676 | } else { |
| 677 | // This is not RTX, the outer seq and time is what we use. The first |
| 678 | // stream.update will have updated the main register. |
| 679 | (receipt_outer, data) |
| 680 | }; |
| 681 | |
| 682 | // Probe packets (SSRC 0) contain only padding, no real media to process |
| 683 | if ssrc.is_probe() { |
| 684 | return; |
| 685 | } |
| 686 | |
| 687 | self.media_bytes_rx += data.len() as u64; |
| 688 | |
| 689 | // One-shot conversion to Arc<[u8]>: a single allocation that copies |
| 690 | // only the trimmed payload bytes out of the SRTP scratch buffer. |
| 691 | let payload: Arc<[u8]> = Arc::from(data); |
| 692 | |
| 693 | let packet = stream.handle_rtp(now, header, payload, seq_no, receipt.time); |
| 694 | |
| 695 | if self.rtp_mode { |
| 696 | // In RTP mode, we store the packet temporarily here for the next poll_output(). |
| 697 | // However only if this is a packet not seen before. This filters out spurious |
| 698 | // resends for padding. |
| 699 | if receipt.is_new_packet { |
| 700 | self.pending_packet = Some(packet); |
| 701 | } |
| 702 | } else { |
| 703 | // In non-RTP mode, we let the Media use a Depayloader. |
| 704 | // unwrap is fine because mid_and_ssrc_for_header guarantees it. |
| 705 | let media = self.medias.iter_mut().find(|m| m.mid() == mid).unwrap(); |
| 706 | media.depayload( |
| 707 | stream.rid(), |
| 708 | packet, |
| 709 | self.reordering_size_audio, |
| 710 | self.reordering_size_video, |
| 711 | &self.codec_config, |
| 712 | ); |
| 713 | } |
| 714 | } |
| 715 | |
| 716 | fn handle_rtcp(&mut self, now: Instant, buf: &[u8]) -> Option<()> { |
| 717 | let srtp: &mut SrtpContext = self.srtp_rx.as_mut()?; |
| 718 | let unprotected = srtp.unprotect_rtcp(buf)?; |
| 719 | |
| 720 | Rtcp::read_packet(&unprotected, &mut self.feedback_rx); |
| 721 | let mut need_configure_pacer = false; |
| 722 | |
| 723 | if let Some(raw_packets) = &mut self.raw_packets { |
| 724 | for fb in &self.feedback_rx { |
| 725 | raw_packets.push_back(Box::new(RawPacket::RtcpRx(fb.clone()))); |
| 726 | } |
| 727 | } |
| 728 | |
| 729 | for fb in RtcpFb::from_rtcp(self.feedback_rx.drain(..)) { |
| 730 | if let RtcpFb::AppSpecificFeedback(asf) = fb { |
| 731 | self.pending_app_feedback = Some(AppSpecificFeedback { |
| 732 | sender_ssrc: asf.sender_ssrc, |
| 733 | media_ssrc: asf.media_ssrc, |
| 734 | payload: asf.payload, |
| 735 | }); |
| 736 | continue; |
| 737 | } |
| 738 | |
| 739 | if let RtcpFb::Twcc(twcc) = fb { |
| 740 | trace!("Handle TWCC: {:?}", twcc); |
| 741 | let maybe_records = self.twcc_tx_register.apply_report(twcc, now); |
| 742 | |
| 743 | if let (Some(maybe_records), Some(bwe)) = (maybe_records, &mut self.bwe) { |
| 744 | bwe.update(maybe_records, now); |
| 745 | } |
| 746 | need_configure_pacer = true; |
| 747 | |
| 748 | // The funky thing about TWCC reports is that they are never stapled |
| 749 | // together with other RTCP packet. If they were though, we want to |
| 750 | // handle more packets. |
| 751 | continue; |
| 752 | } |
| 753 | |
| 754 | if let RtcpFb::Rrtr((rrtr, ssrc)) = fb { |
| 755 | // DLRR responder (RFC 3611 §4.5): store the remote's RRTR so we can |
| 756 | // reply with a DLRR, letting the remote compute RTT. |
| 757 | let lrr = (rrtr.ntp_time.as_ntp_64() >> 16) as u32; |
| 758 | let entry = LastRrtr { |
| 759 | lrr, |
| 760 | received_at: now, |
| 761 | }; |
| 762 | |
| 763 | // Update in place if this SSRC is already pending (preserves queue order). |
| 764 | if let Some(existing) = self.pending_rrtrs.iter_mut().find(|(s, _)| *s == ssrc) { |
| 765 | existing.1 = entry; |
| 766 | } else { |
| 767 | // Cap at 300 entries; drop the oldest if full. |
| 768 | if self.pending_rrtrs.len() >= MAX_PENDING_RRTRS { |
| 769 | self.pending_rrtrs.pop_front(); |
| 770 | } |
| 771 | self.pending_rrtrs.push_back((ssrc, entry)); |
| 772 | } |
| 773 | continue; |
| 774 | } |
| 775 | |
| 776 | if fb.is_for_rx() { |
| 777 | let Some(stream) = self.streams.stream_rx(&fb.ssrc()) else { |
| 778 | continue; |
| 779 | }; |
| 780 | stream.handle_rtcp(now, fb); |
| 781 | } else { |
| 782 | let Some(stream) = self.streams.stream_tx(&fb.ssrc()) else { |
| 783 | continue; |
| 784 | }; |
| 785 | stream.handle_rtcp(now, fb); |
| 786 | } |
| 787 | } |
| 788 | |
| 789 | // Not in the above if due to lifetime issues, still okay because the method |
| 790 | // doesn't do anything when BWE isn't configured. |
| 791 | if need_configure_pacer { |
| 792 | self.configure_pacer(); |
| 793 | } |
| 794 | |
| 795 | Some(()) |
| 796 | } |
| 797 | |
| 798 | pub fn poll_event(&mut self) -> Option<Event> { |
| 799 | #[cfg(feature = "_internal_test_exports")] |
| 800 | { |
| 801 | if let Some(probe) = self.pending_probe.take() { |
| 802 | return Some(Event::Probe(probe)); |
| 803 | } |
| 804 | } |
| 805 | |
| 806 | if let Some(bitrate_estimate) = self.bwe.as_mut().and_then(|bwe| bwe.poll_estimate()) { |
| 807 | return Some(Event::EgressBitrateEstimate(BweKind::Twcc( |
| 808 | bitrate_estimate, |
| 809 | ))); |
| 810 | } |
| 811 | |
| 812 | // If we're not ready to flow media, don't send any events. |
| 813 | if !self.ready_for_srtp() { |
| 814 | return None; |
| 815 | } |
| 816 | |
| 817 | if let Some(raw_packets) = &mut self.raw_packets { |
| 818 | if let Some(p) = raw_packets.pop_front() { |
| 819 | return Some(Event::RawPacket(p)); |
| 820 | } |
| 821 | } |
| 822 | |
| 823 | if let Some(feedback) = self.pending_app_feedback.take() { |
| 824 | return Some(Event::AppSpecificFeedback(feedback)); |
| 825 | } |
| 826 | |
| 827 | // This must be before pending_packet.take() since we need to emit the unpaused event |
| 828 | // before the first packet causing the unpause. |
| 829 | if let Some(paused) = self.streams.poll_stream_paused() { |
| 830 | if paused.paused { |
| 831 | if let Some(media) = self.medias.iter_mut().find(|m| m.mid() == paused.mid) { |
| 832 | // Drop held partial depacketizer state so pre-pause fragments can't |
| 833 | // complete into stale MediaData after the stream resumes. |
| 834 | media.reset_depayloaders_for_rid(paused.rid); |
| 835 | } |
| 836 | } |
| 837 | return Some(Event::StreamPaused(paused)); |
| 838 | } |
| 839 | |
| 840 | if self.rtp_mode { |
| 841 | if let Some(packet) = self.pending_packet.take() { |
| 842 | return Some(Event::RtpPacket(packet)); |
| 843 | } |
| 844 | } |
| 845 | |
| 846 | if let Some(req) = self.streams.poll_keyframe_request() { |
| 847 | return Some(Event::KeyframeRequest(req)); |
| 848 | } |
| 849 | |
| 850 | if let Some(report) = self.streams.poll_sender_feedback() { |
| 851 | return Some(Event::SenderFeedback(report)); |
| 852 | } |
| 853 | |
| 854 | if let Some((mid, bitrate)) = self.streams.poll_remb_request() { |
| 855 | return Some(Event::EgressBitrateEstimate(BweKind::Remb(mid, bitrate))); |
| 856 | } |
| 857 | |
| 858 | for media in &mut self.medias { |
| 859 | if media.need_open_event { |
| 860 | media.need_open_event = false; |
| 861 | |
| 862 | return Some(Event::MediaAdded(MediaAdded { |
| 863 | mid: media.mid(), |
| 864 | kind: media.kind(), |
| 865 | direction: media.direction(), |
| 866 | simulcast: media.simulcast().map(|s| s.clone().into()), |
| 867 | })); |
| 868 | } |
| 869 | |
| 870 | if media.need_changed_event { |
| 871 | media.need_changed_event = false; |
| 872 | return Some(Event::MediaChanged(MediaChanged { |
| 873 | mid: media.mid(), |
| 874 | direction: media.direction(), |
| 875 | })); |
| 876 | } |
| 877 | } |
| 878 | |
| 879 | None |
| 880 | } |
| 881 | |
| 882 | pub fn poll_event_fallible(&mut self) -> Result<Option<Event>, RtcError> { |
| 883 | // Not relevant in rtp_mode, where the packets are picked up by poll_event(). |
| 884 | if self.rtp_mode { |
| 885 | return Ok(None); |
| 886 | } |
| 887 | |
| 888 | for media in &mut self.medias { |
| 889 | if let Some(e) = media.poll_sample(&self.codec_config)? { |
| 890 | return Ok(Some(Event::MediaData(e))); |
| 891 | } |
| 892 | } |
| 893 | |
| 894 | Ok(None) |
| 895 | } |
| 896 | |
| 897 | fn ready_for_srtp(&self) -> bool { |
| 898 | self.srtp_rx.is_some() && self.srtp_tx.is_some() |
| 899 | } |
| 900 | |
| 901 | pub fn poll_datagram(&mut self, now: Instant) -> Option<net::DatagramSend> { |
| 902 | // Time must have progressed forward from start value. |
| 903 | if now == already_happened() { |
| 904 | return None; |
| 905 | } |
| 906 | |
| 907 | let x = if self.rtp_closing { |
| 908 | self.poll_feedback() |
| 909 | } else { |
| 910 | None.or_else(|| self.poll_feedback()) |
| 911 | .or_else(|| self.poll_packet(now)) |
| 912 | }; |
| 913 | |
| 914 | if let Some(x) = &x { |
| 915 | // In RTP mode we trust the API user feeds the RTP packet sizes they |
| 916 | // need for the MTU they are targeting. This warning is only for when |
| 917 | // str0m does the RTP packetization. |
| 918 | if !self.rtp_mode && x.len() > self.mtu_warn() { |
| 919 | warn!("RTP above MTU {}: {}", self.mtu_warn(), x.len()); |
| 920 | } |
| 921 | } |
| 922 | |
| 923 | x |
| 924 | } |
| 925 | |
| 926 | /// To be called in lieu of [`Self::poll_datagram`] when the owner is not in a position to transmit any |
| 927 | /// generated feedback, and thus such feedback should be dropped. |
| 928 | pub fn clear_feedback(&mut self) { |
| 929 | self.feedback_rx.clear(); |
| 930 | self.feedback_tx.clear(); |
| 931 | } |
| 932 | |
| 933 | pub fn close_rtp(&mut self) { |
| 934 | if self.rtp_closing { |
| 935 | return; |
| 936 | } |
| 937 | |
| 938 | self.rtp_closing = true; |
| 939 | self.clear_feedback(); |
| 940 | self.streams.reset_send_buffers(); |
| 941 | |
| 942 | for reports in ReportList::lists_from_iter(self.streams.local_sender_ssrcs()) { |
| 943 | self.feedback_tx |
| 944 | .push_back(Rtcp::Goodbye(Goodbye { reports })); |
| 945 | } |
| 946 | } |
| 947 | |
| 948 | pub fn is_rtp_closed(&self) -> bool { |
| 949 | self.rtp_closing && self.feedback_tx.is_empty() |
| 950 | } |
| 951 | |
| 952 | fn poll_feedback(&mut self) -> Option<net::DatagramSend> { |
| 953 | if self.feedback_tx.is_empty() { |
| 954 | return None; |
| 955 | } |
| 956 | |
| 957 | // Round to nearest multiple of 4 bytes. |
| 958 | let encryptable_mtu: usize = (self.mtu() - SRTCP_OVERHEAD) & !3; |
| 959 | assert!(encryptable_mtu % 4 == 0); |
| 960 | |
| 961 | let mut data = vec![0_u8; encryptable_mtu]; |
| 962 | |
| 963 | let mut raw_packets = self.raw_packets.as_mut(); |
| 964 | let output = move |fb| { |
| 965 | if let Some(raw_packets) = &mut raw_packets { |
| 966 | raw_packets.push_back(Box::new(RawPacket::RtcpTx(fb))); |
| 967 | } |
| 968 | }; |
| 969 | |
| 970 | let len = Rtcp::write_packet(&mut self.feedback_tx, &mut data, output); |
| 971 | |
| 972 | if len == 0 { |
| 973 | return None; |
| 974 | } |
| 975 | |
| 976 | data.truncate(len); |
| 977 | |
| 978 | let srtp = self.srtp_tx.as_mut()?; |
| 979 | let protected = srtp.protect_rtcp(&data); |
| 980 | |
| 981 | assert!( |
| 982 | protected.len() < self.mtu(), |
| 983 | "Encrypted SRTCP should be less than MTU" |
| 984 | ); |
| 985 | |
| 986 | Some(protected.into()) |
| 987 | } |
| 988 | |
| 989 | fn poll_packet(&mut self, now: Instant) -> Option<DatagramSend> { |
| 990 | let srtp_tx = self.srtp_tx.as_mut()?; |
| 991 | |
| 992 | // Figure out which, if any, queue to poll |
| 993 | // The cluster_id is captured by the pacer at poll time, before register_send() might clear it |
| 994 | let (midrid, cluster_id) = self.pacer.poll_queue()?; |
| 995 | let Some(media) = self.medias.iter().find(|m| m.mid() == midrid.mid()) else { |
| 996 | trace!("Pacer pointed to mid {} which has no media", midrid.mid()); |
| 997 | return None; |
| 998 | }; |
| 999 | |
| 1000 | let buf = &mut self.poll_packet_buf; |
| 1001 | let twcc_seq = self.twcc; |
| 1002 | |
| 1003 | let stream = self.streams.stream_tx_by_midrid(midrid)?; |
| 1004 | |
| 1005 | let params = &self.codec_config; |
| 1006 | let exts = media.remote_extmap(); |
| 1007 | |
| 1008 | // TWCC might not be enabled for this m-line. Firefox do use TWCC, but not |
| 1009 | // for audio. This is indiciated via the SDP. |
| 1010 | let twcc_enabled = exts.id_of(Extension::TransportSequenceNumber).is_some(); |
| 1011 | let twcc = twcc_enabled.then_some(&mut self.twcc); |
| 1012 | |
| 1013 | let receipt = stream.poll_packet(now, exts, twcc, params, buf)?; |
| 1014 | |
| 1015 | let PacketReceipt { |
| 1016 | header, |
| 1017 | seq_no, |
| 1018 | is_padding, |
| 1019 | payload_size, |
| 1020 | } = receipt; |
| 1021 | |
| 1022 | trace!(payload_size, is_padding, "Poll RTP: {:?}", header); |
| 1023 | |
| 1024 | #[cfg(feature = "_internal_dont_use_log_stats")] |
| 1025 | { |
| 1026 | let kind = if is_padding { "padding" } else { "media" }; |
| 1027 | |
| 1028 | crate::log_stat!("PACKET_SENT", header.ssrc, payload_size, kind); |
| 1029 | } |
| 1030 | |
| 1031 | self.pacer.register_send(now, payload_size.into(), midrid); |
| 1032 | |
| 1033 | if let Some(raw_packets) = &mut self.raw_packets { |
| 1034 | raw_packets.push_back(Box::new(RawPacket::RtpTx(header.clone(), buf.clone()))); |
| 1035 | } |
| 1036 | |
| 1037 | let protected = srtp_tx.protect_rtp(buf, &header, *seq_no); |
| 1038 | |
| 1039 | if twcc_enabled { |
| 1040 | let packet_id = if let Some(cluster) = cluster_id { |
| 1041 | TwccPacketId::with_cluster(twcc_seq, cluster) |
| 1042 | } else { |
| 1043 | TwccPacketId::new(twcc_seq) |
| 1044 | }; |
| 1045 | self.twcc_tx_register |
| 1046 | .register_seq(packet_id, now, payload_size); |
| 1047 | } |
| 1048 | |
| 1049 | // Update BWE subsystem |
| 1050 | if let Some(bwe) = self.bwe.as_mut() { |
| 1051 | bwe.on_media_sent(payload_size.into(), is_padding, now); |
| 1052 | } |
| 1053 | |
| 1054 | if !is_padding && !header.ssrc.is_probe() { |
| 1055 | self.media_bytes_tx += payload_size as u64; |
| 1056 | } |
| 1057 | |
| 1058 | if !self.packet_first_sent { |
| 1059 | self.packet_first_sent = true; |
| 1060 | } |
| 1061 | |
| 1062 | Some(protected.into()) |
| 1063 | } |
| 1064 | |
| 1065 | pub fn poll_timeout(&mut self) -> (Option<Instant>, Reason) { |
| 1066 | if self.rtp_closing { |
| 1067 | return if self.feedback_tx.is_empty() { |
| 1068 | (None, Reason::NotHappening) |
| 1069 | } else { |
| 1070 | (Some(already_happened()), Reason::Feedback) |
| 1071 | }; |
| 1072 | } |
| 1073 | |
| 1074 | let feedback_at = self.regular_feedback_at(); |
| 1075 | let nack_at = self.nack_at(); |
| 1076 | let twcc_at = self.twcc_at(); |
| 1077 | let pacing_at = self.pacer.poll_timeout(); |
| 1078 | let packetize_at = self.medias.iter().flat_map(|m| m.poll_timeout()).next(); |
| 1079 | let paused_at = self.paused_at(); |
| 1080 | let send_stream_at = self.streams.send_stream(); |
| 1081 | |
| 1082 | // Gives us built-in reason |
| 1083 | let bwe_at = self |
| 1084 | .bwe |
| 1085 | .as_ref() |
| 1086 | .map(|bwe| bwe.poll_timeout()) |
| 1087 | // We should never see this reason. |
| 1088 | .unwrap_or((None, Reason::BweDelayControl)); |
| 1089 | |
| 1090 | (feedback_at, Reason::Feedback) |
| 1091 | .soonest((nack_at, Reason::Nack)) |
| 1092 | .soonest((twcc_at, Reason::Twcc)) |
| 1093 | .soonest(pacing_at) |
| 1094 | .soonest((packetize_at, Reason::Packetize)) |
| 1095 | .soonest((paused_at, Reason::PauseCheck)) |
| 1096 | .soonest((send_stream_at, Reason::SendStream)) |
| 1097 | .soonest(bwe_at) |
| 1098 | } |
| 1099 | |
| 1100 | pub fn has_mid(&self, mid: Mid) -> bool { |
| 1101 | self.medias.iter().any(|m| m.mid() == mid) |
| 1102 | } |
| 1103 | |
| 1104 | fn regular_feedback_at(&self) -> Option<Instant> { |
| 1105 | self.streams.regular_feedback_at(self.intervals) |
| 1106 | } |
| 1107 | |
| 1108 | fn paused_at(&self) -> Option<Instant> { |
| 1109 | self.streams.paused_at() |
| 1110 | } |
| 1111 | |
| 1112 | fn nack_at(&mut self) -> Option<Instant> { |
| 1113 | if !self.streams.any_nack_enabled() { |
| 1114 | return None; |
| 1115 | } |
| 1116 | |
| 1117 | Some(self.last_nack + NACK_MIN_INTERVAL) |
| 1118 | } |
| 1119 | |
| 1120 | fn twcc_at(&self) -> Option<Instant> { |
| 1121 | let is_receiving = self.streams.is_receiving(); |
| 1122 | if is_receiving && self.enable_twcc_feedback && self.twcc_rx_register.has_unreported() { |
| 1123 | Some(self.last_twcc + TWCC_INTERVAL) |
| 1124 | } else { |
| 1125 | None |
| 1126 | } |
| 1127 | } |
| 1128 | |
| 1129 | pub fn enable_twcc_feedback(&mut self) { |
| 1130 | if !self.enable_twcc_feedback { |
| 1131 | debug!("Enable TWCC feedback"); |
| 1132 | self.enable_twcc_feedback = true; |
| 1133 | } |
| 1134 | } |
| 1135 | |
| 1136 | pub fn visit_stats(&mut self, now: Instant, snapshot: &mut StatsSnapshot) { |
| 1137 | for stream in self.streams.streams_tx() { |
| 1138 | stream.visit_stats(snapshot, now); |
| 1139 | } |
| 1140 | |
| 1141 | for stream in self.streams.streams_rx() { |
| 1142 | stream.visit_stats(snapshot, now); |
| 1143 | } |
| 1144 | |
| 1145 | snapshot.tx = self.media_bytes_tx; |
| 1146 | snapshot.rx = self.media_bytes_rx; |
| 1147 | snapshot.bwe_tx = self.bwe.as_ref().and_then(|bwe| bwe.last_estimate()); |
| 1148 | |
| 1149 | snapshot.egress_loss_fraction = self.twcc_tx_register.loss(Duration::from_secs(1), now); |
| 1150 | snapshot.rtt = self.twcc_tx_register.rtt(); |
| 1151 | snapshot.ingress_loss_fraction = self.twcc_rx_register.loss(); |
| 1152 | } |
| 1153 | |
| 1154 | pub fn set_bwe_desired_bitrate(&mut self, desired_bitrate: Bitrate) { |
| 1155 | if let Some(bwe) = self.bwe.as_mut() { |
| 1156 | bwe.set_desired_bitrate(desired_bitrate); |
| 1157 | self.configure_pacer(); |
| 1158 | } |
| 1159 | } |
| 1160 | |
| 1161 | pub fn reset_bwe(&mut self, init_bitrate: Bitrate) { |
| 1162 | if let Some(bwe) = self.bwe.as_mut() { |
| 1163 | bwe.reset(init_bitrate); |
| 1164 | } |
| 1165 | } |
| 1166 | |
| 1167 | #[allow(dead_code)] |
| 1168 | pub fn line_count(&self) -> usize { |
| 1169 | self.medias.len() + if self.app.is_some() { 1 } else { 0 } |
| 1170 | } |
| 1171 | |
| 1172 | pub fn add_media(&mut self, media: Media) { |
| 1173 | self.medias.push(media); |
| 1174 | } |
| 1175 | |
| 1176 | pub fn medias(&self) -> &[Media] { |
| 1177 | &self.medias |
| 1178 | } |
| 1179 | |
| 1180 | pub fn remove_media(&mut self, mid: Mid) { |
| 1181 | self.medias.retain(|media| media.mid() != mid); |
| 1182 | self.streams.remove_streams_by_mid(mid); |
| 1183 | } |
| 1184 | |
| 1185 | fn configure_pacer(&mut self) { |
| 1186 | let Some(bwe) = self.bwe.as_mut() else { |
| 1187 | return; |
| 1188 | }; |
| 1189 | |
| 1190 | let Some(current_estimate) = bwe.last_estimate() else { |
| 1191 | // No estimate yet, no padding |
| 1192 | return; |
| 1193 | }; |
| 1194 | let is_overuse = bwe.is_overusing(); |
| 1195 | |
| 1196 | let has_active_media = self.has_active_outgoing_media(); |
| 1197 | |
| 1198 | // Calculate pacing and padding rates |
| 1199 | let result = self |
| 1200 | .pacer_control |
| 1201 | .calculate(has_active_media, current_estimate, is_overuse); |
| 1202 | |
| 1203 | self.pacer.set_padding_rate(result.padding_rate); |
| 1204 | self.pacer.set_pacing_rate(result.pacing_rate); |
| 1205 | } |
| 1206 | |
| 1207 | fn has_active_outgoing_media(&self) -> bool { |
| 1208 | self.medias |
| 1209 | .iter() |
| 1210 | .any(|m| m.direction().is_sending() && !m.disabled()) |
| 1211 | } |
| 1212 | |
| 1213 | pub fn media_by_mid(&self, mid: Mid) -> Option<&Media> { |
| 1214 | self.medias.iter().find(|m| m.mid() == mid) |
| 1215 | } |
| 1216 | |
| 1217 | pub fn media_by_mid_mut(&mut self, mid: Mid) -> Option<&mut Media> { |
| 1218 | self.medias.iter_mut().find(|m| m.mid() == mid) |
| 1219 | } |
| 1220 | |
| 1221 | fn do_payload(&mut self) -> Result<(), RtcError> { |
| 1222 | let mtu = self.mtu(); |
| 1223 | for m in &mut self.medias { |
| 1224 | m.do_payload( |
| 1225 | &mut self.streams, |
| 1226 | &self.codec_config, |
| 1227 | self.vp9_packetizer_mode, |
| 1228 | mtu, |
| 1229 | )?; |
| 1230 | } |
| 1231 | |
| 1232 | Ok(()) |
| 1233 | } |
| 1234 | |
| 1235 | pub fn set_direction(&mut self, mid: Mid, direction: Direction) -> bool { |
| 1236 | let Some(media) = self.media_by_mid_mut(mid) else { |
| 1237 | return false; |
| 1238 | }; |
| 1239 | let old_dir = media.direction(); |
| 1240 | if old_dir == direction { |
| 1241 | return false; |
| 1242 | } |
| 1243 | |
| 1244 | media.set_direction(direction); |
| 1245 | |
| 1246 | if old_dir.is_sending() && !direction.is_sending() { |
| 1247 | self.streams.reset_buffers_tx(mid); |
| 1248 | } |
| 1249 | |
| 1250 | let max_seq_lookup = make_max_seq_lookup(&self.max_rx_seq_lookup); |
| 1251 | if old_dir.is_receiving() && !direction.is_receiving() { |
| 1252 | self.streams.reset_buffers_rx(mid, max_seq_lookup); |
| 1253 | } |
| 1254 | |
| 1255 | true |
| 1256 | } |
| 1257 | |
| 1258 | pub fn stop_media(&mut self, mid: Mid) -> bool { |
| 1259 | let Some(media) = self.media_by_mid_mut(mid) else { |
| 1260 | return false; |
| 1261 | }; |
| 1262 | if media.disabled() { |
| 1263 | return false; |
| 1264 | } |
| 1265 | |
| 1266 | media.mark_stopped(); |
| 1267 | self.set_direction(mid, Direction::Inactive); |
| 1268 | |
| 1269 | true |
| 1270 | } |
| 1271 | |
| 1272 | pub fn is_request_keyframe_possible(&self, kind: KeyframeRequestKind) -> bool { |
| 1273 | // TODO: It's possible to have different set of feedback enabled for different |
| 1274 | // payload types. I.e. we could have FIR enabled for H264, but not for VP8. |
| 1275 | // We might want to make this check more fine grained by testing which PT is |
| 1276 | // in "active use" right now. |
| 1277 | self.codec_config.iter().any(|r| match kind { |
| 1278 | KeyframeRequestKind::Pli => r.fb_pli, |
| 1279 | KeyframeRequestKind::Fir => r.fb_fir, |
| 1280 | }) |
| 1281 | } |
| 1282 | |
| 1283 | /// Checks whether the SRTP contexts are up. |
| 1284 | pub fn is_connected(&self) -> bool { |
| 1285 | self.srtp_rx.is_some() && self.srtp_tx.is_some() |
| 1286 | } |
| 1287 | } |
| 1288 | |
| 1289 | pub struct PacketReceipt { |
| 1290 | pub header: RtpHeader, |
| 1291 | pub seq_no: SeqNo, |
| 1292 | pub is_padding: bool, |
| 1293 | pub payload_size: usize, |
| 1294 | } |
| 1295 | |
| 1296 | /// Find the PayloadParams for the given Pt, either when the Pt is the main Pt for the Codec or |
| 1297 | /// when it's the RTX Pt. |
| 1298 | fn main_payload_params(c: &CodecConfig, pt: Pt) -> Option<&PayloadParams> { |
| 1299 | c.iter().find(|p| p.pt == pt || p.resend == Some(pt)) |
| 1300 | } |
| 1301 | |
| 1302 | fn make_max_seq_lookup(map: &HashMap<Ssrc, SeqNo>) -> impl Fn(Ssrc) -> Option<SeqNo> + '_ { |
| 1303 | |ssrc| map.get(&ssrc).cloned() |
| 1304 | } |
| 1305 | |
| 1306 | fn update_max_seq(map: &mut HashMap<Ssrc, SeqNo>, ssrc: Ssrc, seq_no: SeqNo) { |
| 1307 | let current = map.entry(ssrc).or_insert(seq_no); |
| 1308 | if seq_no > *current { |
| 1309 | *current = seq_no; |
| 1310 | } |
| 1311 | } |
| 1312 | |
| 1313 | #[cfg(test)] |
| 1314 | mod tests { |
| 1315 | use super::*; |
| 1316 | use crate::RtcConfig; |
| 1317 | use crate::io::DATAGRAM_MTU_TARGET; |
| 1318 | |
| 1319 | #[test] |
| 1320 | fn session_mtu_matches_config() { |
| 1321 | let cfg = RtcConfig::default(); |
| 1322 | let s = Session::new(&cfg); |
| 1323 | assert_eq!(s.mtu(), DATAGRAM_MTU_TARGET); |
| 1324 | |
| 1325 | let cfg = RtcConfig::default().set_mtu(900..=1280); |
| 1326 | let s = Session::new(&cfg); |
| 1327 | assert_eq!(s.mtu(), 900); |
| 1328 | } |
| 1329 | } |