Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 5 additions & 0 deletions crates/ironrdp-client/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,9 @@ test = false
[features]
default = []

# Internal integration-test access; not part of the supported public API.
__test = []

rustls = [
"ironrdp-tls/rustls",
"tokio-tungstenite/rustls-tls-native-roots",
Expand Down Expand Up @@ -50,6 +53,7 @@ webauthn = ["dep:ironrdp-rdpewa", "dep:ironrdp-rdpewa-native", "dvc-com-plugin"]
vmconnect = ["dep:ironrdp-vmconnect"]
location = ["dep:ironrdp-rdpel"]
udp = [
"dep:ironrdp-rdpemt",
"dep:ironrdp-rdpeudp",
"dep:ironrdp-rdpeudp-tokio",
"ironrdp-rdpeudp-tokio/rustls-aws-lc-rs",
Expand Down Expand Up @@ -82,6 +86,7 @@ ironrdp-displaycontrol = { path = "../ironrdp-displaycontrol", version = "0.8" }
ironrdp-echo = { path = "../ironrdp-echo", version = "0.4" }
ironrdp-egfx = { path = "../ironrdp-egfx", version = "0.3" }
ironrdp-rdpei = { path = "../ironrdp-rdpei", version = "0.1" }
ironrdp-rdpemt = { path = "../ironrdp-rdpemt", version = "0.1", optional = true }
ironrdp-rdpeudp = { path = "../ironrdp-rdpeudp", version = "0.1", optional = true }
ironrdp-rdpeudp-tokio = { path = "../ironrdp-rdpeudp-tokio", version = "0.1", optional = true }
ironrdp-tls = { path = "../ironrdp-tls", version = "0.2" } # public
Expand Down
10 changes: 10 additions & 0 deletions crates/ironrdp-client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,3 +16,13 @@ pub mod rdp;
mod clipboard;

mod ws;

#[cfg(all(feature = "udp", feature = "__test"))]
#[doc(hidden)]
pub mod udp;
#[cfg(all(feature = "udp", not(feature = "__test")))]
#[expect(
unreachable_pub,
reason = "the __test feature exposes this module to the shared integration tests"
)]
mod udp;
91 changes: 70 additions & 21 deletions crates/ironrdp-client/src/rdp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ use std::io;
use std::sync::Arc;
#[cfg(feature = "location")]
use std::sync::mpsc as std_mpsc;
#[cfg(feature = "location")]
#[cfg(any(feature = "location", feature = "udp"))]
use std::time::Instant;

#[cfg(feature = "clipboard")]
Expand Down Expand Up @@ -39,6 +39,7 @@ use ironrdp_pdu::input::mouse::PointerFlags;
all(windows, feature = "webauthn")
))]
use ironrdp_pdu::pdu_other_err;
use ironrdp_pdu::rdp::autodetect::AutoDetectRequest;
use ironrdp_pdu::rdp::multitransport::MultitransportResponsePdu;
use ironrdp_pdu::rdp::session_info::ServerAutoReconnect;
#[cfg(feature = "rdpdr")]
Expand Down Expand Up @@ -83,6 +84,8 @@ use ironrdp_rdpsnd_native::{RdpeaiCaptureBackend, cpal};

use crate::config::{Config, RDCleanPathConfig, Transport};
use crate::rail::{RailClient, RailControlEvent, RailEvent, RailInputEvent};
#[cfg(feature = "udp")]
use crate::udp::{disable_failed_tunnel, tunnel_auto_detect_requests, tunnel_auto_detect_sub_header};
use ironrdp_rail::pdu::{ExecutePdu, ExecuteResultPdu};

// ── Public event types ────────────────────────────────────────────────────────
Expand Down Expand Up @@ -3103,6 +3106,9 @@ async fn active_session(
let mut graceful_shutdown_sent = false;
let mut post_logon_redraw_requested = false;
let mut pending_udp_payload: Option<Vec<u8>> = None;
// Auto-detect on the tunnel is timed against its own monotonic clock.
#[cfg(feature = "udp")]
let tunnel_clock = Instant::now();
let mut initial_outputs = if *graceful_close_receiver.borrow_and_update() {
graceful_shutdown_sent = true;
Some(active_stage.graceful_shutdown()?)
Expand Down Expand Up @@ -3184,7 +3190,8 @@ async fn active_session(
Err(error) => return Err(ironrdp_session::custom_err!("read frame", error)),
};
trace!(?action, frame_length = payload.len(), "Frame received");
let mut outputs = active_stage.process(&mut image, action, &payload)?;
let mut outputs =
active_stage.process_with_timestamp(&mut image, action, &payload, reader.last_read_at())?;
#[cfg(feature = "rdpdr")]
if let Some(output) = poll_deferred_rdpdr_output(&mut active_stage)? {
outputs.push(output);
Expand Down Expand Up @@ -3246,42 +3253,84 @@ async fn active_session(
}
ActiveSessionIteration::outputs(outputs)
}
udp_payload = async {
udp_message = async {
#[cfg(feature = "udp")]
{
match (udp_tunnel.transport.as_mut(), pending_udp_payload.is_none()) {
(Some(transport), true) => transport.recv().await,
(Some(transport), true) => transport
.recv_message()
.await
.map(|message| (tunnel_auto_detect_requests(&message.sub_headers), message.data)),
(Some(_), false) | (None, _) => {
core::future::pending::<Option<Vec<u8>>>().await
core::future::pending::<Option<(Vec<AutoDetectRequest>, Vec<u8>)>>().await
}
}
}
#[cfg(not(feature = "udp"))]
{
core::future::pending::<Option<Vec<u8>>>().await
core::future::pending::<Option<(Vec<AutoDetectRequest>, Vec<u8>)>>().await
}
} => {
match udp_payload {
match udp_message {
None => {
if active_stage.reliable_udp_dvc_tunnel_in_use() {
return Ok(RdpControlFlow::TransportFailure(
ironrdp_session::general_err!("reliable UDP tunnel closed"),
));
}
#[cfg(feature = "udp")]
{
udp_tunnel.transport = None;
if let Err(error) = disable_failed_tunnel(
&mut active_stage,
&mut udp_tunnel.transport,
ironrdp_session::general_err!("reliable UDP tunnel closed"),
) {
return Ok(RdpControlFlow::TransportFailure(error));
}
}
active_stage.disable_reliable_udp_dvc_tunnel()?;
warn!("Reliable UDP tunnel closed before Soft-Sync; continuing with TCP");
ActiveSessionIteration::outputs(Vec::new())
}
Some(payload) if payload.is_empty() => {
trace!("Ignoring reliable UDP tunnel PDU without higher-layer data");
ActiveSessionIteration::outputs(Vec::new())
}
Some(payload) => {
if active_stage.reliable_udp_dvc_tunnel_in_use() {
Some((auto_detect_requests, payload)) => {
#[cfg(feature = "udp")]
{
let received_at = ironrdp_core::MonotonicInstant::from_millis(
u64::try_from(tunnel_clock.elapsed().as_millis()).unwrap_or(u64::MAX),
);
let responses = active_stage.process_tunnel_auto_detect(
auto_detect_requests,
payload.len(),
received_at,
);
let sub_headers: Vec<_> = responses.iter().filter_map(tunnel_auto_detect_sub_header).collect();
if !sub_headers.is_empty() && let Some(transport) = udp_tunnel.transport.as_ref() {
let reply = ironrdp_rdpeudp_tokio::TunnelMessage {
sub_headers,
data: Vec::new(),
};
let Some(result) =
cancelable_operation(transport.send_message(reply), close_receiver).await
else {
return Ok(RdpControlFlow::TerminatedGracefully(
GracefulDisconnectReason::UserInitiated,
));
};
if let Err(error) = result {
if let Err(error) = disable_failed_tunnel(
&mut active_stage,
&mut udp_tunnel.transport,
ironrdp_session::custom_err!("answer reliable UDP tunnel auto-detect", error),
) {
return Ok(RdpControlFlow::TransportFailure(error));
}
// This PDU arrived before any channel migrated. Once the
// tunnel fails, it cannot retain data for a future Soft-Sync.
pending_udp_payload = None;
continue;
}
Comment thread
AKolenda marked this conversation as resolved.
}
}
#[cfg(not(feature = "udp"))]
let _ = auto_detect_requests;

if payload.is_empty() {
trace!("Reliable UDP tunnel PDU without higher-layer data");
ActiveSessionIteration::outputs(Vec::new())
} else if active_stage.reliable_udp_dvc_tunnel_in_use() {
ActiveSessionIteration::tunnel(
SoftSyncTunnelType::RELIABLE_UDP,
active_stage.process_dvc_tunnel(
Expand Down
63 changes: 63 additions & 0 deletions crates/ironrdp-client/src/udp.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
//! Internal client handling for a reliable UDP tunnel.

use ironrdp_core::{Decode, Encode};
use ironrdp_pdu::rdp::autodetect::{AutoDetectRequest, AutoDetectResponse};
use ironrdp_rdpemt::{SubHeaderType, TunnelSubHeader};
use ironrdp_rdpeudp_tokio::UdpTransport;
use ironrdp_session::{ActiveStage, SessionError, SessionResult};
use tracing::{debug, warn};

/// Decodes the auto-detect requests among the sub-headers of a Tunnel Data PDU.
///
/// Each sub-header is the request structure itself: SubHeaderLength and SubHeaderType
/// are headerLength and headerTypeId ([MS-RDPEMT] 2.2.1.1.1). Unknown requests are skipped.
/// RTT requests are accepted for Windows interoperability in addition to the bandwidth
/// requests listed by MS-RDPEMT; RTT encapsulation on a tunnel is not specified there.
///
/// [MS-RDPEMT]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpemt/4f538fd7-3aca-4e7d-a213-13eb5f95c1ad
pub fn tunnel_auto_detect_requests(sub_headers: &[TunnelSubHeader]) -> Vec<AutoDetectRequest> {
sub_headers
.iter()
.filter(|sub_header| sub_header.sub_header_type == SubHeaderType::AutoDetectRequest)
.filter_map(|sub_header| match redecode(sub_header) {
Ok(request) => Some(request),
Err(error) => {
debug!(%error, data = ?sub_header.data, "Ignoring an undecodable auto-detect request on the tunnel");
None
}
})
.collect()
}

/// Encodes a response as its wire-compatible tunnel sub-header.
pub fn tunnel_auto_detect_sub_header(response: &AutoDetectResponse) -> Option<TunnelSubHeader> {
match redecode(response) {
Ok(sub_header) => Some(sub_header),
Err(error) => {
debug!(%error, ?response, "Could not encode an auto-detect response for the tunnel");
None
}
}
}

fn redecode<T: for<'de> Decode<'de>>(value: &dyn Encode) -> Result<T, String> {
let bytes = ironrdp_core::encode_vec(value).map_err(|error| error.to_string())?;
ironrdp_core::decode(&bytes).map_err(|error| error.to_string())
}

/// Drops a failed sideband and withdraws it from future Soft-Sync negotiation.
/// A channel already moved to the tunnel cannot resume over TCP without reconnecting.
/// Both receive closure and auto-detect send failures use this policy.
pub fn disable_failed_tunnel(
stage: &mut ActiveStage,
transport: &mut Option<UdpTransport>,
error: SessionError,
) -> SessionResult<()> {
if stage.reliable_udp_dvc_tunnel_in_use() {
return Err(error);
}
*transport = None;
stage.disable_reliable_udp_dvc_tunnel()?;
warn!(%error, "Reliable UDP tunnel failed before channel migration; continuing with TCP");
Ok(())
}
10 changes: 9 additions & 1 deletion crates/ironrdp-rdpeudp-tokio/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,11 @@ pub enum UdpTransportErrorKind {
/// A `send()` payload exceeds the wire `PayloadLength` field's 65535-byte
/// capacity ([MS-RDPEMT] 2.2.2.3, `RDP_TUNNEL_DATA`).
PayloadTooLarge { len: usize },

/// A `send_message()` message's sub-headers, `len` bytes encoded, do not
/// fit beside the 4-byte tunnel header in the one-byte `HeaderLength`
/// field ([MS-RDPEMT] 2.2.1.1).
SubHeadersTooLarge { len: usize },
}

impl fmt::Display for UdpTransportErrorKind {
Expand All @@ -140,6 +145,9 @@ impl fmt::Display for UdpTransportErrorKind {
"send payload of {len} bytes exceeds the 65535-byte tunnel data limit"
)
}
Self::SubHeadersTooLarge { len } => {
write!(f, "{len} bytes of sub-headers exceed the 251 a tunnel header holds")
}
}
}
}
Expand All @@ -157,7 +165,7 @@ impl core::error::Error for UdpTransportErrorKind {
| Self::TunnelTimeout
| Self::TunnelRejected { .. }
| Self::DriverPanic => None,
Self::UnsupportedProtocol { .. } | Self::PayloadTooLarge { .. } => None,
Self::UnsupportedProtocol { .. } | Self::PayloadTooLarge { .. } | Self::SubHeadersTooLarge { .. } => None,
}
}
}
Expand Down
27 changes: 14 additions & 13 deletions crates/ironrdp-rdpeudp-tokio/src/framed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,9 @@ impl FramedRead for UdpTransport {
// puts no minimum on HigherLayerData, and [MS-RDPBCGR] 1.3.9 sends
// the four Continuous Auto-Detection messages "encapsulated in the
// RDP_TUNNEL_SUBHEADER structure ... over the sideband channels
// that are in active use". The tunnel has already taken what it
// needs from those subheaders by the time we get here, leaving a
// payload of nothing to pass on.
// that are in active use". `recv` drops those subheaders (a caller
// that needs them uses `recv_message`), leaving a payload of
// nothing to pass on.
loop {
match self.recv().await {
Some(data) if data.is_empty() => {
Expand Down Expand Up @@ -74,11 +74,12 @@ mod tests {
use tokio::sync::mpsc;

use super::*;
use crate::transport::TunnelMessage;

/// Build a `UdpTransport` backed by test channels (no real network).
fn test_transport() -> (UdpTransport, mpsc::Sender<Vec<u8>>, mpsc::Receiver<Vec<u8>>) {
let (incoming_tx, incoming_rx) = mpsc::channel::<Vec<u8>>(16);
let (outgoing_tx, outgoing_rx) = mpsc::channel::<Vec<u8>>(16);
fn test_transport() -> (UdpTransport, mpsc::Sender<TunnelMessage>, mpsc::Receiver<TunnelMessage>) {
let (incoming_tx, incoming_rx) = mpsc::channel::<TunnelMessage>(16);
let (outgoing_tx, outgoing_rx) = mpsc::channel::<TunnelMessage>(16);

let transport = UdpTransport::from_channels(incoming_rx, outgoing_tx);

Expand All @@ -88,7 +89,7 @@ mod tests {
#[tokio::test]
async fn framed_read_delivers_one_message() {
let (mut transport, feeder, _) = test_transport();
feeder.send(vec![0xDE, 0xAD, 0xBE, 0xEF]).await.unwrap();
feeder.send(vec![0xDE, 0xAD, 0xBE, 0xEF].into()).await.unwrap();

let mut buf = BytesMut::new();
let n = FramedRead::read(&mut transport, &mut buf).await.unwrap();
Expand Down Expand Up @@ -120,8 +121,8 @@ mod tests {
async fn framed_read_does_not_mistake_an_empty_message_for_eof() {
let (mut transport, feeder, _) = test_transport();

feeder.send(Vec::new()).await.unwrap();
feeder.send(vec![0x11, 0x22]).await.unwrap();
feeder.send(Vec::new().into()).await.unwrap();
feeder.send(vec![0x11, 0x22].into()).await.unwrap();

let mut buf = BytesMut::new();
let n = FramedRead::read(&mut transport, &mut buf).await.unwrap();
Expand All @@ -135,7 +136,7 @@ mod tests {
async fn framed_read_still_reports_eof_after_an_empty_message() {
let (mut transport, feeder, _) = test_transport();

feeder.send(Vec::new()).await.unwrap();
feeder.send(Vec::new().into()).await.unwrap();
drop(feeder);

let mut buf = BytesMut::new();
Expand All @@ -154,7 +155,7 @@ mod tests {
.unwrap();

let data = receiver.recv().await.unwrap();
assert_eq!(data, vec![0x01, 0x02, 0x03]);
assert_eq!(data, TunnelMessage::from(vec![0x01, 0x02, 0x03]));
}

#[tokio::test]
Expand All @@ -171,8 +172,8 @@ mod tests {
async fn framed_read_multiple_messages_accumulate() {
let (mut transport, feeder, _) = test_transport();

feeder.send(vec![0xAA, 0xBB]).await.unwrap();
feeder.send(vec![0xCC, 0xDD]).await.unwrap();
feeder.send(vec![0xAA, 0xBB].into()).await.unwrap();
feeder.send(vec![0xCC, 0xDD].into()).await.unwrap();

let mut buf = BytesMut::new();

Expand Down
3 changes: 2 additions & 1 deletion crates/ironrdp-rdpeudp-tokio/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,5 +15,6 @@ pub(crate) mod tunnel;
pub use self::error::{DriverError, DriverErrorKind, UdpTransportError, UdpTransportErrorKind};
pub use self::multitransport::MultitransportBootstrap;
pub use self::transport::{
UdpAcceptConfig, UdpTlsConfig, UdpTransport, UdpTransportConfig, UdpTransportSender, accept_udp, connect_udp,
TunnelMessage, UdpAcceptConfig, UdpTlsConfig, UdpTransport, UdpTransportConfig, UdpTransportSender, accept_udp,
connect_udp,
};
Loading
Loading