diff --git a/Cargo.lock b/Cargo.lock index ea162c2dfd..25eedc4b6d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3317,6 +3317,9 @@ dependencies = [ "ironrdp-rdpdr", "ironrdp-rdpeai", "ironrdp-rdpei", + "ironrdp-rdpemt", + "ironrdp-rdpeudp", + "ironrdp-rdpeudp-tokio", "ironrdp-rdpeusb", "ironrdp-rdpsnd", "ironrdp-svc", @@ -3445,6 +3448,7 @@ dependencies = [ "ironrdp-daemon", "ironrdp-dvc", "ironrdp-dvc-pipe-proxy", + "ironrdp-egfx", "ironrdp-graphics", "ironrdp-input", "ironrdp-pdu", @@ -3458,6 +3462,7 @@ dependencies = [ "ironrdp-rdpsnd", "ironrdp-rdpsnd-native", "ironrdp-rpc", + "ironrdp-server", "ironrdp-svc", "ironrdp-tls", "ironrdp-tokio", diff --git a/crates/ironrdp-acceptor/src/connection.rs b/crates/ironrdp-acceptor/src/connection.rs index e11b47cdf9..c342625d0b 100644 --- a/crates/ironrdp-acceptor/src/connection.rs +++ b/crates/ironrdp-acceptor/src/connection.rs @@ -73,6 +73,11 @@ pub struct Acceptor { /// Source of randomness for the Initiate Multitransport Request's security /// cookie and request ID. See `set_multitransport_security_rng()`. multitransport_security_rng: Box, + /// Whether the Initiate Multitransport Response matching + /// `sent_multitransport_request` was received, and whether it reported + /// success. `None` until a matching response arrives. See + /// [`AcceptorResult::multitransport_response_success`]. + received_multitransport_response: Option, } /// Source of randomness for the security cookie and request ID the acceptor @@ -185,6 +190,11 @@ pub struct AcceptorResult { /// implement UDP multitransport can use it to decide whether to send a /// Server Initiate Multitransport Request. pub multitransport_flags: gcc::MultiTransportFlags, + /// Whether the Initiate Multitransport Response matching the sent + /// request was received during the connection sequence, and whether it + /// reported success. `None` when no request was sent, or a matching + /// response never arrived. + pub multitransport_response_success: Option, /// Credentials received from the client during SecureSettingsExchange. /// /// Present for TLS-mode connections where the client sends credentials @@ -234,6 +244,7 @@ impl Acceptor { advertised_multitransport: None, sent_multitransport_request: None, multitransport_security_rng: Box::new(OsMultitransportSecurityRng), + received_multitransport_response: None, } } @@ -386,7 +397,7 @@ impl Acceptor { /// handling for anything that isn't really a response, mirroring how /// `ClientConnectorState::ConnectTimeAutoDetection` demuxes the same /// channel client-side. - fn is_late_multitransport_response(&self, data: &mcs::SendDataRequest<'_>) -> bool { + fn is_late_multitransport_response(&mut self, data: &mcs::SendDataRequest<'_>) -> bool { let Some(sent) = self.sent_multitransport_request.as_ref() else { return false; }; @@ -411,6 +422,7 @@ impl Acceptor { success = response.is_success(), "Received Initiate Multitransport Response" ); + self.received_multitransport_response = Some(response.is_success()); } else { warn!( response.request_id, @@ -467,6 +479,7 @@ impl Acceptor { advertised_multitransport: consumed.advertised_multitransport, sent_multitransport_request: consumed.sent_multitransport_request, multitransport_security_rng: consumed.multitransport_security_rng, + received_multitransport_response: consumed.received_multitransport_response, }) } @@ -560,6 +573,7 @@ impl Acceptor { multitransport_flags: self .multitransport_flags .unwrap_or_else(gcc::MultiTransportFlags::empty), + multitransport_response_success: self.received_multitransport_response, client_early_capability_flags: self.early_capability_flags, reactivation: self.reactivation, credentials: self.received_credentials.take(), diff --git a/crates/ironrdp-dvc/src/server.rs b/crates/ironrdp-dvc/src/server.rs index ce7e40be32..d3c0b99dd2 100644 --- a/crates/ironrdp-dvc/src/server.rs +++ b/crates/ironrdp-dvc/src/server.rs @@ -311,6 +311,12 @@ impl DrdynvcServer { /// This API emits exactly one `ReliableUdp` channel list and maps every supplied /// channel to that list. A future multi-tunnel request API must establish an /// explicit response-routing mapping before it is exposed. + /// + /// A connection gets one request: The state never returns to idle, so a + /// later call returns an error. The request declares the TCP path flushed for + /// the supplied channels (SOFT_SYNC_TCP_FLUSHED, [MS-RDPEDYC] 2.2.5.1), so the + /// server sends their data over the tunnel from then on ([MS-RDPEDYC] + /// 3.3.5.3.1), whatever the client's response lists. pub fn request_reliable_udp(&mut self, channel_ids: Vec) -> PduResult { if channel_ids.is_empty() { return Err(pdu_other_err!("soft-sync requires at least one dynamic channel")); @@ -352,6 +358,19 @@ impl DrdynvcServer { self.outgoing_tunnel_channels.get(&channel_id).copied() } + /// Returns whether a Soft-Sync request was sent and its response has not + /// arrived yet, the window in which the server must not read tunnel data + /// (MS-RDPEDYC 3.3.5.3.2). + pub const fn soft_sync_awaiting_response(&self) -> bool { + matches!( + self.soft_sync_state, + SoftSyncState::Active { + response_received: false, + .. + } + ) + } + /// Returns whether the client has acknowledged the Soft-Sync request over TCP. pub const fn soft_sync_response_received(&self) -> bool { matches!( @@ -408,12 +427,24 @@ impl DrdynvcServer { return Err(pdu_other_err!("soft-sync response selected an unrequested tunnel")); } } - self.incoming_tunnel_channels = self + let accepted_channels: BTreeMap = self .outgoing_tunnel_channels .iter() .filter(|(_, tunnel_type)| response.tunnels_to_switch().contains(tunnel_type)) .map(|(channel_id, tunnel_type)| (*channel_id, *tunnel_type)) .collect(); + // The response names the tunnels the client will write on + // (MS-RDPEDYC 2.2.5.2), so it sets only what the server reads from a + // tunnel. The server's own sending was fixed by the request: after + // sending it, the server MUST keep using the tunnel it named for + // those channels (3.3.5.3.1), and the request told the client that no + // more of their data comes over TCP. + debug!( + tunnels = ?response.tunnels_to_switch(), + channels = ?accepted_channels.keys().collect::>(), + "Soft-Sync response received" + ); + self.incoming_tunnel_channels = accepted_channels; *response_received = true; Ok(()) } @@ -562,6 +593,39 @@ mod tests { assert!(server.request_reliable_udp(alloc::vec![channel_id]).is_err()); } + #[test] + fn soft_sync_response_does_not_change_the_outgoing_tunnel() { + let mut server = DrdynvcServer::new(); + let channel_id = server.dynamic_channels.insert_channel(TestDvc, ChannelState::Opened); + + server.request_reliable_udp(alloc::vec![channel_id]).unwrap(); + assert_eq!( + server.tunnel_for_outgoing_channel(channel_id), + Some(SoftSyncTunnelType::RELIABLE_UDP) + ); + + // An empty list means the client keeps writing over TCP. It does not + // take back the server's own switch, which the request announced. + server + .process_soft_sync_response(crate::pdu::SoftSyncResponsePdu::new(alloc::vec![])) + .unwrap(); + + assert!(server.soft_sync_response_received()); + assert_eq!( + server.tunnel_for_outgoing_channel(channel_id), + Some(SoftSyncTunnelType::RELIABLE_UDP), + "the server keeps sending on the tunnel its request named" + ); + let tunnel_data = ironrdp_core::encode_vec(&DrdynvcClientPdu::Data(crate::pdu::DrdynvcDataPdu::Data( + crate::pdu::DataPdu::new(channel_id, Vec::new()), + ))) + .unwrap(); + assert!( + server.process_tunnel(&tunnel_data).is_err(), + "the client did not switch its own writing, so tunnel data for the channel is unexpected" + ); + } + #[test] fn soft_sync_accepts_a_response_after_the_selected_channel_closes() { let mut server = DrdynvcServer::new(); diff --git a/crates/ironrdp-rdpeudp-tokio/src/lib.rs b/crates/ironrdp-rdpeudp-tokio/src/lib.rs index b22487412c..517975ccdd 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/lib.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/lib.rs @@ -14,4 +14,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, accept_udp, connect_udp}; +pub use self::transport::{ + UdpAcceptConfig, UdpTlsConfig, UdpTransport, UdpTransportConfig, UdpTransportSender, accept_udp, connect_udp, +}; diff --git a/crates/ironrdp-rdpeudp-tokio/src/transport.rs b/crates/ironrdp-rdpeudp-tokio/src/transport.rs index 7d06e8c20b..d5b9c79605 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/transport.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/transport.rs @@ -211,6 +211,38 @@ impl Drop for AbortOnDrop { } } +/// A cloneable handle for sending data over an established UDP transport, +/// obtained from [`UdpTransport::sender`]. +/// +/// Independent of [`UdpTransport::recv`]'s `&mut self` requirement: Sending +/// and receiving already run over separate channels fed by separate +/// background tasks, so this never contends with a concurrent `recv()`. +#[derive(Clone)] +pub struct UdpTransportSender(mpsc::Sender>); + +impl UdpTransportSender { + /// Send a higher-layer data frame through the tunnel. + /// + /// Identical validation and error semantics to [`UdpTransport::send`], + /// which delegates to this. + /// + /// # Errors + /// + /// Returns `PayloadTooLarge` if `data` exceeds 65535 bytes, the wire + /// `PayloadLength` field's capacity ([MS-RDPEMT] 2.2.2.3). + pub async fn send(&self, data: Vec) -> Result<(), UdpTransportError> { + if data.len() > usize::from(u16::MAX) { + debug!(len = data.len(), "Rejected oversized tunnel payload"); + return Err(UdpTransportError::payload_too_large("send", data.len())); + } + + self.0 + .send(data) + .await + .map_err(|_| UdpTransportError::driver("send", DriverError::connection_closed("send"))) + } +} + /// Handle to an established UDP transport. /// /// Provides bidirectional higher-layer data (DVC frames) over the @@ -261,15 +293,21 @@ impl UdpTransport { /// discover: that task has no way to report a per-payload failure back /// to a caller who already received `Ok(())` from a channel send. pub async fn send(&self, data: Vec) -> Result<(), UdpTransportError> { - if data.len() > usize::from(u16::MAX) { - debug!(len = data.len(), "Rejected oversized tunnel payload"); - return Err(UdpTransportError::payload_too_large("send", data.len())); - } + self.sender().send(data).await + } - self.data_tx - .send(data) - .await - .map_err(|_| UdpTransportError::driver("send", DriverError::connection_closed("send"))) + /// Returns a cloneable handle for sending data, independent of this + /// object's `&mut self`-requiring [`Self::recv`]. + /// + /// Sending and receiving are already independent internally (separate + /// channels fed by separate background tasks), so a caller that shares + /// one `UdpTransport` between a single dedicated receiver (behind a lock + /// reserved for `recv()` alone, since only one caller should ever call + /// it) and one or more senders can send through this handle without + /// contending on that lock, including for the full duration of an idle + /// `recv()` wait. + pub fn sender(&self) -> UdpTransportSender { + UdpTransportSender(self.data_tx.clone()) } /// Shut down the transport, closing the RDPEUDP2 connection. diff --git a/crates/ironrdp-server/Cargo.toml b/crates/ironrdp-server/Cargo.toml index b35b97e14b..ae73ff2379 100644 --- a/crates/ironrdp-server/Cargo.toml +++ b/crates/ironrdp-server/Cargo.toml @@ -60,6 +60,9 @@ ironrdp-rdpei = { path = "../ironrdp-rdpei", version = "0.1" } # public ironrdp-rdpeai = { path = "../ironrdp-rdpeai", version = "0.1" } # public ironrdp-rdpeusb = { path = "../ironrdp-rdpeusb", version = "0.1", optional = true } ironrdp-usb = { path = "../ironrdp-usb", version = "0.1", optional = true } +ironrdp-rdpemt = { path = "../ironrdp-rdpemt", version = "0.1" } # public +ironrdp-rdpeudp = { path = "../ironrdp-rdpeudp", version = "0.1" } # public +ironrdp-rdpeudp-tokio = { path = "../ironrdp-rdpeudp-tokio", version = "0.1" } # public tracing = { version = "0.1", features = ["log"] } x509-cert = { version = "0.3", optional = true } rustls-pemfile = { version = "2.2", optional = true } diff --git a/crates/ironrdp-server/src/builder.rs b/crates/ironrdp-server/src/builder.rs index c61fc3c9d8..a511adb8dd 100644 --- a/crates/ironrdp-server/src/builder.rs +++ b/crates/ironrdp-server/src/builder.rs @@ -66,6 +66,7 @@ pub struct BuilderDone { connection_policy: ConnectionPolicy, remotefx_quant: Quant, remotefx_entropy_coder: Option, + udp_bind_addr: Option, } pub struct RdpServerBuilder { @@ -180,6 +181,7 @@ impl RdpServerBuilder { auto_reconnect_cookie: None, remotefx_quant: Quant::default(), remotefx_entropy_coder: None, + udp_bind_addr: None, }, } } @@ -215,6 +217,7 @@ impl RdpServerBuilder { auto_reconnect_cookie: None, remotefx_quant: Quant::default(), remotefx_entropy_coder: None, + udp_bind_addr: None, }, } } @@ -479,6 +482,38 @@ impl RdpServerBuilder { self } + /// Offer UDP multitransport (MS-RDPBCGR 2.2.1.4.6/2.2.15.1) to clients + /// that support it, binding a fresh UDP socket to `udp_bind_addr` per + /// connection to accept the sideband RDPEUDP2 + TLS + RDPEMT transport. + /// + /// Requires [`RdpServerSecurity::Tls`] or [`RdpServerSecurity::Hybrid`] + /// (multitransport is Enhanced-Security-only, matching the reference + /// client); ignored under [`RdpServerSecurity::None`]. + /// + /// `udp_bind_addr` is a separate, explicit address rather than reusing + /// [`Self`]'s own TCP `addr`: a caller driving [`RdpServer::run_connection`] + /// with its own accept loop (rather than [`RdpServer::run`]) may not have + /// `addr` bound to anything real, so it cannot be inferred. Typically the + /// same host and port as the TCP listener (UDP and TCP occupy independent + /// port spaces at the same number). When its IP is unspecified, each + /// connection's socket binds to the local address that client reached + /// instead, so replies leave from the address the client sent to; see + /// [`RdpServer::set_connection_local_addr`]. + /// + /// Once established, the transport is used to migrate EGFX graphics + /// traffic off TCP; a failure to establish it at any stage falls back to + /// TCP-only rather than failing the connection. This includes + /// [`Self::with_preempt_existing_session`] overlapping a candidate + /// session's own UDP bind with a still-live session's: When + /// `udp_bind_addr` is the same for both, the second bind fails and that + /// connection degrades to TCP-only. + /// + /// `None` (the default): no UDP socket is ever bound, no behavior change. + pub fn with_udp_transport(mut self, udp_bind_addr: SocketAddr) -> Self { + self.state.udp_bind_addr = Some(udp_bind_addr); + self + } + pub fn build(self) -> RdpServer { let mut server = RdpServer::new( RdpServerOptions { @@ -490,6 +525,7 @@ impl RdpServerBuilder { connection_policy: self.state.connection_policy, remotefx_quant: self.state.remotefx_quant, remotefx_entropy_coder: self.state.remotefx_entropy_coder, + udp_bind_addr: self.state.udp_bind_addr, }, self.state.handler, self.state.display, diff --git a/crates/ironrdp-server/src/gfx.rs b/crates/ironrdp-server/src/gfx.rs index 5663768d4e..2347edd595 100644 --- a/crates/ironrdp-server/src/gfx.rs +++ b/crates/ironrdp-server/src/gfx.rs @@ -106,3 +106,39 @@ impl core::fmt::Display for EgfxServerMessage { } } } + +/// The dynamic channel id EGFX is registered under: the bridge when the +/// factory handed out a frame handle, otherwise the server itself. +pub(crate) fn egfx_channel_id(drdynvc: &ironrdp_dvc::DrdynvcServer) -> Option { + drdynvc + .get_channel_id_by_type::() + .or_else(|| drdynvc.get_channel_id_by_type::()) +} + +#[cfg(test)] +mod tests { + use ironrdp_dvc::DrdynvcServer; + use ironrdp_egfx::pdu::{CapabilitiesAdvertisePdu, CapabilitySet}; + + use super::*; + + struct Handler; + + impl GraphicsPipelineHandler for Handler { + fn capabilities_advertise(&mut self, _pdu: &CapabilitiesAdvertisePdu) {} + + fn on_ready(&mut self, _negotiated: &CapabilitySet) {} + } + + #[test] + fn egfx_is_found_whichever_way_it_was_registered() { + let server = Arc::new(Mutex::new(GraphicsPipelineServer::new(Box::new(Handler)))); + let bridged = DrdynvcServer::new().with_dynamic_channel(GfxDvcBridge::new(server)); + assert!(egfx_channel_id(&bridged).is_some()); + + let direct = DrdynvcServer::new().with_dynamic_channel(GraphicsPipelineServer::new(Box::new(Handler))); + assert!(egfx_channel_id(&direct).is_some()); + + assert!(egfx_channel_id(&DrdynvcServer::new()).is_none()); + } +} diff --git a/crates/ironrdp-server/src/lib.rs b/crates/ironrdp-server/src/lib.rs index 94544265cf..b1f9554080 100644 --- a/crates/ironrdp-server/src/lib.rs +++ b/crates/ironrdp-server/src/lib.rs @@ -18,6 +18,7 @@ mod handler; pub mod heartbeat; #[cfg(feature = "helper")] mod helper; +mod multitransport; mod rdpdr; mod rdpeai; mod rdpei; diff --git a/crates/ironrdp-server/src/multitransport.rs b/crates/ironrdp-server/src/multitransport.rs new file mode 100644 index 0000000000..78a03e9d17 --- /dev/null +++ b/crates/ironrdp-server/src/multitransport.rs @@ -0,0 +1,194 @@ +//! Server-side UDP multitransport bootstrapping. +//! +//! [MS-RDPBCGR] 2.2.15.1/2.2.15.2 (`ironrdp_acceptor::Acceptor`) sends the +//! Initiate Multitransport Request and reports it via +//! [`ironrdp_acceptor::accept_finalize_with_multitransport`]. This module +//! establishes the sideband RDPEUDP2 + TLS + RDPEMT transport that request +//! bootstraps, for [`crate::server`] to migrate DVC channels onto via +//! [`ironrdp_dvc::DrdynvcServer::request_reliable_udp`]. +//! +//! [MS-RDPBCGR]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpbcgr/b8e7c588-51cb-455b-bb73-92d480903133 + +use core::net::{SocketAddr, SocketAddrV6}; +use core::time::Duration; +use std::sync::Arc; + +use ironrdp_pdu::rdp::multitransport::MultitransportRequestPdu; +use ironrdp_rdpemt::TunnelConfig; +use ironrdp_rdpeudp::ConnectionConfig; +use ironrdp_rdpeudp_tokio::{UdpAcceptConfig, UdpTransport, UdpTransportSender, accept_udp}; +use tokio::net::UdpSocket; +use tokio::sync::Mutex; +use tokio_rustls::rustls; +use tracing::{debug, warn}; + +/// How long the full UDP accept sequence (initial datagram, RDPEUDP2 +/// handshake, TLS, RDPEMT tunnel creation) may take before giving up and +/// continuing TCP-only. +const UDP_ACCEPT_TIMEOUT: Duration = Duration::from_secs(15); + +/// Shared handle to an established sideband UDP transport. +/// +/// `send()` uses `UdpTransportSender`, cloned once at construction from the +/// underlying `UdpTransport`. That handle is entirely independent of the +/// `Mutex` below, so it never contends with (or blocks behind) `recv()`, +/// which holds that lock for the full duration of its idle wait for the +/// next datagram: Sending and receiving already run over separate channels +/// fed by separate background tasks, so the two were never meant to share +/// one lock. `recv()` is intended for a single dedicated caller (the +/// `client_loop` select arm). +#[derive(Clone)] +pub(crate) struct UdpTransportHandle { + sender: UdpTransportSender, + transport: Arc>, +} + +impl UdpTransportHandle { + fn new(transport: UdpTransport) -> Self { + let sender = transport.sender(); + Self { + sender, + transport: Arc::new(Mutex::new(transport)), + } + } + + /// Sends one higher-layer (raw DVC) payload over the tunnel. + /// + /// A failure is logged here rather than propagated: The RDP session never + /// fails over an optional sideband transport. + pub(crate) async fn send(&self, data: Vec) { + if let Err(error) = self.sender.send(data).await { + warn!(%error, "Failed to send data over UDP transport"); + } + } + + pub(crate) async fn recv(&self) -> Option> { + self.transport.lock().await.recv().await + } +} + +/// The address to bind the sideband UDP socket to: `bind` itself, unless its +/// IP is unspecified and the local address the client reached over TCP is +/// known, in which case that address (at `bind`'s port). Replies must leave +/// from the address the client sent to, which a socket bound to the +/// unspecified address does not guarantee on a host with several addresses. +/// +/// A link-local IPv6 address keeps its scope ID, without which it cannot be +/// bound. +pub(crate) fn sideband_bind_addr(bind: SocketAddr, connection_local: Option) -> SocketAddr { + let Some(local) = connection_local.filter(|_| bind.ip().is_unspecified()) else { + return bind; + }; + match local { + SocketAddr::V6(v6) if v6.ip().to_ipv4_mapped().is_none() => { + SocketAddr::V6(SocketAddrV6::new(*v6.ip(), bind.port(), 0, v6.scope_id())) + } + _ => SocketAddr::new(local.ip().to_canonical(), bind.port()), + } +} + +/// Attempts to establish the sideband UDP transport for one Initiate +/// Multitransport Request. +/// +/// Binds a fresh socket to `udp_bind_addr` for this attempt, avoiding any +/// shared, long-lived UDP socket state to manage. This assumes at most one +/// live attempt at a time; under +/// [`RdpServerOptions::preempt_existing_session`](crate::server::RdpServerOptions::preempt_existing_session) +/// a candidate session's negotiation can overlap the still-live session it +/// is preempting, and `udp_bind_addr` is typically the same host and port for +/// both (see [`RdpServerBuilder::with_udp_transport`](crate::RdpServerBuilder::with_udp_transport)), +/// so the second bind fails with `AddrInUse`. That failure is handled below +/// like any other: The overlapping connection degrades to TCP-only rather +/// than failing. +/// +/// Any failure (bind, RDPEUDP2 handshake, TLS, RDPEMT tunnel) is logged and +/// reported as `None` rather than propagated: An optional sideband transport +/// failing to come up is never a reason to fail the RDP connection, matching +/// the reference client implementation's posture. +pub(crate) async fn accept( + udp_bind_addr: SocketAddr, + tls_config: Arc, + request: &MultitransportRequestPdu, +) -> Option { + let socket = match UdpSocket::bind(udp_bind_addr).await { + Ok(socket) => socket, + Err(error) => { + warn!(%error, %udp_bind_addr, "Failed to bind UDP socket for multitransport, continuing TCP-only"); + return None; + } + }; + + let config = UdpAcceptConfig { + tls_config, + tunnel_config: TunnelConfig { + request_id: request.request_id, + security_cookie: request.security_cookie, + }, + connection_config: ConnectionConfig::default(), + accept_timeout: UDP_ACCEPT_TIMEOUT, + }; + + match accept_udp(socket, config).await { + Ok(transport) => { + debug!("Sideband UDP transport established"); + Some(UdpTransportHandle::new(transport)) + } + Err(error) => { + warn!(%error, "Failed to establish sideband UDP transport, continuing TCP-only"); + None + } + } +} + +#[cfg(test)] +mod tests { + use core::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + + use super::*; + + const V6_LOCAL: Ipv6Addr = Ipv6Addr::new(0x2001, 0xdb8, 0, 0, 0, 0, 0, 0xc518); + + #[test] + fn unspecified_bind_takes_the_connection_address() { + let bind = SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 3389); + let local = SocketAddr::new(IpAddr::V6(V6_LOCAL), 3389); + assert_eq!( + sideband_bind_addr(bind, Some(local)), + SocketAddr::new(IpAddr::V6(V6_LOCAL), 3389) + ); + } + + #[test] + fn ipv4_client_on_a_dual_stack_listener_binds_ipv4() { + let bind = SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 3389); + let mapped = SocketAddr::new(IpAddr::V6(Ipv4Addr::new(192, 0, 2, 7).to_ipv6_mapped()), 3389); + assert_eq!( + sideband_bind_addr(bind, Some(mapped)), + SocketAddr::new(IpAddr::V4(Ipv4Addr::new(192, 0, 2, 7)), 3389) + ); + } + + #[test] + fn a_link_local_address_keeps_its_scope() { + let bind = SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 3389); + let link_local = Ipv6Addr::new(0xfe80, 0, 0, 0, 0, 0, 0, 1); + let local = SocketAddr::V6(SocketAddrV6::new(link_local, 50000, 0, 2)); + assert_eq!( + sideband_bind_addr(bind, Some(local)), + SocketAddr::V6(SocketAddrV6::new(link_local, 3389, 0, 2)) + ); + } + + #[test] + fn an_explicit_bind_address_is_kept() { + let bind = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(192, 0, 2, 1)), 3390); + let local = SocketAddr::new(IpAddr::V6(V6_LOCAL), 3389); + assert_eq!(sideband_bind_addr(bind, Some(local)), bind); + } + + #[test] + fn unknown_connection_address_keeps_the_bind() { + let bind = SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 3389); + assert_eq!(sideband_bind_addr(bind, None), bind); + } +} diff --git a/crates/ironrdp-server/src/server.rs b/crates/ironrdp-server/src/server.rs index 551e5989a5..434cf140df 100644 --- a/crates/ironrdp-server/src/server.rs +++ b/crates/ironrdp-server/src/server.rs @@ -1,9 +1,11 @@ +use core::cell::RefCell; use core::fmt; use core::net::{IpAddr, SocketAddr}; use core::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering}; use core::time::Duration; #[cfg(feature = "usb")] use std::collections::HashMap; +use std::collections::VecDeque; use std::rc::Rc; use std::sync::Arc; use std::sync::LazyLock; @@ -60,6 +62,7 @@ use crate::error::{ServerError, ServerErrorExt as _, ServerErrorKind, ServerResu use crate::gfx::{EgfxServerMessage, GfxServerFactory}; use crate::handler::RdpServerInputHandler; use crate::heartbeat::HeartbeatConfig; +use crate::multitransport; use crate::rdpeai::RdpeaiServerFactory; use crate::rdpei::RdpeiServerFactory; #[cfg(feature = "usb")] @@ -386,6 +389,10 @@ pub enum ConnectionPolicy { Preempt, } +/// Tunnel payloads held while a Soft-Sync response is pending; beyond this the +/// client is sending far more than the handful of messages the race allows. +const MAX_EARLY_TUNNEL_PAYLOADS: usize = 64; + #[derive(Clone)] #[non_exhaustive] pub struct RdpServerOptions { @@ -422,6 +429,10 @@ pub struct RdpServerOptions { /// /// [MS-RDPRFX]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdprfx/ pub remotefx_entropy_coder: Option, + /// UDP bind address for the sideband multitransport socket, when UDP + /// multitransport is offered. `None` disables the feature entirely. Set + /// via [`RdpServerBuilder::with_udp_transport`](crate::RdpServerBuilder::with_udp_transport). + pub udp_bind_addr: Option, } impl RdpServerOptions { @@ -512,6 +523,20 @@ impl RdpServerSecurity { RdpServerSecurity::Hybrid(_) => nego::SecurityProtocol::HYBRID | nego::SecurityProtocol::HYBRID_EX, } } + + /// The [`TlsAcceptor`] backing this security mode, if any. + /// + /// UDP multitransport reuses this connection's own TLS certificate for + /// the sideband transport (`TlsAcceptor::config()`), rather than needing + /// separate cert/key plumbing. `None` under [`RdpServerSecurity::None`], + /// which has no certificate to reuse and cannot offer multitransport + /// anyway (Enhanced Security only). + fn tls_acceptor(&self) -> Option<&TlsAcceptor> { + match self { + RdpServerSecurity::None => None, + RdpServerSecurity::Tls(acceptor) | RdpServerSecurity::Hybrid((acceptor, _)) => Some(acceptor), + } + } } struct AInputHandler { @@ -704,6 +729,9 @@ pub struct RdpServer { creds: Option, credential_validator: Option>, local_addr: Option, + /// The local address the current client reached. See + /// [`Self::set_connection_local_addr`]. + connection_local_addr: Option, autodetect: Option, heartbeat: Option, connection_handler: Option>, @@ -792,6 +820,34 @@ pub struct RdpServer { /// Tracks whether the current cookie has reached a client. Subsequent /// connections and hourly updates replace it with a new random. auto_reconnect_sent: bool, + + /// Abort handle of the current connection's pending UDP multitransport + /// accept, if one is running. A client that could not establish the + /// sideband transport answers with a failure Initiate Multitransport + /// Response, often after finalization has completed; the message-channel + /// handler uses this to stop the accept instead of letting it hold its + /// socket until `multitransport::UDP_ACCEPT_TIMEOUT`. + pending_udp_accept_abort: Option, + /// Whether the current connection negotiated SOFT_SYNC_TCP_TO_UDP, one of + /// the two conditions for migrating a channel onto the sideband + /// transport (see [`Self::udp_migration_allowed`]). + soft_sync_negotiated: bool, + /// Whether the current connection may migrate EGFX onto the sideband + /// transport. MS-RDPEDYC 3.1.5.3/3.3.5.3.1: Soft-Sync MUST NOT be used + /// unless both peers negotiated SOFT_SYNC_TCP_TO_UDP and a successful + /// Initiate Multitransport Response was received. Finalization sets it + /// when the response came during finalization; the message-channel + /// handler sets it when the response comes later, which is the usual + /// case with mstsc, whose UDP bootstrap outlasts the TCP finalization. + udp_migration_allowed: bool, + /// Whether the current connection's EGFX data has started going over the + /// sideband transport, so the switch is logged once. + egfx_on_udp: bool, + /// Tunnel payloads that arrived while a Soft-Sync request was waiting for + /// its response. The client writes on the tunnel right after sending the + /// response over TCP, so its first tunnel data can overtake it; this holds + /// that data until the response is in instead of dropping it. + early_tunnel_payloads: VecDeque>, } /// Cloneable handle for updating the Server Auto-Reconnect Cookie while @@ -1002,9 +1058,24 @@ impl PendingConnection { capabilities: Vec, creds: Option, honor_client_desktop_size: Option, + udp_bind_addr: Option, ) -> Self { let mut acceptor = Acceptor::new(security.flag(), desktop_size, capabilities, creds); acceptor.set_honor_client_desktop_size(honor_client_desktop_size); + // Multitransport requires Enhanced Security: `RdpServerSecurity::None` + // has nothing to authenticate the sideband transport's TLS with, and + // matches the reference client's own Enhanced-Security-only gate. + if udp_bind_addr.is_some() && security.tls_acceptor().is_some() { + // SOFT_SYNC_TCP_TO_UDP is included alongside the transport type + // itself: Without it, Soft-Sync is never negotiated (MS-RDPEDYC + // 3.1.5.3 requires both peers to advertise it), and + // `dispatch_egfx_messages` gates the EGFX-over-tunnel migration + // on `multitransport_flags` reflecting it. + acceptor.set_multitransport_offer(Some( + ironrdp_pdu::gcc::MultiTransportFlags::TRANSPORT_TYPE_UDP_FECR + | ironrdp_pdu::gcc::MultiTransportFlags::SOFT_SYNC_TCP_TO_UDP, + )); + } Self { security, acceptor } } @@ -1071,6 +1142,7 @@ impl PendingConnection { Ok(Some(NegotiatedConnection { transport: NegotiatedTransport::Tls(Box::new(framed)), acceptor, + local_addr: None, })) } // The stream is already past TLS (terminated at a lower @@ -1082,6 +1154,7 @@ impl PendingConnection { Ok(Some(NegotiatedConnection { transport: NegotiatedTransport::Offloaded(framed), acceptor, + local_addr: None, })) } }, @@ -1089,6 +1162,7 @@ impl PendingConnection { BeginResult::Continue(framed) => Ok(Some(NegotiatedConnection { transport: NegotiatedTransport::Continued(framed), acceptor, + local_addr: None, })), } } @@ -1101,6 +1175,9 @@ impl PendingConnection { struct NegotiatedConnection { transport: NegotiatedTransport, acceptor: Acceptor, + /// The local address the client reached, when the stream is a socket + /// the server accepted itself; see [`RdpServer::set_connection_local_addr`]. + local_addr: Option, } /// Advance a stream that is now past the security upgrade: mark the acceptor @@ -1317,6 +1394,7 @@ async fn negotiate_candidate( capabilities, ctx.creds.clone(), ctx.opts.honor_client_desktop_size, + ctx.opts.udp_bind_addr, ); // NOTE: deliberately NO channel attachment here. Building the cliprdr / @@ -1329,9 +1407,11 @@ async fn negotiate_candidate( // the static channel set until it processes the MCS Connect Initial, which // happens there, not in `accept_begin`. + let local_addr = stream.local_addr().ok(); match pending.negotiate_and_authenticate(stream, TransportTls::Managed).await { - Ok(Some(negotiated)) => { + Ok(Some(mut negotiated)) => { debug!(?peer, "candidate authenticated -- eligible to preempt the live session"); + negotiated.local_addr = local_addr; Some((Box::new(negotiated), peer)) } Ok(None) => { @@ -1451,6 +1531,7 @@ impl RdpServer { creds: None, credential_validator: None, local_addr: None, + connection_local_addr: None, autodetect: None, heartbeat: None, connection_handler, @@ -1477,6 +1558,11 @@ impl RdpServer { auto_reconnect_cookie: None, previous_auto_reconnect_cookie: None, auto_reconnect_sent: false, + pending_udp_accept_abort: None, + soft_sync_negotiated: false, + udp_migration_allowed: false, + egfx_on_udp: false, + early_tunnel_payloads: VecDeque::new(), } } @@ -2151,6 +2237,7 @@ impl RdpServer { // connections through this method with the previous session's backends // still live until the next client attached new ones. self.static_channels = StaticChannelSet::new(); + self.connection_local_addr = None; result } @@ -2179,6 +2266,7 @@ impl RdpServer { capabilities, self.creds.clone(), self.opts.honor_client_desktop_size, + self.opts.udp_bind_addr, ); self.attach_channels(pending.acceptor_mut(), monitor_count); @@ -2203,7 +2291,14 @@ impl RdpServer { where S: AsyncRead + AsyncWrite + Sync + Send + Unpin, { - let NegotiatedConnection { transport, acceptor } = negotiated; + let NegotiatedConnection { + transport, + acceptor, + local_addr, + } = negotiated; + if local_addr.is_some() { + self.connection_local_addr = local_addr; + } match transport { // No security upgrade happened, so there is no TLS session to shut // down — matches the pre-existing `BeginResult::Continue` arm. @@ -2250,6 +2345,8 @@ impl RdpServer { Ok(()) } + /// Bind the configured address and serve RDP connections until + /// [`ServerEvent::Quit`] is received or the event channel closes. pub async fn run(&mut self) -> ServerResult<()> { // Create socket with control over options before binding. // Using TcpSocket instead of TcpListener::bind() allows setting @@ -2379,6 +2476,9 @@ impl RdpServer { let peer = match &entry { Entry::Fresh(_, peer) | Entry::Negotiated(_, peer) => *peer, }; + if let Entry::Fresh(stream, _) = &entry { + self.connection_local_addr = stream.local_addr().ok(); + } debug!(?peer, "Received connection"); // A `Negotiated` winner already passed `on_accept` as a candidate, @@ -2639,6 +2739,7 @@ impl RdpServer { // capture, held open until the next client) for every preemption // takeover. self.static_channels = StaticChannelSet::new(); + self.connection_local_addr = None; if let Some(ref mut handler) = self.connection_handler { let action = handler.on_disconnected(peer, duration, result.as_ref().err()); @@ -2662,6 +2763,10 @@ impl RdpServer { self.static_channels.get_channel_id_by_type::() } + #[expect( + clippy::too_many_arguments, + reason = "private per-connection dispatch; the parameters are the connection's negotiated identifiers and transports" + )] async fn dispatch_pdu( &mut self, action: Action, @@ -2670,6 +2775,7 @@ impl RdpServer { io_channel_id: u16, user_channel_id: u16, message_channel_id: Option, + udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult { match action { Action::FastPath => { @@ -2679,7 +2785,14 @@ impl RdpServer { Action::X224 => { if self - .handle_x224(writer, io_channel_id, user_channel_id, message_channel_id, &bytes) + .handle_x224( + writer, + io_channel_id, + user_channel_id, + message_channel_id, + &bytes, + udp_transport, + ) .await? { debug!("Got disconnect request"); @@ -2735,7 +2848,14 @@ impl RdpServer { io_channel_id: u16, user_channel_id: u16, message_channel_id: Option, + udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult { + // Only referenced under `#[cfg(feature = "egfx")]` below (the only + // channel this integration migrates); this keeps the parameter + // itself unconditional so callers don't need their own `egfx` gate. + #[cfg(not(feature = "egfx"))] + let _ = &udp_transport; + // Avoid wave messages queuing up and causing extra delay. When a // batch carries more than `WAVE_KEEP` waves, drop the OLDEST ones // and keep the most recent — playing stale audio just bakes the @@ -3284,15 +3404,8 @@ impl RdpServer { #[cfg(feature = "egfx")] ServerEvent::Egfx(msg) => match msg { EgfxServerMessage::SendMessages { messages } => { - let drdynvc_channel_id = self - .get_channel_id_by_type::() - .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; - let data = server_encode_svc_messages(messages, drdynvc_channel_id, user_channel_id) - .map_err(ServerError::encode)?; - writer - .write_all(&data) - .await - .map_err(|e| ServerError::io("write_all", e))?; + self.dispatch_egfx_messages(messages, writer, user_channel_id, udp_transport) + .await?; } }, ServerEvent::AutoDetectRttRequest => { @@ -3341,6 +3454,261 @@ impl RdpServer { Ok(RunState::Continue) } + /// Routes EGFX graphics data over the sideband UDP transport when one is + /// up and the channel has migrated, otherwise the normal TCP path. + /// + /// Opportunistically requests the migration (Soft-Sync, + /// [`dvc::DrdynvcServer::request_reliable_udp`]) the first time EGFX has + /// data to send after a transport exists: MS-RDPBCGR sequencing does not + /// give an earlier point where EGFX's dynamic channel id is known to be + /// open (dynamic channels open lazily, driven by the client, well after + /// the transport could have already come up during multitransport + /// bootstrapping). The request is made once per connection: A later attempt + /// gets an "already requested" error from + /// [`dvc::DrdynvcServer::request_reliable_udp`], which is expected and only + /// traced. The request declares the TCP path flushed for the channel + /// (SOFT_SYNC_TCP_FLUSHED, MS-RDPEDYC 2.2.5.1), so its data goes over the + /// tunnel from then on (3.3.5.3.1), whatever the client's response lists. + #[cfg(feature = "egfx")] + async fn dispatch_egfx_messages( + &mut self, + messages: Vec, + writer: &mut impl FramedWrite, + user_channel_id: u16, + udp_transport: Option<&multitransport::UdpTransportHandle>, + ) -> ServerResult<()> { + let drdynvc_channel_id = self + .get_channel_id_by_type::() + .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; + + let mut route_over_udp = false; + + // `self.udp_migration_allowed` gates the Soft-Sync Request itself + // (`request_reliable_udp` below): MS-RDPEDYC 3.1.5.3/3.3.5.3.1 forbid + // it unless both peers negotiated SOFT_SYNC_TCP_TO_UDP and a + // successful Initiate Multitransport Response was actually received, + // neither of which the sideband transport's own handshake succeeding + // (`udp_transport.is_some()`) establishes on its own. + let mut newly_on_udp = None; + if self.udp_migration_allowed + && let Some(udp_transport) = udp_transport + && let Some(drdynvc) = self.get_svc_processor::() + { + let Some(egfx_dvc_id) = crate::gfx::egfx_channel_id(drdynvc) else { + trace!("EGFX channel not open yet, staying on TCP"); + return self + .write_egfx_over_tcp(messages, writer, drdynvc_channel_id, user_channel_id) + .await; + }; + + // First time EGFX has data to send after the tunnel exists: ask + // to migrate it. "Already requested" (a previous batch already + // asked) and "channel not open" (raced with the client closing + // it) are both expected steady-state outcomes here, not errors. + match drdynvc.request_reliable_udp(vec![egfx_dvc_id]) { + Ok(request) => { + let data = server_encode_svc_messages(vec![request], drdynvc_channel_id, user_channel_id) + .map_err(ServerError::encode)?; + writer + .write_all(&data) + .await + .map_err(|e| ServerError::io("write_all", e))?; + debug!( + egfx_dvc_id, + "Soft-Sync request sent, asking to move EGFX to reliable UDP" + ); + } + Err(error) => trace!(%error, "No Soft-Sync request sent"), + } + + // From the request on, not from the client's response: The request + // carries SOFT_SYNC_TCP_FLUSHED, "no more data will be sent over + // TCP for the specified DVCs" (MS-RDPEDYC 2.2.5.1), and the server + // MUST keep using the tunnel it named immediately after sending it + // (3.3.5.3.1). The client does not read the tunnel until the + // request has arrived (3.2.5.3.1), so data that overtakes it waits. + route_over_udp = + drdynvc.tunnel_for_outgoing_channel(egfx_dvc_id) == Some(dvc::pdu::SoftSyncTunnelType::RELIABLE_UDP); + + if route_over_udp { + if !self.egfx_on_udp { + newly_on_udp = Some(egfx_dvc_id); + } + for message in &messages { + let payload = message.encode_unframed_pdu().map_err(ServerError::encode)?; + udp_transport.send(payload).await; + } + } + } + if let Some(egfx_dvc_id) = newly_on_udp { + self.egfx_on_udp = true; + debug!(egfx_dvc_id, "EGFX is now sent over the UDP transport"); + } + + if route_over_udp { + return Ok(()); + } + self.write_egfx_over_tcp(messages, writer, drdynvc_channel_id, user_channel_id) + .await + } + + #[cfg(feature = "egfx")] + async fn write_egfx_over_tcp( + &mut self, + messages: Vec, + writer: &mut impl FramedWrite, + drdynvc_channel_id: u16, + user_channel_id: u16, + ) -> ServerResult<()> { + let data = + server_encode_svc_messages(messages, drdynvc_channel_id, user_channel_id).map_err(ServerError::encode)?; + writer + .write_all(&data) + .await + .map_err(|e| ServerError::io("write_all", e)) + } + + /// Writes DRDYNVC output, sending the data of any channel the Soft-Sync + /// request moved over the tunnel and everything else over TCP. + /// + /// Once the request is sent the server MUST keep using the tunnel it named + /// for those channels (MS-RDPEDYC 3.3.5.3.1), and that covers a channel's + /// replies to the client as well as data the server sends unprompted. + /// Without a live tunnel everything goes over TCP. + async fn write_drdynvc_output( + &mut self, + messages: Vec, + writer: &mut impl FramedWrite, + drdynvc_channel_id: u16, + user_channel_id: u16, + udp_transport: Option<&multitransport::UdpTransportHandle>, + ) -> ServerResult<()> { + let mut over_tcp = Vec::with_capacity(messages.len()); + match (udp_transport, self.get_svc_processor::()) { + (Some(udp_transport), Some(drdynvc)) => { + for message in messages { + let payload = message.encode_unframed_pdu().map_err(ServerError::encode)?; + let tunneled = match decode::(&payload) { + Ok(dvc::pdu::DrdynvcServerPdu::Data(data)) => { + drdynvc.tunnel_for_outgoing_channel(data.channel_id()) + == Some(dvc::pdu::SoftSyncTunnelType::RELIABLE_UDP) + } + _ => false, + }; + if tunneled { + udp_transport.send(payload).await; + } else { + over_tcp.push(message); + } + } + } + _ => over_tcp = messages, + } + + if over_tcp.is_empty() { + return Ok(()); + } + let data = + server_encode_svc_messages(over_tcp, drdynvc_channel_id, user_channel_id).map_err(ServerError::encode)?; + writer + .write_all(&data) + .await + .map_err(|e| ServerError::io("write drdynvc output", e))?; + Ok(()) + } + + /// Feeds tunnel payloads held back by [`Self::dispatch_udp_tunnel_payload`] + /// through, in arrival order, once the Soft-Sync response they were + /// waiting on has been processed. + async fn replay_early_tunnel_payloads( + &mut self, + writer: &mut impl FramedWrite, + user_channel_id: u16, + udp_transport: Option<&multitransport::UdpTransportHandle>, + ) -> ServerResult<()> { + if self.early_tunnel_payloads.is_empty() + || !self + .get_svc_processor::() + .is_some_and(|drdynvc| drdynvc.soft_sync_response_received()) + { + return Ok(()); + } + debug!( + count = self.early_tunnel_payloads.len(), + "Soft-Sync response received, processing the tunnel payloads held for it" + ); + while let Some(payload) = self.early_tunnel_payloads.pop_front() { + self.dispatch_udp_tunnel_payload(&payload, writer, user_channel_id, udp_transport) + .await?; + } + Ok(()) + } + + /// Decodes and dispatches one payload received over the sideband UDP + /// transport (MS-RDPEMT TunnelData): raw DRDYNVC bytes, per + /// [`dvc::DrdynvcServer::process_tunnel`]. Any resulting response + /// messages go back over TCP; only outgoing EGFX frames are proactively + /// routed onto the tunnel in this integration (see + /// [`Self::dispatch_egfx_messages`]). + /// + /// A malformed or unexpected tunnel payload is logged and dropped rather + /// than treated as a connection error: nothing about the primary TCP + /// session depends on the sideband transport's correctness. + async fn dispatch_udp_tunnel_payload( + &mut self, + payload: &[u8], + writer: &mut impl FramedWrite, + user_channel_id: u16, + udp_transport: Option<&multitransport::UdpTransportHandle>, + ) -> ServerResult { + // A single guard: `get_channel_id_by_type` and `get_svc_processor` + // both key off the same registered processor's `TypeId`, so one + // succeeding without the other would mean `static_channels` itself + // is inconsistent, not a case this function can meaningfully + // distinguish from "no drdynvc channel". + let (Some(drdynvc_channel_id), Some(drdynvc)) = ( + self.get_channel_id_by_type::(), + self.get_svc_processor::(), + ) else { + warn!("No drdynvc channel, dropping UDP tunnel payload"); + return Ok(RunState::Continue); + }; + + // MS-RDPEDYC 3.3.5.3.2: the server MUST NOT begin to read tunnel data + // until the Soft-Sync response has arrived. Hold it rather than drop + // it; `replay_early_tunnel_payloads` feeds it through once the + // response is in. + if drdynvc.soft_sync_awaiting_response() { + if self.early_tunnel_payloads.len() < MAX_EARLY_TUNNEL_PAYLOADS { + trace!( + len = payload.len(), + "Holding a tunnel payload until the Soft-Sync response arrives" + ); + self.early_tunnel_payloads.push_back(payload.to_vec()); + } else { + warn!("Too many tunnel payloads ahead of the Soft-Sync response, dropping one"); + } + return Ok(RunState::Continue); + } + + let messages = match drdynvc.process_tunnel(payload) { + Ok(messages) => messages, + Err(error) => { + warn!(%error, "Failed to process UDP tunnel payload, dropping it"); + return Ok(RunState::Continue); + } + }; + + if messages.is_empty() { + return Ok(RunState::Continue); + } + + self.write_drdynvc_output(messages, writer, drdynvc_channel_id, user_channel_id, udp_transport) + .await?; + + Ok(RunState::Continue) + } + #[expect( clippy::too_many_arguments, reason = "private per-connection entry point; the parameters are the connection's negotiated identifiers" @@ -3354,6 +3722,8 @@ impl RdpServer { message_channel_id: Option, client_supports_heartbeat: bool, mut encoder: UpdateEncoder, + udp_transport: Rc>>, + pending_udp_accept: Option>>, ) -> ServerResult where R: FramedRead, @@ -3371,6 +3741,9 @@ impl RdpServer { let mut event_writer = writer.clone(); let mut auto_reconnect_writer = writer.clone(); let mut heartbeat_writer = writer.clone(); + let mut udp_tunnel_writer = writer.clone(); + let udp_transport_for_events = Rc::clone(&udp_transport); + let udp_transport_for_pdus = Rc::clone(&udp_transport); let write_counter = writer.write_counter(); let ev_receiver = Arc::clone(&self.ev_receiver); let s = Rc::new(Mutex::new(self)); @@ -3391,6 +3764,7 @@ impl RdpServer { let lock_wait_ms = u64::try_from(lock_start.elapsed().as_millis()).unwrap_or(u64::MAX); let dispatch_start = Instant::now(); + let current_udp_transport = udp_transport_for_pdus.borrow().clone(); let result = this .dispatch_pdu( action, @@ -3399,6 +3773,7 @@ impl RdpServer { io_channel_id, user_channel_id, message_channel_id, + current_udp_transport.as_ref(), ) .await?; let dispatch_ms = u64::try_from(dispatch_start.elapsed().as_millis()).unwrap_or(u64::MAX); @@ -3488,6 +3863,11 @@ impl RdpServer { let lock_wait_ms = u64::try_from(lock_start.elapsed().as_millis()).unwrap_or(u64::MAX); let dispatch_start = Instant::now(); + // Cloned out of the cell before the call, since `dispatch_server_events` + // is async and this must not hold the `RefCell` borrow across + // an await point. `UdpTransportHandle` is cheap to clone (see + // its own doc comment). + let current_udp_transport = udp_transport_for_events.borrow().clone(); let result = this .dispatch_server_events( &mut events, @@ -3495,6 +3875,7 @@ impl RdpServer { io_channel_id, user_channel_id, message_channel_id, + current_udp_transport.as_ref(), ) .await?; let dispatch_ms = u64::try_from(dispatch_start.elapsed().as_millis()).unwrap_or(u64::MAX); @@ -3568,12 +3949,95 @@ impl RdpServer { } }; + let this = Rc::clone(&s); + let dispatch_udp_tunnel = async move { + // Only the pass of `client_loop` that received a fresh + // `pending_udp_accept` waits on it; every other future in this + // `select!` is already dispatching over TCP while this resolves, + // so the session never stalls on it. The accepted transport (if + // any) is stashed in the shared cell so `dispatch_events`, running + // concurrently, can start using it as soon as it lands too. + // + // If a Deactivation-Reactivation cycle drops this future while + // `handle` is still unresolved, the background accept task (and + // whatever socket/handshake it is mid-negotiation on) is orphaned: + // Its eventual result has nowhere left to be delivered, and this + // session continues TCP-only for the rest of its lifetime. That is + // an accepted narrow-window tradeoff; the alternative is + // persisting in-flight `JoinHandle`s across reactivation, which + // this integration does not attempt. The orphaned task also keeps + // `udp_bind_addr` bound for up to `UDP_ACCEPT_TIMEOUT`, so a NEXT + // connection's own bind attempt during that window can collide + // with it too, not only a concurrently preempted session's (see + // `multitransport::accept`'s own doc comment on that failure + // mode); it degrades the same way, gracefully, to TCP-only. + if let Some(handle) = pending_udp_accept { + match handle.await { + Ok(Some(transport)) => { + *udp_transport.borrow_mut() = Some(transport); + } + Ok(None) => { + debug!("UDP transport did not come up, continuing TCP-only for the rest of the session"); + } + Err(error) if error.is_cancelled() => { + debug!( + "UDP transport accept stopped after the client declined, continuing TCP-only for the rest of the session" + ); + } + Err(error) => { + warn!(%error, "UDP transport accept task panicked, continuing TCP-only for the rest of the session"); + } + } + } + + // This future is the cell's only writer, so the handle cannot + // change under the loop below. + let current = udp_transport.borrow().clone(); + let Some(transport) = current else { + return core::future::pending::>().await; + }; + loop { + let Some(payload) = transport.recv().await else { + // Soft-Sync only moves channels onto a tunnel (MS-RDPEDYC + // 2.2.5.1 has no TCP tunnel type), its request promised + // no more of their data over TCP (SOFT_SYNC_TCP_FLUSHED), + // and the tunnel lasts as long as the connection + // (MS-RDPEMT 1.3.3). A client whose channels moved has + // nowhere left to read them, so end the connection and let + // it reconnect rather than keep a session it cannot draw. + if this.lock().await.egfx_on_udp { + warn!("UDP transport lost with EGFX on it, ending the connection"); + return Err(ServerError::reason( + "UDP transport", + "lost with dynamic channels moved onto it", + )); + } + debug!("UDP transport closed, continuing TCP-only for the rest of the session"); + // Without this, `dispatch_egfx_messages` would keep seeing + // `Some(dead_handle)` here and stay on the UDP branch, + // silently dropping every future EGFX batch instead of + // actually falling back to TCP as this log claims. + *udp_transport.borrow_mut() = None; + return core::future::pending::>().await; + }; + let mut this = this.lock().await; + let result = this + .dispatch_udp_tunnel_payload(&payload, &mut udp_tunnel_writer, user_channel_id, Some(&transport)) + .await?; + match result { + RunState::Continue => continue, + state => break Ok(state), + } + } + }; + let state = tokio::select!( state = dispatch_pdu => state, state = dispatch_display => state, state = dispatch_events => state, state = refresh_auto_reconnect_cookie => state, state = send_heartbeats => state, + state = dispatch_udp_tunnel => state, ); debug!("End of client loop: {state:?}"); @@ -3585,6 +4049,8 @@ impl RdpServer { reader: &mut Framed, writer: &mut Framed, result: AcceptorResult, + udp_transport: Rc>>, + pending_udp_accept: Option>>, ) -> ServerResult where R: FramedRead, @@ -3643,12 +4109,15 @@ impl RdpServer { if !result.input_events.is_empty() { debug!("Handling input event backlog from acceptor sequence"); + // Set on a reactivation pass, where EGFX may already be on the tunnel. + let current_udp_transport = udp_transport.borrow().clone(); self.handle_input_backlog( writer, result.io_channel_id, result.user_channel_id, result.message_channel_id, result.input_events, + current_udp_transport.as_ref(), ) .await?; } @@ -3795,6 +4264,20 @@ impl RdpServer { self.send_next_auto_reconnect_cookie(writer, result.io_channel_id, result.user_channel_id) .await?; + let pending_udp_accept = + Self::drop_declined_udp_accept(pending_udp_accept, result.multitransport_response_success); + self.pending_udp_accept_abort = pending_udp_accept.as_ref().map(task::JoinHandle::abort_handle); + + // See `udp_migration_allowed`: a successful response that arrives + // after this point enables migration from the message-channel handler. + // Only ever raised here, never lowered: a deactivation-reactivation + // pass lands here again without repeating the multitransport exchange, + // and must not take back a migration the session may already be using. + self.soft_sync_negotiated |= result + .multitransport_flags + .contains(ironrdp_pdu::gcc::MultiTransportFlags::SOFT_SYNC_TCP_TO_UDP); + self.udp_migration_allowed |= self.soft_sync_negotiated && result.multitransport_response_success == Some(true); + let state = self .client_loop( reader, @@ -3806,12 +4289,34 @@ impl RdpServer { .client_early_capability_flags .contains(ironrdp_pdu::gcc::ClientEarlyCapabilityFlags::SUPPORT_HEART_BEAT_PDU), encoder, + udp_transport, + pending_udp_accept, ) .await?; Ok(state) } + /// A failure hrResponse (E_ABORT) means the client could not establish + /// the multitransport connection (MS-RDPBCGR 2.2.15.2), so no UDP + /// handshake is coming. Stop the accept now instead of holding the + /// socket until `multitransport::UDP_ACCEPT_TIMEOUT` and then reporting a + /// timeout. No response at all leaves it running: the client may still be + /// bringing UDP up. + fn drop_declined_udp_accept( + pending_udp_accept: Option>, + multitransport_response_success: Option, + ) -> Option> { + match pending_udp_accept { + Some(handle) if multitransport_response_success == Some(false) => { + handle.abort(); + debug!("Client could not establish the UDP multitransport connection, continuing TCP-only"); + None + } + other => other, + } + } + async fn handle_input_backlog( &mut self, writer: &mut impl FramedWrite, @@ -3819,6 +4324,7 @@ impl RdpServer { user_channel_id: u16, message_channel_id: Option, frames: Vec>, + udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult<()> { for frame in frames { match Action::from_fp_output_header(frame[0]) { @@ -3829,7 +4335,14 @@ impl RdpServer { Ok(Action::X224) => { let _ = self - .handle_x224(writer, io_channel_id, user_channel_id, message_channel_id, &frame) + .handle_x224( + writer, + io_channel_id, + user_channel_id, + message_channel_id, + &frame, + udp_transport, + ) .await; } @@ -3991,6 +4504,21 @@ impl RdpServer { hr_response = format!("{:#x}", pdu.hr_response), "Received Initiate Multitransport Response" ); + // E_ABORT: the client could not establish the multitransport + // connection (MS-RDPBCGR 2.2.15.2), so no UDP handshake is coming. + if !pdu.is_success() + && let Some(abort) = self.pending_udp_accept_abort.take() + { + abort.abort(); + debug!("Client could not establish the UDP multitransport connection, continuing TCP-only"); + } + // A success after finalization, the usual order with mstsc, + // completes the condition finalization could not see (see + // `udp_migration_allowed`), so EGFX can migrate from here on. + if pdu.is_success() && self.soft_sync_negotiated && !self.udp_migration_allowed { + self.udp_migration_allowed = true; + debug!("Multitransport confirmed after finalization, EGFX may migrate to UDP"); + } } Err(error) => { warn!(error = format!("{error:#}"), "Unhandled MCS message channel PDU"); @@ -4005,6 +4533,7 @@ impl RdpServer { user_channel_id: u16, message_channel_id: Option, frame: &[u8], + udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult { let message = decode::>>(frame).map_err(ServerError::decode)?; match message.0 { @@ -4028,12 +4557,25 @@ impl RdpServer { let response_pdus = svc .process(&data.user_data) .map_err_kind("svc process", ServerErrorKind::Pdu)?; - let response = server_encode_svc_messages(response_pdus, data.channel_id, user_channel_id) - .map_err(ServerError::encode)?; - writer - .write_all(&response) - .await - .map_err(|e| ServerError::io("write svc response", e))?; + if self.get_channel_id_by_type::() == Some(data.channel_id) { + self.write_drdynvc_output( + response_pdus, + writer, + data.channel_id, + user_channel_id, + udp_transport, + ) + .await?; + self.replay_early_tunnel_payloads(writer, user_channel_id, udp_transport) + .await?; + } else { + let response = server_encode_svc_messages(response_pdus, data.channel_id, user_channel_id) + .map_err(ServerError::encode)?; + writer + .write_all(&response) + .await + .map_err(|e| ServerError::io("write svc response", e))?; + } } else { warn!(channel_id = data.channel_id, "Unexpected channel received: ID",); } @@ -4094,6 +4636,38 @@ impl RdpServer { where S: AsyncRead + AsyncWrite + Sync + Send + Unpin, { + // Per-connection: set again once this connection's finalization knows + // its own negotiated flags and response. + self.soft_sync_negotiated = false; + self.udp_migration_allowed = false; + self.egfx_on_udp = false; + self.early_tunnel_payloads.clear(); + + let udp_bind_addr = self + .opts + .udp_bind_addr + .map(|bind| multitransport::sideband_bind_addr(bind, self.connection_local_addr)); + let tls_config = self + .opts + .security + .tls_acceptor() + .map(|acceptor| Arc::clone(acceptor.config())); + // Persists across a Deactivation-Reactivation pass of this loop: The + // sideband transport, once established, is not torn down or + // re-bootstrapped for a resize (see the acceptor's own doc comment + // on `MultitransportBootstrapping` regarding reactivation). Shared + // with `client_loop`, which populates it once `pending_udp_accept` + // (below) resolves. + let udp_transport: Rc>> = Rc::new(RefCell::new(None)); + // Set once, by the handler below, the first time the acceptor + // actually sends a request. Taken (not cloned) when handed to + // `client_loop`, which awaits it in its own select loop: The full + // RDPEUDP2 + TLS + RDPEMT handshake this drives can take up to + // `multitransport::UDP_ACCEPT_TIMEOUT`, and the handler below runs + // synchronously as part of finalize, so establishing it must never + // block the RDP handshake finalize itself is driving. + let mut pending_udp_accept: Option>> = None; + loop { // Bounded: see `FINALIZE_TIMEOUT`. The bound belongs on THIS call // and not on `accept_finalize` itself or its callers — the loop @@ -4102,7 +4676,30 @@ impl RdpServer { // length. Applying it per pass also gives a // deactivation-reactivation its own budget rather than sharing one // with the initial handshake. - let finalize = ironrdp_acceptor::accept_finalize(framed, &mut acceptor); + // + // Always the `_with_multitransport` driver rather than branching + // on whether UDP is configured: when it is not, `acceptor` never + // has a request to report (`set_multitransport_offer` was never + // called, see `PendingConnection::new`), so the handler below + // simply never runs and this is exactly `accept_finalize`'s own + // behavior. + let finalize = ironrdp_acceptor::accept_finalize_with_multitransport( + framed, + &mut acceptor, + async |request, _soft_sync| { + // Defensive only: reached only if the acceptor sent a + // request despite one of these being unset, which + // `PendingConnection::new`'s gate should prevent. + let (Some(udp_bind_addr), Some(tls_config)) = (udp_bind_addr, tls_config.clone()) else { + return; + }; + // Spawned rather than awaited inline: See the comment on + // `pending_udp_accept`'s declaration above. + pending_udp_accept = Some(task::spawn(async move { + multitransport::accept(udp_bind_addr, tls_config, &request).await + })); + }, + ); let (new_framed, result) = match tokio::time::timeout(FINALIZE_TIMEOUT, finalize).await { Ok(res) => res.map_err_kind("failed to accept client during finalize", ServerErrorKind::Connector)?, Err(_) => { @@ -4119,7 +4716,16 @@ impl RdpServer { let (mut reader, mut writer) = split_tokio_framed(new_framed); - match self.client_accepted(&mut reader, &mut writer, result).await? { + match self + .client_accepted( + &mut reader, + &mut writer, + result, + Rc::clone(&udp_transport), + pending_udp_accept.take(), + ) + .await? + { RunState::Continue => { unreachable!(); } @@ -4149,6 +4755,24 @@ impl RdpServer { debug!(?creds, "Changing credentials"); self.creds = creds } + + /// Tell the server which local address the next client reached, for + /// connections driven through [`Self::run_connection`] or + /// [`Self::run_connection_with`]. [`Self::run`] sets it itself. + /// + /// When the UDP transport address (see + /// [`RdpServerBuilder::with_udp_transport`](crate::RdpServerBuilder::with_udp_transport)) + /// has an unspecified IP, the sideband socket binds to this address + /// instead. A socket bound to the unspecified address replies from + /// whichever local address the routing table picks, and on a host with + /// several addresses (IPv6 temporary addresses, for example) that is not + /// always the one the client sent to. The client then drops the replies + /// and the UDP handshake never completes. + /// + /// Cleared when the connection ends. + pub fn set_connection_local_addr(&mut self, addr: Option) { + self.connection_local_addr = addr; + } } /// Encode a server-initiated Auto-Detect Request PDU for the MCS message channel. @@ -4399,6 +5023,7 @@ mod preempt_tests { connection_policy: ConnectionPolicy::Preempt, remotefx_quant: Quant::default(), remotefx_entropy_coder: None, + udp_bind_addr: None, }, creds: None, display: Arc::new(Mutex::new(Box::new(NoDisplay))), @@ -4595,6 +5220,104 @@ mod preempt_tests { /// /// Drive exactly that: a live session, a silent candidate, then end the /// session and require the server to still respond and still shut down. + #[tokio::test] + async fn a_declined_multitransport_response_stops_the_udp_accept() { + let local = task::LocalSet::new(); + local + .run_until(async { + let pending = || Some(task::spawn_local(core::future::pending::<()>())); + + let declined = RdpServer::drop_declined_udp_accept(pending(), Some(false)); + assert!(declined.is_none()); + + let kept = RdpServer::drop_declined_udp_accept(pending(), None).expect("no response keeps the accept"); + assert!(!kept.is_finished()); + kept.abort(); + + let kept = + RdpServer::drop_declined_udp_accept(pending(), Some(true)).expect("success keeps the accept"); + assert!(!kept.is_finished()); + kept.abort(); + + // The aborted task really stops, so its socket is released. + let handle = task::spawn_local(core::future::pending::<()>()); + let probe = handle.abort_handle(); + assert!(RdpServer::drop_declined_udp_accept(Some(handle), Some(false)).is_none()); + task::yield_now().await; + assert!(probe.is_finished()); + }) + .await; + } + + #[test] + fn a_late_multitransport_success_enables_migration_only_with_soft_sync() { + use ironrdp_pdu::rdp::multitransport::MultitransportResponsePdu; + + let mut server = RdpServer::builder() + .with_addr((Ipv4Addr::LOCALHOST, 0)) + .with_no_security() + .with_no_input() + .with_no_display() + .build(); + let response = |pdu: &MultitransportResponsePdu| SendDataRequest { + initiator_id: 1007, + channel_id: 1008, + user_data: encode_vec(pdu).expect("encode response").into(), + }; + + // Soft-Sync not negotiated: MS-RDPEDYC forbids migration whatever the response. + server.handle_message_channel_data(response(&MultitransportResponsePdu::success(1))); + assert!(!server.udp_migration_allowed); + + // Negotiated, but the client could not bring UDP up. + server.soft_sync_negotiated = true; + server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); + assert!(!server.udp_migration_allowed); + + // Negotiated, and the success arrives after finalization. + server.handle_message_channel_data(response(&MultitransportResponsePdu::success(1))); + assert!(server.udp_migration_allowed); + } + + #[tokio::test] + async fn a_late_multitransport_failure_stops_the_pending_udp_accept() { + use ironrdp_pdu::rdp::multitransport::MultitransportResponsePdu; + + let local = task::LocalSet::new(); + local + .run_until(async { + let mut server = RdpServer::builder() + .with_addr((Ipv4Addr::LOCALHOST, 0)) + .with_no_security() + .with_no_input() + .with_no_display() + .build(); + let response = |pdu: &MultitransportResponsePdu| SendDataRequest { + initiator_id: 1007, + channel_id: 1008, + user_data: encode_vec(pdu).expect("encode response").into(), + }; + + // No accept pending: a failure response is only logged. + server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); + + let accept = task::spawn_local(core::future::pending::<()>()); + server.pending_udp_accept_abort = Some(accept.abort_handle()); + + // Success leaves the accept running. + server.handle_message_channel_data(response(&MultitransportResponsePdu::success(1))); + task::yield_now().await; + assert!(!accept.is_finished()); + + server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); + task::yield_now().await; + assert!(accept.is_finished()); + assert!(accept.await.expect_err("accept was aborted").is_cancelled()); + assert!(server.pending_udp_accept_abort.is_none()); + }) + .await; + } + #[tokio::test] async fn a_silent_candidate_cannot_wedge_the_accept_loop() { let local = task::LocalSet::new(); @@ -5112,7 +5835,7 @@ mod cliprdr_error_tests { let mut writer = CapturingWriter::default(); let state = server - .dispatch_server_events(&mut events, &mut writer, 1003, 1002, None) + .dispatch_server_events(&mut events, &mut writer, 1003, 1002, None, None) .await .expect("a refused clipboard message must not surface as a session error"); diff --git a/crates/ironrdp-testsuite-extra/Cargo.toml b/crates/ironrdp-testsuite-extra/Cargo.toml index 18f6b4bfd2..3f5a348296 100644 --- a/crates/ironrdp-testsuite-extra/Cargo.toml +++ b/crates/ironrdp-testsuite-extra/Cargo.toml @@ -34,6 +34,7 @@ ironrdp-client = { path = "../ironrdp-client", features = ["sound", "udp"] } ironrdp-core.path = "../ironrdp-core" ironrdp-dvc.path = "../ironrdp-dvc" ironrdp-dvc-pipe-proxy.path = "../ironrdp-dvc-pipe-proxy" +ironrdp-egfx.path = "../ironrdp-egfx" ironrdp-graphics.path = "../ironrdp-graphics" ironrdp-svc.path = "../ironrdp-svc" ironrdp-input.path = "../ironrdp-input" @@ -47,6 +48,9 @@ ironrdp-rdpeudp-tokio.path = "../ironrdp-rdpeudp-tokio" ironrdp-rdpsnd.path = "../ironrdp-rdpsnd" ironrdp-rdpsnd-native = { path = "../ironrdp-rdpsnd-native", features = ["capture"] } ironrdp-rpc = { path = "../ironrdp-rpc", features = ["__test"] } +# The EGFX-over-UDP migration in ironrdp-server is behind its `egfx` feature, +# which nothing else in the workspace enables. +ironrdp-server = { path = "../ironrdp-server", features = ["egfx"] } ironrdp-viewer.path = "../ironrdp-viewer" ironrdp-tokio.path = "../ironrdp-tokio" ironrdp-tls = { path = "../ironrdp-tls", features = ["rustls"] } diff --git a/crates/ironrdp-testsuite-extra/tests/e2e.rs b/crates/ironrdp-testsuite-extra/tests/e2e.rs index 8e2d7055f9..bf0db21a3e 100644 --- a/crates/ironrdp-testsuite-extra/tests/e2e.rs +++ b/crates/ironrdp-testsuite-extra/tests/e2e.rs @@ -1,5 +1,6 @@ // FIXME: tests in this module can probably be rewritten to be much shorter using the ironrdp-client crate. +use core::net::SocketAddr; use core::sync::atomic::{AtomicBool, Ordering}; use core::time::Duration; use std::path::Path; @@ -26,6 +27,8 @@ use ironrdp::session::{self, ActiveStage, ActiveStageBuilder, ActiveStageOutput} use ironrdp::svc::{StaticChannelSet, SvcMessage, SvcProcessor, SvcServerProcessor}; use ironrdp_async::{Framed, FramedWrite as _}; use ironrdp_bulk::{BulkCompressor, CompressionType as BulkCompressionType, flags as bulk_flags}; +use ironrdp_dvc::{DvcClientProcessor, DvcEncode, DvcMessage, DvcProcessor}; +use ironrdp_egfx::pdu::{CapabilitiesAdvertisePdu, CapabilitiesV8Flags, CapabilitySet, GfxPdu}; use ironrdp_rdpdr::pdu::RdpdrPdu; use ironrdp_rdpdr::pdu::efs::{ Capabilities, CoreCapability, CoreCapabilityKind, DeviceCreateResponse, DeviceIoRequest, DeviceIoResponse, @@ -34,6 +37,7 @@ use ironrdp_rdpdr::pdu::efs::{ }; use ironrdp_rdpdr::pdu::esc::{ScardCall, ScardIoCtlCode}; use ironrdp_rdpdr::{Rdpdr, RdpdrBackend, RdpdrBackendFactory, RdpdrBackendProduct, RdpdrDrive}; +use ironrdp_server::{GfxServerFactory, ServerEventSender}; use ironrdp_testsuite_extra as _; use ironrdp_tls::TlsStream; use ironrdp_tokio::TokioStream; @@ -61,6 +65,24 @@ async fn test_client_server() { .await } +/// Configuring UDP multitransport on the server must not disturb a client +/// that never advertises support for it: `set_multitransport_offer` gates on +/// the client's own GCC `MultiTransportChannelData` reciprocating, so with +/// `multitransport_flags: None` (the default) the acceptor never sends the +/// Initiate Multitransport Request and the connection proceeds exactly as it +/// does with no UDP transport configured at all. +#[tokio::test] +async fn test_client_server_with_udp_transport_configured_but_unused() { + client_server_with_connector( + default_client_config(), + Vec::new(), + Some(([127, 0, 0, 1], 0).into()), + |connector| connector, + |stage, _activation_factory, framed, _display_tx, _echo_handle| async { (stage, framed) }, + ) + .await; +} + /// Advertising the Graphics Pipeline early-capability bit must not disturb connection establishment. /// /// The core testsuite separately decodes the emitted Connect Initial PDU and verifies the capability flag. @@ -249,6 +271,7 @@ async fn test_echo_virtual_channel_end_to_end() { client_server_with_connector( default_client_config(), Vec::new(), + None, |connector| connector.with_static_channel(DrdynvcClient::new().with_dynamic_channel(EchoClient::new())), move |mut stage, _activation_factory, mut framed, display_tx, echo_handle| async move { let _display_tx = display_tx; @@ -304,6 +327,7 @@ async fn rdpdr_static_channel_announces_a_drive_and_completes_an_unsupported_cre client_server_with_connector( default_client_config(), vec![Box::new(fixture)], + None, |connector| connector.with_static_channel(test_rdpdr_channel()), move |stage, _activation_factory, framed, display_tx, _echo_handle| { drive_rdpdr_until_complete(stage, framed, display_tx, fixture_state) @@ -332,6 +356,7 @@ async fn rdpdr_static_channel_creates_a_file_with_the_windows_backend() { client_server_with_connector( default_client_config(), vec![Box::new(fixture)], + None, move |connector| connector.with_static_channel(rdpdr_channel(&factory)), move |stage, _activation_factory, framed, display_tx, _echo_handle| { drive_rdpdr_until_complete(stage, framed, display_tx, fixture_state_for_client) @@ -369,6 +394,7 @@ async fn rdpdr_static_channel_preserves_large_read_response_lengths() { client_server_with_connector( default_client_config(), vec![Box::new(fixture)], + None, move |connector| connector.with_static_channel(rdpdr_channel(&factory)), move |stage, _activation_factory, framed, display_tx, _echo_handle| { drive_rdpdr_until_complete(stage, framed, display_tx, fixture_state_for_client) @@ -884,6 +910,7 @@ where client_server_with_connector( client_config, Vec::new(), + None, |connector| connector, move |stage, connection_activation, framed, display_tx, _echo_handle| { clientfn(stage, connection_activation, framed, display_tx) @@ -895,6 +922,7 @@ where async fn client_server_with_connector( client_config: connector::Config, static_channel_factories: Vec>, + server_udp_addr: Option, connector_factory: C, clientfn: F, ) where @@ -930,6 +958,9 @@ async fn client_server_with_connector( for factory in static_channel_factories { server_builder = server_builder.with_static_channel_factory(factory); } + if let Some(udp_addr) = server_udp_addr { + server_builder = server_builder.with_udp_transport(udp_addr); + } let mut server = server_builder.build(); server.set_credentials(Some(server::Credentials { username: USERNAME.into(), @@ -1026,6 +1057,372 @@ async fn client_server_with_connector( .await; } +/// EGFX moves onto the UDP tunnel with Soft-Sync (MS-RDPEDYC 3.1.5.3), end to end. +/// +/// The server sends the Soft-Sync Request when EGFX first has data to send, and +/// from then on every server message on that channel goes over the tunnel +/// (3.3.5.3.1): the batch that triggered the request, and the channel's reply to +/// a client message. The IronRDP client rejects TCP data for a channel it moved +/// to the tunnel, so any EGFX data sent over TCP after the request fails the test. +/// +/// The client writes on the tunnel as soon as it has sent the Soft-Sync +/// Response, so its first tunnel data can reach the server before the response +/// does. The test forces that order and checks the server still processes it. +/// +/// Soft-Sync has no tunnel type that moves a channel back to TCP (2.2.5.1), +/// and the tunnel lasts as long as the connection (MS-RDPEMT 1.3.3), so when +/// the tunnel closes under EGFX the server has to end the connection. +#[tokio::test] +async fn egfx_moves_onto_the_udp_tunnel_with_soft_sync() { + let _ = tracing_subscriber::fmt() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .try_init(); + + // The server binds its sideband socket to the configured port, so the + // client has to know that port before the server reports it. + let udp_port = std::net::UdpSocket::bind("127.0.0.1:0") + .and_then(|socket| socket.local_addr()) + .expect("pick a free UDP port") + .port(); + let udp_addr: SocketAddr = ([127, 0, 0, 1], udp_port).into(); + + let cert_path = Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/certs/server-cert.pem"); + let key_path = Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/certs/server-key.pem"); + let identity = TlsIdentityCtx::init_from_paths(&cert_path, &key_path).expect("failed to init TLS identity"); + let acceptor = identity.make_acceptor().expect("failed to build TLS acceptor"); + + let (caps_tx, mut caps_rx) = mpsc::unbounded_channel(); + let (_display_tx, display_rx) = mpsc::unbounded_channel(); + let mut server = RdpServer::builder() + .with_addr(([127, 0, 0, 1], 0)) + .with_tls(acceptor) + .with_input_handler(TestInputHandler) + .with_display_handler(TestDisplay { + rx: Arc::new(Mutex::new(display_rx)), + }) + .with_gfx_factory(Some(Box::new(TestGfxFactory { caps_tx }))) + .with_udp_transport(udp_addr) + .build(); + server.set_credentials(Some(server::Credentials { + username: USERNAME.into(), + password: PASSWORD.into(), + domain: None, + })); + let ev = server.event_sender().clone(); + + let local = tokio::task::LocalSet::new(); + local + .run_until(async move { + let server = tokio::task::spawn_local(async move { + server.run().await.unwrap(); + }); + + let client = tokio::task::spawn_local(async move { + let (tx, rx) = oneshot::channel(); + ev.send(ServerEvent::GetLocalAddr(tx)).unwrap(); + let server_addr = rx.await.unwrap().unwrap(); + let tcp_stream = TcpStream::connect(server_addr).await.expect("TCP connect"); + let client_addr = tcp_stream.local_addr().expect("local_addr"); + + let egfx = TestEgfxClient::default(); + let egfx_channel = Arc::clone(&egfx.channel_id); + let egfx_received = Arc::clone(&egfx.received); + let client_config = connector::Config { + support_dyn_vc_gfx_protocol: true, + multitransport_flags: Some( + gcc::MultiTransportFlags::TRANSPORT_TYPE_UDP_FECR + | gcc::MultiTransportFlags::SOFT_SYNC_TCP_TO_UDP, + ), + ..default_client_config() + }; + let mut connector = connector::ClientConnector::new(client_config, client_addr) + .with_static_channel(DrdynvcClient::new().with_dynamic_channel(egfx)); + + let mut framed = ironrdp_tokio::TokioFramed::new(tcp_stream); + let should_upgrade = ironrdp_async::connect_begin(&mut framed, &mut connector) + .await + .expect("begin connection"); + let (upgraded_stream, tls_cert) = ironrdp_tls::upgrade_with_certificate_validation( + framed.into_inner_no_leftover(), + "localhost", + ironrdp_tls::CertificateValidation::DangerouslyAcceptInvalidCertificate, + ) + .await + .expect("TLS upgrade"); + let upgraded = ironrdp_tokio::mark_as_upgraded(should_upgrade, &mut connector); + let mut framed = ironrdp_tokio::TokioFramed::new(upgraded_stream); + let server_public_key = + ironrdp_tls::extract_tls_server_public_key(&tls_cert).expect("extract server public key"); + + let mut tunnel = None; + let connection_result = ironrdp_async::connect_finalize_with_multitransport( + upgraded, + connector, + &mut framed, + &mut ironrdp_tokio::reqwest::ReqwestNetworkClient::new(), + "localhost".into(), + server_public_key.to_owned(), + None, + async |request, soft_sync| { + assert!(soft_sync, "both peers advertised SOFT_SYNC_TCP_TO_UDP"); + let mut bootstrap = ironrdp_rdpeudp_tokio::MultitransportBootstrap::new(request); + bootstrap + .connect( + udp_addr, + "localhost".into(), + ironrdp_rdpeudp::ConnectionConfig::default(), + ironrdp_rdpeudp_tokio::UdpTlsConfig::new("localhost".into()), + ) + .await + .expect("UDP tunnel handshake"); + tunnel = bootstrap.take_transport(); + Ok(connector::MultitransportResult::Success) + }, + ) + .await + .expect("finalize connection"); + let mut tunnel = tunnel.expect("the server sent an Initiate Multitransport Request"); + + let mut stage = ActiveStageBuilder { + static_channels: connection_result.static_channels, + user_channel_id: connection_result.user_channel_id, + io_channel_id: connection_result.io_channel_id, + message_channel_id: connection_result.message_channel_id, + share_id: connection_result.share_id, + compression_type: connection_result.compression_type, + enable_server_pointer: connection_result.enable_server_pointer, + pointer_software_rendering: connection_result.pointer_software_rendering, + } + .build(); + stage.enable_reliable_udp_dvc_tunnel().expect("DRDYNVC is present"); + let mut image = DecodedImage::new(PixelFormat::RgbA32, DESKTOP_WIDTH, DESKTOP_HEIGHT); + + // The client advertises its capabilities over TCP when the channel + // opens, as a real client does; the server handler seeing them means + // the server has the channel open too. + tokio::time::timeout(Duration::from_secs(10), async { + loop { + tokio::select! { + caps = caps_rx.recv() => { + caps.expect("capabilities reached the EGFX handler"); + break; + } + pdu = framed.read_pdu() => { + let (action, frame) = pdu.expect("read PDU"); + for output in stage.process(&mut image, action, &frame).expect("stage process") { + if let ActiveStageOutput::ResponseFrame(frame) = output { + framed.write_all(&frame).await.expect("write response frame"); + } + } + } + } + } + }) + .await + .expect("the server opened the EGFX channel"); + let egfx_channel_id = egfx_channel + .lock() + .unwrap() + .expect("the client opened the EGFX channel"); + + let batch = b"first egfx batch".to_vec(); + let messages = ironrdp_dvc::encode_dvc_messages( + egfx_channel_id, + vec![Box::new(RawDvcPayload(batch.clone()))], + ironrdp::svc::ChannelFlags::empty(), + ) + .expect("encode EGFX batch"); + ev.send(ServerEvent::Egfx(ironrdp_server::EgfxServerMessage::SendMessages { + messages, + })) + .unwrap(); + + // Process TCP until the Soft-Sync Request has moved the channel, + // holding back the Soft-Sync Response it produced. + let held_response = tokio::time::timeout(Duration::from_secs(10), async { + loop { + let (action, frame) = framed.read_pdu().await.expect("read PDU"); + let mut responses = Vec::new(); + for output in stage.process(&mut image, action, &frame).expect("stage process") { + if let ActiveStageOutput::ResponseFrame(frame) = output { + responses.push(frame); + } + } + if stage.dvc_tunnel_for_channel(egfx_channel_id) + == Some(ironrdp_dvc::pdu::SoftSyncTunnelType::RELIABLE_UDP) + { + break responses; + } + for frame in responses { + framed.write_all(&frame).await.expect("write response frame"); + } + } + }) + .await + .expect("the server sent a Soft-Sync Request for EGFX"); + + // Client data on the tunnel that reaches the server ahead of the + // Soft-Sync Response, and that the EGFX channel answers: a second + // capabilities advertisement, which a client sends to reset its + // decoder. + let caps = ironrdp_dvc::pdu::DrdynvcClientPdu::Data(ironrdp_dvc::pdu::DrdynvcDataPdu::Data( + ironrdp_dvc::pdu::DataPdu::new(egfx_channel_id, encode_vec(&v8_capabilities()).unwrap()), + )); + tunnel + .send(encode_vec(&caps).unwrap()) + .await + .expect("send on the tunnel"); + tokio::time::sleep(Duration::from_millis(200)).await; + for frame in held_response { + framed.write_all(&frame).await.expect("write Soft-Sync Response"); + } + + // Over TCP, the confirmation of the first capabilities. Over the + // tunnel, the triggering batch and the second confirmation. + tokio::time::timeout(Duration::from_secs(10), async { + while egfx_received.lock().unwrap().len() < 3 { + tokio::select! { + payload = tunnel.recv() => { + let payload = payload.expect("tunnel open"); + stage + .process_dvc_tunnel(&mut image, ironrdp_dvc::pdu::SoftSyncTunnelType::RELIABLE_UDP, &payload) + .expect("DVC data for the tunneled channel"); + } + pdu = framed.read_pdu() => { + let (action, frame) = pdu.expect("read PDU"); + for output in stage.process(&mut image, action, &frame).expect("stage process") { + if let ActiveStageOutput::ResponseFrame(frame) = output { + framed.write_all(&frame).await.expect("write response frame"); + } + } + } + } + } + }) + .await + .expect("EGFX batch and capabilities reply arrived on the tunnel"); + assert_eq!( + egfx_received.lock().unwrap()[1], + batch, + "the batch that triggered the request goes over the tunnel" + ); + + tokio::time::timeout(Duration::from_secs(10), caps_rx.recv()) + .await + .expect("the server processed the early tunnel data") + .expect("capabilities reached the EGFX handler"); + + // With EGFX on the tunnel there is no way back to TCP, so the + // server ends the connection when the tunnel closes. Graphics + // still being sent is what shows the server the tunnel is gone. + drop(tunnel); + let messages = ironrdp_dvc::encode_dvc_messages( + egfx_channel_id, + vec![Box::new(RawDvcPayload(b"after the tunnel closed".to_vec()))], + ironrdp::svc::ChannelFlags::empty(), + ) + .expect("encode EGFX batch"); + ev.send(ServerEvent::Egfx(ironrdp_server::EgfxServerMessage::SendMessages { + messages, + })) + .unwrap(); + tokio::time::timeout(Duration::from_secs(20), async { + while let Ok(pdu) = framed.read_pdu().await { + debug!(?pdu); + } + }) + .await + .expect("the server ended the connection once the tunnel closed"); + ev.send(ServerEvent::Quit("bye".into())).unwrap(); + }); + + tokio::try_join!(server, client).expect("join"); + }) + .await; +} + +/// Server EGFX handler that reports the client's advertised capabilities. +struct TestGfxFactory { + caps_tx: UnboundedSender<()>, +} + +impl ServerEventSender for TestGfxFactory { + fn set_sender(&mut self, _sender: UnboundedSender) {} +} + +impl GfxServerFactory for TestGfxFactory { + fn build_gfx_handler(&self) -> Box { + Box::new(TestGfxHandler { + caps_tx: self.caps_tx.clone(), + }) + } +} + +struct TestGfxHandler { + caps_tx: UnboundedSender<()>, +} + +impl ironrdp_egfx::server::GraphicsPipelineHandler for TestGfxHandler { + fn capabilities_advertise(&mut self, _pdu: &CapabilitiesAdvertisePdu) { + let _ = self.caps_tx.send(()); + } + + fn on_ready(&mut self, _negotiated: &CapabilitySet) {} +} + +/// Client end of the EGFX channel that records what the server sends on it. +#[derive(Default)] +struct TestEgfxClient { + channel_id: Arc>>, + received: Arc>>>, +} + +impl_as_any!(TestEgfxClient); + +impl DvcProcessor for TestEgfxClient { + fn channel_name(&self) -> &str { + "Microsoft::Windows::RDS::Graphics" + } + + fn start(&mut self, channel_id: u32) -> pdu::PduResult> { + *self.channel_id.lock().unwrap() = Some(channel_id); + Ok(vec![Box::new(v8_capabilities())]) + } + + fn process(&mut self, _channel_id: u32, payload: &[u8]) -> pdu::PduResult> { + self.received.lock().unwrap().push(payload.to_vec()); + Ok(Vec::new()) + } +} + +impl DvcClientProcessor for TestEgfxClient {} + +fn v8_capabilities() -> GfxPdu { + GfxPdu::CapabilitiesAdvertise(CapabilitiesAdvertisePdu::from_typed(&[CapabilitySet::V8 { + flags: CapabilitiesV8Flags::empty(), + }])) +} + +struct RawDvcPayload(Vec); + +impl ironrdp::core::Encode for RawDvcPayload { + fn encode(&self, dst: &mut ironrdp::core::WriteCursor<'_>) -> ironrdp::core::EncodeResult<()> { + ironrdp::core::ensure_size!(in: dst, size: self.0.len()); + dst.write_slice(&self.0); + Ok(()) + } + + fn name(&self) -> &'static str { + "RawDvcPayload" + } + + fn size(&self) -> usize { + self.0.len() + } +} + +impl DvcEncode for RawDvcPayload {} + pub(super) fn default_client_config() -> connector::Config { connector::Config { desktop_size: DesktopSize { diff --git a/crates/ironrdp/examples/server.rs b/crates/ironrdp/examples/server.rs index 58e0468bb2..7e6c9c9e3f 100644 --- a/crates/ironrdp/examples/server.rs +++ b/crates/ironrdp/examples/server.rs @@ -54,7 +54,11 @@ async fn main() -> Result<(), anyhow::Error> { pass, cert, key, - } => run(bind_addr, hybrid, user, pass, cert, key).await, + } => { + // The server's future is large enough to trip clippy::large_futures, + // so keep it on the heap rather than in main's stack frame. + Box::pin(run(bind_addr, hybrid, user, pass, cert, key)).await + } } }