From 6ffb506961a27b2c6f11de24920c422606f93604 Mon Sep 17 00:00:00 2001 From: rustdesk Date: Thu, 27 Aug 2026 11:36:02 +0800 Subject: [PATCH] webrtc: skip the controller's ICE re-send instead of queueing it twice The controller sends every candidate twice, because the server's hop to a peer registered over UDP can lose one. The ICE agent that dedups repeats sits downstream of the answerer's queue, so the answerer paid for both copies: a slot, a JSON parse, and the ICE agent's lock, once per repeat. Remember a digest of what was queued and skip the repeat. Recorded only once queued, so a candidate a full queue refused stays repairable by the re-send. The queue's depth is unchanged. A real peer gathers well under it - four STUN servers, link-local IPv6 filtered, one component - and the drain empties it as candidates trickle in, so what this removes is the redundant work, not an overflow. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_019aokqJuhjvB3kijXtAg5Ns --- src/rendezvous_mediator.rs | 164 ++++++++++++++++++++++++++++++++----- 1 file changed, 143 insertions(+), 21 deletions(-) diff --git a/src/rendezvous_mediator.rs b/src/rendezvous_mediator.rs index 2c09128c4..ed951314d 100644 --- a/src/rendezvous_mediator.rs +++ b/src/rendezvous_mediator.rs @@ -1,5 +1,6 @@ use std::{ - collections::HashMap, + collections::{hash_map::RandomState, HashMap, VecDeque}, + hash::BuildHasher, net::SocketAddr, sync::{ atomic::{AtomicBool, Ordering}, @@ -53,12 +54,16 @@ lazy_static::lazy_static! { static ref SOLVING_PK_MISMATCH: Mutex = Default::default(); static ref LAST_MSG: Mutex<(SocketAddr, Instant)> = Mutex::new((SocketAddr::new([0; 4].into(), 0), Instant::now())); static ref LAST_RELAY_MSG: Mutex<(SocketAddr, Instant)> = Mutex::new((SocketAddr::new([0; 4].into(), 0), Instant::now())); - static ref WEBRTC_ICE_TXS: Mutex>> = Default::default(); + static ref WEBRTC_ICE_TXS: Mutex> = Default::default(); + static ref ICE_DIGEST_STATE: RandomState = Default::default(); } -/// Remote ICE candidates buffered per session while the answerer applies them. Mirrors the -/// controller's own cap: gathering yields host, then srflx, then relay, so a real peer sends -/// well under this, and anything past it is someone deciding how much memory this process holds. +/// Remote ICE candidates buffered per session while the answerer applies them. Same depth as the +/// controller's own buffer (`Client::MAX_PENDING_WEBRTC_ICE`), though that one evicts its oldest +/// where a full channel here refuses the newest. const MAX_PENDING_REMOTE_ICE: usize = 64; +/// Queued candidates remembered so the controller's re-send is skipped instead of taking a slot +/// of its own. Far more than an honest peer gathers, at eight bytes each. +const ICE_DEDUP_WINDOW: usize = 256; // The rendezvous ICE route is reachable without a prior punch and the peer decides how many // candidates it sends, so these sites would let someone else set how much this machine writes to // its log file. One line a minute each, carrying the suppressed count. @@ -70,6 +75,45 @@ static REJECTED_REMOTE_ICE_LOG: hbb_common::log_throttle::LogThrottle = static FULL_ICE_QUEUE_LOG: hbb_common::log_throttle::LogThrottle = hbb_common::log_throttle::LogThrottle::new(ICE_LOG_INTERVAL); +struct IceRoute { + tx: mpsc::Sender, + recent: VecDeque, +} + +impl IceRoute { + fn new(tx: mpsc::Sender) -> Self { + Self { + tx, + recent: VecDeque::new(), + } + } + + /// Keeps `queue` the only way onto the channel, so nothing reaches it unrecorded. + fn is_same_channel(&self, other: &mpsc::Sender) -> bool { + self.tx.same_channel(other) + } + + /// Skip the controller's re-send of a candidate already queued: the ICE agent that dedups + /// repeats is downstream of this queue, so the copy would spend a slot of its own. + /// False means the candidate was dropped. + fn queue(&mut self, candidate: String) -> bool { + let digest = ICE_DIGEST_STATE.hash_one(candidate.as_str()); + if self.recent.contains(&digest) { + // Only honest about the drop if the route is still alive to have taken it. + return !self.tx.is_closed(); + } + // Recorded once queued, never before: a refused candidate stays repairable by the re-send. + if self.tx.try_send(candidate).is_err() { + return false; + } + if self.recent.len() >= ICE_DEDUP_WINDOW { + self.recent.pop_front(); + } + self.recent.push_back(digest); + true + } +} + static SHOULD_EXIT: AtomicBool = AtomicBool::new(false); static MANUAL_RESTARTED: AtomicBool = AtomicBool::new(false); static SENT_REGISTER_PK: AtomicBool = AtomicBool::new(false); @@ -422,19 +466,27 @@ impl RendezvousMediator { }); } Some(rendezvous_message::Union::IceCandidate(ice)) => { - let tx = WEBRTC_ICE_TXS.lock().await.get(&ice.session_key).cloned(); - if let Some(tx) = tx { - if tx.try_send(ice.candidate).is_err() { + let queued = { + let mut txs = WEBRTC_ICE_TXS.lock().await; + txs.get_mut(&ice.session_key) + .map(|route| route.queue(ice.candidate)) + }; + match queued { + Some(false) => { if let Some(n) = FULL_ICE_QUEUE_LOG.due() { log::debug!("dropped {} ICE candidate(s): queue full or closed", n); } } - } else if let Some(n) = UNKNOWN_ICE_SESSION_LOG.due() { - log::debug!( - "dropped {} ICE candidate(s) for unknown WebRTC session key, last: {}", - n, - ice.session_key - ); + None => { + if let Some(n) = UNKNOWN_ICE_SESSION_LOG.due() { + log::debug!( + "dropped {} ICE candidate(s) for unknown WebRTC session key, last: {}", + n, + ice.session_key + ); + } + } + _ => {} } } Some(rendezvous_message::Union::ConfigureUpdate(cu)) => { @@ -714,17 +766,17 @@ impl RendezvousMediator { return Ok(answer); }; - // Bounded, like the controller's own candidate buffer: how many candidates arrive is the - // sender's choice, while draining one costs a JSON parse and the ICE agent's lock, so an - // unbounded queue lets whoever can reach this session's route grow it without limit inside - // a long-lived service process. A full queue drops the newest candidate, which costs at - // most one path; ICE keeps whatever pairs it already has. + // Bounded: how many candidates arrive is the sender's choice, while draining one costs a + // JSON parse and the ICE agent's lock, so an unbounded queue lets whoever can reach this + // session's route grow it without limit inside a long-lived service process. A full queue + // drops the newest candidate, and the controller re-sends it once — the digests beside the + // sender are what keep that re-send from spending a slot of its own. let (remote_ice_tx, mut remote_ice_rx) = mpsc::channel::(MAX_PENDING_REMOTE_ICE); let own_ice_tx = remote_ice_tx.clone(); WEBRTC_ICE_TXS .lock() .await - .insert(session_key.clone(), remote_ice_tx); + .insert(session_key.clone(), IceRoute::new(remote_ice_tx)); let stream_for_remote_ice = stream.clone(); tokio::spawn(async move { @@ -806,7 +858,7 @@ impl RendezvousMediator { let mut txs = WEBRTC_ICE_TXS.lock().await; if txs .get(&session_key_for_cleanup) - .is_some_and(|tx| tx.same_channel(&own_ice_tx)) + .is_some_and(|route| route.is_same_channel(&own_ice_tx)) { txs.remove(&session_key_for_cleanup); } @@ -1236,3 +1288,73 @@ impl Drop for CheckIfResendPk { } } } + +#[cfg(test)] +mod tests { + use super::{mpsc, IceRoute, ICE_DEDUP_WINDOW, MAX_PENDING_REMOTE_ICE}; + + fn queue(route: &mut IceRoute, candidate: &str) -> bool { + route.queue(candidate.to_owned()) + } + + #[test] + fn the_re_sent_copy_does_not_spend_a_queue_slot() { + // Two slots, three sends: without the dedup the re-send takes the second and "relay", + // the one that traverses NAT, is the one refused. + let (tx, mut rx) = mpsc::channel::(2); + let mut route = IceRoute::new(tx); + for _ in 0..2 { + assert!(queue(&mut route, "host")); + } + assert!(queue(&mut route, "relay")); + let mut queued = Vec::new(); + while let Ok(candidate) = rx.try_recv() { + queued.push(candidate); + } + assert_eq!(queued, vec!["host".to_owned(), "relay".to_owned()]); + } + + #[test] + fn a_candidate_the_full_queue_refused_is_not_remembered() { + let (tx, mut rx) = mpsc::channel::(1); + let mut route = IceRoute::new(tx); + assert!(queue(&mut route, "host")); + assert!(!queue(&mut route, "relay")); + // The re-send is the only repair for a refused candidate; remembering it would swallow it. + assert_eq!(rx.try_recv().ok(), Some("host".to_owned())); + assert!(queue(&mut route, "relay")); + assert_eq!(rx.try_recv().ok(), Some("relay".to_owned())); + } + + #[test] + fn a_re_send_is_skipped_while_the_original_is_still_queued() { + let (tx, mut rx) = mpsc::channel::(MAX_PENDING_REMOTE_ICE); + let mut route = IceRoute::new(tx); + for i in 0..MAX_PENDING_REMOTE_ICE { + assert!(queue(&mut route, &format!("candidate-{}", i))); + } + assert!(queue(&mut route, "candidate-0")); + let mut queued = 0; + while rx.try_recv().is_ok() { + queued += 1; + } + assert_eq!(queued, MAX_PENDING_REMOTE_ICE); + } + + #[test] + fn the_window_forgets_in_arrival_order() { + let (tx, mut rx) = mpsc::channel::(MAX_PENDING_REMOTE_ICE); + let mut route = IceRoute::new(tx); + for i in 0..=ICE_DEDUP_WINDOW { + assert!(queue(&mut route, &format!("candidate-{}", i))); + assert!(rx.try_recv().is_ok()); + } + // The oldest digest made room for the newest, so its re-send is admitted again. + assert!(queue(&mut route, "candidate-0")); + assert!(rx.try_recv().is_ok()); + // A recent one is still skipped. + let recent = format!("candidate-{}", ICE_DEDUP_WINDOW); + assert!(queue(&mut route, &recent)); + assert!(rx.try_recv().is_err()); + } +}