From 27bddf6fb21123476e59fc05c836c048d98df765 Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Mon, 28 Sep 2026 00:21:36 -0500 Subject: [PATCH 1/4] feat(rdpeudp): carry tunnel sub-headers through the UDP transport The RDPEMT codec already handled sub-headers, but UdpTransport exposed only a PDU's data, so the Continuous Auto-Detection messages MS-RDPBCGR 1.3.9 carries there could be neither sent nor received. TunnelMessage carries both; send_message and recv_message use it, and send/recv keep working on data alone. --- crates/ironrdp-rdpeudp-tokio/src/error.rs | 10 +- crates/ironrdp-rdpeudp-tokio/src/framed.rs | 21 +- crates/ironrdp-rdpeudp-tokio/src/lib.rs | 3 +- crates/ironrdp-rdpeudp-tokio/src/transport.rs | 204 +++++++++++++++--- crates/ironrdp-rdpeudp-tokio/src/tunnel.rs | 59 ++++- .../tests/rdpeudp_tokio.rs | 41 +++- 6 files changed, 290 insertions(+), 48 deletions(-) diff --git a/crates/ironrdp-rdpeudp-tokio/src/error.rs b/crates/ironrdp-rdpeudp-tokio/src/error.rs index 2ebe906487..2018af9bcc 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/error.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/error.rs @@ -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 { @@ -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") + } } } } @@ -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, } } } diff --git a/crates/ironrdp-rdpeudp-tokio/src/framed.rs b/crates/ironrdp-rdpeudp-tokio/src/framed.rs index f38f5ef361..4fff17426d 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/framed.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/framed.rs @@ -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>, mpsc::Receiver>) { - let (incoming_tx, incoming_rx) = mpsc::channel::>(16); - let (outgoing_tx, outgoing_rx) = mpsc::channel::>(16); + fn test_transport() -> (UdpTransport, mpsc::Sender, mpsc::Receiver) { + let (incoming_tx, incoming_rx) = mpsc::channel::(16); + let (outgoing_tx, outgoing_rx) = mpsc::channel::(16); let transport = UdpTransport::from_channels(incoming_rx, outgoing_tx); @@ -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(); @@ -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(); @@ -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(); @@ -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] @@ -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(); diff --git a/crates/ironrdp-rdpeudp-tokio/src/lib.rs b/crates/ironrdp-rdpeudp-tokio/src/lib.rs index 517975ccdd..20fece5f1a 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/lib.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/lib.rs @@ -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, }; diff --git a/crates/ironrdp-rdpeudp-tokio/src/transport.rs b/crates/ironrdp-rdpeudp-tokio/src/transport.rs index d5b9c79605..f3d42d184d 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/transport.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/transport.rs @@ -29,7 +29,10 @@ use tokio::task::JoinHandle; use tracing::{debug, info, trace, warn}; use crate::driver::Driver; -use crate::error::{DriverError, DriverErrorExt as _, DriverErrorKind, UdpTransportError, UdpTransportErrorExt as _}; +use crate::error::{ + DriverError, DriverErrorExt as _, DriverErrorKind, UdpTransportError, UdpTransportErrorExt as _, + UdpTransportErrorKind, +}; use crate::stream::{RdpeudpStream, SharedIo}; use crate::tls::{tls_accept, tls_upgrade}; use crate::tunnel::{read_tunnel_pdu, tunnel_data_loop, write_tunnel_pdu}; @@ -211,6 +214,28 @@ impl Drop for AbortOnDrop { } } +/// One RDPEMT Tunnel Data PDU's content: the higher-layer data and the +/// sub-headers carried beside it ([MS-RDPEMT] 2.2.2.3). +/// +/// The sub-headers carry the Continuous Auto-Detection messages that +/// [MS-RDPBCGR] 1.3.9 sends over a sideband channel in use, such as a +/// bandwidth measurement's Start and Stop and the client's results. `data` +/// may be empty when a message carries only sub-headers. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct TunnelMessage { + pub sub_headers: Vec, + pub data: Vec, +} + +impl From> for TunnelMessage { + fn from(data: Vec) -> Self { + Self { + sub_headers: Vec::new(), + data, + } + } +} + /// A cloneable handle for sending data over an established UDP transport, /// obtained from [`UdpTransport::sender`]. /// @@ -218,7 +243,7 @@ impl Drop for AbortOnDrop { /// 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>); +pub struct UdpTransportSender(mpsc::Sender); impl UdpTransportSender { /// Send a higher-layer data frame through the tunnel. @@ -231,13 +256,36 @@ impl UdpTransportSender { /// 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.send_message(TunnelMessage::from(data)).await + } + + /// Send higher-layer data through the tunnel with sub-headers beside it + /// in the same Tunnel Data PDU. + /// + /// # Errors + /// + /// As [`Self::send`], for `message.data`, and `SubHeadersTooLarge` if + /// the sub-headers do not fit the one-byte `HeaderLength` field + /// ([MS-RDPEMT] 2.2.1.1), checked here for the same reason. + pub async fn send_message(&self, message: TunnelMessage) -> Result<(), UdpTransportError> { + if message.data.len() > usize::from(u16::MAX) { + debug!(len = message.data.len(), "Rejected oversized tunnel payload"); + return Err(UdpTransportError::payload_too_large("send", message.data.len())); + } + let sub_headers_len = message + .sub_headers + .iter() + .map(ironrdp_rdpemt::TunnelSubHeader::wire_size) + .sum::(); + if 4 /* RDP_TUNNEL_HEADER */ + sub_headers_len > usize::from(u8::MAX) { + return Err(UdpTransportError::new( + "send", + UdpTransportErrorKind::SubHeadersTooLarge { len: sub_headers_len }, + )); } self.0 - .send(data) + .send(message) .await .map_err(|_| UdpTransportError::driver("send", DriverError::connection_closed("send"))) } @@ -252,10 +300,10 @@ impl UdpTransportSender { /// Drop this handle to initiate shutdown of the background tasks. pub struct UdpTransport { /// Receives higher-layer data (DVC frames) from the tunnel. - data_rx: mpsc::Receiver>, + data_rx: mpsc::Receiver, /// Sends higher-layer data into the tunnel for encryption and transmission. - data_tx: mpsc::Sender>, + data_tx: mpsc::Sender, /// Shared I/O bridge between the driver and the TLS/RDPEMT layer. /// Held here so `shutdown()` can signal closure to both the @@ -275,8 +323,16 @@ pub struct UdpTransport { impl UdpTransport { /// Receive the next higher-layer data frame from the tunnel. /// - /// Returns `None` when the tunnel is closed. + /// Returns `None` when the tunnel is closed. Any sub-headers that came + /// with the frame are dropped; use [`Self::recv_message`] to keep them. pub async fn recv(&mut self) -> Option> { + self.recv_message().await.map(|message| message.data) + } + + /// Receive the next Tunnel Data PDU's content, sub-headers included. + /// + /// Returns `None` when the tunnel is closed. + pub async fn recv_message(&mut self) -> Option { self.data_rx.recv().await } @@ -296,6 +352,16 @@ impl UdpTransport { self.sender().send(data).await } + /// Send higher-layer data through the tunnel with sub-headers beside it + /// in the same Tunnel Data PDU. + /// + /// # Errors + /// + /// As [`UdpTransportSender::send_message`]. + pub async fn send_message(&self, message: TunnelMessage) -> Result<(), UdpTransportError> { + self.sender().send_message(message).await + } + /// Returns a cloneable handle for sending data, independent of this /// object's `&mut self`-requiring [`Self::recv`]. /// @@ -380,7 +446,7 @@ impl UdpTransport { /// For unit tests that exercise the channel-based API (FramedRead, /// FramedWrite) without needing a real UDP socket or TLS stack. #[cfg(test)] - pub(crate) fn from_channels(data_rx: mpsc::Receiver>, data_tx: mpsc::Sender>) -> Self { + pub(crate) fn from_channels(data_rx: mpsc::Receiver, data_tx: mpsc::Sender) -> Self { Self { data_rx, data_tx, @@ -548,8 +614,8 @@ pub async fn connect_udp(config: UdpTransportConfig) -> Result>(64); - let (outgoing_tx, mut outgoing_rx) = mpsc::channel::>(64); + let (incoming_tx, incoming_rx) = mpsc::channel::(64); + let (outgoing_tx, mut outgoing_rx) = mpsc::channel::(64); // Read pump: TLS → RDPEMT decode → channel let pump_handle = AbortOnDrop::new(tokio::spawn(async move { @@ -740,8 +806,8 @@ async fn accept_udp_inner(socket: UdpSocket, config: UdpAcceptConfig) -> Result< debug!("RDPEMT server tunnel established, starting data pump"); // Phase 6: Set up data channels and spawn pumps (identical to client side) - let (incoming_tx, incoming_rx) = mpsc::channel::>(64); - let (outgoing_tx, mut outgoing_rx) = mpsc::channel::>(64); + let (incoming_tx, incoming_rx) = mpsc::channel::(64); + let (outgoing_tx, mut outgoing_rx) = mpsc::channel::(64); let pump_handle = AbortOnDrop::new(tokio::spawn(async move { tunnel_data_loop(&mut tls_read, &mut tunnel, &incoming_tx).await @@ -818,15 +884,19 @@ where /// Shared by both `connect_udp` and `accept_udp`. Returns the first failure /// encountered rather than only logging it, so `shutdown()` can surface it /// to the caller instead of the pump silently going quiet. -async fn write_pump(tls_write: &mut W, outgoing_rx: &mut mpsc::Receiver>) -> Result<(), UdpTransportError> +async fn write_pump( + tls_write: &mut W, + outgoing_rx: &mut mpsc::Receiver, +) -> Result<(), UdpTransportError> where W: tokio::io::AsyncWrite + Unpin, { - while let Some(data) = outgoing_rx.recv().await { - let len = data.len(); + while let Some(message) = outgoing_rx.recv().await { + let len = message.data.len(); + let sub_headers = message.sub_headers.len(); let pdu = ironrdp_rdpemt::TunnelData { - sub_headers: Vec::new(), - higher_layer_data: data, + sub_headers: message.sub_headers, + higher_layer_data: message.data, }; let encoded = ironrdp_core::encode_vec(&pdu) .map_err(|error| UdpTransportError::rdpemt("write pump", ironrdp_rdpemt::RdpemtError::encode(error)))?; @@ -838,7 +908,7 @@ where debug!(%error, "Write pump failed to flush tunnel data"); UdpTransportError::tls("write pump", error) })?; - trace!(len, encoded_len = encoded.len(), "Sent tunnel data"); + trace!(len, sub_headers, encoded_len = encoded.len(), "Sent tunnel data"); } debug!("Write pump stopped, send channel closed"); Ok(()) @@ -917,8 +987,8 @@ mod tests { static DRIVER_RUNNING: AtomicBool = AtomicBool::new(false); DRIVER_RUNNING.store(false, Ordering::SeqCst); - let (_incoming_tx, incoming_rx) = mpsc::channel::>(4); - let (outgoing_tx, _outgoing_rx) = mpsc::channel::>(4); + let (_incoming_tx, incoming_rx) = mpsc::channel::(4); + let (outgoing_tx, _outgoing_rx) = mpsc::channel::(4); let transport = UdpTransport { data_rx: incoming_rx, @@ -960,7 +1030,7 @@ mod tests { let error = driver_exit_during_handshake("test", Ok(Err(driver_error))); assert!( - matches!(error.kind(), crate::error::UdpTransportErrorKind::Handshake(_)), + matches!(error.kind(), UdpTransportErrorKind::Handshake(_)), "got {error:?}, expected a Handshake error carrying the driver's own cause" ); } @@ -973,7 +1043,7 @@ mod tests { let error = driver_exit_during_handshake("test", Ok(Ok(()))); assert!( - matches!(error.kind(), crate::error::UdpTransportErrorKind::Handshake(_)), + matches!(error.kind(), UdpTransportErrorKind::Handshake(_)), "got {error:?}, expected a Handshake error for an unexpectedly-clean driver exit" ); } @@ -988,7 +1058,7 @@ mod tests { let error = driver_exit_during_handshake("test", join_result); assert!( - matches!(error.kind(), crate::error::UdpTransportErrorKind::DriverPanic), + matches!(error.kind(), UdpTransportErrorKind::DriverPanic), "got {error:?}, expected DriverPanic" ); } @@ -1031,8 +1101,90 @@ mod tests { .expect("the driver branch should resolve almost immediately, well inside 5 seconds"); assert!( - matches!(outcome.kind(), crate::error::UdpTransportErrorKind::Handshake(_)), + matches!(outcome.kind(), UdpTransportErrorKind::Handshake(_)), "got {outcome:?}, expected the driver's real error to surface" ); } + + /// A message's sub-headers go into the Tunnel Data PDU beside its data. + #[tokio::test] + async fn the_write_pump_encodes_sub_headers() { + use ironrdp_rdpemt::{SubHeaderType, TunnelData, TunnelSubHeader}; + use tokio::io::AsyncReadExt as _; + + let sub_header = TunnelSubHeader { + sub_header_type: SubHeaderType::AutoDetectRequest, + data: vec![0x06, 0x00, 0x07, 0x00, 0x14, 0x00], + }; + let (tx, mut rx) = mpsc::channel(4); + tx.send(TunnelMessage { + sub_headers: vec![sub_header.clone()], + data: vec![0xaa, 0xbb], + }) + .await + .expect("queue message"); + drop(tx); + + let (mut writer, mut reader) = tokio::io::duplex(1024); + write_pump(&mut writer, &mut rx).await.expect("pump drains"); + drop(writer); + + let mut wire = Vec::new(); + reader.read_to_end(&mut wire).await.expect("read back"); + let pdu: TunnelData = ironrdp_core::decode(&wire).expect("decode"); + assert_eq!(pdu.sub_headers, vec![sub_header]); + assert_eq!(pdu.higher_layer_data, vec![0xaa, 0xbb]); + } + + /// `recv` keeps its old shape and drops sub-headers; `recv_message` keeps them. + #[tokio::test] + async fn recv_message_keeps_what_recv_drops() { + use ironrdp_rdpemt::{SubHeaderType, TunnelSubHeader}; + + let (incoming_tx, incoming_rx) = mpsc::channel(4); + let (outgoing_tx, _outgoing_rx) = mpsc::channel(4); + let mut transport = UdpTransport::from_channels(incoming_rx, outgoing_tx); + let message = TunnelMessage { + sub_headers: vec![TunnelSubHeader { + sub_header_type: SubHeaderType::AutoDetectResponse, + data: vec![0x01], + }], + data: vec![0x02], + }; + + incoming_tx.send(message.clone()).await.expect("queue"); + incoming_tx.send(message.clone()).await.expect("queue"); + + assert_eq!(transport.recv().await, Some(vec![0x02])); + assert_eq!(transport.recv_message().await, Some(message)); + } + + /// Sub-headers that would overflow `HeaderLength` are refused at the + /// call, not left to fail the write pump; the largest that fit go out. + #[tokio::test] + async fn send_message_refuses_sub_headers_the_header_cannot_hold() { + use ironrdp_rdpemt::{SubHeaderType, TunnelSubHeader}; + + let (_incoming_tx, incoming_rx) = mpsc::channel(4); + let (outgoing_tx, mut outgoing_rx) = mpsc::channel(4); + let transport = UdpTransport::from_channels(incoming_rx, outgoing_tx); + let message = |data_len| TunnelMessage { + sub_headers: vec![TunnelSubHeader { + sub_header_type: SubHeaderType::AutoDetectRequest, + data: vec![0; data_len], + }], + data: Vec::new(), + }; + + // 4 (tunnel header) + 2 (sub-header header) + 249 = 255. + transport.send_message(message(249)).await.expect("fits"); + let error = transport.send_message(message(250)).await.expect_err("one byte over"); + assert!(matches!( + error.kind(), + UdpTransportErrorKind::SubHeadersTooLarge { len: 252 } + )); + + assert_eq!(outgoing_rx.recv().await, Some(message(249))); + assert!(outgoing_rx.try_recv().is_err(), "the refused message was not queued"); + } } diff --git a/crates/ironrdp-rdpeudp-tokio/src/tunnel.rs b/crates/ironrdp-rdpeudp-tokio/src/tunnel.rs index 89c682bcbe..654c66b4ab 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/tunnel.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/tunnel.rs @@ -15,6 +15,7 @@ use tokio::io::{AsyncRead, AsyncReadExt as _, AsyncWrite, AsyncWriteExt as _}; use tracing::{debug, trace}; use crate::error::{UdpTransportError, UdpTransportErrorExt as _}; +use crate::transport::TunnelMessage; /// Read a complete RDPEMT PDU from the stream using self-framing. /// @@ -128,7 +129,7 @@ where pub(crate) async fn tunnel_data_loop( stream: &mut S, tunnel: &mut RdpemtTunnel, - data_tx: &tokio::sync::mpsc::Sender>, + data_tx: &tokio::sync::mpsc::Sender, ) -> Result<(), UdpTransportError> where S: AsyncRead + Unpin, @@ -154,13 +155,16 @@ where while let Some(event) = tunnel.poll_event() { match event { - // `sub_headers` (e.g. auto-detect bandwidth measurement, MS-RDPBCGR - // 2.2.14) are not consumed here; this driver only wires the DVC - // payload through. A future auto-detect integration would need to - // dispatch them instead of discarding them. - TunnelEvent::Data { data, .. } => { - trace!(len = data.len(), "Forwarding tunnel data"); - if data_tx.send(data).await.is_err() { + // Sub-headers go through with the data: they carry the + // auto-detect messages ([MS-RDPBCGR] 1.3.9) of a sideband + // channel in use, which the application answers. + TunnelEvent::Data { sub_headers, data } => { + trace!( + len = data.len(), + sub_headers = sub_headers.len(), + "Forwarding tunnel data" + ); + if data_tx.send(TunnelMessage { sub_headers, data }).await.is_err() { debug!("Tunnel data receiver dropped, stopping read pump"); // Application dropped the receiver return Ok(()); @@ -263,4 +267,43 @@ mod tests { let mut cursor = io::Cursor::new(vec![0x02, 0x05, 0x00, 0x04, 0x48, 0x65]); assert!(read_tunnel_pdu(&mut cursor).await.is_err()); } + + /// Sub-headers reach the application with the data they came with. They + /// carry the auto-detect messages of a sideband channel in use + /// ([MS-RDPBCGR] 1.3.9), which used to be dropped here. + #[tokio::test] + async fn the_data_loop_forwards_sub_headers() { + use ironrdp_rdpemt::{SubHeaderType, TunnelConfig, TunnelCreateResponse, TunnelData, TunnelSubHeader}; + + let mut tunnel = RdpemtTunnel::client(TunnelConfig { + request_id: 1, + security_cookie: [0; 16], + }); + let response = ironrdp_core::encode_vec(&TunnelCreateResponse { + hr_response: TunnelCreateResponse::S_OK, + }) + .expect("encode response"); + tunnel.handle_pdu(&response).expect("tunnel established"); + while tunnel.poll_event().is_some() {} + + let sub_header = TunnelSubHeader { + sub_header_type: SubHeaderType::AutoDetectResponse, + data: vec![0x0e, 0x00, 0x05, 0x00], + }; + let wire = ironrdp_core::encode_vec(&TunnelData { + sub_headers: vec![sub_header.clone()], + higher_layer_data: Vec::new(), + }) + .expect("encode data"); + + let (tx, mut rx) = tokio::sync::mpsc::channel(4); + let mut stream = io::Cursor::new(wire); + tunnel_data_loop(&mut stream, &mut tunnel, &tx) + .await + .expect("clean EOF"); + + let message = rx.recv().await.expect("one message"); + assert_eq!(message.sub_headers, vec![sub_header]); + assert!(message.data.is_empty()); + } } diff --git a/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs b/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs index b207ba923e..421f38e97a 100644 --- a/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs +++ b/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs @@ -11,10 +11,11 @@ use core::sync::atomic::{AtomicUsize, Ordering}; use core::time::Duration; use std::sync::Arc; -use ironrdp_rdpemt::TunnelConfig; +use ironrdp_rdpemt::{SubHeaderType, TunnelConfig, TunnelSubHeader}; use ironrdp_rdpeudp::ConnectionConfig; use ironrdp_rdpeudp_tokio::{ - MultitransportBootstrap, UdpAcceptConfig, UdpTlsConfig, UdpTransport, UdpTransportConfig, accept_udp, connect_udp, + MultitransportBootstrap, TunnelMessage, UdpAcceptConfig, UdpTlsConfig, UdpTransport, UdpTransportConfig, + accept_udp, connect_udp, }; use ironrdp_tls::{CertificateValidation, CertificateValidationCallback}; use tokio::net::UdpSocket; @@ -222,6 +223,42 @@ async fn full_stack_bidirectional_data() { server.shutdown().await.expect("server shutdown"); } +/// Sub-headers cross the full stack beside their data, and a message may +/// carry sub-headers alone, as an auto-detect Start or Stop does. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn full_stack_sub_headers() { + let (mut client, server) = establish_loopback_pair_on("127.0.0.1:0").await; + + let start = TunnelMessage { + sub_headers: vec![TunnelSubHeader { + sub_header_type: SubHeaderType::AutoDetectRequest, + data: vec![0x06, 0x00, 0x07, 0x00, 0x14, 0x00], + }], + data: Vec::new(), + }; + let graphics = TunnelMessage { + sub_headers: vec![TunnelSubHeader { + sub_header_type: SubHeaderType::AutoDetectRequest, + data: vec![0x08], + }], + data: vec![0x01, 0x02, 0x03], + }; + server + .send_message(start.clone()) + .await + .expect("send sub-headers alone"); + server + .send_message(graphics.clone()) + .await + .expect("send sub-headers with data"); + + assert_eq!(client.recv_message().await.expect("recv"), start); + assert_eq!(client.recv_message().await.expect("recv"), graphics); + + client.shutdown().await.expect("client shutdown"); + server.shutdown().await.expect("server shutdown"); +} + /// A payload over the wire `PayloadLength` field's 65535-byte capacity must /// be rejected synchronously by `send()`, and must not take the write pump /// down with it: a normal-sized send right after must still succeed. From 505914bb67099333a6ced623363fa60110c5286b Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Mon, 28 Sep 2026 01:26:57 -0500 Subject: [PATCH 2/4] test(rdpeudp): use a real Bandwidth Measure Start's sub-header data The sub-header examples carried a whole Bandwidth Measure Start, header bytes included. An auto-detect sub-header is the auto-detect structure itself (MS-RDPEMT 2.2.1.1.1): its SubHeaderLength and SubHeaderType are the structure's headerLength and headerTypeId, so on the wire those examples would repeat the header. The transport treats the data as opaque and the tests passed either way; the examples now carry only what follows the header. --- crates/ironrdp-rdpeudp-tokio/src/transport.rs | 4 +++- crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs | 4 +++- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/crates/ironrdp-rdpeudp-tokio/src/transport.rs b/crates/ironrdp-rdpeudp-tokio/src/transport.rs index f3d42d184d..1ceb0b9d88 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/transport.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/transport.rs @@ -1114,7 +1114,9 @@ mod tests { let sub_header = TunnelSubHeader { sub_header_type: SubHeaderType::AutoDetectRequest, - data: vec![0x06, 0x00, 0x07, 0x00, 0x14, 0x00], + // A Bandwidth Measure Start: sequenceNumber 7, requestType 0x0014. The + // sub-header's own two bytes are its headerLength and headerTypeId. + data: vec![0x07, 0x00, 0x14, 0x00], }; let (tx, mut rx) = mpsc::channel(4); tx.send(TunnelMessage { diff --git a/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs b/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs index 421f38e97a..a0069ef8c0 100644 --- a/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs +++ b/crates/ironrdp-testsuite-extra/tests/rdpeudp_tokio.rs @@ -232,7 +232,9 @@ async fn full_stack_sub_headers() { let start = TunnelMessage { sub_headers: vec![TunnelSubHeader { sub_header_type: SubHeaderType::AutoDetectRequest, - data: vec![0x06, 0x00, 0x07, 0x00, 0x14, 0x00], + // A Bandwidth Measure Start: sequenceNumber 7, requestType 0x0014. The + // sub-header's own two bytes are its headerLength and headerTypeId. + data: vec![0x07, 0x00, 0x14, 0x00], }], data: Vec::new(), }; From 72705d90be8831ff84a5e1ff5764851bf8e9746e Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Mon, 28 Sep 2026 09:58:20 -0500 Subject: [PATCH 3/4] review: drop the unused Default derive and correct the FramedRead comment Nothing builds a default TunnelMessage, and a default one would be an empty Tunnel Data PDU. The FramedRead comment said the tunnel consumes the sub-headers before this point; recv drops them, and a caller that needs them uses recv_message. --- crates/ironrdp-rdpeudp-tokio/src/framed.rs | 6 +++--- crates/ironrdp-rdpeudp-tokio/src/transport.rs | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/ironrdp-rdpeudp-tokio/src/framed.rs b/crates/ironrdp-rdpeudp-tokio/src/framed.rs index 4fff17426d..ea74fbcd8f 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/framed.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/framed.rs @@ -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() => { diff --git a/crates/ironrdp-rdpeudp-tokio/src/transport.rs b/crates/ironrdp-rdpeudp-tokio/src/transport.rs index 1ceb0b9d88..f171da5928 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/transport.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/transport.rs @@ -221,7 +221,7 @@ impl Drop for AbortOnDrop { /// [MS-RDPBCGR] 1.3.9 sends over a sideband channel in use, such as a /// bandwidth measurement's Start and Stop and the client's results. `data` /// may be empty when a message carries only sub-headers. -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, PartialEq, Eq)] pub struct TunnelMessage { pub sub_headers: Vec, pub data: Vec, From efe177750c9b0a40c00e944e13edf448602fb3d6 Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Tue, 29 Sep 2026 16:22:06 -0500 Subject: [PATCH 4/4] review: pin the send_message header bound to the TunnelData encoder The send_message precheck restates the HeaderLength bound that ironrdp-rdpemt's TunnelData encoder owns. The boundary test now also runs the largest accepted message and the smallest refused one through the real encoder, so the two cannot drift apart unnoticed, and the comment on the check cites MS-RDPEMT 2.2.1.1. --- crates/ironrdp-rdpeudp-tokio/src/transport.rs | 24 +++++++++++++++++-- 1 file changed, 22 insertions(+), 2 deletions(-) diff --git a/crates/ironrdp-rdpeudp-tokio/src/transport.rs b/crates/ironrdp-rdpeudp-tokio/src/transport.rs index f171da5928..0a2b86be15 100644 --- a/crates/ironrdp-rdpeudp-tokio/src/transport.rs +++ b/crates/ironrdp-rdpeudp-tokio/src/transport.rs @@ -277,6 +277,9 @@ impl UdpTransportSender { .iter() .map(ironrdp_rdpemt::TunnelSubHeader::wire_size) .sum::(); + // HeaderLength is one byte and counts the 4 byte fixed part of the tunnel + // header plus the sub-headers ([MS-RDPEMT] 2.2.1.1). The write pump's + // TunnelData encoder enforces the same bound; a test pins the two together. if 4 /* RDP_TUNNEL_HEADER */ + sub_headers_len > usize::from(u8::MAX) { return Err(UdpTransportError::new( "send", @@ -1162,10 +1165,12 @@ mod tests { } /// Sub-headers that would overflow `HeaderLength` are refused at the - /// call, not left to fail the write pump; the largest that fit go out. + /// call, not left to fail the write pump; the largest that fit go out. The + /// boundary is also checked against the real `TunnelData` encoder, so the + /// precheck cannot drift from the bound the write pump enforces. #[tokio::test] async fn send_message_refuses_sub_headers_the_header_cannot_hold() { - use ironrdp_rdpemt::{SubHeaderType, TunnelSubHeader}; + use ironrdp_rdpemt::{SubHeaderType, TunnelData, TunnelSubHeader}; let (_incoming_tx, incoming_rx) = mpsc::channel(4); let (outgoing_tx, mut outgoing_rx) = mpsc::channel(4); @@ -1188,5 +1193,20 @@ mod tests { assert_eq!(outgoing_rx.recv().await, Some(message(249))); assert!(outgoing_rx.try_recv().is_err(), "the refused message was not queued"); + + let encode = |message: TunnelMessage| { + ironrdp_core::encode_vec(&TunnelData { + sub_headers: message.sub_headers, + higher_layer_data: message.data, + }) + }; + assert!( + encode(message(249)).is_ok(), + "the encoder accepts what the precheck accepts" + ); + assert!( + encode(message(250)).is_err(), + "the encoder refuses what the precheck refuses" + ); } }