webrtc: reject data channels effectively; hand the permit to a desynced close

- The on_data_channel guards called dc.close() from inside the handler,
  which webrtc-rs runs to completion BEFORE handle_open binds the SCTP
  stream. At that point close() only flips the ready state and returns,
  and handle_open sets it back to Open; with detached channels no read
  loop is spawned either. So a refused channel stayed open and undrained,
  and its queued bytes count against the association-wide receive window
  - about a megabyte written into an ignored channel stalls the one
  carrying the session, from an unauthenticated peer. Close from the
  channel's own on_open instead, where the stream exists. This is also
  what makes the ordered+reliable precondition an actual rejection
  rather than a log line.

- A write that failed part-way through a fragment sequence called
  close_detached() and let the send permit drop, while the timeout path
  hands the permit to the teardown for exactly this reason: the close is
  in flight, state_notify is still Open, and callers wrap sends in
  allow_err! and keep going - so the next message could be written onto
  the orphaned prefix the failure left on the peer. send_bytes_inner now
  reports that it desynced and the caller, which holds the permit, does
  the teardown.

- Drop the catch_unwind around Handle::spawn: release builds set
  panic = "abort", so it can never catch anything, and the comment
  claimed a mitigation that does not exist. The reachable case is having
  no runtime handle at all, which is now logged at debug rather than
  warn - Stream's Drop reaches this on every non-runtime thread.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ExUfAkYbq8UC9pQCiLy8TQ
This commit is contained in:
rustdesk
2026-08-08 19:13:24 +08:00
parent ddea60cd3d
commit 2cb8d0c7d8
+64 -33
View File
@@ -256,6 +256,28 @@ impl WebRTCStream {
) )
} }
/// Reject a data channel from inside `on_data_channel`, where closing it directly does not
/// work.
///
/// webrtc-rs runs this handler to completion BEFORE `handle_open` binds the SCTP stream, so
/// at this point `close()` only flips the ready state and returns — no stream reset — and
/// `handle_open` then sets it back to Open. With detached channels no read loop is spawned
/// either, so a channel refused here would stay open and undrained, and its queued bytes
/// count against the association-wide receive window: about a megabyte written into an
/// ignored channel stalls the one carrying the session. Close from the channel's own
/// `on_open` instead, which runs once the stream exists.
fn close_unbound_channel(dc: Arc<RTCDataChannel>) {
let dc_for_close = dc.clone();
dc.on_open(Box::new(move || {
let dc = dc_for_close.clone();
Box::pin(async move {
if let Err(err) = dc.close().await {
log::debug!("failed to close rejected data channel: {}", err);
}
})
}));
}
/// Whether a `SESSIONS` entry may be handed to a caller asking for `force_relay`. /// Whether a `SESSIONS` entry may be handed to a caller asking for `force_relay`.
/// ///
/// A hit returns a pc built for the FIRST caller, and two of its properties belong to this /// A hit returns a pc built for the FIRST caller, and two of its properties belong to this
@@ -583,7 +605,7 @@ impl WebRTCStream {
"Rejecting WebRTC data channel {}: not ordered and fully reliable", "Rejecting WebRTC data channel {}: not ordered and fully reliable",
d_label d_label
); );
let _ = dc.close().await; Self::close_unbound_channel(dc);
return; return;
} }
// Bind the first channel only. `detached` caches the first detached handle // Bind the first channel only. `detached` caches the first detached handle
@@ -593,7 +615,7 @@ impl WebRTCStream {
// Closed — and that watch gates both `send_bytes_inner` and `next()`. // Closed — and that watch gates both `send_bytes_inner` and `next()`.
if dc_bound.swap(true, Ordering::SeqCst) { if dc_bound.swap(true, Ordering::SeqCst) {
log::warn!("Ignoring extra WebRTC data channel {}", d_label); log::warn!("Ignoring extra WebRTC data channel {}", d_label);
let _ = dc.close().await; Self::close_unbound_channel(dc);
return; return;
} }
log::debug!("Remote data channel {} ready", d_label); log::debug!("Remote data channel {} ready", d_label);
@@ -941,24 +963,20 @@ impl WebRTCStream {
pub fn close_detached_with<T: Send + 'static>(&self, keep: T) { pub fn close_detached_with<T: Send + 'static>(&self, keep: T) {
let pc = self.pc.clone(); let pc = self.pc.clone();
// Take the runtime handle explicitly rather than calling `tokio::spawn`: this also runs // Take the runtime handle explicitly rather than calling `tokio::spawn`: this also runs
// from `Drop` impls, which can execute during runtime teardown where a bare spawn panics. // from `Drop`, which can execute on a thread with no runtime, where a bare spawn panics
// `Handle::spawn` can panic while the runtime is shutting down too, so catch it — a brief // and a panic in a destructor aborts the process. Without a handle the pc cannot be closed
// leak until process exit beats aborting the process from a destructor. // here at all; it is released at process exit. (Catching a panic from `Handle::spawn`
// during runtime shutdown is not an option: release builds set `panic = "abort"`.)
let Ok(handle) = tokio::runtime::Handle::try_current() else { let Ok(handle) = tokio::runtime::Handle::try_current() else {
log::warn!("no tokio runtime available to close the WebRTC peer connection"); log::debug!("no tokio runtime available to close the WebRTC peer connection");
return; return;
}; };
let spawned = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { handle.spawn(async move {
handle.spawn(async move { let _keep = keep;
let _keep = keep; if let Err(err) = pc.close().await {
if let Err(err) = pc.close().await { log::debug!("WebRTC background close failed: {}", err);
log::debug!("WebRTC background close failed: {}", err); }
} });
});
}));
if spawned.is_err() {
log::warn!("failed to spawn the WebRTC close (runtime shutting down)");
}
} }
#[inline] #[inline]
@@ -1088,8 +1106,17 @@ impl WebRTCStream {
// connection down for missing them. The narrower the miss, the more certain the // connection down for missing them. The narrower the miss, the more certain the
// teardown, which is precisely backwards. // teardown, which is precisely backwards.
let write_deadline = Instant::now() + Duration::from_millis(send_timeout); let write_deadline = Instant::now() + Duration::from_millis(send_timeout);
match timeout_at(write_deadline, self.send_bytes_inner(bytes)).await { let mut desynced = false;
Ok(res) => res, match timeout_at(write_deadline, self.send_bytes_inner(bytes, &mut desynced)).await {
// A write that failed part-way leaves the same orphaned prefix a timeout does, so
// it needs the same teardown — and the same handoff of the permit, or a waiting
// clone appends its message to that prefix while the close is still in flight.
Ok(res) => {
if desynced {
self.close_detached_with(send_permit);
}
res
}
Err(_) => { Err(_) => {
// Hand the logical-message permit to the teardown so no waiting clone can // Hand the logical-message permit to the teardown so no waiting clone can
// append a new message after a partially-written fragment sequence. Holding // append a new message after a partially-written fragment sequence. Holding
@@ -1099,17 +1126,24 @@ impl WebRTCStream {
} }
} }
} else { } else {
let _send_permit = send_gate.acquire_owned().await.map_err(|err| { let send_permit = send_gate.acquire_owned().await.map_err(|err| {
Error::new( Error::new(
ErrorKind::BrokenPipe, ErrorKind::BrokenPipe,
format!("WebRTC send gate closed: {}", err), format!("WebRTC send gate closed: {}", err),
) )
})?; })?;
self.send_bytes_inner(bytes).await let mut desynced = false;
let res = self.send_bytes_inner(bytes, &mut desynced).await;
if desynced {
self.close_detached_with(send_permit);
}
res
} }
} }
async fn send_bytes_inner(&mut self, bytes: Bytes) -> ResultType<()> { /// `desynced` is set when the failure left a partial fragment sequence on the wire, i.e. when
/// the stream can no longer be made consistent and the caller must close it.
async fn send_bytes_inner(&mut self, bytes: Bytes, desynced: &mut bool) -> ResultType<()> {
// Same bound the receiver enforces, so we never emit a message a same-version peer // Same bound the receiver enforces, so we never emit a message a same-version peer
// would have to kill the connection over. // would have to kill the connection over.
if bytes.len() > MAX_RECV_MESSAGE { if bytes.len() > MAX_RECV_MESSAGE {
@@ -1131,17 +1165,14 @@ impl WebRTCStream {
framed.put_u8(if is_last { FRAG_END } else { FRAG_MORE }); framed.put_u8(if is_last { FRAG_END } else { FRAG_MORE });
framed.put_slice(chunk); framed.put_slice(chunk);
if let Err(err) = dc.write(&framed.freeze()).await { if let Err(err) = dc.write(&framed.freeze()).await {
if wrote_any { // A sequence that stops with fragments already on the wire leaves the peer's
// The sequence stops with fragments already on the wire, so the peer's // accumulator holding a prefix that will never be terminated — and this framing
// accumulator holds a prefix that will never be terminated — and this framing // carries no length or sequence number for it to notice with, so the next message
// carries no length or sequence number for it to notice with, so the next // is appended to the orphan and parsed as one corrupt frame. Report it so the
// message is appended to the orphan and parsed as one corrupt frame. The // caller can tear the stream down while still holding the send permit; callers
// stream cannot be made consistent again from this side; kill it. (Callers // wrap sends in allow_err! and keep going, so returning the error alone would
// wrap sends in allow_err! and keep going, so returning the error alone would // leave the corruption in place.
// leave the corruption in place.) *desynced = wrote_any;
log::warn!("WebRTC send failed mid-message, closing: {}", err);
self.close_detached();
}
return Err(err.into()); return Err(err.into());
} }
wrote_any = true; wrote_any = true;