mirror of
https://github.com/rustdesk/rustdesk.git
synced 2026-09-11 15:01:02 +03:00
fix: preserve WebRTC transport preference
This commit is contained in:
Submodule libs/hbb_common updated: f98f3e8732...0952f18b8e
157
src/client.rs
157
src/client.rs
@@ -214,9 +214,10 @@ impl Drop for OffererGuard {
|
|||||||
/// - a WebRTC success is committed immediately;
|
/// - a WebRTC success is committed immediately;
|
||||||
/// - an `is_p2p` result from `others` (e.g. IPv6 direct) is committed immediately;
|
/// - an `is_p2p` result from `others` (e.g. IPv6 direct) is committed immediately;
|
||||||
/// - a relayed result inside the preference window is *held*, giving WebRTC until the window
|
/// - a relayed result inside the preference window is *held*, giving WebRTC until the window
|
||||||
/// expires to connect and take priority, while unfinished `others` keep racing so a slower
|
/// expires to connect and take priority. The window starts when the first relay is ready, not
|
||||||
/// direct transport can still win; when the window ends (or WebRTC fails) the held connection
|
/// when setup begins, so rendezvous latency cannot consume the preference budget. Unfinished
|
||||||
/// is committed;
|
/// `others` keep racing so a slower direct transport can still win; when the window ends (or
|
||||||
|
/// WebRTC fails) the held connection is committed;
|
||||||
/// - a failure on one side commits the survivor as soon as it succeeds; two failures compose
|
/// - a failure on one side commits the survivor as soon as it succeeds; two failures compose
|
||||||
/// into one error.
|
/// into one error.
|
||||||
///
|
///
|
||||||
@@ -234,10 +235,15 @@ async fn race_transports_prefer_webrtc<'a, T: 'a>(
|
|||||||
let mut others_err: Option<hbb_common::anyhow::Error> = None;
|
let mut others_err: Option<hbb_common::anyhow::Error> = None;
|
||||||
let window = tokio::time::sleep(Duration::from_millis(window_ms));
|
let window = tokio::time::sleep(Duration::from_millis(window_ms));
|
||||||
tokio::pin!(window);
|
tokio::pin!(window);
|
||||||
let mut window_over = false;
|
let mut window_started = false;
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
res = async { webrtc_fut.as_mut().unwrap().await }, if webrtc_fut.is_some() => {
|
res = async {
|
||||||
|
match webrtc_fut.as_mut() {
|
||||||
|
Some(fut) => fut.await,
|
||||||
|
None => std::future::pending().await,
|
||||||
|
}
|
||||||
|
}, if webrtc_fut.is_some() => {
|
||||||
webrtc_fut = None;
|
webrtc_fut = None;
|
||||||
match res {
|
match res {
|
||||||
// WebRTC connected: preferred outright; a held relay conn just drops.
|
// WebRTC connected: preferred outright; a held relay conn just drops.
|
||||||
@@ -254,17 +260,26 @@ async fn race_transports_prefer_webrtc<'a, T: 'a>(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
res = async { others_fut.as_mut().unwrap().await }, if others_fut.is_some() => {
|
res = async {
|
||||||
|
match others_fut.as_mut() {
|
||||||
|
Some(fut) => fut.await,
|
||||||
|
None => std::future::pending().await,
|
||||||
|
}
|
||||||
|
}, if others_fut.is_some() => {
|
||||||
others_fut = None;
|
others_fut = None;
|
||||||
match res {
|
match res {
|
||||||
Ok((conn, unfinished)) => {
|
Ok((conn, unfinished)) => {
|
||||||
if is_p2p(&conn) || window_over || webrtc_fut.is_none() {
|
if is_p2p(&conn) || webrtc_fut.is_none() {
|
||||||
return Ok(conn);
|
return Ok(conn);
|
||||||
}
|
}
|
||||||
// Hold the first relay, but keep polling unfinished alternatives: an IPv6
|
// Hold the first relay, but keep polling unfinished alternatives: an IPv6
|
||||||
// direct attempt may complete inside the same WebRTC preference window.
|
// direct attempt may complete inside the same WebRTC preference window.
|
||||||
if held.is_none() {
|
if held.is_none() {
|
||||||
held = Some(conn);
|
held = Some(conn);
|
||||||
|
window.as_mut().reset(
|
||||||
|
Instant::now() + Duration::from_millis(window_ms),
|
||||||
|
);
|
||||||
|
window_started = true;
|
||||||
}
|
}
|
||||||
if !unfinished.is_empty() {
|
if !unfinished.is_empty() {
|
||||||
others_fut = Some(select_ok(unfinished));
|
others_fut = Some(select_ok(unfinished));
|
||||||
@@ -277,16 +292,24 @@ async fn race_transports_prefer_webrtc<'a, T: 'a>(
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
_ = &mut window, if !window_over => {
|
_ = &mut window, if window_started => {
|
||||||
window_over = true;
|
|
||||||
if let Some(conn) = held.take() {
|
if let Some(conn) = held.take() {
|
||||||
return Ok(conn);
|
return Ok(conn);
|
||||||
}
|
}
|
||||||
|
window_started = false;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn request_can_carry_webrtc(udp_port: u16, force_relay: bool) -> bool {
|
||||||
|
// A normal TCP punch request must close its rendezvous socket before reusing that local
|
||||||
|
// address, while WebRTC trickle ICE retains the socket as its signaling bridge. Therefore a
|
||||||
|
// request with udp_port=0 must never carry WebRTC. UDP requests can carry it because their
|
||||||
|
// punch socket is separate; force-relay requests can too because they never enter TCP punching.
|
||||||
|
udp_port > 0 || force_relay
|
||||||
|
}
|
||||||
|
|
||||||
impl Client {
|
impl Client {
|
||||||
const CLIENT_CLIPBOARD_NAME: &'static str = "client-clipboard";
|
const CLIENT_CLIPBOARD_NAME: &'static str = "client-clipboard";
|
||||||
|
|
||||||
@@ -438,7 +461,13 @@ impl Client {
|
|||||||
} else {
|
} else {
|
||||||
None
|
None
|
||||||
};
|
};
|
||||||
let webrtc_offerer = if Self::should_create_webrtc_offerer(&interface) {
|
// Prepare WebRTC only for a possible UDP request, or for force-relay where TURN is the
|
||||||
|
// only WebRTC path. `_start_inner` applies the stricter wire invariant after the UDP NAT
|
||||||
|
// test: an actual request with udp_port=0 never includes the offer.
|
||||||
|
let may_prepare_webrtc = udp.0.is_some() || interface.is_force_relay();
|
||||||
|
let webrtc_offerer = if may_prepare_webrtc
|
||||||
|
&& Self::should_create_webrtc_offerer(&interface)
|
||||||
|
{
|
||||||
match WebRTCStream::new("", interface.is_force_relay(), CONNECT_TIMEOUT).await {
|
match WebRTCStream::new("", interface.is_force_relay(), CONNECT_TIMEOUT).await {
|
||||||
Ok(stream) => Some(stream),
|
Ok(stream) => Some(stream),
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
@@ -449,6 +478,7 @@ impl Client {
|
|||||||
} else {
|
} else {
|
||||||
None
|
None
|
||||||
};
|
};
|
||||||
|
let has_webrtc_offerer = webrtc_offerer.is_some();
|
||||||
let fut = Self::_start_inner(
|
let fut = Self::_start_inner(
|
||||||
peer.to_owned(),
|
peer.to_owned(),
|
||||||
key.to_owned(),
|
key.to_owned(),
|
||||||
@@ -466,9 +496,12 @@ impl Client {
|
|||||||
if udp.0.is_none() {
|
if udp.0.is_none() {
|
||||||
return fut.await;
|
return fut.await;
|
||||||
}
|
}
|
||||||
let mut connect_futures = Vec::new();
|
let preferred_fut = fut.boxed();
|
||||||
connect_futures.push(fut.boxed());
|
// This is deliberately a pure TCP punch request: its WebRTC argument must stay `None`.
|
||||||
let fut = Self::_start_inner(
|
// TCP punching closes the rendezvous socket before binding a new connection to the same
|
||||||
|
// local address; a WebRTC ICE bridge would retain that socket and break the port reuse.
|
||||||
|
// The preferred request above may carry WebRTC only because it owns a separate UDP socket.
|
||||||
|
let fallback_fut = Self::_start_inner(
|
||||||
peer.to_owned(),
|
peer.to_owned(),
|
||||||
key.to_owned(),
|
key.to_owned(),
|
||||||
token.to_owned(),
|
token.to_owned(),
|
||||||
@@ -481,8 +514,18 @@ impl Client {
|
|||||||
rendezvous_server,
|
rendezvous_server,
|
||||||
servers,
|
servers,
|
||||||
contained,
|
contained,
|
||||||
);
|
)
|
||||||
connect_futures.push(fut.boxed());
|
.boxed();
|
||||||
|
if has_webrtc_offerer {
|
||||||
|
return race_transports_prefer_webrtc(
|
||||||
|
preferred_fut,
|
||||||
|
vec![fallback_fut],
|
||||||
|
Self::WEBRTC_PREFER_WINDOW_MS,
|
||||||
|
|result| result.0 .1,
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
let connect_futures = vec![preferred_fut, fallback_fut];
|
||||||
match select_ok(connect_futures).await {
|
match select_ok(connect_futures).await {
|
||||||
Ok(conn) => Ok((conn.0 .0, conn.0 .1, conn.0 .2)),
|
Ok(conn) => Ok((conn.0 .0, conn.0 .1, conn.0 .2)),
|
||||||
Err(e) => Err(e),
|
Err(e) => Err(e),
|
||||||
@@ -531,7 +574,9 @@ impl Client {
|
|||||||
/// too.
|
/// too.
|
||||||
/// When force_relay *and* TURN are configured, WebRTC via TURN is a valid "relayed" path:
|
/// When force_relay *and* TURN are configured, WebRTC via TURN is a valid "relayed" path:
|
||||||
/// `connect` keeps a WebRTC win instead of replacing it with the RustDesk relay, and the
|
/// `connect` keeps a WebRTC win instead of replacing it with the RustDesk relay, and the
|
||||||
/// RelayResponse path prefers it within the WebRTC preference window.
|
/// RelayResponse path races it without a P2P preference delay. The caller additionally
|
||||||
|
/// requires either a usable UDP request or force_relay; a normal TCP punch request must never
|
||||||
|
/// carry a WebRTC offer because trickle ICE retains the socket TCP punching needs to reuse.
|
||||||
fn should_create_webrtc_offerer(interface: &impl Interface) -> bool {
|
fn should_create_webrtc_offerer(interface: &impl Interface) -> bool {
|
||||||
if !crate::get_udp_punch_enabled() || use_ws() || Config::is_proxy() {
|
if !crate::get_udp_punch_enabled() || use_ws() || Config::is_proxy() {
|
||||||
return false;
|
return false;
|
||||||
@@ -752,20 +797,26 @@ impl Client {
|
|||||||
.take()
|
.take()
|
||||||
.map(|(socket, addr)| (Some(socket), Some(addr)))
|
.map(|(socket, addr)| (Some(socket), Some(addr)))
|
||||||
.unwrap_or((None, None));
|
.unwrap_or((None, None));
|
||||||
let webrtc_sdp_offer = if let Some(stream) =
|
let udp_nat_port = udp.1.map(|x| *x.lock().unwrap()).unwrap_or(0);
|
||||||
webrtc_offerer.as_ref().and_then(|g| g.stream())
|
let webrtc_sdp_offer =
|
||||||
{
|
if request_can_carry_webrtc(udp_nat_port, interface.is_force_relay()) {
|
||||||
match stream.get_local_endpoint().await {
|
if let Some(stream) = webrtc_offerer.as_ref().and_then(|g| g.stream()) {
|
||||||
Ok(endpoint) => endpoint,
|
match stream.get_local_endpoint_trickle().await {
|
||||||
Err(err) => {
|
Ok(endpoint) => endpoint,
|
||||||
log::warn!("failed to read local WebRTC offer: {}", err);
|
Err(err) => {
|
||||||
|
log::warn!("failed to read local WebRTC offer: {}", err);
|
||||||
|
String::new()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} else {
|
||||||
String::new()
|
String::new()
|
||||||
}
|
}
|
||||||
}
|
} else {
|
||||||
} else {
|
// Hard protocol invariant: a normal request with udp_port=0 is a TCP punch
|
||||||
String::new()
|
// request. It must not start WebRTC trickle on the rendezvous socket that TCP
|
||||||
};
|
// punching needs to close and reuse by local address.
|
||||||
let udp_nat_port = udp.1.map(|x| *x.lock().unwrap()).unwrap_or(0);
|
String::new()
|
||||||
|
};
|
||||||
let punch_type = if udp_nat_port > 0 { "UDP" } else { "TCP" };
|
let punch_type = if udp_nat_port > 0 { "UDP" } else { "TCP" };
|
||||||
msg_out.set_punch_hole_request(PunchHoleRequest {
|
msg_out.set_punch_hole_request(PunchHoleRequest {
|
||||||
id: peer.to_owned(),
|
id: peer.to_owned(),
|
||||||
@@ -960,16 +1011,24 @@ impl Client {
|
|||||||
Ok((Stream::WebRTC(raced), None, "WebRTC"))
|
Ok((Stream::WebRTC(raced), None, "WebRTC"))
|
||||||
}
|
}
|
||||||
.boxed();
|
.boxed();
|
||||||
// The peer answered WebRTC: prefer P2P. The relay result is held for
|
if interface.is_force_relay() {
|
||||||
// the preference window so WebRTC can win even though a relay TCP
|
// Relay-only WebRTC can use only TURN, so it has no P2P advantage
|
||||||
// connect completes much faster than ICE + DTLS setup.
|
// over the RustDesk relay. Take the first successful relay instead
|
||||||
race_transports_prefer_webrtc(
|
// of delaying an already-ready result for the preference window.
|
||||||
webrtc_fut,
|
connect_futures.push(webrtc_fut);
|
||||||
connect_futures,
|
select_ok(connect_futures).await.map(|r| r.0)
|
||||||
Self::WEBRTC_PREFER_WINDOW_MS,
|
} else {
|
||||||
|result| result.2 == "IPv6",
|
// The peer answered WebRTC: prefer P2P. The relay result is held
|
||||||
)
|
// for the preference window so WebRTC can win even though a relay
|
||||||
.await
|
// TCP connect completes much faster than ICE + DTLS setup.
|
||||||
|
race_transports_prefer_webrtc(
|
||||||
|
webrtc_fut,
|
||||||
|
connect_futures,
|
||||||
|
Self::WEBRTC_PREFER_WINDOW_MS,
|
||||||
|
|result| result.2 == "IPv6",
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
// Run all connection attempts concurrently, take the first success.
|
// Run all connection attempts concurrently, take the first success.
|
||||||
select_ok(connect_futures).await.map(|r| r.0)
|
select_ok(connect_futures).await.map(|r| r.0)
|
||||||
@@ -5093,7 +5152,7 @@ async fn udp_nat_connect(
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod webrtc_race_tests {
|
mod webrtc_race_tests {
|
||||||
use super::race_transports_prefer_webrtc;
|
use super::{race_transports_prefer_webrtc, request_can_carry_webrtc};
|
||||||
use hbb_common::{
|
use hbb_common::{
|
||||||
anyhow::anyhow,
|
anyhow::anyhow,
|
||||||
futures::future::{BoxFuture, FutureExt},
|
futures::future::{BoxFuture, FutureExt},
|
||||||
@@ -5120,6 +5179,13 @@ mod webrtc_race_tests {
|
|||||||
|
|
||||||
const NOT_P2P: fn(&&'static str) -> bool = |_| false;
|
const NOT_P2P: fn(&&'static str) -> bool = |_| false;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn tcp_punch_request_never_carries_webrtc() {
|
||||||
|
assert!(!request_can_carry_webrtc(0, false));
|
||||||
|
assert!(request_can_carry_webrtc(1, false));
|
||||||
|
assert!(request_can_carry_webrtc(0, true));
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn webrtc_preferred_over_faster_relay_within_window() {
|
async fn webrtc_preferred_over_faster_relay_within_window() {
|
||||||
let got = race_transports_prefer_webrtc(
|
let got = race_transports_prefer_webrtc(
|
||||||
@@ -5133,6 +5199,19 @@ mod webrtc_race_tests {
|
|||||||
assert_eq!(got, "webrtc");
|
assert_eq!(got, "webrtc");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn preference_window_starts_when_relay_is_ready() {
|
||||||
|
let got = race_transports_prefer_webrtc(
|
||||||
|
ok_after(350, "webrtc"),
|
||||||
|
vec![ok_after(250, "relay")],
|
||||||
|
200,
|
||||||
|
NOT_P2P,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(got, "webrtc");
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn unfinished_ipv6_still_beats_held_relay() {
|
async fn unfinished_ipv6_still_beats_held_relay() {
|
||||||
let start = Instant::now();
|
let start = Instant::now();
|
||||||
|
|||||||
@@ -697,7 +697,7 @@ impl RendezvousMediator {
|
|||||||
) -> ResultType<String> {
|
) -> ResultType<String> {
|
||||||
let mut stream =
|
let mut stream =
|
||||||
WebRTCStream::new(&ph.webrtc_sdp_offer, force_relay, CONNECT_TIMEOUT).await?;
|
WebRTCStream::new(&ph.webrtc_sdp_offer, force_relay, CONNECT_TIMEOUT).await?;
|
||||||
let answer = match stream.get_local_endpoint().await {
|
let answer = match stream.get_local_endpoint_trickle().await {
|
||||||
Ok(answer) => answer,
|
Ok(answer) => answer,
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
// Close the freshly-created pc so a failure here doesn't leak it in SESSIONS.
|
// Close the freshly-created pc so a failure here doesn't leak it in SESSIONS.
|
||||||
|
|||||||
Reference in New Issue
Block a user