From 69cb54ae7588ca1dfa57fe78ffa1349ba247f595 Mon Sep 17 00:00:00 2001 From: AKolenda Date: Fri, 25 Sep 2026 23:41:09 -0600 Subject: [PATCH 1/4] feat(session): answer bandwidth measurements during the session The session answered RTT requests on the message channel but ignored Bandwidth Measure Start and Stop, so a server running continuous network auto-detection (MS-RDPBCGR 2.2.14, 3.2.5.14) never got a Bandwidth Measure Results PDU back. Windows runs continuous detection by default. The x224 processor now keeps a bandwidth window: - A continuous Start (0x0014) opens a window, and every server byte received until the Stop (0x0429) is counted. Following 3.2.5.14, only the bytes after a Basic Security Header count on the message channel. On the IO and static channels the MCS user data is counted, and on the fast path the data after the fast-path header, which is what FreeRDP counts too. - A connect-time Start (0x1014) counts Bandwidth Measure Payload PDUs and the 0x002B Stop the way the connector does: payloadLength plus the eight header bytes. - A repeated Start restarts the window, and the Results report the elapsed time floored at 1 ms (the server divides by it) and the byte count. - Lossy (0x0114/0x0629) requests belong to a lossy tunnel and are not answered on this channel. Timing follows the connector: the caller passes the frame's arrival time through the new `ActiveStage::process_with_timestamp` and `x224::Processor::process_with_timestamp`, and the client passes the time its framed reader filled the buffer. Without a timestamp no window opens and a Stop is answered with a zero-byte, 1 ms result instead of a figure the client never measured. `process` keeps working that way. --- crates/ironrdp-client/src/rdp.rs | 3 +- crates/ironrdp-session/src/active_stage.rs | 32 +++- crates/ironrdp-session/src/x224/mod.rs | 159 +++++++++++++++--- .../tests/session/autodetect.rs | 126 +++++++++++++- 4 files changed, 288 insertions(+), 32 deletions(-) diff --git a/crates/ironrdp-client/src/rdp.rs b/crates/ironrdp-client/src/rdp.rs index a67d49c9af..44061e2792 100644 --- a/crates/ironrdp-client/src/rdp.rs +++ b/crates/ironrdp-client/src/rdp.rs @@ -3165,7 +3165,8 @@ async fn active_session( Err(error) => return Err(ironrdp_session::custom_err!("read frame", error)), }; trace!(?action, frame_length = payload.len(), "Frame received"); - let mut outputs = active_stage.process(&mut image, action, &payload)?; + let mut outputs = + active_stage.process_with_timestamp(&mut image, action, &payload, reader.last_read_at())?; #[cfg(feature = "rdpdr")] if let Some(output) = poll_deferred_rdpdr_output(&mut active_stage)? { outputs.push(output); diff --git a/crates/ironrdp-session/src/active_stage.rs b/crates/ironrdp-session/src/active_stage.rs index f53e0aedac..5d815bc573 100644 --- a/crates/ironrdp-session/src/active_stage.rs +++ b/crates/ironrdp-session/src/active_stage.rs @@ -1,12 +1,13 @@ use std::sync::Arc; use ironrdp_bulk::{BulkCompressor, CompressionType as BulkCompressionType}; -use ironrdp_core::{ReadCursor, WriteBuf}; +use ironrdp_core::{Decode as _, MonotonicInstant, ReadCursor, WriteBuf}; use ironrdp_displaycontrol::client::DisplayControlClient; use ironrdp_dvc::pdu::SoftSyncTunnelType; use ironrdp_dvc::{DrdynvcClient, DvcClientProcessor, DvcMessageBatch, DynamicChannelMut, DynamicChannelRef}; use ironrdp_egfx::client::GraphicsPipelineClient; use ironrdp_graphics::pointer::DecodedPointer; +use ironrdp_pdu::fast_path::FastPathHeader; use ironrdp_pdu::gcc::{ChannelName, Monitor}; use ironrdp_pdu::geometry::{ExclusiveRectangle, InclusiveRectangle, Rectangle as _}; use ironrdp_pdu::input::fast_path::{FastPathInput, FastPathInputEvent}; @@ -183,15 +184,38 @@ impl ActiveStage { } /// Process a frame received from the server. + /// + /// Without an arrival time, bandwidth measurements are answered with an untimed, + /// zero-byte result; see [`Self::process_with_timestamp`]. pub fn process( &mut self, image: &mut DecodedImage, action: Action, frame: &[u8], + ) -> SessionResult> { + self.process_with_timestamp(image, action, frame, None) + } + + /// Process a frame received from the server, together with its arrival time. + /// + /// The clock stays outside this state machine: `received_at` is the time the transport + /// read the frame, from the same monotonic clock for every frame. Frames that were + /// buffered must keep their read time so a bandwidth measurement reflects network arrival + /// rather than the time spent decoding earlier frames. + pub fn process_with_timestamp( + &mut self, + image: &mut DecodedImage, + action: Action, + frame: &[u8], + received_at: Option, ) -> SessionResult> { self.damage_regions.clear(); let (mut stage_outputs, processor_updates) = match action { Action::FastPath => { + // A continuous bandwidth measurement counts what follows the fast-path header. + let mut header = ReadCursor::new(frame); + FastPathHeader::decode(&mut header).map_err(SessionError::decode)?; + self.x224_processor.record_bandwidth_bytes(header.len()); let mut output = WriteBuf::new(); let processor_updates = self.fast_path_processor @@ -202,7 +226,9 @@ impl ActiveStage { ) } Action::X224 => { - let x224_outputs = self.x224_processor.process(frame, &mut self.bulk_decompressor)?; + let x224_outputs = + self.x224_processor + .process_with_timestamp(frame, &mut self.bulk_decompressor, received_at)?; let mut stage_outputs = Vec::new(); let mut processor_updates = Vec::new(); @@ -1040,7 +1066,7 @@ mod tests { use core::any::TypeId; use super::*; - use ironrdp_core::{Decode as _, encode_vec}; + use ironrdp_core::encode_vec; use ironrdp_displaycontrol::pdu::{DisplayControlCapabilities, DisplayControlPdu}; use ironrdp_dvc::pdu::{ CreateRequestPdu, DataPdu, DrdynvcDataPdu, DrdynvcServerPdu, SoftSyncChannelList, SoftSyncRequestPdu, diff --git a/crates/ironrdp-session/src/x224/mod.rs b/crates/ironrdp-session/src/x224/mod.rs index d006393ba4..08f3076807 100644 --- a/crates/ironrdp-session/src/x224/mod.rs +++ b/crates/ironrdp-session/src/x224/mod.rs @@ -1,9 +1,12 @@ use ironrdp_bulk::BulkCompressor; -use ironrdp_core::{Decode as _, ReadCursor, WriteBuf, decode}; +use ironrdp_core::{Decode as _, MonotonicInstant, ReadCursor, WriteBuf, decode}; use ironrdp_dvc::{DrdynvcClient, DvcClientProcessor, DynamicChannelMut, DynamicChannelRef}; use ironrdp_pdu::gcc::{ChannelName, Monitor}; use ironrdp_pdu::mcs::{DisconnectProviderUltimatum, DisconnectReason, McsMessage, SendDataIndicationCtx}; -use ironrdp_pdu::rdp::autodetect::{AutoDetectReqPdu, AutoDetectRequest, AutoDetectResponse, AutoDetectRspPdu}; +use ironrdp_pdu::rdp::autodetect::{ + AutoDetectReqPdu, AutoDetectRequest, AutoDetectResponse, AutoDetectRspPdu, BW_RESULTS_CONNECT_TIME, + BW_RESULTS_CONTINUOUS, BW_START_CONNECT_TIME, BW_START_RELIABLE_UDP, BW_STOP_CONNECT_TIME, BW_STOP_RELIABLE_UDP, +}; use ironrdp_pdu::rdp::client_info::CompressionType; use ironrdp_pdu::rdp::headers::{ BasicSecurityHeader, BasicSecurityHeaderFlags, CompressionFlags, IoChannelPdu, ShareDataCtx, ShareDataPdu, @@ -72,7 +75,7 @@ pub enum ProcessorOutput { /// Auto-detect network characteristics from server ([\[MS-RDPBCGR\] 2.2.14]). /// /// Currently only surfaces [`AutoDetectRequest::NetworkCharacteristicsResult`]. - /// RTT requests are handled internally with automatic responses. + /// RTT and bandwidth measurement requests are answered internally. /// /// [\[MS-RDPBCGR\] 2.2.14]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpbcgr/dc672839-4f4e-40b1-a71c-cd6a959baa38 AutoDetect(AutoDetectRequest), @@ -107,6 +110,13 @@ pub struct Processor { io_channel_id: u16, message_channel_id: Option, share_id: u32, + bandwidth: Option, +} + +struct BandwidthMeasurement { + started_at: MonotonicInstant, + bytes: u32, + continuous: bool, } impl Processor { @@ -123,6 +133,7 @@ impl Processor { io_channel_id, message_channel_id, share_id, + bandwidth: None, } } @@ -212,6 +223,19 @@ impl Processor { &mut self, frame: &[u8], bulk_decompressor: &mut Option, + ) -> SessionResult> { + self.process_with_timestamp(frame, bulk_decompressor, None) + } + + /// Processes a frame with its driver-observed arrival time. + /// + /// Supply the same monotonic clock for every frame. Without a timestamp, + /// bandwidth replies use a conservative zero-byte measurement. + pub fn process_with_timestamp( + &mut self, + frame: &[u8], + bulk_decompressor: &mut Option, + received_at: Option, ) -> SessionResult> { let data_ctx: SendDataIndicationCtx<'_> = match ironrdp_pdu::mcs::decode_send_data_indication(frame) { Ok(data_ctx) => data_ctx, @@ -232,10 +256,15 @@ impl Processor { }; let channel_id = data_ctx.channel_id; + // Message-channel PDUs carry a Basic Security Header, excluded below. + // Ordinary session data on TLS-protected IO/SVC channels does not. + if self.message_channel_id != Some(channel_id) { + self.record_bandwidth_bytes(data_ctx.user_data.len()); + } if channel_id == self.io_channel_id { self.process_io_channel_data_indication(data_ctx, bulk_decompressor) } else if self.message_channel_id == Some(channel_id) { - self.process_message_channel(data_ctx) + self.process_message_channel(data_ctx, received_at) } else { let maximum_chunk_size = self.static_channels.maximum_chunk_size(); if let Some(svc) = self.static_channels.get_by_channel_id_mut(channel_id) { @@ -445,6 +474,16 @@ impl Processor { Ok(decompressed) } + /// Counts session bytes after transport/security headers while a continuous + /// bandwidth window is open ([MS-RDPBCGR] 3.2.5.14). + pub(crate) fn record_bandwidth_bytes(&mut self, bytes: usize) { + if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| measurement.continuous) { + measurement.bytes = measurement + .bytes + .saturating_add(u32::try_from(bytes).unwrap_or(u32::MAX)); + } + } + /// Process a PDU received on the MCS message channel: auto-detect /// ([MS-RDPBCGR] 2.2.14), multitransport ([MS-RDPBCGR] 2.2.15), or /// Heartbeat ([MS-RDPBCGR] 2.2.16.1). @@ -457,13 +496,18 @@ impl Processor { /// session-fatal decode error: this channel is forward-safe for future /// message-channel PDU types the same way the connect-time demux /// (`ironrdp-connector`) already is. - fn process_message_channel(&self, data_ctx: SendDataIndicationCtx<'_>) -> SessionResult> { + fn process_message_channel( + &mut self, + data_ctx: SendDataIndicationCtx<'_>, + received_at: Option, + ) -> SessionResult> { let Some(message_channel_id) = self.message_channel_id else { return Err(reason_err!("message channel", "no message channel negotiated")); }; let mut peek = ReadCursor::new(data_ctx.user_data); let security_header = BasicSecurityHeader::decode(&mut peek).map_err(SessionError::decode)?; + self.record_bandwidth_bytes(peek.len()); let flags = security_header .flags .difference(BasicSecurityHeaderFlags::RESET_SEQNO | BasicSecurityHeaderFlags::IGNORE_SEQNO); @@ -495,29 +539,87 @@ impl Processor { let req = decode::(data_ctx.user_data).map_err(SessionError::decode)?; - match req.request { + let response = match req.request { AutoDetectRequest::RttRequest { sequence_number, .. } => { - let response = AutoDetectRspPdu::new(AutoDetectResponse::RttResponse { sequence_number }); - let mut frame = WriteBuf::new(); - ironrdp_pdu::mcs::encode_send_data_request( - self.user_channel_id, - message_channel_id, - &response, - &mut frame, - ) - .map_err(SessionError::encode)?; - debug!(sequence_number, "Responded to auto-detect RTT request"); - Ok(vec![ProcessorOutput::ResponseFrame(frame.into_inner())]) + AutoDetectResponse::RttResponse { sequence_number } + } + AutoDetectRequest::BandwidthMeasureStart { request_type, .. } + if matches!(request_type, BW_START_CONNECT_TIME | BW_START_RELIABLE_UDP) => + { + self.bandwidth = received_at.map(|started_at| BandwidthMeasurement { + started_at, + bytes: 0, + continuous: request_type == BW_START_RELIABLE_UDP, + }); + return Ok(Vec::new()); + } + AutoDetectRequest::BandwidthMeasurePayload { payload, .. } => { + if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| !measurement.continuous) { + // [MS-RDPBCGR] 3.2.5.14 counts the eight-byte auto-detect + // header as well as payloadLength, but not the security header. + measurement.bytes = measurement + .bytes + .saturating_add(u32::try_from(payload.len()).unwrap_or(u32::MAX).saturating_add(8)); + } + return Ok(Vec::new()); + } + AutoDetectRequest::BandwidthMeasureStop { + sequence_number, + request_type, + payload, + } if matches!(request_type, BW_STOP_CONNECT_TIME | BW_STOP_RELIABLE_UDP) => { + let continuous = request_type == BW_STOP_RELIABLE_UDP; + let measurement = self + .bandwidth + .take() + .filter(|measurement| measurement.continuous == continuous); + let stop_bytes = if continuous { + 0 + } else { + u32::try_from(payload.as_ref().map_or(0, Vec::len)) + .unwrap_or(u32::MAX) + .saturating_add(8) + }; + let (time_delta_ms, byte_count) = match (measurement, received_at) { + (Some(measurement), Some(stopped_at)) => ( + u32::try_from(stopped_at.duration_since(measurement.started_at).as_millis()) + .unwrap_or(u32::MAX) + .max(1), + measurement.bytes.saturating_add(stop_bytes), + ), + _ => (1, stop_bytes), + }; + AutoDetectResponse::BandwidthMeasureResults { + sequence_number, + response_type: if continuous { + BW_RESULTS_CONTINUOUS + } else { + BW_RESULTS_CONNECT_TIME + }, + time_delta_ms, + byte_count, + } } req @ AutoDetectRequest::NetworkCharacteristicsResult { .. } => { debug!(?req, "Received network characteristics from server"); - Ok(vec![ProcessorOutput::AutoDetect(req)]) + return Ok(vec![ProcessorOutput::AutoDetect(req)]); } req => { - debug!(?req, "Auto-detect request not yet implemented"); - Ok(Vec::new()) + // Lossy variants belong to a lossy tunnel, never this channel. + debug!(?req, "Ignoring auto-detect request for another transport"); + return Ok(Vec::new()); } - } + }; + debug!(?response, "Responding to an auto-detect request"); + let mut frame = WriteBuf::new(); + ironrdp_pdu::mcs::encode_send_data_request( + self.user_channel_id, + message_channel_id, + &AutoDetectRspPdu::new(response), + &mut frame, + ) + .map_err(SessionError::encode)?; + Ok(vec![ProcessorOutput::ResponseFrame(frame.into_inner())]) } /// Encodes an Initiate Multitransport Response on the MCS message channel. @@ -638,14 +740,17 @@ mod tests { fn processor_surfaces_multitransport_request_on_message_channel() { let request = multitransport_request(); let encoded = encode_vec(&request).expect("encode multitransport request"); - let processor = Processor::new(StaticChannelSet::new(), 1002, 1003, Some(1004), 0); + let mut processor = Processor::new(StaticChannelSet::new(), 1002, 1003, Some(1004), 0); let outputs = processor - .process_message_channel(SendDataIndicationCtx { - initiator_id: 1002, - channel_id: 1004, - user_data: &encoded, - }) + .process_message_channel( + SendDataIndicationCtx { + initiator_id: 1002, + channel_id: 1004, + user_data: &encoded, + }, + None, + ) .expect("surface multitransport request"); assert!(matches!( diff --git a/crates/ironrdp-testsuite-core/tests/session/autodetect.rs b/crates/ironrdp-testsuite-core/tests/session/autodetect.rs index 4591f3fe92..00f93d0285 100644 --- a/crates/ironrdp-testsuite-core/tests/session/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/session/autodetect.rs @@ -130,7 +130,7 @@ fn bandwidth_measure_stop_does_not_crash() { let frame = encode_server_autodetect(request); let outputs = process_frame(&mut processor, &frame); - assert!(outputs.is_empty(), "BW stop should produce no output"); + assert_eq!(bandwidth_result(&outputs), (200, 0x000b, 1, 0)); } #[test] @@ -142,3 +142,127 @@ fn bandwidth_measure_payload_does_not_crash() { let outputs = process_frame(&mut processor, &frame); assert!(outputs.is_empty(), "BW payload should produce no output"); } + +fn timed_request( + processor: &mut Processor, + request: AutoDetectRequest, + millis: u64, +) -> Vec { + processor + .process_with_timestamp( + &encode_server_autodetect(request), + &mut None, + Some(ironrdp_core::MonotonicInstant::from_millis(millis)), + ) + .expect("process timed auto-detect request") +} + +fn bandwidth_result(outputs: &[ironrdp_session::x224::ProcessorOutput]) -> (u16, u16, u32, u32) { + let [ironrdp_session::x224::ProcessorOutput::ResponseFrame(frame)] = outputs else { + panic!("expected exactly one bandwidth response"); + }; + let X224(McsMessage::SendDataRequest(message)) = ironrdp_core::decode::>>(frame).unwrap() + else { + panic!("expected main-channel response"); + }; + assert_eq!(message.channel_id, MESSAGE_CHANNEL_ID); + let response = ironrdp_core::decode::(&message.user_data).unwrap(); + let AutoDetectResponse::BandwidthMeasureResults { + sequence_number, + response_type, + time_delta_ms, + byte_count, + } = response.response + else { + panic!("expected bandwidth results"); + }; + (sequence_number, response_type, time_delta_ms, byte_count) +} + +#[test] +fn continuous_measurement_counts_data_without_security_headers() { + let mut processor = make_processor(); + assert!(timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 10).is_empty()); + timed_request(&mut processor, AutoDetectRequest::rtt_continuous(2), 20); + let result = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(3), 35); + // RTT and Stop have six-byte auto-detect headers. Neither four-byte + // security header, nor TPKT/X224/MCS framing, belongs to the count. + assert_eq!(bandwidth_result(&result), (3, 0x000b, 25, 12)); +} + +#[test] +fn connect_time_measurement_includes_payload_headers_once() { + let mut processor = make_processor(); + timed_request(&mut processor, AutoDetectRequest::bw_start_connect_time(1), 100); + timed_request(&mut processor, AutoDetectRequest::bw_payload(2, vec![0xaa; 64]), 110); + timed_request(&mut processor, AutoDetectRequest::rtt_continuous(3), 112); + let result = timed_request( + &mut processor, + AutoDetectRequest::bw_stop_connect_time(4, vec![0xbb; 16]), + 140, + ); + assert_eq!(bandwidth_result(&result), (4, 0x0003, 40, 64 + 8 + 16 + 8)); +} + +#[test] +fn repeated_start_resets_count_and_timer_and_stop_ends_window() { + let mut processor = make_processor(); + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 100); + timed_request(&mut processor, AutoDetectRequest::rtt_continuous(2), 110); + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(3), 200); + let result = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(4), 210); + assert_eq!(bandwidth_result(&result), (4, 0x000b, 10, 6)); + let repeated = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(5), 220); + assert_eq!(bandwidth_result(&repeated), (5, 0x000b, 1, 0)); +} + +#[test] +fn measurement_timing_saturates_and_never_reports_zero_divisor() { + for (start, stop, expected) in [(100, 100, 1), (100, 90, 1), (0, u64::MAX, u32::MAX)] { + let mut processor = make_processor(); + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), start); + let result = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(2), stop); + assert_eq!(bandwidth_result(&result), (2, 0x000b, expected, 6)); + } +} + +#[test] +fn lossy_requests_on_main_channel_are_not_answered() { + use ironrdp_pdu::rdp::autodetect::{BW_START_LOSSY_UDP, BW_STOP_LOSSY_UDP}; + let mut processor = make_processor(); + assert!( + timed_request( + &mut processor, + AutoDetectRequest::BandwidthMeasureStart { + sequence_number: 1, + request_type: BW_START_LOSSY_UDP, + }, + 10 + ) + .is_empty() + ); + assert!( + timed_request( + &mut processor, + AutoDetectRequest::BandwidthMeasureStop { + sequence_number: 1, + request_type: BW_STOP_LOSSY_UDP, + payload: None, + }, + 20 + ) + .is_empty() + ); +} + +#[test] +fn untimed_driver_does_not_report_accumulated_bytes_as_a_real_measurement() { + let mut processor = make_processor(); + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 10); + timed_request(&mut processor, AutoDetectRequest::rtt_continuous(2), 20); + let outputs = process_frame( + &mut processor, + &encode_server_autodetect(AutoDetectRequest::bw_stop_continuous(3)), + ); + assert_eq!(bandwidth_result(&outputs), (3, 0x000b, 1, 0)); +} From ec764862def3b5c04454e54fc17dd4c7a0fc4f05 Mon Sep 17 00:00:00 2001 From: AKolenda Date: Mon, 28 Sep 2026 00:33:42 -0600 Subject: [PATCH 2/4] refactor(session): answer auto-detect requests through one responder Move the handling of RTT and bandwidth measurement requests out of the x224 processor into an `AutoDetectResponder`, so the message channel and a multitransport tunnel can answer them with the same code. Each transport keeps its own responder, because a continuous measurement counts the data received on its own transport ([MS-RDPBCGR] 3.2.5.14). No behavior change on the message channel. --- crates/ironrdp-session/src/autodetect.rs | 122 +++++++++++++++++++++++ crates/ironrdp-session/src/lib.rs | 1 + crates/ironrdp-session/src/x224/mod.rs | 98 ++---------------- 3 files changed, 134 insertions(+), 87 deletions(-) create mode 100644 crates/ironrdp-session/src/autodetect.rs diff --git a/crates/ironrdp-session/src/autodetect.rs b/crates/ironrdp-session/src/autodetect.rs new file mode 100644 index 0000000000..ccc0a3105a --- /dev/null +++ b/crates/ironrdp-session/src/autodetect.rs @@ -0,0 +1,122 @@ +//! Client answers to the network auto-detect requests of [MS-RDPBCGR] 2.2.14. +//! +//! [MS-RDPBCGR]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpbcgr/dc672839-4f4e-40b1-a71c-cd6a959baa38 + +use ironrdp_core::MonotonicInstant; +use ironrdp_pdu::rdp::autodetect::{ + AutoDetectRequest, AutoDetectResponse, BW_RESULTS_CONNECT_TIME, BW_RESULTS_CONTINUOUS, BW_START_CONNECT_TIME, + BW_START_RELIABLE_UDP, BW_STOP_CONNECT_TIME, BW_STOP_RELIABLE_UDP, +}; +use tracing::debug; + +/// Answers the RTT and bandwidth measurement requests that arrive on one transport. +/// +/// Each transport keeps its own, because a continuous measurement counts the data received on +/// the transport it runs on ([MS-RDPBCGR] 3.2.5.14). Timestamps come from the caller, from the +/// same monotonic clock for every request on that transport. +/// +/// [MS-RDPBCGR]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpbcgr/16ffa852-8aa7-481c-99a0-36c1a9a198f6 +#[derive(Debug, Default)] +pub(crate) struct AutoDetectResponder { + bandwidth: Option, +} + +#[derive(Debug)] +struct BandwidthMeasurement { + started_at: MonotonicInstant, + bytes: u32, + continuous: bool, +} + +impl AutoDetectResponder { + /// Counts received bytes while a continuous bandwidth window is open. + pub(crate) fn record_bytes(&mut self, bytes: usize) { + if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| measurement.continuous) { + measurement.bytes = measurement + .bytes + .saturating_add(u32::try_from(bytes).unwrap_or(u32::MAX)); + } + } + + /// Returns the response `request` calls for, if any. + /// + /// Without an arrival time, a bandwidth measurement is answered with an untimed, zero-byte + /// result. + pub(crate) fn respond( + &mut self, + request: AutoDetectRequest, + received_at: Option, + ) -> Option { + match request { + AutoDetectRequest::RttRequest { sequence_number, .. } => { + Some(AutoDetectResponse::RttResponse { sequence_number }) + } + AutoDetectRequest::BandwidthMeasureStart { request_type, .. } + if matches!(request_type, BW_START_CONNECT_TIME | BW_START_RELIABLE_UDP) => + { + self.bandwidth = received_at.map(|started_at| BandwidthMeasurement { + started_at, + bytes: 0, + continuous: request_type == BW_START_RELIABLE_UDP, + }); + None + } + AutoDetectRequest::BandwidthMeasurePayload { payload, .. } => { + if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| !measurement.continuous) { + // [MS-RDPBCGR] 3.2.5.14 counts the eight-byte auto-detect + // header as well as payloadLength, but not the security header. + measurement.bytes = measurement + .bytes + .saturating_add(u32::try_from(payload.len()).unwrap_or(u32::MAX).saturating_add(8)); + } + None + } + AutoDetectRequest::BandwidthMeasureStop { + sequence_number, + request_type, + payload, + } if matches!(request_type, BW_STOP_CONNECT_TIME | BW_STOP_RELIABLE_UDP) => { + let continuous = request_type == BW_STOP_RELIABLE_UDP; + let measurement = self + .bandwidth + .take() + .filter(|measurement| measurement.continuous == continuous); + let stop_bytes = if continuous { + 0 + } else { + u32::try_from(payload.as_ref().map_or(0, Vec::len)) + .unwrap_or(u32::MAX) + .saturating_add(8) + }; + let (time_delta_ms, byte_count) = match (measurement, received_at) { + (Some(measurement), Some(stopped_at)) => ( + u32::try_from(stopped_at.duration_since(measurement.started_at).as_millis()) + .unwrap_or(u32::MAX) + .max(1), + measurement.bytes.saturating_add(stop_bytes), + ), + _ => (1, stop_bytes), + }; + Some(AutoDetectResponse::BandwidthMeasureResults { + sequence_number, + response_type: if continuous { + BW_RESULTS_CONTINUOUS + } else { + BW_RESULTS_CONNECT_TIME + }, + time_delta_ms, + byte_count, + }) + } + request @ AutoDetectRequest::NetworkCharacteristicsResult { .. } => { + debug!(?request, "Received network characteristics from server"); + None + } + request => { + // Measurements for a lossy transport are not answered here. + debug!(?request, "Ignoring auto-detect request for another transport"); + None + } + } + } +} diff --git a/crates/ironrdp-session/src/lib.rs b/crates/ironrdp-session/src/lib.rs index 3103da0a37..84420484ed 100644 --- a/crates/ironrdp-session/src/lib.rs +++ b/crates/ironrdp-session/src/lib.rs @@ -11,6 +11,7 @@ pub mod rfx; // FIXME: maybe this module should not be in this crate pub mod x224; mod active_stage; +mod autodetect; mod palette; use core::fmt; diff --git a/crates/ironrdp-session/src/x224/mod.rs b/crates/ironrdp-session/src/x224/mod.rs index 08f3076807..a85b55ac6b 100644 --- a/crates/ironrdp-session/src/x224/mod.rs +++ b/crates/ironrdp-session/src/x224/mod.rs @@ -3,10 +3,7 @@ use ironrdp_core::{Decode as _, MonotonicInstant, ReadCursor, WriteBuf, decode}; use ironrdp_dvc::{DrdynvcClient, DvcClientProcessor, DynamicChannelMut, DynamicChannelRef}; use ironrdp_pdu::gcc::{ChannelName, Monitor}; use ironrdp_pdu::mcs::{DisconnectProviderUltimatum, DisconnectReason, McsMessage, SendDataIndicationCtx}; -use ironrdp_pdu::rdp::autodetect::{ - AutoDetectReqPdu, AutoDetectRequest, AutoDetectResponse, AutoDetectRspPdu, BW_RESULTS_CONNECT_TIME, - BW_RESULTS_CONTINUOUS, BW_START_CONNECT_TIME, BW_START_RELIABLE_UDP, BW_STOP_CONNECT_TIME, BW_STOP_RELIABLE_UDP, -}; +use ironrdp_pdu::rdp::autodetect::{AutoDetectReqPdu, AutoDetectRequest, AutoDetectRspPdu}; use ironrdp_pdu::rdp::client_info::CompressionType; use ironrdp_pdu::rdp::headers::{ BasicSecurityHeader, BasicSecurityHeaderFlags, CompressionFlags, IoChannelPdu, ShareDataCtx, ShareDataPdu, @@ -21,6 +18,7 @@ use ironrdp_svc::{ }; use tracing::debug; +use crate::autodetect::AutoDetectResponder; use crate::{SessionError, SessionErrorExt as _, SessionResult, reason_err}; /// X224 Processor output @@ -110,13 +108,7 @@ pub struct Processor { io_channel_id: u16, message_channel_id: Option, share_id: u32, - bandwidth: Option, -} - -struct BandwidthMeasurement { - started_at: MonotonicInstant, - bytes: u32, - continuous: bool, + auto_detect: AutoDetectResponder, } impl Processor { @@ -133,7 +125,7 @@ impl Processor { io_channel_id, message_channel_id, share_id, - bandwidth: None, + auto_detect: AutoDetectResponder::default(), } } @@ -477,11 +469,7 @@ impl Processor { /// Counts session bytes after transport/security headers while a continuous /// bandwidth window is open ([MS-RDPBCGR] 3.2.5.14). pub(crate) fn record_bandwidth_bytes(&mut self, bytes: usize) { - if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| measurement.continuous) { - measurement.bytes = measurement - .bytes - .saturating_add(u32::try_from(bytes).unwrap_or(u32::MAX)); - } + self.auto_detect.record_bytes(bytes); } /// Process a PDU received on the MCS message channel: auto-detect @@ -539,76 +527,12 @@ impl Processor { let req = decode::(data_ctx.user_data).map_err(SessionError::decode)?; - let response = match req.request { - AutoDetectRequest::RttRequest { sequence_number, .. } => { - AutoDetectResponse::RttResponse { sequence_number } - } - AutoDetectRequest::BandwidthMeasureStart { request_type, .. } - if matches!(request_type, BW_START_CONNECT_TIME | BW_START_RELIABLE_UDP) => - { - self.bandwidth = received_at.map(|started_at| BandwidthMeasurement { - started_at, - bytes: 0, - continuous: request_type == BW_START_RELIABLE_UDP, - }); - return Ok(Vec::new()); - } - AutoDetectRequest::BandwidthMeasurePayload { payload, .. } => { - if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| !measurement.continuous) { - // [MS-RDPBCGR] 3.2.5.14 counts the eight-byte auto-detect - // header as well as payloadLength, but not the security header. - measurement.bytes = measurement - .bytes - .saturating_add(u32::try_from(payload.len()).unwrap_or(u32::MAX).saturating_add(8)); - } - return Ok(Vec::new()); - } - AutoDetectRequest::BandwidthMeasureStop { - sequence_number, - request_type, - payload, - } if matches!(request_type, BW_STOP_CONNECT_TIME | BW_STOP_RELIABLE_UDP) => { - let continuous = request_type == BW_STOP_RELIABLE_UDP; - let measurement = self - .bandwidth - .take() - .filter(|measurement| measurement.continuous == continuous); - let stop_bytes = if continuous { - 0 - } else { - u32::try_from(payload.as_ref().map_or(0, Vec::len)) - .unwrap_or(u32::MAX) - .saturating_add(8) - }; - let (time_delta_ms, byte_count) = match (measurement, received_at) { - (Some(measurement), Some(stopped_at)) => ( - u32::try_from(stopped_at.duration_since(measurement.started_at).as_millis()) - .unwrap_or(u32::MAX) - .max(1), - measurement.bytes.saturating_add(stop_bytes), - ), - _ => (1, stop_bytes), - }; - AutoDetectResponse::BandwidthMeasureResults { - sequence_number, - response_type: if continuous { - BW_RESULTS_CONTINUOUS - } else { - BW_RESULTS_CONNECT_TIME - }, - time_delta_ms, - byte_count, - } - } - req @ AutoDetectRequest::NetworkCharacteristicsResult { .. } => { - debug!(?req, "Received network characteristics from server"); - return Ok(vec![ProcessorOutput::AutoDetect(req)]); - } - req => { - // Lossy variants belong to a lossy tunnel, never this channel. - debug!(?req, "Ignoring auto-detect request for another transport"); - return Ok(Vec::new()); - } + if let request @ AutoDetectRequest::NetworkCharacteristicsResult { .. } = req.request { + debug!(?request, "Received network characteristics from server"); + return Ok(vec![ProcessorOutput::AutoDetect(request)]); + } + let Some(response) = self.auto_detect.respond(req.request, received_at) else { + return Ok(Vec::new()); }; debug!(?response, "Responding to an auto-detect request"); let mut frame = WriteBuf::new(); From 3d578e6db98cda259ff85d267b10d38eca1c6245 Mon Sep 17 00:00:00 2001 From: AKolenda Date: Mon, 28 Sep 2026 13:17:37 -0600 Subject: [PATCH 3/4] refactor(session): tidy the bandwidth responder Decode the fast-path header for the byte count only while a continuous measurement is running, since `fast_path::Processor::process` decodes it again. Keep the payloadLength plus header rule in one `counted_len` helper, and log a timed window that is dropped because its Stop has no arrival time, as the connector does. Adds a test for fast-path counting. --- crates/ironrdp-session/src/active_stage.rs | 10 ++- crates/ironrdp-session/src/autodetect.rs | 50 +++++++++--- crates/ironrdp-session/src/x224/mod.rs | 5 ++ .../tests/session/autodetect.rs | 79 +++++++++++++++++++ 4 files changed, 131 insertions(+), 13 deletions(-) diff --git a/crates/ironrdp-session/src/active_stage.rs b/crates/ironrdp-session/src/active_stage.rs index 5d815bc573..8e3ff668d3 100644 --- a/crates/ironrdp-session/src/active_stage.rs +++ b/crates/ironrdp-session/src/active_stage.rs @@ -213,9 +213,13 @@ impl ActiveStage { let (mut stage_outputs, processor_updates) = match action { Action::FastPath => { // A continuous bandwidth measurement counts what follows the fast-path header. - let mut header = ReadCursor::new(frame); - FastPathHeader::decode(&mut header).map_err(SessionError::decode)?; - self.x224_processor.record_bandwidth_bytes(header.len()); + // The header is decoded again by `process`, so only pay for it here while a + // measurement is running. + if self.x224_processor.is_counting_bandwidth() { + let mut header = ReadCursor::new(frame); + FastPathHeader::decode(&mut header).map_err(SessionError::decode)?; + self.x224_processor.record_bandwidth_bytes(header.len()); + } let mut output = WriteBuf::new(); let processor_updates = self.fast_path_processor diff --git a/crates/ironrdp-session/src/autodetect.rs b/crates/ironrdp-session/src/autodetect.rs index ccc0a3105a..e1f59a5f3d 100644 --- a/crates/ironrdp-session/src/autodetect.rs +++ b/crates/ironrdp-session/src/autodetect.rs @@ -21,6 +21,24 @@ pub(crate) struct AutoDetectResponder { bandwidth: Option, } +/// Size of the auto-detect header fields [MS-RDPBCGR] 3.2.5.14 counts along with payloadLength: +/// headerLength, headerTypeId, sequenceNumber, requestType and payloadLength itself. +/// +/// [MS-RDPBCGR]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpbcgr/16ffa852-8aa7-481c-99a0-36c1a9a198f6 +const AUTO_DETECT_HEADER_LEN: u32 = 8; + +/// Reported as `timeDelta` for a window that was not timed, and the floor for one that was: +/// a server computing `byteCount * 8 / timeDelta` divides by it ([MS-RDPBCGR] 3.3.5.14). +const UNMEASURABLE_INTERVAL_MS: u32 = 1; + +/// Bytes a connect-time Payload or Stop adds to the count: payloadLength plus the auto-detect +/// header, but not the security header. This is the rule the connector applies to the same PDUs. +fn counted_len(payload_len: usize) -> u32 { + u32::try_from(payload_len) + .unwrap_or(u32::MAX) + .saturating_add(AUTO_DETECT_HEADER_LEN) +} + #[derive(Debug)] struct BandwidthMeasurement { started_at: MonotonicInstant, @@ -29,6 +47,13 @@ struct BandwidthMeasurement { } impl AutoDetectResponder { + /// Whether a continuous bandwidth window is open, so received bytes are being counted. + pub(crate) fn is_counting(&self) -> bool { + self.bandwidth + .as_ref() + .is_some_and(|measurement| measurement.continuous) + } + /// Counts received bytes while a continuous bandwidth window is open. pub(crate) fn record_bytes(&mut self, bytes: usize) { if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| measurement.continuous) { @@ -63,11 +88,7 @@ impl AutoDetectResponder { } AutoDetectRequest::BandwidthMeasurePayload { payload, .. } => { if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| !measurement.continuous) { - // [MS-RDPBCGR] 3.2.5.14 counts the eight-byte auto-detect - // header as well as payloadLength, but not the security header. - measurement.bytes = measurement - .bytes - .saturating_add(u32::try_from(payload.len()).unwrap_or(u32::MAX).saturating_add(8)); + measurement.bytes = measurement.bytes.saturating_add(counted_len(payload.len())); } None } @@ -84,18 +105,27 @@ impl AutoDetectResponder { let stop_bytes = if continuous { 0 } else { - u32::try_from(payload.as_ref().map_or(0, Vec::len)) - .unwrap_or(u32::MAX) - .saturating_add(8) + counted_len(payload.as_ref().map_or(0, Vec::len)) }; let (time_delta_ms, byte_count) = match (measurement, received_at) { (Some(measurement), Some(stopped_at)) => ( u32::try_from(stopped_at.duration_since(measurement.started_at).as_millis()) .unwrap_or(u32::MAX) - .max(1), + .max(UNMEASURABLE_INTERVAL_MS), measurement.bytes.saturating_add(stop_bytes), ), - _ => (1, stop_bytes), + (Some(measurement), None) => { + // The window was timed but this Stop was not, so there is nothing to + // divide the count by. Log the drop so it does not look like the + // ordinary no-window case. + debug!( + dropped_bytes = measurement.bytes, + "Bandwidth Measure Stop arrived with no arrival time although its window was open; \ + dropping the accumulated count" + ); + (UNMEASURABLE_INTERVAL_MS, stop_bytes) + } + (None, _) => (UNMEASURABLE_INTERVAL_MS, stop_bytes), }; Some(AutoDetectResponse::BandwidthMeasureResults { sequence_number, diff --git a/crates/ironrdp-session/src/x224/mod.rs b/crates/ironrdp-session/src/x224/mod.rs index a85b55ac6b..735a2968cf 100644 --- a/crates/ironrdp-session/src/x224/mod.rs +++ b/crates/ironrdp-session/src/x224/mod.rs @@ -472,6 +472,11 @@ impl Processor { self.auto_detect.record_bytes(bytes); } + /// Whether a continuous bandwidth window is open, so session bytes are being counted. + pub(crate) fn is_counting_bandwidth(&self) -> bool { + self.auto_detect.is_counting() + } + /// Process a PDU received on the MCS message channel: auto-detect /// ([MS-RDPBCGR] 2.2.14), multitransport ([MS-RDPBCGR] 2.2.15), or /// Heartbeat ([MS-RDPBCGR] 2.2.16.1). diff --git a/crates/ironrdp-testsuite-core/tests/session/autodetect.rs b/crates/ironrdp-testsuite-core/tests/session/autodetect.rs index 00f93d0285..ef531d3475 100644 --- a/crates/ironrdp-testsuite-core/tests/session/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/session/autodetect.rs @@ -1,10 +1,15 @@ use std::borrow::Cow; use ironrdp_core::encode_vec; +use ironrdp_graphics::image_processing::PixelFormat; +use ironrdp_pdu::Action; +use ironrdp_pdu::fast_path::{EncryptionFlags, FastPathHeader, FastPathUpdatePdu, Fragmentation, UpdateCode}; use ironrdp_pdu::mcs::{McsMessage, SendDataIndication}; use ironrdp_pdu::rdp::autodetect::{AutoDetectReqPdu, AutoDetectRequest, AutoDetectResponse, AutoDetectRspPdu}; use ironrdp_pdu::x224::X224; +use ironrdp_session::image::DecodedImage; use ironrdp_session::x224::Processor; +use ironrdp_session::{ActiveStage, ActiveStageBuilder, ActiveStageOutput}; use ironrdp_svc::StaticChannelSet; const USER_CHANNEL_ID: u16 = 1002; @@ -161,6 +166,10 @@ fn bandwidth_result(outputs: &[ironrdp_session::x224::ProcessorOutput]) -> (u16, let [ironrdp_session::x224::ProcessorOutput::ResponseFrame(frame)] = outputs else { panic!("expected exactly one bandwidth response"); }; + bandwidth_result_frame(frame) +} + +fn bandwidth_result_frame(frame: &[u8]) -> (u16, u16, u32, u32) { let X224(McsMessage::SendDataRequest(message)) = ironrdp_core::decode::>>(frame).unwrap() else { panic!("expected main-channel response"); @@ -266,3 +275,73 @@ fn untimed_driver_does_not_report_accumulated_bytes_as_a_real_measurement() { ); assert_eq!(bandwidth_result(&outputs), (3, 0x000b, 1, 0)); } + +fn make_active_stage() -> ActiveStage { + ActiveStageBuilder { + static_channels: StaticChannelSet::new(), + user_channel_id: USER_CHANNEL_ID, + io_channel_id: IO_CHANNEL_ID, + message_channel_id: Some(MESSAGE_CHANNEL_ID), + share_id: SHARE_ID, + compression_type: None, + enable_server_pointer: false, + pointer_software_rendering: false, + } + .build() +} + +/// A fast-path frame carrying one Synchronize update with `data_len` bytes of data. +fn fast_path_frame(data_len: usize) -> Vec { + let data = vec![0; data_len]; + let update = encode_vec(&FastPathUpdatePdu { + fragmentation: Fragmentation::Single, + update_code: UpdateCode::Synchronize, + compression_flags: None, + compression_type: None, + data: &data, + }) + .unwrap(); + let mut frame = encode_vec(&FastPathHeader::new(EncryptionFlags::empty(), update.len())).unwrap(); + frame.extend_from_slice(&update); + frame +} + +fn process_stage_frame( + stage: &mut ActiveStage, + image: &mut DecodedImage, + action: Action, + frame: &[u8], + millis: u64, +) -> Vec { + stage + .process_with_timestamp( + image, + action, + frame, + Some(ironrdp_core::MonotonicInstant::from_millis(millis)), + ) + .expect("process timed frame") +} + +#[test] +fn continuous_measurement_counts_fast_path_data_only_inside_the_window() { + let mut stage = make_active_stage(); + let mut image = DecodedImage::new(PixelFormat::RgbA32, 64, 64); + let frame = fast_path_frame(100); + + process_stage_frame(&mut stage, &mut image, Action::FastPath, &frame, 5); + let start = encode_server_autodetect(AutoDetectRequest::bw_start_continuous(1)); + process_stage_frame(&mut stage, &mut image, Action::X224, &start, 10); + process_stage_frame(&mut stage, &mut image, Action::FastPath, &frame, 20); + let stop = encode_server_autodetect(AutoDetectRequest::bw_stop_continuous(2)); + let outputs = process_stage_frame(&mut stage, &mut image, Action::X224, &stop, 40); + process_stage_frame(&mut stage, &mut image, Action::FastPath, &frame, 50); + + let [ActiveStageOutput::ResponseFrame(response)] = outputs.as_slice() else { + panic!("expected exactly one bandwidth response, got {outputs:?}"); + }; + // The frame between Start and Stop counts without its fast-path header: a + // three-byte update header and 100 bytes of update data. The Stop adds its + // six-byte auto-detect header, as in the message-channel tests above. + assert_eq!(bandwidth_result_frame(response), (2, 0x000b, 30, 3 + 100 + 6)); +} From 5a0e3117defe6015b7a4969a2fb67915569a6cbe Mon Sep 17 00:00:00 2001 From: AKolenda Date: Fri, 2 Oct 2026 00:11:40 -0600 Subject: [PATCH 4/4] fix(session): count complete bandwidth frames and reject untimed estimates --- crates/ironrdp-session/src/active_stage.rs | 16 +-- crates/ironrdp-session/src/autodetect.rs | 14 +- crates/ironrdp-session/src/x224/mod.rs | 29 ++-- .../tests/session/autodetect.rs | 130 +++++++++++++++++- 4 files changed, 150 insertions(+), 39 deletions(-) diff --git a/crates/ironrdp-session/src/active_stage.rs b/crates/ironrdp-session/src/active_stage.rs index 8e3ff668d3..e9068eda8a 100644 --- a/crates/ironrdp-session/src/active_stage.rs +++ b/crates/ironrdp-session/src/active_stage.rs @@ -1,13 +1,12 @@ use std::sync::Arc; use ironrdp_bulk::{BulkCompressor, CompressionType as BulkCompressionType}; -use ironrdp_core::{Decode as _, MonotonicInstant, ReadCursor, WriteBuf}; +use ironrdp_core::{MonotonicInstant, ReadCursor, WriteBuf}; use ironrdp_displaycontrol::client::DisplayControlClient; use ironrdp_dvc::pdu::SoftSyncTunnelType; use ironrdp_dvc::{DrdynvcClient, DvcClientProcessor, DvcMessageBatch, DynamicChannelMut, DynamicChannelRef}; use ironrdp_egfx::client::GraphicsPipelineClient; use ironrdp_graphics::pointer::DecodedPointer; -use ironrdp_pdu::fast_path::FastPathHeader; use ironrdp_pdu::gcc::{ChannelName, Monitor}; use ironrdp_pdu::geometry::{ExclusiveRectangle, InclusiveRectangle, Rectangle as _}; use ironrdp_pdu::input::fast_path::{FastPathInput, FastPathInputEvent}; @@ -212,14 +211,9 @@ impl ActiveStage { self.damage_regions.clear(); let (mut stage_outputs, processor_updates) = match action { Action::FastPath => { - // A continuous bandwidth measurement counts what follows the fast-path header. - // The header is decoded again by `process`, so only pay for it here while a - // measurement is running. - if self.x224_processor.is_counting_bandwidth() { - let mut header = ReadCursor::new(frame); - FastPathHeader::decode(&mut header).map_err(SessionError::decode)?; - self.x224_processor.record_bandwidth_bytes(header.len()); - } + // TLS-protected fast-path frames have no RDP Security Header, so the + // continuous bandwidth count includes the entire frame. + self.x224_processor.record_bandwidth_bytes(frame.len()); let mut output = WriteBuf::new(); let processor_updates = self.fast_path_processor @@ -1070,7 +1064,7 @@ mod tests { use core::any::TypeId; use super::*; - use ironrdp_core::encode_vec; + use ironrdp_core::{Decode as _, encode_vec}; use ironrdp_displaycontrol::pdu::{DisplayControlCapabilities, DisplayControlPdu}; use ironrdp_dvc::pdu::{ CreateRequestPdu, DataPdu, DrdynvcDataPdu, DrdynvcServerPdu, SoftSyncChannelList, SoftSyncRequestPdu, diff --git a/crates/ironrdp-session/src/autodetect.rs b/crates/ironrdp-session/src/autodetect.rs index e1f59a5f3d..8c43f8df31 100644 --- a/crates/ironrdp-session/src/autodetect.rs +++ b/crates/ironrdp-session/src/autodetect.rs @@ -47,13 +47,6 @@ struct BandwidthMeasurement { } impl AutoDetectResponder { - /// Whether a continuous bandwidth window is open, so received bytes are being counted. - pub(crate) fn is_counting(&self) -> bool { - self.bandwidth - .as_ref() - .is_some_and(|measurement| measurement.continuous) - } - /// Counts received bytes while a continuous bandwidth window is open. pub(crate) fn record_bytes(&mut self, bytes: usize) { if let Some(measurement) = self.bandwidth.as_mut().filter(|measurement| measurement.continuous) { @@ -123,9 +116,9 @@ impl AutoDetectResponder { "Bandwidth Measure Stop arrived with no arrival time although its window was open; \ dropping the accumulated count" ); - (UNMEASURABLE_INTERVAL_MS, stop_bytes) + (UNMEASURABLE_INTERVAL_MS, 0) } - (None, _) => (UNMEASURABLE_INTERVAL_MS, stop_bytes), + (None, _) => (UNMEASURABLE_INTERVAL_MS, 0), }; Some(AutoDetectResponse::BandwidthMeasureResults { sequence_number, @@ -139,6 +132,9 @@ impl AutoDetectResponder { }) } request @ AutoDetectRequest::NetworkCharacteristicsResult { .. } => { + // The TCP message-channel processor surfaces this request itself. Keep this + // arm for the UDP tunnel responder introduced in PR #2009, which handles + // auto-detect requests without passing through that processor. debug!(?request, "Received network characteristics from server"); None } diff --git a/crates/ironrdp-session/src/x224/mod.rs b/crates/ironrdp-session/src/x224/mod.rs index 735a2968cf..7e64c7ebf8 100644 --- a/crates/ironrdp-session/src/x224/mod.rs +++ b/crates/ironrdp-session/src/x224/mod.rs @@ -248,16 +248,13 @@ impl Processor { }; let channel_id = data_ctx.channel_id; - // Message-channel PDUs carry a Basic Security Header, excluded below. - // Ordinary session data on TLS-protected IO/SVC channels does not. - if self.message_channel_id != Some(channel_id) { - self.record_bandwidth_bytes(data_ctx.user_data.len()); - } if channel_id == self.io_channel_id { - self.process_io_channel_data_indication(data_ctx, bulk_decompressor) + self.process_io_channel_data_indication(data_ctx, frame.len(), bulk_decompressor) } else if self.message_channel_id == Some(channel_id) { self.process_message_channel(data_ctx, received_at) } else { + // TLS-protected SVC data has no RDP Security Header: include TPKT/X224/MCS. + self.record_bandwidth_bytes(frame.len()); let maximum_chunk_size = self.static_channels.maximum_chunk_size(); if let Some(svc) = self.static_channels.get_by_channel_id_mut(channel_id) { let response_pdus = svc.process(data_ctx.user_data).map_err(SessionError::pdu)?; @@ -272,6 +269,7 @@ impl Processor { fn process_io_channel_data_indication( &mut self, data_ctx: SendDataIndicationCtx<'_>, + frame_len: usize, bulk_decompressor: &mut Option, ) -> SessionResult> { debug_assert_eq!(data_ctx.channel_id, self.io_channel_id); @@ -282,9 +280,13 @@ impl Processor { ironrdp_pdu::rdp::headers::decode_io_channel(data_ctx), Ok(IoChannelPdu::MultitransportRequest(_)) ) { + self.record_bandwidth_bytes(data_ctx.user_data.len() - BasicSecurityHeader::FIXED_PART_SIZE); return self.process_io_channel(data_ctx, bulk_decompressor); } + // Ordinary TLS-protected IO data has no RDP Security Header. Count the + // entire frame once, including framing and any concatenated Share Control PDUs. + self.record_bandwidth_bytes(frame_len); let mut outputs = Vec::new(); let mut offset = 0usize; let data = data_ctx.user_data; @@ -466,17 +468,13 @@ impl Processor { Ok(decompressed) } - /// Counts session bytes after transport/security headers while a continuous - /// bandwidth window is open ([MS-RDPBCGR] 3.2.5.14). + /// Counts received session bytes while a continuous bandwidth window is open. + /// Exclude framing only when an RDP Security Header is present, counting just + /// the bytes after that header ([MS-RDPBCGR] 3.2.5.14). pub(crate) fn record_bandwidth_bytes(&mut self, bytes: usize) { self.auto_detect.record_bytes(bytes); } - /// Whether a continuous bandwidth window is open, so session bytes are being counted. - pub(crate) fn is_counting_bandwidth(&self) -> bool { - self.auto_detect.is_counting() - } - /// Process a PDU received on the MCS message channel: auto-detect /// ([MS-RDPBCGR] 2.2.14), multitransport ([MS-RDPBCGR] 2.2.15), or /// Heartbeat ([MS-RDPBCGR] 2.2.16.1). @@ -701,6 +699,7 @@ mod tests { channel_id: 1003, user_data: &encoded, }, + encoded.len(), &mut None, ) .expect("ignore a misrouted optional multitransport request"); @@ -918,6 +917,7 @@ mod tests { channel_id: 1003, user_data: &user_data, }, + user_data.len(), &mut None, ) .expect("concatenated Share Control PDUs should be split"); @@ -939,6 +939,7 @@ mod tests { channel_id: 1003, user_data: &[0x06], }, + 1, &mut None, ) .expect_err("a truncated totalLength field is invalid"); @@ -959,6 +960,7 @@ mod tests { channel_id: 1003, user_data: &user_data, }, + user_data.len(), &mut None, ) .expect_err("an overrunning concatenated totalLength is invalid"); @@ -984,6 +986,7 @@ mod tests { channel_id: 1003, user_data: &user_data, }, + user_data.len(), &mut None, ) .expect("invalid first totalLength should fall back to whole-buffer decode"); diff --git a/crates/ironrdp-testsuite-core/tests/session/autodetect.rs b/crates/ironrdp-testsuite-core/tests/session/autodetect.rs index ef531d3475..1622d9ae1f 100644 --- a/crates/ironrdp-testsuite-core/tests/session/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/session/autodetect.rs @@ -1,11 +1,18 @@ use std::borrow::Cow; use ironrdp_core::encode_vec; +use ironrdp_dvc::DrdynvcClient; +use ironrdp_dvc::pdu::{CapabilitiesRequestPdu, CapsVersion, DrdynvcServerPdu}; use ironrdp_graphics::image_processing::PixelFormat; use ironrdp_pdu::Action; use ironrdp_pdu::fast_path::{EncryptionFlags, FastPathHeader, FastPathUpdatePdu, Fragmentation, UpdateCode}; use ironrdp_pdu::mcs::{McsMessage, SendDataIndication}; use ironrdp_pdu::rdp::autodetect::{AutoDetectReqPdu, AutoDetectRequest, AutoDetectResponse, AutoDetectRspPdu}; +use ironrdp_pdu::rdp::headers::{ + BasicSecurityHeader, BasicSecurityHeaderFlags, ServerDeactivateAll, ShareControlHeader, ShareControlPdu, +}; +use ironrdp_pdu::rdp::multitransport::{MultitransportRequestPdu, RequestedProtocol}; +use ironrdp_pdu::rdp::vc::{ChannelControlFlags, ChannelPduHeader}; use ironrdp_pdu::x224::X224; use ironrdp_session::image::DecodedImage; use ironrdp_session::x224::Processor; @@ -38,10 +45,13 @@ fn process_frame(processor: &mut Processor, frame: &[u8]) -> Vec Vec { let pdu = AutoDetectReqPdu::new(request); let user_data = encode_vec(&pdu).unwrap(); + encode_server_data(MESSAGE_CHANNEL_ID, user_data) +} +fn encode_server_data(channel_id: u16, user_data: Vec) -> Vec { let indication = McsMessage::SendDataIndication(SendDataIndication { initiator_id: USER_CHANNEL_ID, - channel_id: MESSAGE_CHANNEL_ID, + channel_id, user_data: Cow::Owned(user_data), }); @@ -189,7 +199,7 @@ fn bandwidth_result_frame(frame: &[u8]) -> (u16, u16, u32, u32) { } #[test] -fn continuous_measurement_counts_data_without_security_headers() { +fn continuous_measurement_excludes_message_channel_security_headers() { let mut processor = make_processor(); assert!(timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 10).is_empty()); timed_request(&mut processor, AutoDetectRequest::rtt_continuous(2), 20); @@ -199,6 +209,82 @@ fn continuous_measurement_counts_data_without_security_headers() { assert_eq!(bandwidth_result(&result), (3, 0x000b, 25, 12)); } +#[test] +fn continuous_measurement_counts_io_framing_once_for_concatenated_pdus() { + let mut processor = make_processor(); + let pdu = encode_vec(&ShareControlHeader { + share_control_pdu: ShareControlPdu::ServerDeactivateAll(ServerDeactivateAll), + pdu_source: USER_CHANNEL_ID, + share_id: SHARE_ID, + }) + .unwrap(); + let frame = encode_server_data(IO_CHANNEL_ID, [pdu.as_slice(), pdu.as_slice()].concat()); + + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 10); + assert_eq!(process_frame(&mut processor, &frame).len(), 2); + let result = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(2), 35); + // There is no RDP Security Header: include TPKT, X224, MCS and both PDUs. + assert_eq!( + bandwidth_result(&result), + (2, 0x000b, 25, u32::try_from(frame.len()).unwrap() + 6) + ); +} + +#[test] +fn continuous_measurement_counts_svc_framing() { + let mut channels = StaticChannelSet::new(); + channels.insert(DrdynvcClient::new()); + channels.attach_channel_id(core::any::TypeId::of::(), 1005); + let mut processor = Processor::new( + channels, + USER_CHANNEL_ID, + IO_CHANNEL_ID, + Some(MESSAGE_CHANNEL_ID), + SHARE_ID, + ); + let capabilities = encode_vec(&DrdynvcServerPdu::Capabilities(CapabilitiesRequestPdu::new( + CapsVersion::V1, + None, + ))) + .unwrap(); + let mut chunk = encode_vec(&ChannelPduHeader { + length: u32::try_from(capabilities.len()).unwrap(), + flags: ChannelControlFlags::FLAG_FIRST | ChannelControlFlags::FLAG_LAST, + }) + .unwrap(); + chunk.extend_from_slice(&capabilities); + let frame = encode_server_data(1005, chunk); + + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 10); + process_frame(&mut processor, &frame); + let result = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(2), 35); + assert_eq!( + bandwidth_result(&result), + (2, 0x000b, 25, u32::try_from(frame.len()).unwrap() + 6) + ); +} + +#[test] +fn continuous_measurement_excludes_io_multitransport_security_header() { + let mut processor = make_processor(); + let data = encode_vec(&MultitransportRequestPdu { + security_header: BasicSecurityHeader { + flags: BasicSecurityHeaderFlags::TRANSPORT_REQ, + }, + request_id: 42, + requested_protocol: RequestedProtocol::UdpFecR, + security_cookie: [0xab; 16], + }) + .unwrap(); + let counted = u32::try_from(data.len() - BasicSecurityHeader::FIXED_PART_SIZE).unwrap(); + let frame = encode_server_data(IO_CHANNEL_ID, data); + + timed_request(&mut processor, AutoDetectRequest::bw_start_continuous(1), 10); + assert!(process_frame(&mut processor, &frame).is_empty()); + let result = timed_request(&mut processor, AutoDetectRequest::bw_stop_continuous(2), 35); + assert_eq!(bandwidth_result(&result), (2, 0x000b, 25, counted + 6)); +} + #[test] fn connect_time_measurement_includes_payload_headers_once() { let mut processor = make_processor(); @@ -276,6 +362,36 @@ fn untimed_driver_does_not_report_accumulated_bytes_as_a_real_measurement() { assert_eq!(bandwidth_result(&outputs), (3, 0x000b, 1, 0)); } +#[test] +fn untimed_connect_time_stop_does_not_report_stop_payload_bytes() { + let mut processor = make_processor(); + timed_request(&mut processor, AutoDetectRequest::bw_start_connect_time(1), 10); + timed_request(&mut processor, AutoDetectRequest::bw_payload(2, vec![0xaa; 64]), 20); + let outputs = process_frame( + &mut processor, + &encode_server_autodetect(AutoDetectRequest::bw_stop_connect_time(3, vec![0xbb; 16])), + ); + assert_eq!(bandwidth_result(&outputs), (3, 0x0003, 1, 0)); +} + +#[test] +fn connect_time_stop_without_a_timed_start_reports_zero_bytes() { + for send_untimed_start in [false, true] { + for stop_time in [None, Some(ironrdp_core::MonotonicInstant::from_millis(20))] { + let mut processor = make_processor(); + if send_untimed_start { + process_frame( + &mut processor, + &encode_server_autodetect(AutoDetectRequest::bw_start_connect_time(1)), + ); + } + let frame = encode_server_autodetect(AutoDetectRequest::bw_stop_connect_time(2, vec![0xbb; 16])); + let outputs = processor.process_with_timestamp(&frame, &mut None, stop_time).unwrap(); + assert_eq!(bandwidth_result(&outputs), (2, 0x0003, 1, 0)); + } + } +} + fn make_active_stage() -> ActiveStage { ActiveStageBuilder { static_channels: StaticChannelSet::new(), @@ -340,8 +456,10 @@ fn continuous_measurement_counts_fast_path_data_only_inside_the_window() { let [ActiveStageOutput::ResponseFrame(response)] = outputs.as_slice() else { panic!("expected exactly one bandwidth response, got {outputs:?}"); }; - // The frame between Start and Stop counts without its fast-path header: a - // three-byte update header and 100 bytes of update data. The Stop adds its - // six-byte auto-detect header, as in the message-channel tests above. - assert_eq!(bandwidth_result_frame(response), (2, 0x000b, 30, 3 + 100 + 6)); + // The whole fast-path frame counts because it has no RDP Security Header. + // The Stop adds only its six-byte auto-detect header, after its Security Header. + assert_eq!( + bandwidth_result_frame(response), + (2, 0x000b, 30, u32::try_from(frame.len()).unwrap() + 6) + ); }