diff --git a/src/webrtc.rs b/src/webrtc.rs index bd4ccced6..65c3fa942 100644 --- a/src/webrtc.rs +++ b/src/webrtc.rs @@ -82,6 +82,10 @@ pub struct WebRTCStream { // Whether anyone still wants this pc — see `HandoffState`. Shared by every clone, the // `SESSIONS` entry included, because that question is about the pc, not about one handle. handoff: Arc, + // The offer/answer envelope as created, kept because `pc.local_description()` grows with + // every gathered candidate — see `get_local_endpoint_trickle`. Shared, not copied per clone: + // a session clones this struct several times and reads the envelope once. + local_endpoint: Arc, } /// Whether a pc `new_inner` built is still wanted, shared by every `NewStreamHandoff` handed out @@ -228,6 +232,7 @@ impl Clone for WebRTCStream { recv_state: self.recv_state.clone(), peer_verified: self.peer_verified.clone(), handoff: self.handoff.clone(), + local_endpoint: self.local_endpoint.clone(), } } } @@ -754,34 +759,58 @@ impl WebRTCStream { // Trickle ICE: local-only work, no gathering wait — the controlled side awaits answer // creation inline on its punch-reply critical path. A failure below leaves a live pc whose // state handler only fires on a terminal ICE state, so close before propagating. - let offer_answer: ResultType = async { + // Encode each endpoint before `set_local_description`, which is what starts gathering, so + // there is nothing for `trickle_endpoint` to strip in the first place. + let offer_answer: ResultType<(String, String)> = async { if start_local_offer { let sdp = pc.create_offer(None).await?; + let endpoint = Self::trickle_endpoint(&sdp, force_relay)?; pc.set_local_description(sdp.clone()).await?; // SDP carries host/srflx IPs and ICE ufrag/pwd; log only its size, not the body. - log::debug!("local offer SDP built ({} bytes)", sdp.sdp.len()); + log::debug!( + "local offer SDP built ({} bytes, endpoint {} bytes)", + sdp.sdp.len(), + endpoint.len() + ); let k = Self::get_key_for_sdp(&sdp)?; log::debug!("Start webrtc with local key: {}", k); - Ok(k) + Ok((k, endpoint)) } else { let sdp = serde_json::from_str::(&remote_offer)?; pc.set_remote_description(sdp.clone()).await?; let answer = pc.create_answer(None).await?; + let endpoint = Self::trickle_endpoint(&answer, force_relay)?; pc.set_local_description(answer).await?; - log::debug!("remote offer SDP received ({} bytes)", sdp.sdp.len()); + log::debug!( + "remote offer SDP received ({} bytes), local answer endpoint {} bytes", + sdp.sdp.len(), + endpoint.len() + ); let k = Self::get_key_for_sdp(&sdp)?; log::debug!("Start webrtc with remote key: {}", k); - Ok(k) + Ok((k, endpoint)) } } .await; - key = match offer_answer { - Ok(k) => k, + let (new_key, local_endpoint) = match offer_answer { + Ok(x) => x, Err(e) => { pc.close().await.ok(); return Err(e); } }; + // Only the UDP leg has a ceiling, and which leg carries this is the rendezvous server's + // to decide, so this reports rather than refuses: over a TCP/WS route a larger endpoint is + // delivered fine, and failing here would cost WebRTC for no reason. + if local_endpoint.len() > Self::UDP_ENDPOINT_BUDGET { + log::warn!( + "WebRTC endpoint is {} bytes, over the {} byte UDP budget; a peer reached over UDP \ + may never receive it", + local_endpoint.len(), + Self::UDP_ENDPOINT_BUDGET + ); + } + key = new_key; let webrtc_stream = Self { pc, @@ -796,6 +825,7 @@ impl WebRTCStream { recv_state: Arc::new(Mutex::new(RecvState::default())), peer_verified: Arc::new(AtomicBool::new(false)), handoff: Default::default(), + local_endpoint: Arc::new(local_endpoint), }; // Insert into the session cache, but never `await pc.close()` while holding this lock: // `close()` fires the peer-connection-state handler inline, which itself locks SESSIONS, @@ -841,29 +871,73 @@ impl WebRTCStream { // gathered host/srflx/relay candidates. let mut gather_complete = self.pc.gathering_complete_promise().await; let _gathering_channel_closed = gather_complete.recv().await; - self.get_local_endpoint_trickle().await + let Some(local_desc) = self.pc.local_description().await else { + return Err(anyhow::anyhow!("Local desc is not set")); + }; + Self::encode_endpoint(&local_desc, self.relay_only) } - /// Return the current local description immediately for callers that signal candidates via - /// `take_local_ice_rx`. Unlike `get_local_endpoint`, this does not wait for ICE gathering. + /// The offer/answer exactly as it was created, for callers that signal candidates via + /// `take_local_ice_rx`. + /// + /// Deliberately not `pc.local_description()`: that one runs `populate_local_candidates`, so + /// it returns the SDP plus every candidate gathered up to that moment — hundreds of bytes on + /// a multi-homed host, and the caller reads it a network round trip after `new`, by which + /// time gathering has filled it in. The rendezvous server forwards this blob to the peer as a + /// single UDP datagram, and one that needs IP fragmentation is dropped outright on paths that + /// discard fragments, silently and every time. Candidates belong on the trickle channel. #[inline] pub async fn get_local_endpoint_trickle(&self) -> ResultType { - if let Some(local_desc) = self.pc.local_description().await { - let sdp = if self.relay_only { - serde_json::to_string(&local_desc)? - } else { - // Rides in the envelope, not a proto field the rendezvous server would forward, - // because it is the offer's own property: the receiver must know whether - // force_relay was policy (stay Relay-only) or transport. Older peers ignore it. - let mut v = serde_json::to_value(&local_desc)?; - v[Self::ICE_POLICY_KEY] = serde_json::Value::from(Self::ICE_POLICY_ALL); - serde_json::to_string(&v)? - }; - let endpoint = Self::sdp_to_endpoint(&sdp); - Ok(endpoint) - } else { - Err(anyhow::anyhow!("Local desc is not set")) + Ok(self.local_endpoint.as_ref().clone()) + } + + /// How large the offer/answer may be before the UDP leg of its route stops being safe. It is + /// the one signaling message that can grow: a peer registered over UDP receives it inside + /// `PunchHole` — with the mangled addresses and the permission blobs — as a single datagram, + /// and one that needs IP fragmentation is dropped without a trace on paths that discard + /// fragments. The candidates that follow are one small message each, so they never approach + /// this. Nothing enforces it: a TCP/WS leg has no packet ceiling, and `trickle_endpoint` + /// already pins the value near 700 bytes either way. It is what the warning and the tests + /// measure against, so a regression shows up as a log line instead of a silent drop. + const UDP_ENDPOINT_BUDGET: usize = 1024; + + /// The endpoint a trickling peer signals: the session parameters needed to start ICE and DTLS + /// (ice-ufrag, ice-pwd, fingerprint, setup, the sctp m-line), without the candidates — those + /// ride `take_local_ice_rx` and arrive as individual `IceCandidate` messages. + /// + /// The split is what keeps this endpoint a fixed ~700 bytes: candidates are the only part of a + /// local description that grows, and `pc.local_description()` appends every one gathered so + /// far. Stripping them here makes the bound a property of the value rather than of when it was + /// taken, so the size cannot drift with gathering however this is later refactored. + fn trickle_endpoint(local_desc: &RTCSessionDescription, relay_only: bool) -> ResultType { + let mut sdp = String::with_capacity(local_desc.sdp.len()); + for line in local_desc.sdp.lines() { + // `end-of-candidates` would tell the remote agent to stop waiting for the trickle. + if line.starts_with("a=candidate:") || line.starts_with("a=end-of-candidates") { + continue; + } + sdp.push_str(line); + sdp.push_str("\r\n"); } + let mut desc = local_desc.clone(); + desc.sdp = sdp; + Self::encode_endpoint(&desc, relay_only) + } + + /// Base64 envelope of a local description, carrying its ICE transport policy unless the pc is + /// Relay-only (see `ICE_POLICY_KEY`). + fn encode_endpoint(local_desc: &RTCSessionDescription, relay_only: bool) -> ResultType { + let sdp = if relay_only { + serde_json::to_string(local_desc)? + } else { + // Rides in the envelope, not a proto field the rendezvous server would forward, + // because it is the offer's own property: the receiver must know whether + // force_relay was policy (stay Relay-only) or transport. Older peers ignore it. + let mut v = serde_json::to_value(local_desc)?; + v[Self::ICE_POLICY_KEY] = serde_json::Value::from(Self::ICE_POLICY_ALL); + serde_json::to_string(&v)? + }; + Ok(Self::sdp_to_endpoint(&sdp)) } /// Whether the peer's endpoint declares it was built with ICE transport policy `all` @@ -1476,6 +1550,88 @@ mod tests { .expect("extra envelope key must not break RTCSessionDescription parsing"); } + /// The rendezvous server hands a trickle endpoint to the peer inside one UDP datagram, so it + /// must stay small however far gathering has got. `pc.local_description()` appends every + /// candidate gathered so far — the stored endpoint must not, or the datagram fragments and is + /// dropped without a trace on paths that discard fragments. + #[tokio::test] + async fn test_trickle_endpoint_never_grows_with_candidates() { + let offerer = WebRTCStream::new("", false, 20000).await.unwrap(); + let first = offerer.get_local_endpoint_trickle().await.unwrap(); + + let mut gather_complete = offerer.pc.gathering_complete_promise().await; + let _ = timeout(Duration::from_secs(10), gather_complete.recv()).await; + + let after = offerer.get_local_endpoint_trickle().await.unwrap(); + assert_eq!(first, after, "trickle endpoint changed while ICE gathered"); + assert!( + after.len() <= WebRTCStream::UDP_ENDPOINT_BUDGET, + "trickle endpoint is {} bytes, over the {} byte budget", + after.len(), + WebRTCStream::UDP_ENDPOINT_BUDGET + ); + assert!( + !WebRTCStream::get_remote_offer(&after) + .unwrap() + .contains("a=candidate"), + "trickle endpoint must carry no ICE candidates" + ); + + // The one-shot endpoint keeps its old contract, which is also what makes the check above + // meaningful: candidates really were gathered by now. + let gathered = offerer.get_local_endpoint().await.unwrap(); + assert!( + gathered.len() > after.len(), + "gathered endpoint {} is not larger than the trickle one {}", + gathered.len(), + after.len() + ); + offerer.close().await; + } + + /// The split itself: candidates out, everything ICE and DTLS need to start in. A stripped + /// endpoint the peer cannot consume would trade a silent drop for a silent failure. + #[tokio::test] + async fn test_trickle_endpoint_splits_candidates_off() { + let offerer = WebRTCStream::new("", false, 20000).await.unwrap(); + let mut gather_complete = offerer.pc.gathering_complete_promise().await; + let _ = timeout(Duration::from_secs(10), gather_complete.recv()).await; + + let gathered = offerer.pc.local_description().await.unwrap(); + assert!( + gathered.sdp.contains("a=candidate:"), + "nothing to strip: the pc gathered no candidate" + ); + + let endpoint = WebRTCStream::trickle_endpoint(&gathered, false).unwrap(); + let json = WebRTCStream::get_remote_offer(&endpoint).unwrap(); + assert!(!json.contains("a=candidate:"), "candidate survived the split"); + assert!( + !json.contains("a=end-of-candidates"), + "end-of-candidates would stop the remote agent waiting for the trickle" + ); + for keep in [ + "m=application", + "a=ice-ufrag:", + "a=ice-pwd:", + "a=fingerprint:", + "a=setup:", + "a=mid:", + "a=sctp-port:", + ] { + assert!(json.contains(keep), "trickle endpoint dropped {keep}"); + } + assert!( + endpoint.len() <= WebRTCStream::UDP_ENDPOINT_BUDGET, + "stripped endpoint is {} bytes", + endpoint.len() + ); + serde_json::from_str::(&json) + .expect("a stripped endpoint must still parse as a session description"); + + offerer.close().await; + } + #[test] fn test_webrtc_session_key() { let mut sdp_str = "".to_owned();