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..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() => { @@ -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..0a2b86be15 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, 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,39 @@ 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::(); + // 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", + UdpTransportErrorKind::SubHeadersTooLarge { len: sub_headers_len }, + )); } self.0 - .send(data) + .send(message) .await .map_err(|_| UdpTransportError::driver("send", DriverError::connection_closed("send"))) } @@ -252,10 +303,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 +326,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 +355,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 +449,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 +617,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 +809,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 +887,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 +911,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 +990,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 +1033,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 +1046,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 +1061,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 +1104,109 @@ 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, + // 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 { + 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. 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, TunnelData, 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"); + + 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" + ); + } } 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..a0069ef8c0 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,44 @@ 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, + // 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(), + }; + 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.