File
Blob: firmware/vendor/str0m/src/change/sdp.rs
| 1 | //! Strategy that amends the [`Rtc`] via SDP OFFER/ANSWER negotiation. |
| 2 | |
| 3 | use std::fmt; |
| 4 | use std::ops::{Deref, DerefMut}; |
| 5 | use std::slice::Iter; |
| 6 | |
| 7 | use crate::Rtc; |
| 8 | use crate::RtcError; |
| 9 | use crate::channel::ChannelId; |
| 10 | use crate::crypto::Fingerprint; |
| 11 | use crate::format::CodecConfig; |
| 12 | use crate::format::PayloadParams; |
| 13 | use crate::media::{Media, Rids, Simulcast}; |
| 14 | use crate::packet::MediaKind; |
| 15 | use crate::rtp_::MidRid; |
| 16 | use crate::rtp_::Rid; |
| 17 | use crate::rtp_::{Direction, Extension, ExtensionMap, Mid, Pt, Ssrc}; |
| 18 | use crate::sctp::ChannelConfig; |
| 19 | use crate::sctp::RtcSctp; |
| 20 | use crate::sdp::{self, MediaAttribute, MediaLine, MediaType, Msid, Sdp}; |
| 21 | use crate::sdp::{Proto, SessionAttribute, Setup, SimulcastGroups}; |
| 22 | use crate::session::Session; |
| 23 | use crate::{Candidate, IceCreds}; |
| 24 | use str0m_proto::Id; |
| 25 | |
| 26 | pub use crate::sdp::{SdpAnswer, SdpOffer}; |
| 27 | use crate::streams::{DEFAULT_RTX_CACHE_DURATION, DEFAULT_RTX_RATIO_CAP, Streams}; |
| 28 | |
| 29 | /// Changes to the Rtc via SDP Offer/Answer dance. |
| 30 | pub struct SdpApi<'a> { |
| 31 | rtc: &'a mut Rtc, |
| 32 | changes: Changes, |
| 33 | } |
| 34 | |
| 35 | impl<'a> SdpApi<'a> { |
| 36 | pub(crate) fn new(rtc: &'a mut Rtc) -> Self { |
| 37 | SdpApi { |
| 38 | rtc, |
| 39 | changes: Changes::default(), |
| 40 | } |
| 41 | } |
| 42 | |
| 43 | /// Accept an [`SdpOffer`] from the remote peer. If this call returns successfully, the |
| 44 | /// changes will have been made to the session. The resulting [`SdpAnswer`] should be |
| 45 | /// sent to the remote peer. |
| 46 | /// |
| 47 | /// <b>Note. Pending changes from a previous non-completed [`SdpApi`][super::SdpApi] will be |
| 48 | /// considered rolled back when calling this function.</b> |
| 49 | /// |
| 50 | /// The incoming SDP is validated in various ways which can cause this call to fail. |
| 51 | /// Example of such problems would be an SDP without any m-lines, missing `a=fingerprint` |
| 52 | /// or if `a=group` doesn't match the number of m-lines. |
| 53 | /// |
| 54 | /// ```no_run |
| 55 | /// # use std::time::Instant; |
| 56 | /// # use str0m::Rtc; |
| 57 | /// # use str0m::change::{SdpOffer}; |
| 58 | /// // obtain offer from remote peer. |
| 59 | /// let json_offer: &[u8] = todo!(); |
| 60 | /// let offer: SdpOffer = serde_json::from_slice(json_offer).unwrap(); |
| 61 | /// |
| 62 | /// let mut rtc = Rtc::new(Instant::now()); |
| 63 | /// let answer = rtc.sdp_api().accept_offer(offer).unwrap(); |
| 64 | /// |
| 65 | /// // send json_answer to remote peer. |
| 66 | /// let json_answer = serde_json::to_vec(&answer).unwrap(); |
| 67 | /// ``` |
| 68 | pub fn accept_offer(self, offer: SdpOffer) -> Result<SdpAnswer, RtcError> { |
| 69 | debug!("Accept offer"); |
| 70 | |
| 71 | // Invalidate any outstanding PendingOffer. |
| 72 | self.rtc.next_change_id(); |
| 73 | |
| 74 | if offer.media_lines.is_empty() { |
| 75 | return Err(RtcError::RemoteSdp("No m-lines in offer".into())); |
| 76 | } |
| 77 | |
| 78 | if self.rtc.ice.ice_lite() && offer.session.ice_lite() { |
| 79 | return Err(RtcError::RemoteSdp( |
| 80 | "Both peers being ICE-Lite not supported".into(), |
| 81 | )); |
| 82 | } |
| 83 | |
| 84 | add_ice_details(self.rtc, &offer, None)?; |
| 85 | |
| 86 | if self.rtc.remote_fingerprint.is_none() { |
| 87 | if let Some(f) = offer.fingerprint() { |
| 88 | self.rtc.remote_fingerprint = Some(f); |
| 89 | } else { |
| 90 | self.rtc.disconnect(); |
| 91 | return Err(RtcError::RemoteSdp("missing a=fingerprint".into())); |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | if !self.rtc.dtls.is_inited() { |
| 96 | // The side that makes the first offer is the controlling side, unless they |
| 97 | // are ICE Lite, in which case the roles are reversed (see RFC 5245). |
| 98 | self.rtc.ice.set_controlling(offer.session.ice_lite()); |
| 99 | } |
| 100 | |
| 101 | // Ensure setup=active/passive is corresponding remote and init dtls. |
| 102 | init_dtls(self.rtc, &offer)?; |
| 103 | |
| 104 | let remote_max_message_size = extract_max_message_size(offer.media_lines.iter()); |
| 105 | |
| 106 | // Extract a=sctp-init from remote offer before apply_offer consumes it. |
| 107 | let remote_sctp_init = offer.sctp_init().map(|v| v.to_owned()); |
| 108 | |
| 109 | let has_snap = process_remote_sctp_init(&mut self.rtc.sctp, remote_sctp_init.as_deref())?; |
| 110 | |
| 111 | // Modify session with offer. |
| 112 | apply_offer(&mut self.rtc.session, offer)?; |
| 113 | |
| 114 | // Handle potentially new m=application line. |
| 115 | let client = self.rtc.dtls.is_active().expect("DTLS active to be set"); |
| 116 | |
| 117 | // Generate local sctp-init for the answer: |
| 118 | // When the remote included a=sctp-init, we reciprocate (§5.4). |
| 119 | if has_snap && !self.rtc.sctp.is_inited() && !self.rtc.sctp.ensure_local_snap_init() { |
| 120 | warn!("Failed to generate SNAP INIT chunk, degrading to non-SNAP"); |
| 121 | } |
| 122 | |
| 123 | if self.rtc.session.app().is_some() { |
| 124 | let init_data = self.rtc.sctp.build_snap_init_data(); |
| 125 | self.rtc |
| 126 | .try_init_sctp(client, init_data, remote_max_message_size)?; |
| 127 | } |
| 128 | |
| 129 | let params = AsSdpParams::new(self.rtc, None); |
| 130 | let sdp = as_sdp(&self.rtc.session, params); |
| 131 | |
| 132 | debug!("Create answer"); |
| 133 | Ok(sdp.into()) |
| 134 | } |
| 135 | |
| 136 | /// Accept an answer to a previously created [`SdpOffer`]. |
| 137 | /// |
| 138 | /// This function returns an [`RtcError::ChangesOutOfOrder`] if we have created and applied another |
| 139 | /// [`SdpApi`][super::SdpApi] before calling this. The same also happens if we use |
| 140 | /// [`SdpApi::accept_offer()`] before using this pending instance. |
| 141 | /// |
| 142 | /// ```no_run |
| 143 | /// # use std::time::Instant; |
| 144 | /// # use str0m::Rtc; |
| 145 | /// # use str0m::media::{MediaKind, Direction}; |
| 146 | /// # use str0m::change::SdpAnswer; |
| 147 | /// let mut rtc = Rtc::new(Instant::now()); |
| 148 | /// |
| 149 | /// let mut changes = rtc.sdp_api(); |
| 150 | /// let mid = changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 151 | /// let (offer, pending) = changes.apply().unwrap(); |
| 152 | /// |
| 153 | /// // send offer to remote peer, receive answer back |
| 154 | /// let answer: SdpAnswer = todo!(); |
| 155 | /// |
| 156 | /// rtc.sdp_api().accept_answer(pending, answer).unwrap(); |
| 157 | /// ``` |
| 158 | pub fn accept_answer( |
| 159 | self, |
| 160 | mut pending: SdpPendingOffer, |
| 161 | answer: SdpAnswer, |
| 162 | ) -> Result<(), RtcError> { |
| 163 | debug!("Accept answer"); |
| 164 | |
| 165 | // Ensure we don't use the wrong changes below. We must use that of pending. |
| 166 | drop(self.changes); |
| 167 | |
| 168 | if !self.rtc.is_correct_change_id(pending.change_id) { |
| 169 | return Err(RtcError::ChangesOutOfOrder); |
| 170 | } |
| 171 | |
| 172 | if self.rtc.ice.ice_lite() && answer.session.ice_lite() { |
| 173 | return Err(RtcError::RemoteSdp( |
| 174 | "Both peers being ICE-Lite not supported".into(), |
| 175 | )); |
| 176 | } |
| 177 | |
| 178 | add_ice_details(self.rtc, &answer, Some(&pending))?; |
| 179 | |
| 180 | // Ensure setup=active/passive is corresponding remote and init dtls. |
| 181 | init_dtls(self.rtc, &answer)?; |
| 182 | |
| 183 | if self.rtc.remote_fingerprint.is_none() { |
| 184 | if let Some(f) = answer.fingerprint() { |
| 185 | self.rtc.remote_fingerprint = Some(f); |
| 186 | } else { |
| 187 | self.rtc.disconnect(); |
| 188 | return Err(RtcError::RemoteSdp("missing a=fingerprint".into())); |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | // Extract a=sctp-init from remote answer before apply_answer consumes it. |
| 193 | let remote_sctp_init = answer.sctp_init().map(|v| v.to_owned()); |
| 194 | |
| 195 | let expected_snap_answer = |
| 196 | !self.rtc.sctp.is_inited() && self.rtc.sctp.local_sctp_init_for_sdp().is_some(); |
| 197 | |
| 198 | // Validate or cache the remote value before mutating the session. This |
| 199 | // keeps a bad re-offer/re-answer from partially applying local state. |
| 200 | let has_snap = process_remote_sctp_init(&mut self.rtc.sctp, remote_sctp_init.as_deref())?; |
| 201 | |
| 202 | // Split out new channels, since that is not handled by the Session. |
| 203 | let new_channels = pending.changes.take_new_channels(); |
| 204 | |
| 205 | let remote_max_message_size = extract_max_message_size(answer.media_lines.iter()); |
| 206 | |
| 207 | // Modify session with answer |
| 208 | apply_answer(&mut self.rtc.session, pending.changes, answer)?; |
| 209 | |
| 210 | // Handle potentially new m=application line. |
| 211 | let client = self.rtc.dtls.is_active().expect("DTLS to be inited"); |
| 212 | |
| 213 | if expected_snap_answer && !has_snap { |
| 214 | debug!("Remote answer did not accept SNAP, falling back to regular SCTP handshake"); |
| 215 | self.rtc.sctp.disable_pending_snap(); |
| 216 | } |
| 217 | |
| 218 | if self.rtc.session.app().is_some() { |
| 219 | let init_data = self.rtc.sctp.build_snap_init_data(); |
| 220 | self.rtc |
| 221 | .try_init_sctp(client, init_data, remote_max_message_size)?; |
| 222 | } |
| 223 | |
| 224 | for (id, config) in new_channels { |
| 225 | self.rtc.chan.confirm(id, config); |
| 226 | } |
| 227 | |
| 228 | Ok(()) |
| 229 | } |
| 230 | |
| 231 | /// Test if any changes have been made. |
| 232 | /// |
| 233 | /// If changes have been made, nothing happens until we call [`SdpApi::apply()`]. |
| 234 | /// |
| 235 | /// ``` |
| 236 | /// # #[cfg(feature = "openssl")] { |
| 237 | /// # use std::time::Instant; |
| 238 | /// # use str0m::{Rtc, media::MediaKind, media::Direction}; |
| 239 | /// let mut rtc = Rtc::new(Instant::now()); |
| 240 | /// |
| 241 | /// let mut changes = rtc.sdp_api(); |
| 242 | /// assert!(!changes.has_changes()); |
| 243 | /// |
| 244 | /// let mid = changes.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None); |
| 245 | /// assert!(changes.has_changes()); |
| 246 | /// # } |
| 247 | /// ``` |
| 248 | pub fn has_changes(&self) -> bool { |
| 249 | !self.changes.0.is_empty() |
| 250 | } |
| 251 | |
| 252 | /// Add audio or video media and get the `mid` that will be used. |
| 253 | /// |
| 254 | /// Each call will result in a new m-line in the offer identified by the [`Mid`]. |
| 255 | /// |
| 256 | /// The mid is not valid to use until the SDP offer-answer dance is complete and |
| 257 | /// the mid been advertised via [`Event::MediaAdded`][crate::Event::MediaAdded]. |
| 258 | /// |
| 259 | /// * `stream_id` is used to synchronize media. It is `a=msid-semantic: WMS <streamId>` line in SDP. |
| 260 | /// * `track_id` is becomes both the track id in `a=msid <streamId> <trackId>` as well as the |
| 261 | /// CNAME in the RTP SDES. |
| 262 | /// |
| 263 | /// ``` |
| 264 | /// # #[cfg(feature = "openssl")] { |
| 265 | /// # use std::time::Instant; |
| 266 | /// # use str0m::{Rtc, media::MediaKind, media::Direction}; |
| 267 | /// let mut rtc = Rtc::new(Instant::now()); |
| 268 | /// |
| 269 | /// let mut changes = rtc.sdp_api(); |
| 270 | /// |
| 271 | /// let mid = changes.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None); |
| 272 | /// # } |
| 273 | /// ``` |
| 274 | pub fn add_media( |
| 275 | &mut self, |
| 276 | kind: MediaKind, |
| 277 | dir: Direction, |
| 278 | stream_id: Option<String>, |
| 279 | track_id: Option<String>, |
| 280 | simulcast: Option<crate::media::Simulcast>, |
| 281 | ) -> Mid { |
| 282 | let mid = self.rtc.new_mid(); |
| 283 | |
| 284 | // https://www.rfc-editor.org/rfc/rfc8830 |
| 285 | // msid-id = 1*64token-char |
| 286 | fn is_token_char(c: &char) -> bool { |
| 287 | // token-char = %x21 / %x23-27 / %x2A-2B / %x2D-2E / %x30-39 |
| 288 | // / %x41-5A / %x5E-7E |
| 289 | let u = *c as u32; |
| 290 | u == 0x21 |
| 291 | || (0x23..=0x27).contains(&u) |
| 292 | || (0x2a..=0x2b).contains(&u) |
| 293 | || (0x2d..=0x2e).contains(&u) |
| 294 | || (0x30..=0x39).contains(&u) |
| 295 | || (0x41..=0x5a).contains(&u) |
| 296 | || (0x5e..0x7e).contains(&u) |
| 297 | } |
| 298 | |
| 299 | let stream_id = if let Some(stream_id) = stream_id { |
| 300 | stream_id.chars().filter(is_token_char).take(64).collect() |
| 301 | } else { |
| 302 | Id::<20>::random().to_string() |
| 303 | }; |
| 304 | |
| 305 | let track_id = if let Some(track_id) = track_id { |
| 306 | track_id.chars().filter(is_token_char).take(64).collect() |
| 307 | } else { |
| 308 | Id::<20>::random().to_string() |
| 309 | }; |
| 310 | |
| 311 | let mut ssrcs = Vec::new(); |
| 312 | |
| 313 | // Main SSRC, not counting RTX. |
| 314 | let main_ssrc_count = simulcast.as_ref().map(|s| s.send.len()).unwrap_or(1); |
| 315 | |
| 316 | for _ in 0..main_ssrc_count { |
| 317 | let rtx = kind.is_video().then(|| self.rtc.session.streams.new_ssrc()); |
| 318 | ssrcs.push((self.rtc.session.streams.new_ssrc(), rtx)); |
| 319 | } |
| 320 | |
| 321 | // TODO: let user configure stream/track name. |
| 322 | let msid = Msid { |
| 323 | stream_id, |
| 324 | track_id: track_id.clone(), |
| 325 | }; |
| 326 | |
| 327 | let add = AddMedia { |
| 328 | mid, |
| 329 | cname: track_id, |
| 330 | msid, |
| 331 | kind, |
| 332 | dir, |
| 333 | ssrcs, |
| 334 | simulcast, |
| 335 | |
| 336 | // Added later |
| 337 | pts: vec![], |
| 338 | exts: ExtensionMap::empty(), |
| 339 | index: 0, |
| 340 | }; |
| 341 | |
| 342 | self.changes.0.push(Change::AddMedia(add)); |
| 343 | mid |
| 344 | } |
| 345 | |
| 346 | /// Change the direction of an already existing media. |
| 347 | /// |
| 348 | /// All media have a direction. The media can be added by this side via |
| 349 | /// [`SdpApi::add_media()`] or by the remote peer. Either way, the direction |
| 350 | /// of the line can be changed at any time. |
| 351 | /// |
| 352 | /// It's possible to set the direction [`Direction::Inactive`] for media that |
| 353 | /// will not be used by the session anymore. |
| 354 | /// |
| 355 | /// If the direction is set for media that doesn't exist, or if the direction is |
| 356 | /// the same that's already set [`SdpApi::apply()`] not require a negotiation. |
| 357 | pub fn set_direction(&mut self, mid: Mid, dir: Direction) { |
| 358 | let changed = self.rtc.session.set_direction(mid, dir); |
| 359 | |
| 360 | if changed { |
| 361 | self.changes.0.push(Change::Direction(mid, dir)); |
| 362 | } |
| 363 | } |
| 364 | |
| 365 | /// Stop an already existing media. |
| 366 | /// |
| 367 | /// The next generated offer emits the m-line with port 0 and excludes |
| 368 | /// it from the BUNDLE group, per [RFC 8843] §7.5.3. The remote |
| 369 | /// transceiver transitions to the "stopped" state and the m-line slot |
| 370 | /// becomes eligible for recycling. |
| 371 | /// |
| 372 | /// Unlike [`SdpApi::set_direction()`] with [`Direction::Inactive`], |
| 373 | /// a stopped m-line cannot be reactivated. |
| 374 | /// |
| 375 | /// If the media doesn't exist, or is already stopped, [`SdpApi::apply()`] |
| 376 | /// will not require a negotiation. |
| 377 | /// |
| 378 | /// [RFC 8843]: https://datatracker.ietf.org/doc/html/rfc8843#section-7.5.3 |
| 379 | pub fn stop_media(&mut self, mid: Mid) { |
| 380 | let changed = self.rtc.session.stop_media(mid); |
| 381 | |
| 382 | if changed { |
| 383 | self.changes |
| 384 | .0 |
| 385 | .push(Change::Direction(mid, Direction::Inactive)); |
| 386 | } |
| 387 | } |
| 388 | |
| 389 | /// Add a new reliable ordered data channel and get the `id` that will be used. |
| 390 | /// |
| 391 | /// Use `add_channel_with_config` when unreliable or unordered data channels are preferred. |
| 392 | /// |
| 393 | /// The first ever data channel added to a WebRTC session results in a media |
| 394 | /// of a special "application" type in the SDP. The m-line is for a SCTP association over |
| 395 | /// DTLS, and all data channels are multiplexed over this single association. |
| 396 | /// |
| 397 | /// That means only the first ever `add_channel` will result in an [`SdpOffer`]. |
| 398 | /// Consecutive channels will be opened without needing a negotiation. |
| 399 | /// |
| 400 | /// The label is used to identify the data channel to the remote peer. This is mostly |
| 401 | /// useful when multiple channels are in use at the same time. |
| 402 | /// |
| 403 | /// ``` |
| 404 | /// # #[cfg(feature = "openssl")] { |
| 405 | /// # use std::time::Instant; |
| 406 | /// # use str0m::Rtc; |
| 407 | /// let mut rtc = Rtc::new(Instant::now()); |
| 408 | /// |
| 409 | /// let mut changes = rtc.sdp_api(); |
| 410 | /// |
| 411 | /// let cid = changes.add_channel("my special channel".to_string()); |
| 412 | /// # } |
| 413 | /// ``` |
| 414 | pub fn add_channel(&mut self, label: String) -> ChannelId { |
| 415 | self.add_channel_with_config(ChannelConfig { |
| 416 | label, |
| 417 | ..Default::default() |
| 418 | }) |
| 419 | } |
| 420 | |
| 421 | /// Add a new data channel with a given configuration and get the `id` that will be used. |
| 422 | /// |
| 423 | /// Refer to `add_channel` for more details. |
| 424 | /// |
| 425 | /// ``` |
| 426 | /// # #[cfg(feature = "openssl")] { |
| 427 | /// # use std::time::Instant; |
| 428 | /// # use str0m::{channel::{ChannelConfig, Reliability}, Rtc}; |
| 429 | /// let mut rtc = Rtc::new(Instant::now()); |
| 430 | /// |
| 431 | /// let mut changes = rtc.sdp_api(); |
| 432 | /// |
| 433 | /// let cid = changes.add_channel_with_config(ChannelConfig { |
| 434 | /// label: "my special channel".to_string(), |
| 435 | /// reliability: Reliability::MaxRetransmits{ retransmits: 0 }, |
| 436 | /// ordered: false, |
| 437 | /// ..Default::default() |
| 438 | /// }); |
| 439 | /// # } |
| 440 | /// ``` |
| 441 | pub fn add_channel_with_config(&mut self, config: ChannelConfig) -> ChannelId { |
| 442 | let has_media = self.rtc.session.app().is_some(); |
| 443 | let changes_contains_add_app = self.changes.contains_add_app(); |
| 444 | |
| 445 | if !has_media && !changes_contains_add_app { |
| 446 | let mid = self.rtc.new_mid(); |
| 447 | self.changes.0.push(Change::AddApp(mid)); |
| 448 | } |
| 449 | |
| 450 | let id = self.rtc.chan.new_channel(&config); |
| 451 | |
| 452 | self.changes.0.push(Change::AddChannel((id, config))); |
| 453 | |
| 454 | id |
| 455 | } |
| 456 | |
| 457 | /// Perform an ICE restart. |
| 458 | /// |
| 459 | /// Only one ICE restart can be pending at the time. Calling this repeatedly removes any other |
| 460 | /// pending ICE restart. |
| 461 | /// |
| 462 | /// The local ICE candidates can be kept as is, or be cleared out, in which case new ice |
| 463 | /// candidates must be added via [`Rtc::add_local_candidate`] before connectivity can be |
| 464 | /// re-established. |
| 465 | /// |
| 466 | /// Returns the new ICE credentials that will be used going forward. |
| 467 | pub fn ice_restart(&mut self, keep_local_candidates: bool) -> IceCreds { |
| 468 | self.changes |
| 469 | .retain(|c| !matches!(c, Change::IceRestart(_, _))); |
| 470 | |
| 471 | let new_creds = IceCreds::new(); |
| 472 | self.changes |
| 473 | .push(Change::IceRestart(new_creds.clone(), keep_local_candidates)); |
| 474 | |
| 475 | new_creds |
| 476 | } |
| 477 | |
| 478 | /// Attempt to apply the changes made. |
| 479 | /// |
| 480 | /// If this returns [`SdpOffer`], the caller the changes are |
| 481 | /// not happening straight away, and the caller is expected to do a negotiation with the remote |
| 482 | /// peer and apply the answer using [`SdpPendingOffer`]. |
| 483 | /// |
| 484 | /// In case this returns `None`, there either were no changes, or the changes could be applied |
| 485 | /// without doing a negotiation. Specifically for additional [`SdpApi::add_channel()`] |
| 486 | /// after the first, there is no negotiation needed. |
| 487 | /// |
| 488 | /// The [`SdpPendingOffer`] is valid until the next time we call this function, at which |
| 489 | /// point using it will raise an error. Using [`SdpApi::accept_offer()`] will also invalidate |
| 490 | /// the current [`SdpPendingOffer`]. |
| 491 | /// |
| 492 | /// ``` |
| 493 | /// # #[cfg(feature = "openssl")] { |
| 494 | /// # use std::time::Instant; |
| 495 | /// # use str0m::Rtc; |
| 496 | /// let mut rtc = Rtc::new(Instant::now()); |
| 497 | /// |
| 498 | /// let changes = rtc.sdp_api(); |
| 499 | /// assert!(changes.apply().is_none()); |
| 500 | /// # } |
| 501 | /// ``` |
| 502 | pub fn apply(self) -> Option<(SdpOffer, SdpPendingOffer)> { |
| 503 | if self.changes.is_empty() { |
| 504 | return None; |
| 505 | } |
| 506 | |
| 507 | let change_id = self.rtc.next_change_id(); |
| 508 | |
| 509 | let requires_negotiation = self.changes.0.iter().any(requires_negotiation); |
| 510 | |
| 511 | if requires_negotiation { |
| 512 | let offer = create_offer(self.rtc, &self.changes); |
| 513 | let pending = SdpPendingOffer { |
| 514 | change_id, |
| 515 | changes: self.changes, |
| 516 | }; |
| 517 | debug!("Create offer"); |
| 518 | Some((offer, pending)) |
| 519 | } else { |
| 520 | debug!("Apply direct changes"); |
| 521 | apply_direct_changes(self.rtc, self.changes); |
| 522 | None |
| 523 | } |
| 524 | } |
| 525 | |
| 526 | /// Combines the modifications made in [`SdpApi`] with those in [`SdpPendingOffer`]. |
| 527 | /// |
| 528 | /// This function merges the changes present in [`SdpApi`] with the changes |
| 529 | /// in [`SdpPendingOffer`]. In result this [`SdpApi`] will incorporate modifications |
| 530 | /// from both the previous [`SdpPendingOffer`] and any newly added changes. |
| 531 | /// |
| 532 | /// ## Example |
| 533 | /// |
| 534 | /// ```no_run |
| 535 | /// # use std::time::Instant; |
| 536 | /// # use str0m::media::{Direction, MediaKind}; |
| 537 | /// # use str0m::Rtc; |
| 538 | /// let mut rtc = Rtc::new(Instant::now()); |
| 539 | /// let mut changes = rtc.sdp_api(); |
| 540 | /// changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 541 | /// let (_offer, pending) = changes.apply().unwrap(); |
| 542 | /// |
| 543 | /// let mut changes = rtc.sdp_api(); |
| 544 | /// changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 545 | /// changes.merge(pending); |
| 546 | /// |
| 547 | /// // This `SdpOffer` will have changes from the first `SdpPendingChanges` |
| 548 | /// // and new changes from `SdpApi` |
| 549 | /// let (_offer, pending) = changes.apply().unwrap(); |
| 550 | /// ``` |
| 551 | pub fn merge(&mut self, mut pending_offer: SdpPendingOffer) { |
| 552 | pending_offer.retain_relevant(self.rtc); |
| 553 | |
| 554 | // Prepend the original pending changes before the current SdpApi's own changes. |
| 555 | // |
| 556 | // AddMedia / AddApp entries from the pending offer carry already-allocated MIDs |
| 557 | // that the remote peer may have committed to specific m-line positions. They must |
| 558 | // appear before any new changes so that as_new_medias() assigns them the same |
| 559 | // indices as in the original offer. |
| 560 | pending_offer.changes.0.append(&mut self.changes.0); |
| 561 | self.changes.0 = pending_offer.changes.0; |
| 562 | } |
| 563 | } |
| 564 | |
| 565 | /// Pending offer from a previous [`Rtc::sdp_api()`] call. |
| 566 | /// |
| 567 | /// This allows us to accept a remote answer. No changes have been made to the session |
| 568 | /// before we call [`SdpApi::accept_answer()`], which means that rolling back a |
| 569 | /// change is as simple as dropping this instance. |
| 570 | /// |
| 571 | /// ```no_run |
| 572 | /// # use std::time::Instant; |
| 573 | /// # use str0m::Rtc; |
| 574 | /// # use str0m::media::{MediaKind, Direction}; |
| 575 | /// # use str0m::change::SdpAnswer; |
| 576 | /// let mut rtc = Rtc::new(Instant::now()); |
| 577 | /// |
| 578 | /// let mut changes = rtc.sdp_api(); |
| 579 | /// let mid = changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 580 | /// let (offer, pending) = changes.apply().unwrap(); |
| 581 | /// |
| 582 | /// // send offer to remote peer, receive answer back |
| 583 | /// let answer: SdpAnswer = todo!(); |
| 584 | /// |
| 585 | /// rtc.sdp_api().accept_answer(pending, answer).unwrap(); |
| 586 | /// ``` |
| 587 | pub struct SdpPendingOffer { |
| 588 | change_id: usize, |
| 589 | changes: Changes, |
| 590 | } |
| 591 | |
| 592 | impl SdpPendingOffer { |
| 593 | /// Retains only the relevant changes in the `changes` vector based on the provided `Rtc` instance. |
| 594 | /// |
| 595 | /// This function filters the vector of `Change` instances stored in the current object and retains |
| 596 | /// only those changes that are considered relevant with respect to the provided `Rtc` instance. |
| 597 | fn retain_relevant(&mut self, rtc: &Rtc) { |
| 598 | fn is_relevant(rtc: &Rtc, c: &Change) -> bool { |
| 599 | match c { |
| 600 | Change::AddMedia(v) => rtc.media(v.mid).is_none(), |
| 601 | Change::AddApp(_) => rtc.session.app().is_none(), |
| 602 | Change::AddChannel(v) => rtc.chan.stream_id_by_channel_id(v.0).is_none(), |
| 603 | Change::Direction(m, d) => { |
| 604 | // If mid is missing, this is not relevant. |
| 605 | rtc.media(*m).map(|m| m.direction() != *d).unwrap_or(false) |
| 606 | } |
| 607 | Change::IceRestart(v, _) => rtc.ice.local_credentials() != v, |
| 608 | } |
| 609 | } |
| 610 | |
| 611 | self.changes.retain(|c| is_relevant(rtc, c)); |
| 612 | } |
| 613 | } |
| 614 | |
| 615 | #[derive(Default)] |
| 616 | pub(crate) struct Changes(pub Vec<Change>); |
| 617 | |
| 618 | impl Changes { |
| 619 | /// Details of the active ICE restart, if any. |
| 620 | /// |
| 621 | /// Returns the new local ICE credentials and the whether to keep local ICE candidates if an |
| 622 | /// ICE restart has been initiated in the offer, otherwise [`None`]. |
| 623 | fn ice_restart(&self) -> Option<(IceCreds, bool)> { |
| 624 | self.iter().find_map(|c| match c { |
| 625 | Change::IceRestart(creds, keep_local_candidates) => { |
| 626 | Some((creds.clone(), *keep_local_candidates)) |
| 627 | } |
| 628 | _ => None, |
| 629 | }) |
| 630 | } |
| 631 | } |
| 632 | |
| 633 | #[derive(Debug)] |
| 634 | #[allow(clippy::large_enum_variant)] |
| 635 | pub(crate) enum Change { |
| 636 | AddMedia(AddMedia), |
| 637 | AddApp(Mid), |
| 638 | AddChannel((ChannelId, ChannelConfig)), |
| 639 | Direction(Mid, Direction), |
| 640 | IceRestart(IceCreds, bool), |
| 641 | } |
| 642 | |
| 643 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 644 | pub(crate) struct AddMedia { |
| 645 | pub mid: Mid, |
| 646 | pub cname: String, |
| 647 | pub msid: Msid, |
| 648 | pub kind: MediaKind, |
| 649 | pub dir: Direction, |
| 650 | pub ssrcs: Vec<(Ssrc, Option<Ssrc>)>, |
| 651 | pub simulcast: Option<Simulcast>, |
| 652 | |
| 653 | // pts and index are filled in when creating the SDP OFFER. |
| 654 | // The default PT order is set by the Session (BUNDLE). |
| 655 | // TODO: We can make this configurable here too. |
| 656 | pub pts: Vec<Pt>, |
| 657 | pub exts: ExtensionMap, |
| 658 | pub index: usize, |
| 659 | } |
| 660 | |
| 661 | impl Deref for Changes { |
| 662 | type Target = Vec<Change>; |
| 663 | |
| 664 | fn deref(&self) -> &Self::Target { |
| 665 | &self.0 |
| 666 | } |
| 667 | } |
| 668 | |
| 669 | impl DerefMut for Changes { |
| 670 | fn deref_mut(&mut self) -> &mut Self::Target { |
| 671 | &mut self.0 |
| 672 | } |
| 673 | } |
| 674 | |
| 675 | fn requires_negotiation(c: &Change) -> bool { |
| 676 | match c { |
| 677 | Change::IceRestart(_, _) => true, |
| 678 | Change::AddMedia(_) => true, |
| 679 | Change::AddApp(_) => true, |
| 680 | Change::AddChannel(_) => false, |
| 681 | Change::Direction(_, _) => true, |
| 682 | } |
| 683 | } |
| 684 | |
| 685 | fn apply_direct_changes(rtc: &mut Rtc, mut changes: Changes) { |
| 686 | // Split out new channels, since that is not handled by the Session. |
| 687 | let new_channels = changes.take_new_channels(); |
| 688 | |
| 689 | for (id, config) in new_channels { |
| 690 | rtc.chan.confirm(id, config); |
| 691 | } |
| 692 | } |
| 693 | |
| 694 | fn create_offer(rtc: &mut Rtc, changes: &Changes) -> SdpOffer { |
| 695 | if !rtc.dtls.is_inited() { |
| 696 | // The side that makes the first offer is the controlling side, unless they |
| 697 | // are ICE Lite, in which case the roles are reversed (see RFC 5245). |
| 698 | rtc.ice.set_controlling(!rtc.ice.ice_lite()); |
| 699 | } |
| 700 | |
| 701 | // Generate local sctp-init for SNAP if enabled and we have/will have an app m-line. |
| 702 | if rtc.sctp.snap_enabled() |
| 703 | && (rtc.session.app().is_some() || changes.contains_add_app()) |
| 704 | && !rtc.sctp.is_inited() |
| 705 | && !rtc.sctp.ensure_local_snap_init() |
| 706 | { |
| 707 | warn!("Failed to generate SNAP INIT chunk, degrading to non-SNAP"); |
| 708 | } |
| 709 | |
| 710 | let params = AsSdpParams::new(rtc, Some(changes)); |
| 711 | let sdp = as_sdp(&rtc.session, params); |
| 712 | |
| 713 | sdp.into() |
| 714 | } |
| 715 | |
| 716 | fn add_ice_details( |
| 717 | rtc: &mut Rtc, |
| 718 | sdp: &Sdp, |
| 719 | pending: Option<&SdpPendingOffer>, |
| 720 | ) -> Result<(), RtcError> { |
| 721 | let Some(creds) = sdp.ice_creds() else { |
| 722 | return Err(RtcError::RemoteSdp("missing a=ice-ufrag/pwd".into())); |
| 723 | }; |
| 724 | |
| 725 | // If we are handling an **offer** from the remote, differing ICE credentials indicate an ICE |
| 726 | // restart initiated by the remote. |
| 727 | // |
| 728 | // If we are handling an **answer** from the remote, differing ICE credentials indicate an |
| 729 | // acceptance of an ICE restart we requested. |
| 730 | let ice_restart = match rtc.ice.remote_credentials() { |
| 731 | Some(v) => *v != creds, |
| 732 | None => false, |
| 733 | }; |
| 734 | if ice_restart { |
| 735 | let (new_local_creds, keep_local_candidates) = if let Some(pending) = pending { |
| 736 | // Since we have a pending, this is an answer to our offer. |
| 737 | pending.changes.ice_restart().ok_or_else(|| |
| 738 | // Answer contained changed remote creds, indicating an ice restart |
| 739 | // but since we have no pending ice-creds, we didn't initiate it |
| 740 | // Ice restart in an ANSWER breaks spec. |
| 741 | RtcError::RemoteSdp( |
| 742 | "Ice restart in answer without one in the preceeding offer".into(), |
| 743 | ))? |
| 744 | } else { |
| 745 | // The remote OFFER had an ice restart, and we need to respond with |
| 746 | // new credentials in the ANSWER. |
| 747 | (IceCreds::new(), true) |
| 748 | }; |
| 749 | |
| 750 | rtc.ice |
| 751 | .ice_restart(new_local_creds.clone(), keep_local_candidates); |
| 752 | } |
| 753 | |
| 754 | rtc.ice.set_remote_credentials(creds); |
| 755 | |
| 756 | for r in sdp.ice_candidates() { |
| 757 | rtc.ice.add_remote_candidate(r.clone()); |
| 758 | } |
| 759 | |
| 760 | Ok(()) |
| 761 | } |
| 762 | |
| 763 | fn init_dtls(rtc: &mut Rtc, remote_sdp: &Sdp) -> Result<(), RtcError> { |
| 764 | let setup = match remote_sdp.setup() { |
| 765 | Some(v) => match v { |
| 766 | // Remote being ActPass, we take Passive role. |
| 767 | Setup::ActPass => Setup::Passive, |
| 768 | _ => v.invert(), |
| 769 | }, |
| 770 | |
| 771 | None => { |
| 772 | warn!("Missing a=setup line"); |
| 773 | Setup::Passive |
| 774 | } |
| 775 | }; |
| 776 | |
| 777 | let active = setup == Setup::Active; |
| 778 | rtc.init_dtls(active)?; |
| 779 | |
| 780 | Ok(()) |
| 781 | } |
| 782 | |
| 783 | /// Shared logic for processing a remote `a=sctp-init` attribute from an offer or answer. |
| 784 | /// |
| 785 | /// Returns `true` if this is a new SNAP negotiation (remote init was accepted), |
| 786 | /// `false` otherwise. |
| 787 | /// |
| 788 | /// The remote init bytes are stored in `RtcSctp.snap_init` for §5.6 |
| 789 | /// re-offer validation. |
| 790 | /// |
| 791 | /// When the remote includes `a=sctp-init`, we always accept and reciprocate — |
| 792 | /// Section 5.4 of draft-hancke-tsvwg-snap says the answerer MAY include the |
| 793 | /// attribute, and doing so is beneficial. |
| 794 | /// |
| 795 | /// When an initial answer omits or rejects `a=sctp-init`, callers are expected |
| 796 | /// to fall back to a regular SCTP handshake. |
| 797 | fn process_remote_sctp_init( |
| 798 | sctp: &mut RtcSctp, |
| 799 | remote_init: Option<&str>, |
| 800 | ) -> Result<bool, RtcError> { |
| 801 | if let Some(remote_init_str) = remote_init { |
| 802 | if sctp.is_inited() { |
| 803 | // §5.6: SCTP already established — remote MUST re-send the |
| 804 | // same sctp-init value on subsequent offers/answers. |
| 805 | match sctp.snap_remote_init_string() { |
| 806 | Some(cached) if cached == remote_init_str => { |
| 807 | debug!("Remote re-sent expected a=sctp-init for established association"); |
| 808 | } |
| 809 | Some(_) => { |
| 810 | return Err(RtcError::RemoteSdp( |
| 811 | "Changed a=sctp-init for existing SCTP association".into(), |
| 812 | )); |
| 813 | } |
| 814 | None => { |
| 815 | // SCTP was established without SNAP but the remote is now |
| 816 | // sending a=sctp-init. We can't transition to SNAP |
| 817 | // mid-session, so just ignore it. |
| 818 | debug!("Ignoring a=sctp-init for non-SNAP established SCTP association"); |
| 819 | } |
| 820 | } |
| 821 | Ok(false) |
| 822 | } else if sctp.set_remote_snap_init_string(remote_init_str) { |
| 823 | Ok(true) |
| 824 | } else { |
| 825 | debug!("Ignoring malformed a=sctp-init"); |
| 826 | Ok(false) |
| 827 | } |
| 828 | } else if sctp.is_snap_established() { |
| 829 | Err(RtcError::RemoteSdp( |
| 830 | "Missing a=sctp-init for established SNAP SCTP association".into(), |
| 831 | )) |
| 832 | } else { |
| 833 | Ok(false) |
| 834 | } |
| 835 | } |
| 836 | |
| 837 | fn as_sdp(session: &Session, params: AsSdpParams) -> Sdp { |
| 838 | let (media_lines, mids, stream_ids) = { |
| 839 | let mut v = as_media_lines(session); |
| 840 | |
| 841 | let mut new_lines = vec![]; |
| 842 | |
| 843 | // When creating new m-lines from the pending changes, the m-line index starts from this. |
| 844 | let new_index_start = v.len(); |
| 845 | |
| 846 | // If there are additions in the pending changes, prepend them now. |
| 847 | if let Some(pending) = params.pending { |
| 848 | new_lines = pending |
| 849 | .as_new_medias(new_index_start, &session.codec_config, &session.exts) |
| 850 | .collect(); |
| 851 | } |
| 852 | |
| 853 | // Add potentially new m-lines to the existing ones. |
| 854 | v.extend(new_lines.iter().map(|n| n as &dyn AsSdpMediaLine)); |
| 855 | |
| 856 | // Turn into sdp::MediaLine (m-line). |
| 857 | let mut lines = v |
| 858 | .iter() |
| 859 | .map(|m| { |
| 860 | // Candidates should only be in the first BUNDLE mid |
| 861 | let include_candidates = m.index() == 0; |
| 862 | |
| 863 | let attrs = params.media_attributes(include_candidates); |
| 864 | |
| 865 | // Already made send stream SSRCs |
| 866 | let mut ssrcs = session.streams.ssrcs_tx(m.mid()); |
| 867 | |
| 868 | // Merged with pending stream SSRCs |
| 869 | if let Some(pending) = params.pending { |
| 870 | ssrcs.extend(pending.ssrcs_for_mid(m.mid())) |
| 871 | } |
| 872 | |
| 873 | let params: Vec<_> = session |
| 874 | .codec_config |
| 875 | .all_for_kind(m.kind()) |
| 876 | .cloned() |
| 877 | .collect(); |
| 878 | |
| 879 | m.as_media_line(attrs, &ssrcs, &session.exts, ¶ms) |
| 880 | }) |
| 881 | .collect::<Vec<_>>(); |
| 882 | |
| 883 | // Use the same local limit for SDP and the SCTP reassembly policy. |
| 884 | for line in &mut lines { |
| 885 | if line.typ.is_channel() { |
| 886 | line.attrs |
| 887 | .retain(|a| !matches!(a, MediaAttribute::MaxMessageSize(_))); |
| 888 | line.attrs.push(MediaAttribute::MaxMessageSize( |
| 889 | params.local_max_message_size as usize, |
| 890 | )); |
| 891 | } |
| 892 | } |
| 893 | |
| 894 | // Add a=sctp-init to the application m-line if SNAP is configured. |
| 895 | if let Some(sctp_init) = ¶ms.local_sctp_init { |
| 896 | for line in &mut lines { |
| 897 | if line.typ.is_channel() { |
| 898 | line.attrs |
| 899 | .push(sdp::MediaAttribute::SctpInit(sctp_init.clone())); |
| 900 | } |
| 901 | } |
| 902 | } |
| 903 | |
| 904 | if let Some(pending) = params.pending { |
| 905 | pending.apply_to(&mut lines); |
| 906 | } |
| 907 | |
| 908 | // Mids go into the session part of the SDP. |
| 909 | // Rejected (disabled) m-lines must not be part of the BUNDLE group. |
| 910 | let mids = lines |
| 911 | .iter() |
| 912 | .filter(|l| !l.disabled) |
| 913 | .map(|l| l.mid()) |
| 914 | .collect(); |
| 915 | |
| 916 | let mut stream_ids = vec![]; |
| 917 | for msid in v.iter().filter_map(|v| v.msid()) { |
| 918 | if !stream_ids.contains(&msid.stream_id) { |
| 919 | stream_ids.push(msid.stream_id.clone()); |
| 920 | } |
| 921 | } |
| 922 | |
| 923 | (lines, mids, stream_ids) |
| 924 | }; |
| 925 | |
| 926 | // AllowMixedExts adds "a=extmap-allow-mixed" at session level to signal |
| 927 | // support for mixing one-byte and two-byte RTP header extensions. |
| 928 | // TODO: It would make sense to perform an actual negotiation, however |
| 929 | // just adding this line should work fine: |
| 930 | // https://github.com/meetecho/janus-gateway/blob/d2e74fdf9bb8aa7a39ed68ed28394afe1e0cd22d/src/sdp.c#L1519 |
| 931 | let mut attrs = vec![ |
| 932 | SessionAttribute::Group { |
| 933 | typ: "BUNDLE".into(), |
| 934 | mids, |
| 935 | }, |
| 936 | SessionAttribute::AllowMixedExts, |
| 937 | SessionAttribute::MsidSemantic { |
| 938 | semantic: "WMS".to_string(), |
| 939 | stream_ids, |
| 940 | }, |
| 941 | ]; |
| 942 | |
| 943 | if session.ice_lite { |
| 944 | attrs.push(SessionAttribute::IceLite); |
| 945 | } |
| 946 | |
| 947 | Sdp { |
| 948 | session: sdp::Session { |
| 949 | id: session.id(), |
| 950 | bw: None, |
| 951 | attrs, |
| 952 | }, |
| 953 | media_lines, |
| 954 | } |
| 955 | } |
| 956 | |
| 957 | fn apply_offer(session: &mut Session, offer: SdpOffer) -> Result<(), RtcError> { |
| 958 | offer.assert_consistency()?; |
| 959 | |
| 960 | update_session(session, &offer); |
| 961 | |
| 962 | let bundle_mids = offer.bundle_mids(); |
| 963 | let new_lines = sync_medias(session, &offer, true).map_err(RtcError::RemoteSdp)?; |
| 964 | |
| 965 | add_new_lines(session, &new_lines, true, bundle_mids).map_err(RtcError::RemoteSdp)?; |
| 966 | |
| 967 | ensure_stream_tx(session); |
| 968 | |
| 969 | Ok(()) |
| 970 | } |
| 971 | |
| 972 | fn apply_answer( |
| 973 | session: &mut Session, |
| 974 | pending: Changes, |
| 975 | answer: SdpAnswer, |
| 976 | ) -> Result<(), RtcError> { |
| 977 | answer.assert_consistency()?; |
| 978 | |
| 979 | update_session(session, &answer); |
| 980 | |
| 981 | let bundle_mids = answer.bundle_mids(); |
| 982 | let new_lines = sync_medias(session, &answer, false).map_err(RtcError::RemoteSdp)?; |
| 983 | |
| 984 | // The new_lines from the answer must correspond to what we sent in the offer. |
| 985 | if let Some(err) = pending.ensure_correct_answer(&new_lines) { |
| 986 | return Err(RtcError::RemoteSdp(err)); |
| 987 | } |
| 988 | |
| 989 | add_new_lines(session, &new_lines, false, bundle_mids).map_err(RtcError::RemoteSdp)?; |
| 990 | |
| 991 | // Add all pending changes (since we pre-allocated SSRC communicated in the Offer). |
| 992 | add_pending_changes(session, pending); |
| 993 | |
| 994 | ensure_stream_tx(session); |
| 995 | |
| 996 | Ok(()) |
| 997 | } |
| 998 | |
| 999 | fn ensure_stream_tx(session: &mut Session) { |
| 1000 | for media in &session.medias { |
| 1001 | // Only make send streams when we have to. |
| 1002 | if !media.direction().is_sending() { |
| 1003 | continue; |
| 1004 | } |
| 1005 | |
| 1006 | let mut rids: Vec<Option<Rid>> = vec![]; |
| 1007 | |
| 1008 | if let Some(sim) = media.simulcast() { |
| 1009 | for layer in &*sim.send { |
| 1010 | let rid: Rid = layer.restriction_id.0.as_str().into(); |
| 1011 | rids.push(Some(rid)); |
| 1012 | } |
| 1013 | } else { |
| 1014 | rids.push(None); |
| 1015 | } |
| 1016 | |
| 1017 | // If any payload param has RTX, we need to prepare for RTX. This is because we always |
| 1018 | // communicate a=ssrc lines, which need to be complete with main and RTX SSRC. |
| 1019 | let has_rtx = session |
| 1020 | .codec_config |
| 1021 | .iter() |
| 1022 | .filter(|p| media.remote_pts().contains(&p.pt)) |
| 1023 | .any(|p| p.resend().is_some()); |
| 1024 | |
| 1025 | for rid in rids { |
| 1026 | let midrid = MidRid(media.mid(), rid); |
| 1027 | |
| 1028 | // If we already have the stream, we don't make any new one. |
| 1029 | let has_stream = session.streams.stream_tx_by_midrid(midrid).is_some(); |
| 1030 | |
| 1031 | if has_stream { |
| 1032 | continue; |
| 1033 | } |
| 1034 | |
| 1035 | let (ssrc, rtx) = if has_rtx { |
| 1036 | let (ssrc, rtx) = session.streams.new_ssrc_pair(); |
| 1037 | (ssrc, Some(rtx)) |
| 1038 | } else { |
| 1039 | let ssrc = session.streams.new_ssrc(); |
| 1040 | (ssrc, None) |
| 1041 | }; |
| 1042 | |
| 1043 | let stream = session.streams.declare_stream_tx(ssrc, rtx, midrid); |
| 1044 | |
| 1045 | // Configure cache size |
| 1046 | let size = if media.kind().is_audio() { |
| 1047 | session.send_buffer_audio |
| 1048 | } else { |
| 1049 | session.send_buffer_video |
| 1050 | }; |
| 1051 | |
| 1052 | stream.set_rtx_cache(size, DEFAULT_RTX_CACHE_DURATION, DEFAULT_RTX_RATIO_CAP); |
| 1053 | } |
| 1054 | } |
| 1055 | } |
| 1056 | |
| 1057 | fn add_pending_changes(session: &mut Session, pending: Changes) { |
| 1058 | // For pending AddMedia, we have outgoing SSRC communicated that needs to be added. |
| 1059 | for change in pending.0 { |
| 1060 | let add_media = match change { |
| 1061 | Change::AddMedia(v) => v, |
| 1062 | _ => continue, |
| 1063 | }; |
| 1064 | |
| 1065 | let media = session |
| 1066 | .medias |
| 1067 | .iter_mut() |
| 1068 | .find(|m| m.mid() == add_media.mid) |
| 1069 | .expect("Media to be added for pending mid"); |
| 1070 | |
| 1071 | // the cname/msid has already been communicated in the offer, we need to kep |
| 1072 | // it the same once the m-line is created. |
| 1073 | media.set_cname(add_media.cname); |
| 1074 | media.set_msid(add_media.msid); |
| 1075 | |
| 1076 | // If there are RIDs, the SSRC order matches that of the rid order. |
| 1077 | let layers = add_media.simulcast.map(|x| x.send).unwrap_or(vec![]); |
| 1078 | |
| 1079 | for (i, (ssrc, rtx)) in add_media.ssrcs.into_iter().enumerate() { |
| 1080 | let maybe_layer = layers.get(i).cloned(); |
| 1081 | let midrid = MidRid(add_media.mid, maybe_layer.map(|layer| layer.rid)); |
| 1082 | |
| 1083 | let stream = session.streams.declare_stream_tx(ssrc, rtx, midrid); |
| 1084 | |
| 1085 | let size = if media.kind().is_audio() { |
| 1086 | session.send_buffer_audio |
| 1087 | } else { |
| 1088 | session.send_buffer_video |
| 1089 | }; |
| 1090 | |
| 1091 | stream.set_rtx_cache(size, DEFAULT_RTX_CACHE_DURATION, DEFAULT_RTX_RATIO_CAP); |
| 1092 | } |
| 1093 | } |
| 1094 | } |
| 1095 | |
| 1096 | /// Compares m-lines in Sdp with that already in the session. |
| 1097 | /// |
| 1098 | /// * Existing m-lines can apply changes (such as direction change). |
| 1099 | /// * New m-lines are returned to the caller paired with the session |
| 1100 | /// index they should occupy. The index is normally the next free slot, |
| 1101 | /// but can be a recycled slot (RFC 8829 §5.2.2) when the remote has |
| 1102 | /// replaced a previously disabled Media with a new mid. |
| 1103 | fn sync_medias<'a>( |
| 1104 | session: &mut Session, |
| 1105 | sdp: &'a Sdp, |
| 1106 | is_offer: bool, |
| 1107 | ) -> Result<Vec<(usize, &'a MediaLine)>, String> { |
| 1108 | let mut new_lines = Vec::with_capacity(sdp.media_lines.len()); |
| 1109 | let bundle_mids = sdp.bundle_mids(); |
| 1110 | |
| 1111 | for (idx, m) in sdp.media_lines.iter().enumerate() { |
| 1112 | // First, match existing m-lines. |
| 1113 | match m.typ { |
| 1114 | MediaType::Application => { |
| 1115 | if let Some((_, index)) = session.app() { |
| 1116 | if idx != *index { |
| 1117 | return index_err(m.mid()); |
| 1118 | } |
| 1119 | continue; |
| 1120 | } |
| 1121 | } |
| 1122 | MediaType::Audio | MediaType::Video => { |
| 1123 | if let Some(media) = session.medias.iter_mut().find(|l| l.mid() == m.mid()) { |
| 1124 | if idx != media.index() { |
| 1125 | return index_err(m.mid()); |
| 1126 | } |
| 1127 | |
| 1128 | update_media( |
| 1129 | media, |
| 1130 | m, |
| 1131 | &session.codec_config, |
| 1132 | &session.exts, |
| 1133 | &mut session.streams, |
| 1134 | bundle_mids, |
| 1135 | ); |
| 1136 | |
| 1137 | continue; |
| 1138 | } |
| 1139 | |
| 1140 | // Unknown mid at an index held by a disabled Media means |
| 1141 | // the remote has recycled the slot (RFC 8829 §5.2.2). |
| 1142 | // Recycling is only permitted in offers; an answer that |
| 1143 | // rewires a mid at a stopped slot is malformed. |
| 1144 | if let Some(pos) = session.medias.iter().position(|l| l.index() == idx) { |
| 1145 | if !is_offer { |
| 1146 | return Err(format!( |
| 1147 | "Answer recycles stopped m-line (not permitted per RFC 8829 §5.2.2): {}", |
| 1148 | m.mid() |
| 1149 | )); |
| 1150 | } |
| 1151 | if !session.medias[pos].disabled() { |
| 1152 | return index_err(m.mid()); |
| 1153 | } |
| 1154 | let retired_mid = session.medias[pos].mid(); |
| 1155 | session.medias.swap_remove(pos); |
| 1156 | session.streams.remove_streams_by_mid(retired_mid); |
| 1157 | } |
| 1158 | } |
| 1159 | _ => { |
| 1160 | continue; |
| 1161 | } |
| 1162 | } |
| 1163 | |
| 1164 | // Second, discover new m-lines. |
| 1165 | new_lines.push((idx, m)); |
| 1166 | } |
| 1167 | |
| 1168 | fn index_err<T>(mid: Mid) -> Result<T, String> { |
| 1169 | Err(format!("Changed order for m-line with mid: {mid}")) |
| 1170 | } |
| 1171 | |
| 1172 | Ok(new_lines) |
| 1173 | } |
| 1174 | |
| 1175 | /// Adds new m-lines as found in an offer or answer. |
| 1176 | /// |
| 1177 | /// Each entry in `new_lines` pairs an m-line with the session index it |
| 1178 | /// should occupy. |
| 1179 | fn add_new_lines( |
| 1180 | session: &mut Session, |
| 1181 | new_lines: &[(usize, &MediaLine)], |
| 1182 | is_offer: bool, |
| 1183 | bundle_mids: Option<&[Mid]>, |
| 1184 | ) -> Result<(), String> { |
| 1185 | for (idx, m) in new_lines { |
| 1186 | let idx = *idx; |
| 1187 | |
| 1188 | if m.typ.is_media() { |
| 1189 | let mut media = Media::from_remote_media_line(m, idx, is_offer); |
| 1190 | |
| 1191 | // For disabled (rejected) m-lines, don't fire open event |
| 1192 | // and the direction will be set to Inactive by update_media. |
| 1193 | // In max-bundle, port=0 with the MID in BUNDLE group is NOT rejected. |
| 1194 | let is_in_bundle = bundle_mids |
| 1195 | .map(|mids| mids.contains(&m.mid())) |
| 1196 | .unwrap_or(false); |
| 1197 | let is_rejected = m.disabled && !is_in_bundle; |
| 1198 | media.need_open_event = is_offer && !is_rejected; |
| 1199 | |
| 1200 | // Match/remap remote params. |
| 1201 | session |
| 1202 | .codec_config |
| 1203 | .update_params(&m.rtp_params(), m.direction()); |
| 1204 | |
| 1205 | // Remap the extension to that of the answer. |
| 1206 | session.exts.remap(&m.extmaps()); |
| 1207 | |
| 1208 | update_media( |
| 1209 | &mut media, |
| 1210 | m, |
| 1211 | &session.codec_config, |
| 1212 | &session.exts, |
| 1213 | &mut session.streams, |
| 1214 | bundle_mids, |
| 1215 | ); |
| 1216 | |
| 1217 | session.add_media(media); |
| 1218 | } else if m.typ.is_channel() { |
| 1219 | session.set_app(m.mid(), idx)?; |
| 1220 | } else { |
| 1221 | return Err(format!( |
| 1222 | "New m-line is neither media nor channel: {}", |
| 1223 | m.mid() |
| 1224 | )); |
| 1225 | } |
| 1226 | } |
| 1227 | |
| 1228 | Ok(()) |
| 1229 | } |
| 1230 | |
| 1231 | /// Update session level properties like |
| 1232 | /// Extensions from offer or answer. |
| 1233 | fn update_session(session: &mut Session, sdp: &Sdp) { |
| 1234 | // Does any m-line contain a a=rtcp-fb:xx transport-cc? |
| 1235 | let has_transport_cc = sdp |
| 1236 | .media_lines |
| 1237 | .iter() |
| 1238 | .any(|m| m.rtp_params().iter().any(|p| p.fb_transport_cc)); |
| 1239 | |
| 1240 | // Is the session level sequence number enabled? |
| 1241 | let has_twcc_header = session |
| 1242 | .exts |
| 1243 | .id_of(Extension::TransportSequenceNumber) |
| 1244 | .is_some(); |
| 1245 | |
| 1246 | // Since twcc feedback is session wide we enable it if there are _any_ |
| 1247 | // m-line with a a=rtcp-fb transport-cc parameter and the sequence number |
| 1248 | // header is enabled. It can later be disabled for specific m-lines based |
| 1249 | // on the extensions map. |
| 1250 | if has_transport_cc && has_twcc_header { |
| 1251 | session.enable_twcc_feedback(); |
| 1252 | } |
| 1253 | } |
| 1254 | |
| 1255 | /// Returns all media/channels as `AsMediaLine` trait. |
| 1256 | fn as_media_lines(session: &Session) -> Vec<&dyn AsSdpMediaLine> { |
| 1257 | let mut v = vec![]; |
| 1258 | |
| 1259 | if let Some(app) = session.app() { |
| 1260 | v.push(app as &dyn AsSdpMediaLine); |
| 1261 | } |
| 1262 | v.extend(session.medias().iter().map(|m| m as &dyn AsSdpMediaLine)); |
| 1263 | v.sort_by_key(|f| f.index()); |
| 1264 | v |
| 1265 | } |
| 1266 | |
| 1267 | fn update_media( |
| 1268 | media: &mut Media, |
| 1269 | m: &MediaLine, |
| 1270 | config: &CodecConfig, |
| 1271 | exts: &ExtensionMap, |
| 1272 | streams: &mut Streams, |
| 1273 | bundle_mids: Option<&[Mid]>, |
| 1274 | ) { |
| 1275 | // If the m-line has port=0, it could mean: |
| 1276 | // 1. The m-line is rejected (not in BUNDLE group) |
| 1277 | // 2. The m-line is bundled with another m-line (max-bundle format, RFC 8843) |
| 1278 | // |
| 1279 | // In max-bundle, secondary m-lines use port=0 to indicate they share transport |
| 1280 | // with the first m-line. This is NOT a rejection. |
| 1281 | let is_in_bundle = bundle_mids |
| 1282 | .map(|mids| mids.contains(&m.mid())) |
| 1283 | .unwrap_or(false); |
| 1284 | let is_rejected = m.disabled && !is_in_bundle; |
| 1285 | |
| 1286 | if is_rejected { |
| 1287 | if !media.disabled() { |
| 1288 | debug!( |
| 1289 | "Mid ({}) is rejected (port=0, not in BUNDLE), setting to Inactive", |
| 1290 | media.mid() |
| 1291 | ); |
| 1292 | } |
| 1293 | media.mark_stopped(); |
| 1294 | media.set_direction(Direction::Inactive); |
| 1295 | return; |
| 1296 | } |
| 1297 | |
| 1298 | // Direction changes |
| 1299 | // |
| 1300 | // All changes come from the other side, either via an incoming OFFER |
| 1301 | // or a ANSWER from our OFFER. Either way, the direction is inverted to |
| 1302 | // how we have it locally. |
| 1303 | let new_dir = m.direction().invert(); |
| 1304 | // |
| 1305 | let change_direction_disallowed = !media.remote_created() |
| 1306 | && media.direction() == Direction::Inactive |
| 1307 | && new_dir == Direction::SendOnly; |
| 1308 | |
| 1309 | if change_direction_disallowed { |
| 1310 | debug!( |
| 1311 | "Ignore attempt to change inactive to recvonly by remote peer for locally created mid: {}", |
| 1312 | media.mid() |
| 1313 | ); |
| 1314 | } else { |
| 1315 | media.set_direction(new_dir); |
| 1316 | } |
| 1317 | |
| 1318 | if new_dir.is_sending() { |
| 1319 | // The other side has declared how it EXPECTING to receive. We must only send |
| 1320 | // the RIDs declared in the answer. |
| 1321 | let rids = m.rids(); |
| 1322 | let rid_tx = if rids.is_empty() { |
| 1323 | Rids::None |
| 1324 | } else { |
| 1325 | Rids::Specific(rids) |
| 1326 | }; |
| 1327 | media.set_rid_tx(rid_tx); |
| 1328 | } |
| 1329 | if new_dir.is_receiving() { |
| 1330 | // The other side has declared what it proposes to send. We are accepting it. |
| 1331 | let rids = m.rids(); |
| 1332 | let rid_rx = if rids.is_empty() { |
| 1333 | Rids::Any |
| 1334 | } else { |
| 1335 | Rids::Specific(rids) |
| 1336 | }; |
| 1337 | media.set_rid_rx(rid_rx); |
| 1338 | } |
| 1339 | |
| 1340 | // Narrowing/ordering of of PT |
| 1341 | let pts: Vec<Pt> = m |
| 1342 | .rtp_params() |
| 1343 | .into_iter() |
| 1344 | .filter_map(|p| config.sdp_match_remote(p, m.direction())) |
| 1345 | .collect(); |
| 1346 | media.set_remote_pts(pts); |
| 1347 | |
| 1348 | let mut remote_extmap = ExtensionMap::empty(); |
| 1349 | for (id, ext) in m.extmaps().into_iter() { |
| 1350 | // The remapping of extensions should already have happened, which |
| 1351 | // means the ID are matching in the session to the remote. |
| 1352 | |
| 1353 | // Does the ID exist in session? |
| 1354 | let in_session = match exts.lookup(id) { |
| 1355 | Some(v) => v, |
| 1356 | None => continue, |
| 1357 | }; |
| 1358 | |
| 1359 | if in_session != ext { |
| 1360 | // Don't set any extensions that aren't enabled in Session. |
| 1361 | continue; |
| 1362 | } |
| 1363 | |
| 1364 | // Use the Extension from session, since there might be a special |
| 1365 | // serializer for cases like VLA. |
| 1366 | remote_extmap.set(id, in_session.clone()); |
| 1367 | } |
| 1368 | media.set_remote_extmap(remote_extmap); |
| 1369 | |
| 1370 | // SSRC changes |
| 1371 | // This will always be for ReceiverSource since any incoming a=ssrc line will be |
| 1372 | // about the remote side's SSRC. |
| 1373 | if !new_dir.is_receiving() { |
| 1374 | return; |
| 1375 | } |
| 1376 | |
| 1377 | // Simulcast configuration |
| 1378 | if let Some(s) = m.simulcast() { |
| 1379 | if s.is_munged { |
| 1380 | warn!("Not supporting simulcast via munging SDP"); |
| 1381 | } else if media.simulcast().is_none() { |
| 1382 | // Invert before setting, since it has a recv and send config. |
| 1383 | media.set_simulcast(s.invert()); |
| 1384 | } |
| 1385 | } |
| 1386 | |
| 1387 | // Only use pre-communicated SSRC if we are running without simulcast. |
| 1388 | // We found a bug in FF where the order of the simulcast lines does not |
| 1389 | // correspond to the order of the simulcast declarations. In this case |
| 1390 | // it's better to fall back on mid/rid dynamic mapping. |
| 1391 | if m.simulcast().is_some() { |
| 1392 | return; |
| 1393 | } |
| 1394 | |
| 1395 | let infos = m.ssrc_info(); |
| 1396 | let main = infos.iter().filter(|i| i.repairs.is_none()); |
| 1397 | |
| 1398 | for i in main { |
| 1399 | // TODO: If the remote is communicating _BOTH_ rid and a=ssrc this will fail. |
| 1400 | debug!("Adding pre-communicated SSRC: {:?}", i); |
| 1401 | let repair_ssrc = infos |
| 1402 | .iter() |
| 1403 | .find(|r| r.repairs == Some(i.ssrc)) |
| 1404 | .map(|r| r.ssrc); |
| 1405 | |
| 1406 | // If remote communicated a main a=ssrc, but no RTX, we will not send nacks. |
| 1407 | let midrid = MidRid(media.mid(), None); |
| 1408 | let suppress_nack = repair_ssrc.is_none(); |
| 1409 | streams.expect_stream_rx(i.ssrc, repair_ssrc, midrid, suppress_nack); |
| 1410 | } |
| 1411 | } |
| 1412 | |
| 1413 | fn extract_max_message_size(mut media_lines: Iter<MediaLine>) -> Option<u32> { |
| 1414 | if let Some(app_line) = media_lines.find(|m| m.typ.is_channel()) { |
| 1415 | if let Some(max_size) = app_line.max_message_size() { |
| 1416 | // RFC 8841 §6.1: a value of 0 means the peer can receive a message of |
| 1417 | // any size, subject to local capacity. Represent that as the maximum. |
| 1418 | let max_size = if max_size == 0 { |
| 1419 | u32::MAX |
| 1420 | } else { |
| 1421 | u32::try_from(max_size).unwrap_or(u32::MAX) |
| 1422 | }; |
| 1423 | return Some(max_size); |
| 1424 | } |
| 1425 | } |
| 1426 | |
| 1427 | None |
| 1428 | } |
| 1429 | |
| 1430 | trait AsSdpMediaLine { |
| 1431 | fn mid(&self) -> Mid; |
| 1432 | fn msid(&self) -> Option<&Msid>; |
| 1433 | fn index(&self) -> usize; |
| 1434 | fn kind(&self) -> MediaKind; |
| 1435 | fn as_media_line( |
| 1436 | &self, |
| 1437 | attrs: Vec<MediaAttribute>, |
| 1438 | ssrcs_tx: &[(Ssrc, Option<Ssrc>)], |
| 1439 | exts: &ExtensionMap, |
| 1440 | params: &[PayloadParams], |
| 1441 | ) -> MediaLine; |
| 1442 | } |
| 1443 | |
| 1444 | impl AsSdpMediaLine for (Mid, usize) { |
| 1445 | fn mid(&self) -> Mid { |
| 1446 | self.0 |
| 1447 | } |
| 1448 | fn msid(&self) -> Option<&Msid> { |
| 1449 | None |
| 1450 | } |
| 1451 | fn index(&self) -> usize { |
| 1452 | self.1 |
| 1453 | } |
| 1454 | fn kind(&self) -> MediaKind { |
| 1455 | MediaKind::Audio // doesn't matter for App |
| 1456 | } |
| 1457 | fn as_media_line( |
| 1458 | &self, |
| 1459 | mut attrs: Vec<MediaAttribute>, |
| 1460 | _ssrcs_tx: &[(Ssrc, Option<Ssrc>)], |
| 1461 | _exts: &ExtensionMap, |
| 1462 | _params: &[PayloadParams], |
| 1463 | ) -> MediaLine { |
| 1464 | attrs.push(MediaAttribute::Mid(self.0)); |
| 1465 | attrs.push(MediaAttribute::SctpPort(5000)); |
| 1466 | attrs.push(MediaAttribute::MaxMessageSize( |
| 1467 | crate::sctp::LOCAL_MAX_MESSAGE_SIZE as usize, |
| 1468 | )); |
| 1469 | |
| 1470 | MediaLine { |
| 1471 | typ: sdp::MediaType::Application, |
| 1472 | disabled: false, |
| 1473 | proto: Proto::Sctp, |
| 1474 | pts: vec![], |
| 1475 | bw: None, |
| 1476 | attrs, |
| 1477 | } |
| 1478 | } |
| 1479 | } |
| 1480 | |
| 1481 | impl AsSdpMediaLine for Media { |
| 1482 | fn mid(&self) -> Mid { |
| 1483 | Media::mid(self) |
| 1484 | } |
| 1485 | fn msid(&self) -> Option<&Msid> { |
| 1486 | Some(Media::msid(self)) |
| 1487 | } |
| 1488 | fn index(&self) -> usize { |
| 1489 | Media::index(self) |
| 1490 | } |
| 1491 | fn kind(&self) -> MediaKind { |
| 1492 | Media::kind(self) |
| 1493 | } |
| 1494 | fn as_media_line( |
| 1495 | &self, |
| 1496 | mut attrs: Vec<MediaAttribute>, |
| 1497 | ssrcs_tx: &[(Ssrc, Option<Ssrc>)], |
| 1498 | exts: &ExtensionMap, |
| 1499 | params: &[PayloadParams], |
| 1500 | ) -> MediaLine { |
| 1501 | if self.app_tmp { |
| 1502 | let app = (self.mid(), self.index()); |
| 1503 | return app.as_media_line(attrs, ssrcs_tx, exts, params); |
| 1504 | } |
| 1505 | |
| 1506 | attrs.push(MediaAttribute::Mid(self.mid())); |
| 1507 | |
| 1508 | let audio = self.kind() == MediaKind::Audio; |
| 1509 | for (id, ext) in self.remote_extmap().iter_by_media_type(audio) { |
| 1510 | attrs.push(MediaAttribute::ExtMap { |
| 1511 | id, |
| 1512 | ext: ext.clone(), |
| 1513 | }); |
| 1514 | } |
| 1515 | |
| 1516 | attrs.push(self.direction().into()); |
| 1517 | attrs.push(MediaAttribute::Msid(self.msid().clone())); |
| 1518 | attrs.push(MediaAttribute::RtcpMux); |
| 1519 | |
| 1520 | // The effective params start from the Session::codec_config to retain the |
| 1521 | // user's configured preferred order, however they are narrowed only include |
| 1522 | // those the remote peer wants. |
| 1523 | let effective_params = params.iter().filter(|p| self.remote_pts().contains(&p.pt)); |
| 1524 | |
| 1525 | let mut pts = vec![]; |
| 1526 | |
| 1527 | for p in effective_params { |
| 1528 | p.as_media_attrs(&mut attrs); |
| 1529 | |
| 1530 | // The pts that will be advertised in the SDP |
| 1531 | pts.push(p.pt()); |
| 1532 | if let Some(rtx) = p.resend() { |
| 1533 | pts.push(rtx); |
| 1534 | } |
| 1535 | } |
| 1536 | |
| 1537 | if let Some(s) = self.simulcast() { |
| 1538 | fn to_rids<'a>( |
| 1539 | gs: &'a SimulcastGroups, |
| 1540 | direction: &'static str, |
| 1541 | ) -> impl Iterator<Item = MediaAttribute> + 'a { |
| 1542 | gs.iter().map(move |layer| MediaAttribute::Rid { |
| 1543 | id: layer.restriction_id.clone(), |
| 1544 | direction, |
| 1545 | pt: vec![], |
| 1546 | restriction: layer.attributes.clone().unwrap_or_default(), |
| 1547 | }) |
| 1548 | } |
| 1549 | attrs.extend(to_rids(&s.recv, "recv")); |
| 1550 | attrs.extend(to_rids(&s.send, "send")); |
| 1551 | attrs.push(MediaAttribute::Simulcast(s.clone())); |
| 1552 | } |
| 1553 | |
| 1554 | // Outgoing SSRCs |
| 1555 | let msid = format!("{} {}", self.msid().stream_id, self.msid().track_id); |
| 1556 | for (ssrc, ssrc_rtx) in ssrcs_tx { |
| 1557 | attrs.push(MediaAttribute::Ssrc { |
| 1558 | ssrc: *ssrc, |
| 1559 | attr: "cname".to_string(), |
| 1560 | value: self.cname().to_string(), |
| 1561 | }); |
| 1562 | attrs.push(MediaAttribute::Ssrc { |
| 1563 | ssrc: *ssrc, |
| 1564 | attr: "msid".to_string(), |
| 1565 | value: msid.clone(), |
| 1566 | }); |
| 1567 | if let Some(ssrc_rtx) = ssrc_rtx { |
| 1568 | attrs.push(MediaAttribute::Ssrc { |
| 1569 | ssrc: *ssrc_rtx, |
| 1570 | attr: "cname".to_string(), |
| 1571 | value: self.cname().to_string(), |
| 1572 | }); |
| 1573 | attrs.push(MediaAttribute::Ssrc { |
| 1574 | ssrc: *ssrc_rtx, |
| 1575 | attr: "msid".to_string(), |
| 1576 | value: msid.clone(), |
| 1577 | }); |
| 1578 | } |
| 1579 | } |
| 1580 | |
| 1581 | for (ssrc, ssrc_rtx) in ssrcs_tx { |
| 1582 | if let Some(ssrc_rtx) = ssrc_rtx { |
| 1583 | attrs.push(MediaAttribute::SsrcGroup { |
| 1584 | semantics: "FID".to_string(), |
| 1585 | ssrcs: vec![*ssrc, *ssrc_rtx], |
| 1586 | }); |
| 1587 | } |
| 1588 | } |
| 1589 | |
| 1590 | MediaLine { |
| 1591 | typ: self.kind().into(), |
| 1592 | disabled: self.disabled(), |
| 1593 | proto: Proto::Srtp, |
| 1594 | pts, |
| 1595 | bw: None, |
| 1596 | attrs, |
| 1597 | } |
| 1598 | } |
| 1599 | } |
| 1600 | |
| 1601 | impl From<MediaKind> for MediaType { |
| 1602 | fn from(value: MediaKind) -> Self { |
| 1603 | match value { |
| 1604 | MediaKind::Audio => MediaType::Audio, |
| 1605 | MediaKind::Video => MediaType::Video, |
| 1606 | } |
| 1607 | } |
| 1608 | } |
| 1609 | |
| 1610 | struct AsSdpParams<'a, 'b> { |
| 1611 | pub candidates: Vec<Candidate>, |
| 1612 | pub creds: IceCreds, |
| 1613 | pub fingerprint: &'a Fingerprint, |
| 1614 | pub setup: Setup, |
| 1615 | pub pending: Option<&'b Changes>, |
| 1616 | pub local_sctp_init: Option<String>, |
| 1617 | pub local_max_message_size: u32, |
| 1618 | } |
| 1619 | |
| 1620 | impl<'a, 'b> AsSdpParams<'a, 'b> { |
| 1621 | pub fn new(rtc: &'a Rtc, pending: Option<&'b Changes>) -> Self { |
| 1622 | let (creds, candidates) = if let Some((new_creds, keep_local_candidates)) = |
| 1623 | pending.and_then(|p| p.ice_restart()) |
| 1624 | { |
| 1625 | if keep_local_candidates { |
| 1626 | // If we are performing an ICE restart and we are keeping the same |
| 1627 | // candidates we need to use ufrag from the new ICE credentials |
| 1628 | // in our offer. |
| 1629 | let mut new_candidates = rtc.ice.local_candidates().collect::<Vec<_>>(); |
| 1630 | for c in &mut new_candidates { |
| 1631 | c.set_ufrag(&new_creds.ufrag); |
| 1632 | } |
| 1633 | |
| 1634 | (new_creds, new_candidates) |
| 1635 | } else { |
| 1636 | (new_creds, vec![]) |
| 1637 | } |
| 1638 | } else { |
| 1639 | ( |
| 1640 | rtc.ice.local_credentials().clone(), |
| 1641 | rtc.ice.local_candidates().collect::<Vec<_>>(), |
| 1642 | ) |
| 1643 | }; |
| 1644 | |
| 1645 | // RFC 8842 Section 5.2/5.5: Offerers MUST always use a=setup:actpass. |
| 1646 | // Answerers use the negotiated role (active/passive). |
| 1647 | // We distinguish offer vs answer by the presence of `pending` (Some = offer). |
| 1648 | let setup = if pending.is_some() { |
| 1649 | // This is an offer — always actpass per RFC 8842 |
| 1650 | Setup::ActPass |
| 1651 | } else { |
| 1652 | // This is an answer — use the negotiated DTLS role |
| 1653 | match rtc.dtls.is_active() { |
| 1654 | Some(true) => Setup::Active, |
| 1655 | Some(false) => Setup::Passive, |
| 1656 | None => Setup::ActPass, |
| 1657 | } |
| 1658 | }; |
| 1659 | |
| 1660 | AsSdpParams { |
| 1661 | candidates, |
| 1662 | creds, |
| 1663 | fingerprint: rtc.dtls.local_fingerprint(), |
| 1664 | setup, |
| 1665 | pending, |
| 1666 | local_sctp_init: rtc.sctp.local_sctp_init_for_sdp(), |
| 1667 | local_max_message_size: rtc.sctp.local_max_message_size(), |
| 1668 | } |
| 1669 | } |
| 1670 | |
| 1671 | fn media_attributes(&self, include_candidates: bool) -> Vec<MediaAttribute> { |
| 1672 | use MediaAttribute::*; |
| 1673 | |
| 1674 | let mut v = if include_candidates { |
| 1675 | self.candidates |
| 1676 | .iter() |
| 1677 | .map(|c| Candidate(c.clone())) |
| 1678 | .collect() |
| 1679 | } else { |
| 1680 | vec![] |
| 1681 | }; |
| 1682 | |
| 1683 | v.push(IceUfrag(self.creds.ufrag.clone())); |
| 1684 | v.push(IcePwd(self.creds.pass.clone())); |
| 1685 | v.push(IceOptions("trickle".into())); |
| 1686 | v.push(Fingerprint(self.fingerprint.clone())); |
| 1687 | v.push(Setup(self.setup)); |
| 1688 | |
| 1689 | v |
| 1690 | } |
| 1691 | } |
| 1692 | |
| 1693 | impl fmt::Debug for SdpPendingOffer { |
| 1694 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 1695 | f.debug_struct("SdpPendingOffer").finish() |
| 1696 | } |
| 1697 | } |
| 1698 | |
| 1699 | impl Changes { |
| 1700 | pub fn contains_add_app(&self) -> bool { |
| 1701 | for i in 0..self.0.len() { |
| 1702 | if matches!(&self.0[i], Change::AddApp(_)) { |
| 1703 | return true; |
| 1704 | } |
| 1705 | } |
| 1706 | false |
| 1707 | } |
| 1708 | |
| 1709 | pub fn take_new_channels(&mut self) -> Vec<(ChannelId, ChannelConfig)> { |
| 1710 | let mut v = vec![]; |
| 1711 | |
| 1712 | if self.0.is_empty() { |
| 1713 | return v; |
| 1714 | } |
| 1715 | |
| 1716 | for i in (0..self.0.len()).rev() { |
| 1717 | if matches!(&self.0[i], Change::AddChannel(_)) { |
| 1718 | if let Change::AddChannel(id) = self.0.remove(i) { |
| 1719 | v.push(id); |
| 1720 | } |
| 1721 | } |
| 1722 | } |
| 1723 | |
| 1724 | v |
| 1725 | } |
| 1726 | |
| 1727 | /// Tests the given lines (from answer) corresponds to changes. |
| 1728 | fn ensure_correct_answer(&self, lines: &[(usize, &MediaLine)]) -> Option<String> { |
| 1729 | if self.count_new_medias() != lines.len() { |
| 1730 | return Some(format!( |
| 1731 | "Differing m-line count in offer vs answer: {} != {}", |
| 1732 | self.count_new_medias(), |
| 1733 | lines.len() |
| 1734 | )); |
| 1735 | } |
| 1736 | |
| 1737 | 'next: for (_, l) in lines { |
| 1738 | let mid = l.mid(); |
| 1739 | |
| 1740 | for m in &self.0 { |
| 1741 | use Change::*; |
| 1742 | match m { |
| 1743 | AddMedia(v) if v.mid == mid => { |
| 1744 | if !l.typ.is_media() { |
| 1745 | return Some(format!( |
| 1746 | "Answer m-line for mid ({}) is not of media type: {:?}", |
| 1747 | mid, l.typ |
| 1748 | )); |
| 1749 | } |
| 1750 | continue 'next; |
| 1751 | } |
| 1752 | AddApp(v) if *v == mid => { |
| 1753 | if !l.typ.is_channel() { |
| 1754 | return Some(format!( |
| 1755 | "Answer m-line for mid ({}) is not a data channel: {:?}", |
| 1756 | mid, l.typ |
| 1757 | )); |
| 1758 | } |
| 1759 | continue 'next; |
| 1760 | } |
| 1761 | _ => {} |
| 1762 | } |
| 1763 | } |
| 1764 | |
| 1765 | return Some(format!("Mid in answer is not in offer: {mid}")); |
| 1766 | } |
| 1767 | |
| 1768 | None |
| 1769 | } |
| 1770 | |
| 1771 | fn count_new_medias(&self) -> usize { |
| 1772 | self.0 |
| 1773 | .iter() |
| 1774 | .filter(|c| matches!(c, Change::AddMedia(_) | Change::AddApp(_))) |
| 1775 | .count() |
| 1776 | } |
| 1777 | |
| 1778 | pub fn as_new_medias<'a, 'b: 'a>( |
| 1779 | &'a self, |
| 1780 | index_start: usize, |
| 1781 | config: &'b CodecConfig, |
| 1782 | exts: &'b ExtensionMap, |
| 1783 | ) -> impl Iterator<Item = Media> + 'a { |
| 1784 | // Use a separate counter that only advances for entries that actually produce |
| 1785 | // an m-line (AddMedia / AddApp). Non-media entries (Direction, IceRestart, |
| 1786 | // AddChannel) must not consume an index slot, otherwise the resulting m-line |
| 1787 | // indices would be non-contiguous and out of sync with their array positions. |
| 1788 | let mut media_idx = 0usize; |
| 1789 | self.0.iter().filter_map(move |c| { |
| 1790 | let result = c.as_new_media(index_start + media_idx, config, exts); |
| 1791 | if result.is_some() { |
| 1792 | media_idx += 1; |
| 1793 | } |
| 1794 | result |
| 1795 | }) |
| 1796 | } |
| 1797 | |
| 1798 | pub(crate) fn apply_to(&self, lines: &mut [MediaLine]) { |
| 1799 | for change in &self.0 { |
| 1800 | if let Change::Direction(mid, dir) = change { |
| 1801 | if let Some(line) = lines.iter_mut().find(|l| l.mid() == *mid) { |
| 1802 | if let Some(dir_pos) = line.attrs.iter().position(|a| a.is_direction()) { |
| 1803 | line.attrs[dir_pos] = (*dir).into(); |
| 1804 | } |
| 1805 | } |
| 1806 | } |
| 1807 | } |
| 1808 | } |
| 1809 | |
| 1810 | fn ssrcs_for_mid(&self, mid: Mid) -> &[(Ssrc, Option<Ssrc>)] { |
| 1811 | let maybe_add_media = self |
| 1812 | .0 |
| 1813 | .iter() |
| 1814 | .filter_map(|c| { |
| 1815 | if let Change::AddMedia(m) = c { |
| 1816 | Some(m) |
| 1817 | } else { |
| 1818 | None |
| 1819 | } |
| 1820 | }) |
| 1821 | .find(|m| m.mid == mid); |
| 1822 | |
| 1823 | let Some(m) = maybe_add_media else { |
| 1824 | return &[]; |
| 1825 | }; |
| 1826 | |
| 1827 | &m.ssrcs |
| 1828 | } |
| 1829 | } |
| 1830 | |
| 1831 | impl Change { |
| 1832 | fn as_new_media( |
| 1833 | &self, |
| 1834 | index: usize, |
| 1835 | config: &CodecConfig, |
| 1836 | exts: &ExtensionMap, |
| 1837 | ) -> Option<Media> { |
| 1838 | use Change::*; |
| 1839 | match self { |
| 1840 | AddMedia(v) => { |
| 1841 | // TODO can we avoid all this cloning? |
| 1842 | let mut add = v.clone(); |
| 1843 | add.pts = config.all_for_kind(v.kind).map(|p| p.pt()).collect(); |
| 1844 | add.exts = exts.cloned_with_type(v.kind.is_audio()); |
| 1845 | add.index = index; |
| 1846 | |
| 1847 | Some(Media::from_add_media(add)) |
| 1848 | } |
| 1849 | AddApp(mid) => Some(Media::from_app_tmp(*mid, index)), |
| 1850 | _ => None, |
| 1851 | } |
| 1852 | } |
| 1853 | } |
| 1854 | |
| 1855 | #[cfg(test)] |
| 1856 | mod test { |
| 1857 | use std::time::Instant; |
| 1858 | |
| 1859 | use sdp::RestrictionId; |
| 1860 | use sdp::SimulcastLayer as SdpSimulcastLayer; |
| 1861 | |
| 1862 | use crate::format::Codec; |
| 1863 | use crate::media::{Simulcast, SimulcastLayer}; |
| 1864 | use crate::sdp::RtpMap; |
| 1865 | |
| 1866 | use super::*; |
| 1867 | |
| 1868 | fn resolve_pt(m_line: &MediaLine, needle: Pt) -> RtpMap { |
| 1869 | m_line |
| 1870 | .attrs |
| 1871 | .iter() |
| 1872 | .find_map(|attr| match attr { |
| 1873 | MediaAttribute::RtpMap { pt, value } if *pt == needle => Some(*value), |
| 1874 | _ => None, |
| 1875 | }) |
| 1876 | .unwrap_or_else(|| panic!("Expected to find RtpMap for {needle}")) |
| 1877 | } |
| 1878 | |
| 1879 | fn count_lines(lines: &str, what: &str) -> usize { |
| 1880 | lines.lines().filter(|l| l == &what).count() |
| 1881 | } |
| 1882 | |
| 1883 | fn get_setup_from_media_line(line: &MediaLine) -> Setup { |
| 1884 | line.setup().expect("Expected a=setup attribute in SDP") |
| 1885 | } |
| 1886 | |
| 1887 | /// RFC 8842 §5.2: an offer MUST use a=setup:actpass regardless of DTLS role. |
| 1888 | #[test] |
| 1889 | fn offer_uses_actpass() { |
| 1890 | crate::init_crypto_default(); |
| 1891 | |
| 1892 | let now = Instant::now(); |
| 1893 | let mut rtc = Rtc::new(now); |
| 1894 | |
| 1895 | let mut change = rtc.sdp_api(); |
| 1896 | change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None); |
| 1897 | let (offer, _pending) = change.apply().unwrap(); |
| 1898 | |
| 1899 | let setup = get_setup_from_media_line(&offer.media_lines[0]); |
| 1900 | assert_eq!( |
| 1901 | setup, |
| 1902 | Setup::ActPass, |
| 1903 | "Offer must use a=setup:actpass per RFC 8842 §5.2" |
| 1904 | ); |
| 1905 | } |
| 1906 | |
| 1907 | /// RFC 8842 §5.5: the answerer uses the negotiated DTLS role (active or passive). |
| 1908 | /// When the offerer uses actpass, the answerer picks passive (DTLS server) and the |
| 1909 | /// offerer becomes active (DTLS client). The answerer's SDP must reflect passive. |
| 1910 | #[test] |
| 1911 | fn answer_uses_negotiated_dtls_role() { |
| 1912 | crate::init_crypto_default(); |
| 1913 | |
| 1914 | let now = Instant::now(); |
| 1915 | let mut offerer = Rtc::new(now); |
| 1916 | let mut answerer = Rtc::new(now); |
| 1917 | |
| 1918 | // Offerer creates the offer. |
| 1919 | let mut change = offerer.sdp_api(); |
| 1920 | change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None); |
| 1921 | let (offer, pending) = change.apply().unwrap(); |
| 1922 | |
| 1923 | // Verify the offer itself is actpass. |
| 1924 | assert_eq!( |
| 1925 | get_setup_from_media_line(&offer.media_lines[0]), |
| 1926 | Setup::ActPass, |
| 1927 | "Offer must use a=setup:actpass" |
| 1928 | ); |
| 1929 | |
| 1930 | // Answerer accepts and generates the answer. |
| 1931 | let answer = answerer.sdp_api().accept_offer(offer).unwrap(); |
| 1932 | |
| 1933 | // When the offerer sends actpass, the answerer takes the passive DTLS role |
| 1934 | // (see init_dtls: ActPass from remote → local takes Passive). |
| 1935 | // The answerer's SDP must reflect that with a=setup:passive. |
| 1936 | assert_eq!( |
| 1937 | get_setup_from_media_line(&answer.media_lines[0]), |
| 1938 | Setup::Passive, |
| 1939 | "Answerer's SDP must use a=setup:passive when responding to an actpass offer" |
| 1940 | ); |
| 1941 | |
| 1942 | // Complete the exchange on the offerer side. |
| 1943 | offerer.sdp_api().accept_answer(pending, answer).unwrap(); |
| 1944 | |
| 1945 | // Now generate a subsequent offer from the offerer (re-offer). |
| 1946 | // Even after the DTLS role is settled (offerer is now passive), a new offer |
| 1947 | // must still carry a=setup:actpass per RFC 8842 §5.2. |
| 1948 | let mut change = offerer.sdp_api(); |
| 1949 | change.add_media(MediaKind::Video, Direction::SendRecv, None, None, None); |
| 1950 | let (reoffer, _) = change.apply().unwrap(); |
| 1951 | |
| 1952 | assert_eq!( |
| 1953 | get_setup_from_media_line(&reoffer.media_lines[0]), |
| 1954 | Setup::ActPass, |
| 1955 | "Re-offer must still use a=setup:actpass even after DTLS role is settled" |
| 1956 | ); |
| 1957 | } |
| 1958 | |
| 1959 | #[test] |
| 1960 | fn test_out_of_order_error() { |
| 1961 | crate::init_crypto_default(); |
| 1962 | |
| 1963 | let now = Instant::now(); |
| 1964 | let mut rtc1 = Rtc::new(now); |
| 1965 | let mut rtc2 = Rtc::new(now); |
| 1966 | |
| 1967 | let mut change1 = rtc1.sdp_api(); |
| 1968 | change1.add_channel("ch1".into()); |
| 1969 | let (offer1, pending1) = change1.apply().unwrap(); |
| 1970 | |
| 1971 | let mut change2 = rtc2.sdp_api(); |
| 1972 | change2.add_channel("ch2".into()); |
| 1973 | let (offer2, _) = change2.apply().unwrap(); |
| 1974 | |
| 1975 | // invalidates pending1 |
| 1976 | let _ = rtc1.sdp_api().accept_offer(offer2).unwrap(); |
| 1977 | let answer2 = rtc2.sdp_api().accept_offer(offer1).unwrap(); |
| 1978 | |
| 1979 | let r = rtc1.sdp_api().accept_answer(pending1, answer2); |
| 1980 | |
| 1981 | assert!(matches!(r, Err(RtcError::ChangesOutOfOrder))); |
| 1982 | } |
| 1983 | |
| 1984 | #[test] |
| 1985 | fn sdp_api_merge_works() { |
| 1986 | crate::init_crypto_default(); |
| 1987 | |
| 1988 | let mut rtc = Rtc::new(Instant::now()); |
| 1989 | let mut changes = rtc.sdp_api(); |
| 1990 | changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 1991 | let (offer, pending) = changes.apply().unwrap(); |
| 1992 | |
| 1993 | let mut changes = rtc.sdp_api(); |
| 1994 | changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 1995 | changes.merge(pending); |
| 1996 | let (new_offer, _) = changes.apply().unwrap(); |
| 1997 | |
| 1998 | // After merge(), the original AddMedia (audio) is prepended before the new AddMedia |
| 1999 | // (video) to preserve its original m-line position. So audio is at index 0 and video at 1. |
| 2000 | assert_eq!(offer.media_lines[0], new_offer.media_lines[0]); |
| 2001 | assert_eq!(new_offer.media_lines[0].typ, MediaType::Audio); |
| 2002 | assert_eq!(new_offer.media_lines[1].typ, MediaType::Video); |
| 2003 | assert_eq!(new_offer.media_lines.len(), 2); |
| 2004 | } |
| 2005 | |
| 2006 | // The indices assigned to Media objects by as_new_medias() must be contiguous starting |
| 2007 | // from index_start. Only entries that produce m-lines (AddMedia / AddApp) should |
| 2008 | // consume an index slot – non-media entries (AddChannel, Direction, IceRestart) must |
| 2009 | // not advance the counter. |
| 2010 | // |
| 2011 | // We test this by building a Changes vector directly and inspecting the indices |
| 2012 | // returned by as_new_medias(). |
| 2013 | #[test] |
| 2014 | fn as_new_medias_contiguous_indices_with_non_media_changes() { |
| 2015 | crate::init_crypto_default(); |
| 2016 | |
| 2017 | let now = Instant::now(); |
| 2018 | let mut rtc = Rtc::new(now); |
| 2019 | |
| 2020 | // Build a Changes vector that interleaves non-media entries with media entries. |
| 2021 | // The non-media AddChannel entry must not consume an index slot. |
| 2022 | // |
| 2023 | // Changes = [AddApp(mid0), AddChannel(id, cfg), AddMedia(audio)] |
| 2024 | // |
| 2025 | // index_start = 0 (fresh session) |
| 2026 | // Expected indices: AddApp → 0, AddMedia(audio) → 1 |
| 2027 | // Buggy indices: AddApp → 0, AddMedia(audio) → 2 |
| 2028 | let mut changes = Changes::default(); |
| 2029 | let mid_app = rtc.new_mid(); |
| 2030 | changes.0.push(Change::AddApp(mid_app)); |
| 2031 | let mid_audio = rtc.new_mid(); |
| 2032 | let chan_id = rtc.chan.new_channel(&ChannelConfig { |
| 2033 | label: "ch".into(), |
| 2034 | ..Default::default() |
| 2035 | }); |
| 2036 | changes.0.push(Change::AddChannel(( |
| 2037 | chan_id, |
| 2038 | ChannelConfig { |
| 2039 | label: "ch".into(), |
| 2040 | ..Default::default() |
| 2041 | }, |
| 2042 | ))); |
| 2043 | changes.0.push(Change::AddMedia(AddMedia { |
| 2044 | mid: mid_audio, |
| 2045 | cname: "test".into(), |
| 2046 | msid: Msid { |
| 2047 | stream_id: "stream".into(), |
| 2048 | track_id: "track".into(), |
| 2049 | }, |
| 2050 | kind: MediaKind::Audio, |
| 2051 | dir: Direction::SendOnly, |
| 2052 | ssrcs: vec![], |
| 2053 | simulcast: None, |
| 2054 | pts: vec![], |
| 2055 | exts: ExtensionMap::empty(), |
| 2056 | index: 0, // will be overwritten by as_new_media() |
| 2057 | })); |
| 2058 | |
| 2059 | let config = CodecConfig::new_with_defaults(); |
| 2060 | let exts = ExtensionMap::standard(); |
| 2061 | |
| 2062 | let medias: Vec<Media> = changes.as_new_medias(0, &config, &exts).collect(); |
| 2063 | |
| 2064 | // Should produce two Media objects: one for AddApp, one for AddMedia(audio). |
| 2065 | assert_eq!(medias.len(), 2, "Expected 2 media-producing entries"); |
| 2066 | |
| 2067 | // The first m-line is AddApp at index 0. |
| 2068 | assert_eq!(medias[0].index(), 0, "AddApp should have index 0"); |
| 2069 | |
| 2070 | // The second m-line is Audio. It must get the contiguous index 1, not 2 (which |
| 2071 | // would result from counting the non-media AddChannel entry as a slot). |
| 2072 | assert_eq!( |
| 2073 | medias[1].index(), |
| 2074 | 1, |
| 2075 | "AddMedia(audio) should have contiguous index 1, not 2" |
| 2076 | ); |
| 2077 | } |
| 2078 | |
| 2079 | // AddMedia/AddApp entries from a pending offer carry already-allocated MIDs whose |
| 2080 | // m-line positions the remote peer may have committed to. When merging a pending offer |
| 2081 | // into a new SdpApi, those original entries must be prepended so they keep their |
| 2082 | // original lower-index positions rather than being displaced by new changes. |
| 2083 | #[test] |
| 2084 | fn sdp_api_merge_stale_media_keeps_original_position() { |
| 2085 | crate::init_crypto_default(); |
| 2086 | |
| 2087 | let mut rtc = Rtc::new(Instant::now()); |
| 2088 | |
| 2089 | // First offer: add audio. Audio will be at index 0. |
| 2090 | let mut changes = rtc.sdp_api(); |
| 2091 | let _mid_audio = changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 2092 | let (offer1, pending) = changes.apply().unwrap(); |
| 2093 | |
| 2094 | // The original offer has audio at position 0. |
| 2095 | assert_eq!(offer1.media_lines[0].typ, MediaType::Audio); |
| 2096 | |
| 2097 | // Simulate glare: a new SdpApi adds video, then merges the pending offer |
| 2098 | // containing the original audio AddMedia. |
| 2099 | let mut changes = rtc.sdp_api(); |
| 2100 | changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 2101 | changes.merge(pending); |
| 2102 | let (offer2, _) = changes.apply().unwrap(); |
| 2103 | |
| 2104 | // The original audio must be prepended before the new video, keeping it at |
| 2105 | // position 0 — the same position as in the original offer. |
| 2106 | assert_eq!( |
| 2107 | offer2.media_lines[0].typ, |
| 2108 | MediaType::Audio, |
| 2109 | "Original audio from pending should be at position 0 (same as original offer)" |
| 2110 | ); |
| 2111 | assert_eq!( |
| 2112 | offer2.media_lines[1].typ, |
| 2113 | MediaType::Video, |
| 2114 | "New video added to current SdpApi should be at position 1" |
| 2115 | ); |
| 2116 | assert_eq!( |
| 2117 | offer1.media_lines[0], offer2.media_lines[0], |
| 2118 | "Audio m-line must be identical between original and merged offers" |
| 2119 | ); |
| 2120 | } |
| 2121 | |
| 2122 | // When a new SdpApi has its own AddMedia changes and also merges a pending offer that |
| 2123 | // has AddMedia changes, the original pending entries must come before the new ones so |
| 2124 | // their m-line positions are stable: |
| 2125 | // 1. First offer: [AddMedia(audio), AddMedia(video)] → audio at 0, video at 1 |
| 2126 | // 2. Glare – pending is saved |
| 2127 | // 3. New SdpApi: [AddMedia(screen)] |
| 2128 | // 4. After merge: [AddMedia(audio), AddMedia(video), AddMedia(screen)] |
| 2129 | // → audio at 0, video at 1, screen at 2 (CORRECT) |
| 2130 | #[test] |
| 2131 | fn sdp_api_merge_with_direction_and_new_media_preserves_positions() { |
| 2132 | crate::init_crypto_default(); |
| 2133 | |
| 2134 | let now = Instant::now(); |
| 2135 | let mut rtc1 = Rtc::new(now); |
| 2136 | let mut rtc2 = Rtc::new(now); |
| 2137 | |
| 2138 | // Establish one audio m-line so we can change its direction later. |
| 2139 | let mut changes = rtc1.sdp_api(); |
| 2140 | let mid_audio = changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 2141 | let (offer0, pending0) = changes.apply().unwrap(); |
| 2142 | let answer0 = rtc2.sdp_api().accept_offer(offer0).unwrap(); |
| 2143 | rtc1.sdp_api().accept_answer(pending0, answer0).unwrap(); |
| 2144 | // Session now: audio at index 0. |
| 2145 | |
| 2146 | // Create an offer with two new medias: video1 and video2. |
| 2147 | // video1 → session position 1, video2 → session position 2. |
| 2148 | let mut changes = rtc1.sdp_api(); |
| 2149 | let mid_v1 = changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 2150 | let mid_v2 = changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 2151 | let (offer1, pending1) = changes.apply().unwrap(); |
| 2152 | |
| 2153 | assert_eq!(offer1.media_lines[1].mid(), mid_v1, "video1 at idx 1"); |
| 2154 | assert_eq!(offer1.media_lines[2].mid(), mid_v2, "video2 at idx 2"); |
| 2155 | |
| 2156 | // Glare: create a new SdpApi that: |
| 2157 | // (a) changes direction of existing audio m-line (non-media Change) |
| 2158 | // (b) adds a new screen-share video m-line |
| 2159 | // Then merge the original pending1 into it. |
| 2160 | let mut changes = rtc1.sdp_api(); |
| 2161 | changes.set_direction(mid_audio, Direction::Inactive); |
| 2162 | let mid_screen = changes.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 2163 | changes.merge(pending1); |
| 2164 | let (offer2, _) = changes.apply().unwrap(); |
| 2165 | |
| 2166 | // Session already has audio at index 0. |
| 2167 | // The merged offer must have: |
| 2168 | // position 0: audio (existing, unchanged) |
| 2169 | // position 1: video1 (from pending — same position as offer1) |
| 2170 | // position 2: video2 (from pending — same position as offer1) |
| 2171 | // position 3: screen (newly added in the new SdpApi) |
| 2172 | assert_eq!( |
| 2173 | offer2.media_lines.len(), |
| 2174 | 4, |
| 2175 | "Merged offer should have 4 m-lines: audio + video1 + video2 + screen" |
| 2176 | ); |
| 2177 | |
| 2178 | // Original video m-lines from the pending offer must keep their positions (1 and 2). |
| 2179 | assert_eq!( |
| 2180 | offer2.media_lines[1].mid(), |
| 2181 | mid_v1, |
| 2182 | "video1 from pending should remain at position 1" |
| 2183 | ); |
| 2184 | assert_eq!( |
| 2185 | offer1.media_lines[1], offer2.media_lines[1], |
| 2186 | "video1 m-line should be identical in both offers" |
| 2187 | ); |
| 2188 | assert_eq!( |
| 2189 | offer2.media_lines[2].mid(), |
| 2190 | mid_v2, |
| 2191 | "video2 from pending should remain at position 2" |
| 2192 | ); |
| 2193 | assert_eq!( |
| 2194 | offer1.media_lines[2], offer2.media_lines[2], |
| 2195 | "video2 m-line should be identical in both offers" |
| 2196 | ); |
| 2197 | |
| 2198 | // New screen share is at position 3, after the original pending m-lines. |
| 2199 | assert_eq!( |
| 2200 | offer2.media_lines[3].mid(), |
| 2201 | mid_screen, |
| 2202 | "New screen-share should be at position 3" |
| 2203 | ); |
| 2204 | } |
| 2205 | |
| 2206 | #[test] |
| 2207 | fn test_rtp_payload_priority() { |
| 2208 | crate::init_crypto_default(); |
| 2209 | |
| 2210 | let now = Instant::now(); |
| 2211 | let mut rtc1 = Rtc::builder() |
| 2212 | .clear_codecs() |
| 2213 | .enable_h264(true) |
| 2214 | .enable_vp8(true) |
| 2215 | .enable_vp9(true) |
| 2216 | .build(now); |
| 2217 | let mut rtc2 = Rtc::builder() |
| 2218 | .clear_codecs() |
| 2219 | .enable_vp8(true) |
| 2220 | .enable_h264(true) |
| 2221 | .build(now); |
| 2222 | |
| 2223 | let mut change1 = rtc1.sdp_api(); |
| 2224 | change1.add_media(MediaKind::Video, Direction::SendOnly, None, None, None); |
| 2225 | let (offer1, _) = change1.apply().unwrap(); |
| 2226 | |
| 2227 | let answer = rtc2.sdp_api().accept_offer(offer1).unwrap(); |
| 2228 | assert_eq!( |
| 2229 | answer.media_lines.len(), |
| 2230 | 1, |
| 2231 | "There should be one mline only" |
| 2232 | ); |
| 2233 | |
| 2234 | let first_mline = &answer.media_lines[0]; |
| 2235 | let first_pt = resolve_pt(first_mline, first_mline.pts[0]); |
| 2236 | |
| 2237 | assert_eq!( |
| 2238 | first_pt.codec, |
| 2239 | Codec::Vp8, |
| 2240 | "The first PT returned should be the highest priority PT from the answer that is supported." |
| 2241 | ); |
| 2242 | |
| 2243 | let vp9_unsupported = first_mline |
| 2244 | .pts |
| 2245 | .iter() |
| 2246 | .any(|pt| resolve_pt(first_mline, *pt).codec == Codec::Vp9); |
| 2247 | |
| 2248 | assert!( |
| 2249 | !vp9_unsupported, |
| 2250 | "VP9 was not offered, so it should not be present in the answer" |
| 2251 | ); |
| 2252 | } |
| 2253 | |
| 2254 | #[test] |
| 2255 | fn non_simulcast_rids() { |
| 2256 | crate::init_crypto_default(); |
| 2257 | |
| 2258 | let now = Instant::now(); |
| 2259 | let mut rtc1 = Rtc::new(now); |
| 2260 | let mut rtc2 = Rtc::new(now); |
| 2261 | |
| 2262 | // Test initial media creation |
| 2263 | let mid = { |
| 2264 | let mut changes = rtc1.sdp_api(); |
| 2265 | let mid = changes.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 2266 | let (offer, pending) = changes.apply().unwrap(); |
| 2267 | let answer = rtc2.sdp_api().accept_offer(offer).unwrap(); |
| 2268 | rtc1.sdp_api().accept_answer(pending, answer).unwrap(); |
| 2269 | |
| 2270 | assert!(matches!(rtc1.media(mid).unwrap().rids_rx(), Rids::Any)); |
| 2271 | assert!(matches!(rtc1.media(mid).unwrap().rids_tx(), Rids::None)); |
| 2272 | assert!(matches!(rtc2.media(mid).unwrap().rids_rx(), Rids::Any)); |
| 2273 | assert!(matches!(rtc2.media(mid).unwrap().rids_tx(), Rids::None)); |
| 2274 | |
| 2275 | mid |
| 2276 | }; |
| 2277 | |
| 2278 | // Test later updates to that media |
| 2279 | { |
| 2280 | let mut changes = rtc1.sdp_api(); |
| 2281 | changes.set_direction(mid, Direction::Inactive); |
| 2282 | let (offer, pending) = changes.apply().unwrap(); |
| 2283 | let answer = rtc2.sdp_api().accept_offer(offer).unwrap(); |
| 2284 | rtc1.sdp_api().accept_answer(pending, answer).unwrap(); |
| 2285 | |
| 2286 | assert!(matches!(rtc1.media(mid).unwrap().rids_rx(), Rids::Any)); |
| 2287 | assert!(matches!(rtc1.media(mid).unwrap().rids_tx(), Rids::None)); |
| 2288 | assert!(matches!(rtc2.media(mid).unwrap().rids_rx(), Rids::Any)); |
| 2289 | assert!(matches!(rtc2.media(mid).unwrap().rids_tx(), Rids::None)); |
| 2290 | } |
| 2291 | } |
| 2292 | |
| 2293 | #[test] |
| 2294 | fn simulcast_ssrc_allocation() { |
| 2295 | crate::init_crypto_default(); |
| 2296 | |
| 2297 | let mut rtc1 = Rtc::new(Instant::now()); |
| 2298 | |
| 2299 | let mut simulcast = Simulcast::new(); |
| 2300 | |
| 2301 | simulcast.add_send_layer(SimulcastLayer::new("h")); |
| 2302 | simulcast.add_send_layer(SimulcastLayer::new("m")); |
| 2303 | simulcast.add_send_layer(SimulcastLayer::new("l")); |
| 2304 | |
| 2305 | let mut change = rtc1.sdp_api(); |
| 2306 | change.add_media( |
| 2307 | MediaKind::Video, |
| 2308 | Direction::SendOnly, |
| 2309 | None, |
| 2310 | None, |
| 2311 | Some(simulcast), |
| 2312 | ); |
| 2313 | |
| 2314 | let Change::AddMedia(am) = &change.changes[0] else { |
| 2315 | panic!("Not AddMedia?!"); |
| 2316 | }; |
| 2317 | |
| 2318 | // these should be organized in order: m, h, l |
| 2319 | let pending_ssrcs = am.ssrcs.clone(); |
| 2320 | assert_eq!(pending_ssrcs.len(), 3); |
| 2321 | |
| 2322 | for p in &pending_ssrcs { |
| 2323 | assert!(p.1.is_some()); // all should have rtx |
| 2324 | } |
| 2325 | |
| 2326 | let (offer, _) = change.apply().unwrap(); |
| 2327 | let sdp = offer.into_inner(); |
| 2328 | let line = &sdp.media_lines[0]; |
| 2329 | |
| 2330 | assert_eq!( |
| 2331 | line.simulcast().unwrap().send, |
| 2332 | SimulcastGroups(vec![ |
| 2333 | SdpSimulcastLayer { |
| 2334 | restriction_id: RestrictionId("h".into(), true), |
| 2335 | attributes: None, |
| 2336 | }, |
| 2337 | SdpSimulcastLayer { |
| 2338 | restriction_id: RestrictionId("m".into(), true), |
| 2339 | attributes: None, |
| 2340 | }, |
| 2341 | SdpSimulcastLayer { |
| 2342 | restriction_id: RestrictionId("l".into(), true), |
| 2343 | attributes: None, |
| 2344 | }, |
| 2345 | ]) |
| 2346 | ); |
| 2347 | |
| 2348 | // Each SSRC, both regular and RTX get their own a=ssrc line. |
| 2349 | assert_eq!(line.ssrc_info().len(), pending_ssrcs.len() * 2); |
| 2350 | |
| 2351 | let fids: Vec<_> = line |
| 2352 | .attrs |
| 2353 | .iter() |
| 2354 | .filter_map(|a| { |
| 2355 | if let MediaAttribute::SsrcGroup { semantics, ssrcs } = a { |
| 2356 | // We don't have any other semantics right now. |
| 2357 | assert_eq!(semantics, "FID"); |
| 2358 | assert_eq!(ssrcs.len(), 2); |
| 2359 | Some((ssrcs[0], ssrcs[1])) |
| 2360 | } else { |
| 2361 | None |
| 2362 | } |
| 2363 | }) |
| 2364 | .collect(); |
| 2365 | |
| 2366 | assert_eq!(fids.len(), pending_ssrcs.len()); |
| 2367 | |
| 2368 | for (a, b) in fids.iter().zip(pending_ssrcs.iter()) { |
| 2369 | assert_eq!(a.0, b.0); |
| 2370 | assert_eq!(Some(a.1), b.1); |
| 2371 | } |
| 2372 | |
| 2373 | let line_string = line.to_string(); |
| 2374 | |
| 2375 | // The SDP offer should contain layers without any attributes |
| 2376 | assert_eq!(count_lines(&line_string, "a=rid:h send"), 1); |
| 2377 | assert_eq!(count_lines(&line_string, "a=rid:m send"), 1); |
| 2378 | assert_eq!(count_lines(&line_string, "a=rid:l send"), 1); |
| 2379 | } |
| 2380 | |
| 2381 | #[test] |
| 2382 | fn simulcast_attributes() { |
| 2383 | crate::init_crypto_default(); |
| 2384 | |
| 2385 | let mut rtc1 = Rtc::new(Instant::now()); |
| 2386 | |
| 2387 | let mut simulcast = Simulcast::new(); |
| 2388 | |
| 2389 | // High layer |
| 2390 | simulcast.add_send_layer( |
| 2391 | SimulcastLayer::new_with_attributes("high") |
| 2392 | .max_width(1280) |
| 2393 | .max_height(720) |
| 2394 | .max_br(1100000) |
| 2395 | .max_br(1300000) |
| 2396 | .max_br(1500000) // the last one wins |
| 2397 | .max_fps(30) |
| 2398 | .build(), |
| 2399 | ); |
| 2400 | |
| 2401 | // Medium layer |
| 2402 | simulcast.add_send_layer( |
| 2403 | SimulcastLayer::new_with_attributes("medium") |
| 2404 | .max_width(640) |
| 2405 | .max_height(360) |
| 2406 | .max_br(600000) |
| 2407 | // No max_fps |
| 2408 | .build(), |
| 2409 | ); |
| 2410 | |
| 2411 | // Low layer |
| 2412 | simulcast.add_send_layer( |
| 2413 | SimulcastLayer::new_with_attributes("low") |
| 2414 | // No max_width |
| 2415 | .max_height(180) |
| 2416 | .max_br(200000) |
| 2417 | .max_fps(15) |
| 2418 | .build(), |
| 2419 | ); |
| 2420 | |
| 2421 | // Custom attribute |
| 2422 | simulcast.add_send_layer( |
| 2423 | SimulcastLayer::new_with_attributes("custom") |
| 2424 | .custom("foo", "bar") |
| 2425 | .build(), |
| 2426 | ); |
| 2427 | |
| 2428 | // No attributes |
| 2429 | simulcast.add_send_layer(SimulcastLayer::new_with_attributes("no_attrs").build()); |
| 2430 | |
| 2431 | let mut change = rtc1.sdp_api(); |
| 2432 | change.add_media( |
| 2433 | MediaKind::Video, |
| 2434 | Direction::SendOnly, |
| 2435 | None, |
| 2436 | None, |
| 2437 | Some(simulcast), |
| 2438 | ); |
| 2439 | |
| 2440 | let (offer, _) = change.apply().unwrap(); |
| 2441 | let sdp = offer.into_inner(); |
| 2442 | let line = &sdp.media_lines[0]; |
| 2443 | |
| 2444 | fn pairs_to_some_vec(pairs: &[(&str, &str)]) -> Option<Vec<(String, String)>> { |
| 2445 | Some( |
| 2446 | pairs |
| 2447 | .iter() |
| 2448 | .map(|(k, v)| (k.to_string(), v.to_string())) |
| 2449 | .collect(), |
| 2450 | ) |
| 2451 | } |
| 2452 | |
| 2453 | assert_eq!( |
| 2454 | line.simulcast().unwrap().send, |
| 2455 | SimulcastGroups(vec![ |
| 2456 | SdpSimulcastLayer { |
| 2457 | restriction_id: RestrictionId("high".into(), true), |
| 2458 | attributes: pairs_to_some_vec(&[ |
| 2459 | ("max-width", "1280"), |
| 2460 | ("max-height", "720"), |
| 2461 | ("max-br", "1500000"), |
| 2462 | ("max-fps", "30"), |
| 2463 | ]), |
| 2464 | }, |
| 2465 | SdpSimulcastLayer { |
| 2466 | restriction_id: RestrictionId("medium".into(), true), |
| 2467 | attributes: pairs_to_some_vec(&[ |
| 2468 | ("max-width", "640"), |
| 2469 | ("max-height", "360"), |
| 2470 | ("max-br", "600000"), |
| 2471 | ]), |
| 2472 | }, |
| 2473 | SdpSimulcastLayer { |
| 2474 | restriction_id: RestrictionId("low".into(), true), |
| 2475 | attributes: pairs_to_some_vec(&[ |
| 2476 | ("max-height", "180"), |
| 2477 | ("max-br", "200000"), |
| 2478 | ("max-fps", "15"), |
| 2479 | ]), |
| 2480 | }, |
| 2481 | SdpSimulcastLayer { |
| 2482 | restriction_id: RestrictionId("custom".into(), true), |
| 2483 | attributes: pairs_to_some_vec(&[("foo", "bar"),]), |
| 2484 | }, |
| 2485 | SdpSimulcastLayer { |
| 2486 | restriction_id: RestrictionId("no_attrs".into(), true), |
| 2487 | attributes: None, |
| 2488 | }, |
| 2489 | ]) |
| 2490 | ); |
| 2491 | |
| 2492 | // The SDP offer should contain layers with our attributes |
| 2493 | let line_string = line.to_string(); |
| 2494 | assert_eq!( |
| 2495 | count_lines( |
| 2496 | &line_string, |
| 2497 | "a=rid:high send max-width=1280;max-height=720;max-br=1500000;max-fps=30" |
| 2498 | ), |
| 2499 | 1 |
| 2500 | ); |
| 2501 | assert_eq!( |
| 2502 | count_lines( |
| 2503 | &line_string, |
| 2504 | "a=rid:medium send max-width=640;max-height=360;max-br=600000" |
| 2505 | ), |
| 2506 | 1 |
| 2507 | ); |
| 2508 | assert_eq!( |
| 2509 | count_lines( |
| 2510 | &line_string, |
| 2511 | "a=rid:low send max-height=180;max-br=200000;max-fps=15" |
| 2512 | ), |
| 2513 | 1 |
| 2514 | ); |
| 2515 | assert_eq!(count_lines(&line_string, "a=rid:custom send foo=bar"), 1); |
| 2516 | // No space at the end |
| 2517 | assert_eq!(count_lines(&line_string, "a=rid:no_attrs send"), 1); |
| 2518 | } |
| 2519 | |
| 2520 | #[test] |
| 2521 | fn test_local_max_message_size_advertised() { |
| 2522 | crate::init_crypto_default(); |
| 2523 | |
| 2524 | let now = Instant::now(); |
| 2525 | let mut rtc = Rtc::new(now); |
| 2526 | |
| 2527 | // Create an offer with a data channel |
| 2528 | let mut change = rtc.sdp_api(); |
| 2529 | change.add_channel("test-channel".into()); |
| 2530 | let (offer, _) = change.apply().unwrap(); |
| 2531 | |
| 2532 | // Find the application m-line in the offer |
| 2533 | let app_line = offer |
| 2534 | .media_lines |
| 2535 | .iter() |
| 2536 | .find(|m| m.typ.is_channel()) |
| 2537 | .expect("should have application m-line"); |
| 2538 | |
| 2539 | // Verify that max-message-size attribute is present and matches LOCAL_MAX_MESSAGE_SIZE |
| 2540 | let max_size = app_line |
| 2541 | .max_message_size() |
| 2542 | .expect("max-message-size attribute should be present"); |
| 2543 | |
| 2544 | assert_eq!( |
| 2545 | max_size, |
| 2546 | crate::sctp::LOCAL_MAX_MESSAGE_SIZE as usize, |
| 2547 | "max-message-size should match LOCAL_MAX_MESSAGE_SIZE constant" |
| 2548 | ); |
| 2549 | |
| 2550 | // Also verify it's in the SDP string output |
| 2551 | let sdp_string = offer.to_sdp_string(); |
| 2552 | let expected_line = format!("a=max-message-size:{}", crate::sctp::LOCAL_MAX_MESSAGE_SIZE); |
| 2553 | assert!( |
| 2554 | sdp_string.contains(&expected_line), |
| 2555 | "SDP should contain max-message-size attribute with LOCAL_MAX_MESSAGE_SIZE value" |
| 2556 | ); |
| 2557 | } |
| 2558 | |
| 2559 | #[test] |
| 2560 | fn test_configured_max_message_size_advertised() { |
| 2561 | crate::init_crypto_default(); |
| 2562 | let limits = crate::channel::SctpReceiveLimits::new(8192, 32768, 64, 8); |
| 2563 | let mut rtc = Rtc::builder() |
| 2564 | .set_sctp_receive_limits(limits) |
| 2565 | .build(Instant::now()); |
| 2566 | let mut change = rtc.sdp_api(); |
| 2567 | change.add_channel("control".into()); |
| 2568 | let (offer, _) = change.apply().unwrap(); |
| 2569 | let app = offer |
| 2570 | .media_lines |
| 2571 | .iter() |
| 2572 | .find(|m| m.typ.is_channel()) |
| 2573 | .unwrap(); |
| 2574 | assert_eq!(app.max_message_size(), Some(8192)); |
| 2575 | assert_eq!( |
| 2576 | offer.to_sdp_string().matches("a=max-message-size:").count(), |
| 2577 | 1 |
| 2578 | ); |
| 2579 | } |
| 2580 | |
| 2581 | #[test] |
| 2582 | fn test_remote_max_message_size_parsing() { |
| 2583 | // Parse SDP with max-message-size attribute and verify value is extracted correctly |
| 2584 | let sdp = "v=0\r\n\ |
| 2585 | o=- 0 0 IN IP4 172.17.0.1\r\n\ |
| 2586 | s=-\r\n\ |
| 2587 | c=IN IP4 172.17.0.1\r\n\ |
| 2588 | t=0 0\r\n\ |
| 2589 | a=group:BUNDLE 0\r\n\ |
| 2590 | a=fingerprint:sha-256 B4:12:1C:7C:7D:ED:F1:FA:61:07:57:9C:29:BE:58:E3:BC:41:E7:13:8E:7D\ |
| 2591 | :D3:9D:1F:94:6E:A5:23:46:94:23\r\n\ |
| 2592 | m=application 9999 UDP/DTLS/SCTP webrtc-datachannel\r\n\ |
| 2593 | a=mid:0\r\n\ |
| 2594 | a=ice-ufrag:test\r\n\ |
| 2595 | a=ice-pwd:testpassword1234\r\n\ |
| 2596 | a=setup:actpass\r\n\ |
| 2597 | a=sctp-port:5000\r\n\ |
| 2598 | a=max-message-size:131072\r\n\ |
| 2599 | "; |
| 2600 | |
| 2601 | let offer = SdpOffer::from_sdp_string(sdp).expect("should parse"); |
| 2602 | |
| 2603 | // Find the application m-line |
| 2604 | let app_line = offer |
| 2605 | .media_lines |
| 2606 | .iter() |
| 2607 | .find(|m| m.typ.is_channel()) |
| 2608 | .expect("should have application m-line"); |
| 2609 | |
| 2610 | assert_eq!( |
| 2611 | app_line.max_message_size(), |
| 2612 | Some(131072), |
| 2613 | "max-message-size should be parsed as 131072" |
| 2614 | ); |
| 2615 | } |
| 2616 | |
| 2617 | #[test] |
| 2618 | fn test_remote_max_message_size_applied() { |
| 2619 | crate::init_crypto_default(); |
| 2620 | |
| 2621 | let now = Instant::now(); |
| 2622 | let mut rtc1 = Rtc::new(now); |
| 2623 | let mut rtc2 = Rtc::new(now); |
| 2624 | |
| 2625 | // Create an offer from rtc1 with a channel |
| 2626 | let mut change1 = rtc1.sdp_api(); |
| 2627 | change1.add_channel("test-channel".into()); |
| 2628 | let (offer1, pending1) = change1.apply().unwrap(); |
| 2629 | |
| 2630 | // Get the offer SDP string and modify it to have a custom max-message-size |
| 2631 | let custom_max_size1 = 98304u32; |
| 2632 | let sdp_string = offer1.to_sdp_string(); |
| 2633 | let modified_sdp = sdp_string.replace( |
| 2634 | &format!("a=max-message-size:{}", crate::sctp::LOCAL_MAX_MESSAGE_SIZE), |
| 2635 | &format!("a=max-message-size:{}", custom_max_size1), |
| 2636 | ); |
| 2637 | |
| 2638 | let modified_offer = |
| 2639 | SdpOffer::from_sdp_string(&modified_sdp).expect("modified SDP should parse"); |
| 2640 | |
| 2641 | let answer = rtc2.sdp_api().accept_offer(modified_offer).unwrap(); |
| 2642 | |
| 2643 | // Verify that rtc2's SCTP send limit is set to the custom value from rtc1's offer |
| 2644 | assert_eq!( |
| 2645 | rtc2.sctp.remote_max_message_size(), |
| 2646 | custom_max_size1, |
| 2647 | "rtc2 should have remote max message size set to rtc1's advertised value" |
| 2648 | ); |
| 2649 | |
| 2650 | // Now verify the reverse: rtc1 accepts rtc2's answer and applies its max-message-size |
| 2651 | let custom_max_size2 = 131072u32; |
| 2652 | let sdp_string = answer.to_sdp_string(); |
| 2653 | let modified_sdp = sdp_string.replace( |
| 2654 | &format!("a=max-message-size:{}", crate::sctp::LOCAL_MAX_MESSAGE_SIZE), |
| 2655 | &format!("a=max-message-size:{}", custom_max_size2), |
| 2656 | ); |
| 2657 | |
| 2658 | let modified_answer = |
| 2659 | SdpAnswer::from_sdp_string(&modified_sdp).expect("modified SDP should parse"); |
| 2660 | |
| 2661 | rtc1.sdp_api() |
| 2662 | .accept_answer(pending1, modified_answer) |
| 2663 | .unwrap(); |
| 2664 | |
| 2665 | assert_eq!( |
| 2666 | rtc1.sctp.remote_max_message_size(), |
| 2667 | custom_max_size2, |
| 2668 | "rtc1 should have remote max message size set to rtc2's advertised value" |
| 2669 | ); |
| 2670 | } |
| 2671 | |
| 2672 | #[test] |
| 2673 | fn test_remote_max_message_size_zero_means_unbounded() { |
| 2674 | // RFC 8841 §6.1: a=max-message-size:0 means the peer can receive a message |
| 2675 | // of any size, so we represent it as u32::MAX rather than passing 0 through |
| 2676 | // (which would reject every non-empty outbound message in sctp-proto). |
| 2677 | crate::init_crypto_default(); |
| 2678 | |
| 2679 | let now = Instant::now(); |
| 2680 | let mut rtc1 = Rtc::new(now); |
| 2681 | let mut rtc2 = Rtc::new(now); |
| 2682 | |
| 2683 | let mut change1 = rtc1.sdp_api(); |
| 2684 | change1.add_channel("test-channel".into()); |
| 2685 | let (offer1, _pending1) = change1.apply().unwrap(); |
| 2686 | |
| 2687 | let sdp_string = offer1.to_sdp_string(); |
| 2688 | let modified_sdp = sdp_string.replace( |
| 2689 | &format!("a=max-message-size:{}", crate::sctp::LOCAL_MAX_MESSAGE_SIZE), |
| 2690 | "a=max-message-size:0", |
| 2691 | ); |
| 2692 | |
| 2693 | let modified_offer = |
| 2694 | SdpOffer::from_sdp_string(&modified_sdp).expect("modified SDP should parse"); |
| 2695 | |
| 2696 | rtc2.sdp_api().accept_offer(modified_offer).unwrap(); |
| 2697 | |
| 2698 | assert_eq!( |
| 2699 | rtc2.sctp.remote_max_message_size(), |
| 2700 | u32::MAX, |
| 2701 | "a=max-message-size:0 should be treated as unbounded (u32::MAX)" |
| 2702 | ); |
| 2703 | } |
| 2704 | |
| 2705 | #[test] |
| 2706 | fn test_remote_max_message_size_too_large_clamps_to_unbounded() { |
| 2707 | crate::init_crypto_default(); |
| 2708 | |
| 2709 | let now = Instant::now(); |
| 2710 | let mut rtc1 = Rtc::new(now); |
| 2711 | let mut rtc2 = Rtc::new(now); |
| 2712 | |
| 2713 | let mut change1 = rtc1.sdp_api(); |
| 2714 | change1.add_channel("test-channel".into()); |
| 2715 | let (offer1, _pending1) = change1.apply().unwrap(); |
| 2716 | |
| 2717 | let sdp_string = offer1.to_sdp_string(); |
| 2718 | let modified_sdp = sdp_string.replace( |
| 2719 | &format!("a=max-message-size:{}", crate::sctp::LOCAL_MAX_MESSAGE_SIZE), |
| 2720 | "a=max-message-size:4294967296", |
| 2721 | ); |
| 2722 | |
| 2723 | let modified_offer = |
| 2724 | SdpOffer::from_sdp_string(&modified_sdp).expect("modified SDP should parse"); |
| 2725 | |
| 2726 | rtc2.sdp_api().accept_offer(modified_offer).unwrap(); |
| 2727 | |
| 2728 | assert_eq!( |
| 2729 | rtc2.sctp.remote_max_message_size(), |
| 2730 | u32::MAX, |
| 2731 | "a=max-message-size larger than u32::MAX should be treated as unbounded" |
| 2732 | ); |
| 2733 | } |
| 2734 | } |