From 253e73405a20d70eafdcb8070b6315af613775d3 Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Sat, 26 Sep 2026 23:03:09 -0500 Subject: [PATCH 1/2] fix(server): measure bandwidth across a single large graphics write The continuous measurement opened a fixed window of RTT ticks, and the client counts only the traffic that happens to pass between Start and Stop, so a window over a quiet desktop timed idle time and reported a fast link as a few hundred kbps. Bracket one EGFX write of at least 10 KiB with Start and Stop instead, at most once a second, so the client times a burst the link carried. Sessions without EGFX keep the tick window until the first bracketed measurement. A zero time delta now counts as one millisecond: the client's timer has millisecond resolution and a burst on a fast link completes within one. A result with no bytes counted still clears the stored figure. --- crates/ironrdp-server/src/autodetect.rs | 95 +++++++++++++++- crates/ironrdp-server/src/server.rs | 78 +++++++++++-- .../tests/server/autodetect.rs | 103 +++++++++++++++++- 3 files changed, 255 insertions(+), 21 deletions(-) diff --git a/crates/ironrdp-server/src/autodetect.rs b/crates/ironrdp-server/src/autodetect.rs index fadf157d52..ecac2932a9 100644 --- a/crates/ironrdp-server/src/autodetect.rs +++ b/crates/ironrdp-server/src/autodetect.rs @@ -1,9 +1,11 @@ //! Server-side auto-detect (RTT and bandwidth measurement) per [MS-RDPBCGR 2.2.14]. //! //! The server periodically sends RTT Measure Request PDUs and records the -//! round-trip time from the client's response, and periodically brackets -//! ordinary traffic with a Bandwidth Measure Start/Stop pair so the client -//! can report the bandwidth it observed. RTT results are exposed via +//! round-trip time from the client's response, and brackets ordinary traffic +//! with a Bandwidth Measure Start/Stop pair so the client can report the +//! bandwidth it observed: around a single large graphics write where there is +//! one (see [`AutoDetectManager::begin_bandwidth_measure()`]), otherwise over +//! a fixed window of RTT ticks. RTT results are exposed via //! [`AutoDetectManager::snapshot()`]; both RTT and bandwidth are reported to //! the client via [`AutoDetectManager::build_netchar_result()`]. //! @@ -34,6 +36,15 @@ const BW_MEASURE_INTERVAL_TICKS: u32 = 8; /// immediately the way a synthetic payload would allow. const BW_WINDOW_TICKS: u32 = 4; +/// Smallest graphics write worth bracketing with a bandwidth measurement. +/// +/// A smaller write crosses even a slow link within the client's millisecond +/// timer resolution, so its figure says little about capacity. +pub const BW_BRACKET_MIN_BYTES: usize = 10 * 1024; + +/// Minimum spacing between bandwidth measurements bracketed around a write. +const BW_BRACKET_MIN_INTERVAL_MS: u64 = 1_000; + /// Minimum spacing between Network Characteristics Result PDUs. /// /// The reported figures are averaged over a window and change far more slowly @@ -68,6 +79,9 @@ pub struct AutoDetectManager { bandwidth_kbps: Option, /// Counts RTT ticks to pace bandwidth measurements (see [`BW_MEASURE_INTERVAL_TICKS`]). bw_tick_count: u32, + /// When the last bracketed measurement began. Once set, the session has + /// graphics writes to measure across and the tick window stops opening. + last_bracket_start_ms: Option, /// When the last Network Characteristics Result was built, for pacing. last_netchar_result_ms: Option, /// Lowest RTT observed over the whole session, for the wire `baseRTT`. @@ -114,6 +128,8 @@ enum PendingBandwidth { /// Start has been sent; the window stays open for this many more ticks /// before Stop is sent, so ordinary traffic can fill it. Open { sequence: u16, ticks_remaining: u32 }, + /// Start has been sent ahead of a single write; Stop follows it. + Bracketing { sequence: u16 }, /// Stop has been sent; awaiting the client's Bandwidth Measure Results. AwaitingResults { sequence: u16 }, } @@ -127,6 +143,7 @@ impl AutoDetectManager { pending_bw: None, bandwidth_kbps: None, bw_tick_count: 0, + last_bracket_start_ms: None, last_netchar_result_ms: None, rtt_is_fresh: false, bandwidth_is_fresh: false, @@ -175,7 +192,8 @@ impl AutoDetectManager { *ticks_remaining -= 1; None } - Some(PendingBandwidth::AwaitingResults { .. }) => None, + Some(PendingBandwidth::AwaitingResults { .. } | PendingBandwidth::Bracketing { .. }) => None, + None if self.last_bracket_start_ms.is_some() => None, None => { self.bw_tick_count = self.bw_tick_count.wrapping_add(1); if !self.bw_tick_count.is_multiple_of(BW_MEASURE_INTERVAL_TICKS) { @@ -192,6 +210,51 @@ impl AutoDetectManager { } } + /// Start a bandwidth measurement around the write the caller is about to + /// make, returning the Bandwidth Measure Start to send immediately before + /// it, or `None` when no measurement should start now. + /// + /// Continuous Auto-Detection counts the ordinary traffic between Start and + /// Stop ([MS-RDPBCGR] 3.2.5.14), so a window that happens to span idle time + /// measures how little the server sent, not what the link can carry. + /// Bracketing one large write, and sending Stop right after it with + /// [`end_bandwidth_measure()`](Self::end_bandwidth_measure), times a burst + /// the link actually had to carry. All three must go out on the same + /// ordered transport. + /// + /// Returns `None` for a write smaller than [`BW_BRACKET_MIN_BYTES`], while + /// another measurement is outstanding, and within + /// [`BW_BRACKET_MIN_INTERVAL_MS`] of the previous bracketed Start. The + /// first successful call also retires the tick window of + /// [`build_bandwidth_measure()`](Self::build_bandwidth_measure) for the + /// rest of the session. + pub fn begin_bandwidth_measure(&mut self, write_len: usize, now_ms: u64) -> Option { + if write_len < BW_BRACKET_MIN_BYTES || self.pending_bw.is_some() { + return None; + } + if let Some(last) = self.last_bracket_start_ms { + if now_ms.saturating_sub(last) < BW_BRACKET_MIN_INTERVAL_MS { + return None; + } + } + self.last_bracket_start_ms = Some(now_ms); + let sequence = self.next_sequence; + self.next_sequence = sequence.wrapping_add(1); + self.pending_bw = Some(PendingBandwidth::Bracketing { sequence }); + Some(AutoDetectRequest::bw_start_continuous(sequence)) + } + + /// Finish the measurement [`begin_bandwidth_measure()`](Self::begin_bandwidth_measure) + /// started, returning the Bandwidth Measure Stop to send right after the + /// bracketed write, or `None` if no bracketed measurement is open. + pub fn end_bandwidth_measure(&mut self) -> Option { + let Some(PendingBandwidth::Bracketing { sequence }) = self.pending_bw else { + return None; + }; + self.pending_bw = Some(PendingBandwidth::AwaitingResults { sequence }); + Some(AutoDetectRequest::bw_stop_continuous(sequence)) + } + /// Build a Network Characteristics Result reporting the measured network. /// /// Returns `None` until the network has actually been characterised, which @@ -299,7 +362,7 @@ impl AutoDetectManager { // previous one rather than leaving it in place, so a run of failures // withholds the result (see `build_netchar_result`) instead of // reporting an increasingly stale bandwidth. - self.bandwidth_kbps = response.computed_bandwidth_kbps(); + self.bandwidth_kbps = measured_bandwidth_kbps(response); if self.bandwidth_kbps.is_some() { self.bandwidth_is_fresh = true; } @@ -367,6 +430,28 @@ impl AutoDetectManager { } } +/// Bandwidth from a Bandwidth Measure Results, in kilobits per second, or +/// `None` when the client counted no bytes. +/// +/// The client's timer has millisecond resolution, so a burst that crosses a +/// fast link within one millisecond reports a zero time. That is a real +/// measurement, bounded below by one millisecond, not a failed one. +fn measured_bandwidth_kbps(response: &AutoDetectResponse) -> Option { + let AutoDetectResponse::BandwidthMeasureResults { + time_delta_ms, + byte_count, + .. + } = response + else { + return None; + }; + if *byte_count == 0 { + return None; + } + let kbps = u64::from(*byte_count) * 8 / u64::from((*time_delta_ms).max(1)); + Some(u32::try_from(kbps).unwrap_or(u32::MAX)) +} + impl Default for AutoDetectManager { fn default() -> Self { Self::new() diff --git a/crates/ironrdp-server/src/server.rs b/crates/ironrdp-server/src/server.rs index 084c2070ac..b4d8acc1b3 100644 --- a/crates/ironrdp-server/src/server.rs +++ b/crates/ironrdp-server/src/server.rs @@ -3411,8 +3411,14 @@ impl RdpServer { #[cfg(feature = "egfx")] ServerEvent::Egfx(msg) => match msg { EgfxServerMessage::SendMessages { messages } => { - self.dispatch_egfx_messages(messages, writer, user_channel_id, udp_transport) - .await?; + self.dispatch_egfx_messages( + messages, + writer, + user_channel_id, + message_channel_id, + udp_transport, + ) + .await?; } }, ServerEvent::AutoDetectRttRequest => { @@ -3442,10 +3448,11 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } - // Periodically measure bandwidth: Start on one tick, Stop several - // ticks later, with ordinary traffic in between counted by the - // client, then a Bandwidth Measure Results PDU in reply. Until one - // has completed there is no characteristics result to send at all. + // Until a large graphics write has been bracketed (see the Egfx + // arm), measure bandwidth over a window of ticks: Start on one, + // Stop several later, with ordinary traffic in between counted by + // the client, then a Bandwidth Measure Results PDU in reply. Until + // one has completed there is no characteristics result to send. if let Some(pdu) = ad.build_bandwidth_measure() { let data = encode_autodetect_request(pdu, message_channel_id, user_channel_id)?; writer @@ -3482,6 +3489,7 @@ impl RdpServer { messages: Vec, writer: &mut impl FramedWrite, user_channel_id: u16, + message_channel_id: Option, udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult<()> { let drdynvc_channel_id = self @@ -3504,7 +3512,13 @@ impl RdpServer { let Some(egfx_dvc_id) = crate::gfx::egfx_channel_id(drdynvc) else { trace!("EGFX channel not open yet, staying on TCP"); return self - .write_egfx_over_tcp(messages, writer, drdynvc_channel_id, user_channel_id) + .write_egfx_over_tcp( + messages, + writer, + drdynvc_channel_id, + user_channel_id, + message_channel_id, + ) .await; }; @@ -3555,8 +3569,14 @@ impl RdpServer { if route_over_udp { return Ok(()); } - self.write_egfx_over_tcp(messages, writer, drdynvc_channel_id, user_channel_id) - .await + self.write_egfx_over_tcp( + messages, + writer, + drdynvc_channel_id, + user_channel_id, + message_channel_id, + ) + .await } #[cfg(feature = "egfx")] @@ -3566,13 +3586,51 @@ impl RdpServer { writer: &mut impl FramedWrite, drdynvc_channel_id: u16, user_channel_id: u16, + message_channel_id: Option, ) -> ServerResult<()> { let data = server_encode_svc_messages(messages, drdynvc_channel_id, user_channel_id).map_err(ServerError::encode)?; + + // Measure bandwidth across this write when it is large enough to + // say something about the link: Start, the graphics data and Stop + // go out back to back on the same stream, so the client times a + // burst rather than whatever idle stretch a fixed window happens + // to span. + let bracket = match (self.autodetect.as_mut(), message_channel_id) { + (Some(ad), Some(message_channel_id)) => ad + .begin_bandwidth_measure(data.len(), monotonic_now_ms()) + .map(|start| (start, message_channel_id)), + _ => None, + }; + let bracket_channel_id = match bracket { + Some((start, message_channel_id)) => { + let start = encode_autodetect_request(start, message_channel_id, user_channel_id)?; + writer + .write_all(&start) + .await + .map_err(|e| ServerError::io("write_all", e))?; + Some(message_channel_id) + } + None => None, + }; writer .write_all(&data) .await - .map_err(|e| ServerError::io("write_all", e)) + .map_err(|e| ServerError::io("write_all", e))?; + if let Some(message_channel_id) = bracket_channel_id { + if let Some(stop) = self + .autodetect + .as_mut() + .and_then(AutoDetectManager::end_bandwidth_measure) + { + let stop = encode_autodetect_request(stop, message_channel_id, user_channel_id)?; + writer + .write_all(&stop) + .await + .map_err(|e| ServerError::io("write_all", e))?; + } + } + Ok(()) } /// Writes DRDYNVC output, sending the data of any channel the Soft-Sync diff --git a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs index 8b1e6a6391..9d5d0599df 100644 --- a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs @@ -1,5 +1,5 @@ use ironrdp_pdu::rdp::autodetect::{AutoDetectRequest, AutoDetectResponse}; -use ironrdp_server::autodetect::{AutoDetectManager, AutoDetectOutcome}; +use ironrdp_server::autodetect::{AutoDetectManager, AutoDetectOutcome, BW_BRACKET_MIN_BYTES}; /// Upper bound on ticks to drive while waiting for a bandwidth transaction: /// generous enough to cover the pacing plus the window, tight enough that a @@ -296,12 +296,12 @@ fn mismatched_bandwidth_sequence_is_ignored() { ); } -/// A Bandwidth Measure Results with `time_delta_ms: 0` cannot compute a +/// A Bandwidth Measure Results in which the client counted no bytes carries no /// figure. It must not be reported, and it must not leave a stale /// `bandwidth_kbps` from an earlier successful measurement on the wire as if /// it were current. #[test] -fn zero_time_delta_ages_out_a_previous_bandwidth_figure() { +fn zero_byte_count_ages_out_a_previous_bandwidth_figure() { let mut mgr = AutoDetectManager::new(); let req = mgr.send_rtt_request(0); let _ = mgr.handle_response( @@ -328,7 +328,7 @@ fn zero_time_delta_ages_out_a_previous_bandwidth_figure() { "a real bandwidth figure is known" ); - // Second measurement fails: timeDelta 0. + // Second measurement fails: no bytes counted. let req = mgr.send_rtt_request(1_000); let _ = mgr.handle_response( &AutoDetectResponse::RttResponse { @@ -340,8 +340,8 @@ fn zero_time_delta_ages_out_a_previous_bandwidth_figure() { let zero_delta_results = AutoDetectResponse::BandwidthMeasureResults { sequence_number: bw_seq, response_type: ironrdp_pdu::rdp::autodetect::BW_RESULTS_CONTINUOUS, - time_delta_ms: 0, - byte_count: 100_000, + time_delta_ms: 10, + byte_count: 0, }; // Matched (it completes the outstanding transaction), but unusable: distinct from // an unmatched reply, and distinct from `Bandwidth(Some(_))`. A caller collapsing @@ -360,6 +360,97 @@ fn zero_time_delta_ages_out_a_previous_bandwidth_figure() { ); } +/// The client times the window in whole milliseconds, so a burst that crosses a +/// fast link within one reports zero. That is a measurement bounded by one +/// millisecond, not a failure. +#[test] +fn zero_time_delta_counts_as_one_millisecond() { + let mut mgr = AutoDetectManager::new(); + let bw_seq = drive_bandwidth_start_and_stop(&mut mgr); + let results = AutoDetectResponse::BandwidthMeasureResults { + sequence_number: bw_seq, + response_type: ironrdp_pdu::rdp::autodetect::BW_RESULTS_CONTINUOUS, + time_delta_ms: 0, + byte_count: 100_000, + }; + assert_eq!( + mgr.handle_response(&results, 20), + AutoDetectOutcome::Bandwidth(Some(800_000)) + ); +} + +/// Completes the bracketed measurement `start` opened, with a usable figure. +fn complete_bracket(mgr: &mut AutoDetectManager, start: &AutoDetectRequest, now_ms: u64) { + let stop = mgr.end_bandwidth_measure().expect("Stop after a bracketed Start"); + assert_eq!(stop.sequence_number(), start.sequence_number()); + let results = AutoDetectResponse::BandwidthMeasureResults { + sequence_number: stop.sequence_number(), + response_type: ironrdp_pdu::rdp::autodetect::BW_RESULTS_CONTINUOUS, + time_delta_ms: 2, + byte_count: 20_000, + }; + assert_eq!( + mgr.handle_response(&results, now_ms), + AutoDetectOutcome::Bandwidth(Some(80_000)) + ); +} + +#[test] +fn a_large_write_is_bracketed_by_start_and_stop() { + let mut mgr = AutoDetectManager::new(); + let start = mgr + .begin_bandwidth_measure(BW_BRACKET_MIN_BYTES, 0) + .expect("a write at the threshold is measured"); + assert!(matches!(start, AutoDetectRequest::BandwidthMeasureStart { .. })); + let stop = mgr.end_bandwidth_measure().expect("Stop follows the write"); + assert!(matches!(stop, AutoDetectRequest::BandwidthMeasureStop { .. })); + assert_eq!(stop.sequence_number(), start.sequence_number()); + assert!(mgr.end_bandwidth_measure().is_none(), "one Stop per Start"); +} + +#[test] +fn a_small_write_is_not_bracketed() { + let mut mgr = AutoDetectManager::new(); + assert!(mgr.begin_bandwidth_measure(BW_BRACKET_MIN_BYTES - 1, 0).is_none()); + assert!(mgr.end_bandwidth_measure().is_none()); +} + +#[test] +fn bracketed_measurements_are_paced_and_never_overlap() { + let mut mgr = AutoDetectManager::new(); + let start = mgr.begin_bandwidth_measure(64 * 1024, 0).expect("first write measured"); + let _ = mgr.end_bandwidth_measure(); + assert!( + mgr.begin_bandwidth_measure(64 * 1024, 5_000).is_none(), + "no new measurement while the client's results are outstanding" + ); + let results = AutoDetectResponse::BandwidthMeasureResults { + sequence_number: start.sequence_number(), + response_type: ironrdp_pdu::rdp::autodetect::BW_RESULTS_CONTINUOUS, + time_delta_ms: 2, + byte_count: 20_000, + }; + let _ = mgr.handle_response(&results, 10); + assert!( + mgr.begin_bandwidth_measure(64 * 1024, 999).is_none(), + "within a second of the previous Start" + ); + let next = mgr.begin_bandwidth_measure(64 * 1024, 1_000).expect("a second later"); + complete_bracket(&mut mgr, &next, 1_010); +} + +/// Once the session has large writes to measure across, the tick window (which +/// times idle stretches as well as traffic) stops opening. +#[test] +fn the_tick_window_retires_after_a_bracketed_measurement() { + let mut mgr = AutoDetectManager::new(); + let start = mgr.begin_bandwidth_measure(64 * 1024, 0).expect("measured"); + complete_bracket(&mut mgr, &start, 10); + for _ in 0..MAX_BANDWIDTH_TICKS { + assert!(mgr.build_bandwidth_measure().is_none()); + } +} + /// Two consecutive measurements landing on the identical kbps figure (plausible on a /// stable link) must both be reported as a match. A caller that derives "did this /// complete a measurement" from comparing the value before and after the call, rather From 5cdf387f05d48662fe108161a991319396e52020 Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Mon, 28 Sep 2026 11:34:26 -0500 Subject: [PATCH 2/2] review: keep bandwidth measurement going when brackets stop The first bracketed measurement no longer turns the tick window off for good: each bracketed Start resets the tick count, so the window takes over again once eight ticks pass without one. A bracketed measurement, or one waiting for the client's results, is now dropped by expire_stale_probes once it is older than the RTT probe age, so a Stop that is never sent or a client that never answers no longer stops measurement for the rest of the session. A tick window that is still open is not timed, since it sends its own Stop. Also folds the Start and Stop writes into let chains, passes the measured fields to measured_bandwidth_kbps, and allocates sequence numbers through one helper. --- crates/ironrdp-server/src/autodetect.rs | 116 ++++++++++++------ crates/ironrdp-server/src/server.rs | 59 ++++----- .../tests/server/autodetect.rs | 100 ++++++++++++++- 3 files changed, 203 insertions(+), 72 deletions(-) diff --git a/crates/ironrdp-server/src/autodetect.rs b/crates/ironrdp-server/src/autodetect.rs index ecac2932a9..ddb49cae62 100644 --- a/crates/ironrdp-server/src/autodetect.rs +++ b/crates/ironrdp-server/src/autodetect.rs @@ -40,7 +40,7 @@ const BW_WINDOW_TICKS: u32 = 4; /// /// A smaller write crosses even a slow link within the client's millisecond /// timer resolution, so its figure says little about capacity. -pub const BW_BRACKET_MIN_BYTES: usize = 10 * 1024; +pub(crate) const BW_BRACKET_MIN_BYTES: usize = 10 * 1024; /// Minimum spacing between bandwidth measurements bracketed around a write. const BW_BRACKET_MIN_INTERVAL_MS: u64 = 1_000; @@ -79,8 +79,7 @@ pub struct AutoDetectManager { bandwidth_kbps: Option, /// Counts RTT ticks to pace bandwidth measurements (see [`BW_MEASURE_INTERVAL_TICKS`]). bw_tick_count: u32, - /// When the last bracketed measurement began. Once set, the session has - /// graphics writes to measure across and the tick window stops opening. + /// When the last bracketed measurement began, for pacing brackets. last_bracket_start_ms: Option, /// When the last Network Characteristics Result was built, for pacing. last_netchar_result_ms: Option, @@ -129,9 +128,16 @@ enum PendingBandwidth { /// before Stop is sent, so ordinary traffic can fill it. Open { sequence: u16, ticks_remaining: u32 }, /// Start has been sent ahead of a single write; Stop follows it. - Bracketing { sequence: u16 }, + /// `since_ms` is when it began, so [`AutoDetectManager::expire_stale_probes()`] + /// can give up on a Stop that is never sent. + Bracketing { sequence: u16, since_ms: u64 }, /// Stop has been sent; awaiting the client's Bandwidth Measure Results. - AwaitingResults { sequence: u16 }, + /// `since_ms` is when the wait began, so + /// [`AutoDetectManager::expire_stale_probes()`] can give up on a client + /// that never answers. A bracketed measurement carries over the time its + /// Start was sent. A tick window has no clock when its Stop goes out, so + /// it stays `None` until the first expiry pass that sees it stamps it. + AwaitingResults { sequence: u16, since_ms: Option }, } impl AutoDetectManager { @@ -158,8 +164,7 @@ impl AutoDetectManager { /// header ([MS-RDPBCGR] 2.2.14.3). `now_ms` is recorded as the send time /// and is what [`handle_response()`](Self::handle_response) measures against. pub fn send_rtt_request(&mut self, now_ms: u64) -> AutoDetectRequest { - let seq = self.next_sequence; - self.next_sequence = seq.wrapping_add(1); + let seq = self.allocate_sequence(); self.pending_probes.push((seq, now_ms)); AutoDetectRequest::rtt_continuous(seq) } @@ -186,21 +191,22 @@ impl AutoDetectManager { }) => { if *ticks_remaining == 0 { let sequence = *sequence; - self.pending_bw = Some(PendingBandwidth::AwaitingResults { sequence }); + self.pending_bw = Some(PendingBandwidth::AwaitingResults { + sequence, + since_ms: None, + }); return Some(AutoDetectRequest::bw_stop_continuous(sequence)); } *ticks_remaining -= 1; None } Some(PendingBandwidth::AwaitingResults { .. } | PendingBandwidth::Bracketing { .. }) => None, - None if self.last_bracket_start_ms.is_some() => None, None => { self.bw_tick_count = self.bw_tick_count.wrapping_add(1); if !self.bw_tick_count.is_multiple_of(BW_MEASURE_INTERVAL_TICKS) { return None; } - let seq = self.next_sequence; - self.next_sequence = seq.wrapping_add(1); + let seq = self.allocate_sequence(); self.pending_bw = Some(PendingBandwidth::Open { sequence: seq, ticks_remaining: BW_WINDOW_TICKS, @@ -222,12 +228,13 @@ impl AutoDetectManager { /// the link actually had to carry. All three must go out on the same /// ordered transport. /// - /// Returns `None` for a write smaller than [`BW_BRACKET_MIN_BYTES`], while + /// Returns `None` for a write smaller than 10 KiB, while /// another measurement is outstanding, and within - /// [`BW_BRACKET_MIN_INTERVAL_MS`] of the previous bracketed Start. The - /// first successful call also retires the tick window of - /// [`build_bandwidth_measure()`](Self::build_bandwidth_measure) for the - /// rest of the session. + /// `BW_BRACKET_MIN_INTERVAL_MS` of the previous bracketed Start. Each + /// successful call also holds off the tick window of + /// [`build_bandwidth_measure()`](Self::build_bandwidth_measure) for another + /// `BW_MEASURE_INTERVAL_TICKS` ticks, so it only measures while no large + /// writes are coming. pub fn begin_bandwidth_measure(&mut self, write_len: usize, now_ms: u64) -> Option { if write_len < BW_BRACKET_MIN_BYTES || self.pending_bw.is_some() { return None; @@ -238,9 +245,12 @@ impl AutoDetectManager { } } self.last_bracket_start_ms = Some(now_ms); - let sequence = self.next_sequence; - self.next_sequence = sequence.wrapping_add(1); - self.pending_bw = Some(PendingBandwidth::Bracketing { sequence }); + self.bw_tick_count = 0; + let sequence = self.allocate_sequence(); + self.pending_bw = Some(PendingBandwidth::Bracketing { + sequence, + since_ms: now_ms, + }); Some(AutoDetectRequest::bw_start_continuous(sequence)) } @@ -248,10 +258,13 @@ impl AutoDetectManager { /// started, returning the Bandwidth Measure Stop to send right after the /// bracketed write, or `None` if no bracketed measurement is open. pub fn end_bandwidth_measure(&mut self) -> Option { - let Some(PendingBandwidth::Bracketing { sequence }) = self.pending_bw else { + let Some(PendingBandwidth::Bracketing { sequence, since_ms }) = self.pending_bw else { return None; }; - self.pending_bw = Some(PendingBandwidth::AwaitingResults { sequence }); + self.pending_bw = Some(PendingBandwidth::AwaitingResults { + sequence, + since_ms: Some(since_ms), + }); Some(AutoDetectRequest::bw_stop_continuous(sequence)) } @@ -305,8 +318,7 @@ impl AutoDetectManager { self.rtt_is_fresh = false; self.bandwidth_is_fresh = false; - let seq = self.next_sequence; - self.next_sequence = seq.wrapping_add(1); + let seq = self.allocate_sequence(); Some(AutoDetectRequest::netchar_result( seq, base_rtt_ms, @@ -349,10 +361,15 @@ impl AutoDetectManager { AutoDetectOutcome::Rtt(rtt_ms) } - AutoDetectResponse::BandwidthMeasureResults { sequence_number, .. } => { + AutoDetectResponse::BandwidthMeasureResults { + sequence_number, + time_delta_ms, + byte_count, + .. + } => { let awaiting = matches!( self.pending_bw, - Some(PendingBandwidth::AwaitingResults { sequence }) if sequence == *sequence_number + Some(PendingBandwidth::AwaitingResults { sequence, .. }) if sequence == *sequence_number ); if !awaiting { return AutoDetectOutcome::Unmatched; @@ -362,7 +379,7 @@ impl AutoDetectManager { // previous one rather than leaving it in place, so a run of failures // withholds the result (see `build_netchar_result`) instead of // reporting an increasingly stale bandwidth. - self.bandwidth_kbps = measured_bandwidth_kbps(response); + self.bandwidth_kbps = measured_bandwidth_kbps(*time_delta_ms, *byte_count); if self.bandwidth_kbps.is_some() { self.bandwidth_is_fresh = true; } @@ -420,13 +437,44 @@ impl AutoDetectManager { self.bandwidth_kbps } - /// Discard probes older than `max_age_ms` to prevent unbounded growth. + /// Discard probes older than `max_age_ms` to prevent unbounded growth, and + /// give up on a bandwidth measurement that has waited that long. + /// + /// A measurement is otherwise only finished by the client's results, or for + /// a bracketed one by [`end_bandwidth_measure()`](Self::end_bandwidth_measure) + /// first, so a client that never answers or a Stop that is never sent would + /// stop bandwidth measurement for the rest of the session. A tick window + /// that is still open is left alone, since it sends its own Stop. Results + /// that arrive after this are unmatched. /// /// `now_ms` is on the same clock passed to /// [`send_rtt_request()`](Self::send_rtt_request). pub fn expire_stale_probes(&mut self, now_ms: u64, max_age_ms: u64) { self.pending_probes .retain(|(_, sent_at_ms)| now_ms.saturating_sub(*sent_at_ms) < max_age_ms); + + let expired = match &mut self.pending_bw { + Some(PendingBandwidth::Bracketing { since_ms, .. }) => now_ms.saturating_sub(*since_ms) >= max_age_ms, + Some(PendingBandwidth::AwaitingResults { since_ms, .. }) => match *since_ms { + None => { + *since_ms = Some(now_ms); + false + } + Some(since) => now_ms.saturating_sub(since) >= max_age_ms, + }, + Some(PendingBandwidth::Open { .. }) | None => false, + }; + if expired { + self.pending_bw = None; + } + } + + /// Returns the sequence number for the next request and advances the + /// counter, wrapping at `u16::MAX`. + fn allocate_sequence(&mut self) -> u16 { + let sequence = self.next_sequence; + self.next_sequence = sequence.wrapping_add(1); + sequence } } @@ -436,19 +484,11 @@ impl AutoDetectManager { /// The client's timer has millisecond resolution, so a burst that crosses a /// fast link within one millisecond reports a zero time. That is a real /// measurement, bounded below by one millisecond, not a failed one. -fn measured_bandwidth_kbps(response: &AutoDetectResponse) -> Option { - let AutoDetectResponse::BandwidthMeasureResults { - time_delta_ms, - byte_count, - .. - } = response - else { - return None; - }; - if *byte_count == 0 { +fn measured_bandwidth_kbps(time_delta_ms: u32, byte_count: u32) -> Option { + if byte_count == 0 { return None; } - let kbps = u64::from(*byte_count) * 8 / u64::from((*time_delta_ms).max(1)); + let kbps = u64::from(byte_count) * 8 / u64::from(time_delta_ms.max(1)); Some(u32::try_from(kbps).unwrap_or(u32::MAX)) } diff --git a/crates/ironrdp-server/src/server.rs b/crates/ironrdp-server/src/server.rs index b4d8acc1b3..dac83b6e18 100644 --- a/crates/ironrdp-server/src/server.rs +++ b/crates/ironrdp-server/src/server.rs @@ -3448,10 +3448,10 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } - // Until a large graphics write has been bracketed (see the Egfx - // arm), measure bandwidth over a window of ticks: Start on one, - // Stop several later, with ordinary traffic in between counted by - // the client, then a Bandwidth Measure Results PDU in reply. Until + // While no large graphics write has been bracketed recently (see + // `write_egfx_over_tcp`), measure bandwidth over a window of ticks: Start on + // one, Stop several later, with ordinary traffic in between counted + // by the client, then a Bandwidth Measure Results PDU in reply. Until // one has completed there is no characteristics result to send. if let Some(pdu) = ad.build_bandwidth_measure() { let data = encode_autodetect_request(pdu, message_channel_id, user_channel_id)?; @@ -3596,39 +3596,40 @@ impl RdpServer { // go out back to back on the same stream, so the client times a // burst rather than whatever idle stretch a fixed window happens // to span. - let bracket = match (self.autodetect.as_mut(), message_channel_id) { - (Some(ad), Some(message_channel_id)) => ad - .begin_bandwidth_measure(data.len(), monotonic_now_ms()) - .map(|start| (start, message_channel_id)), - _ => None, - }; - let bracket_channel_id = match bracket { - Some((start, message_channel_id)) => { - let start = encode_autodetect_request(start, message_channel_id, user_channel_id)?; - writer - .write_all(&start) - .await - .map_err(|e| ServerError::io("write_all", e))?; - Some(message_channel_id) - } - None => None, + // + // The Stop goes out only for the Start this write sent, which + // `bracket_channel_id` records. The manager outlives the + // connection, so after a failed Start a stale bracket can still be + // pending when the next connection writes, and asking the manager + // alone would send that client a Stop it has no Start for. + let bracket_channel_id = if let (Some(ad), Some(message_channel_id)) = + (self.autodetect.as_mut(), message_channel_id) + && let Some(start) = ad.begin_bandwidth_measure(data.len(), monotonic_now_ms()) + { + let start = encode_autodetect_request(start, message_channel_id, user_channel_id)?; + writer + .write_all(&start) + .await + .map_err(|e| ServerError::io("write_all", e))?; + Some(message_channel_id) + } else { + None }; writer .write_all(&data) .await .map_err(|e| ServerError::io("write_all", e))?; - if let Some(message_channel_id) = bracket_channel_id { - if let Some(stop) = self + if let Some(message_channel_id) = bracket_channel_id + && let Some(stop) = self .autodetect .as_mut() .and_then(AutoDetectManager::end_bandwidth_measure) - { - let stop = encode_autodetect_request(stop, message_channel_id, user_channel_id)?; - writer - .write_all(&stop) - .await - .map_err(|e| ServerError::io("write_all", e))?; - } + { + let stop = encode_autodetect_request(stop, message_channel_id, user_channel_id)?; + writer + .write_all(&stop) + .await + .map_err(|e| ServerError::io("write_all", e))?; } Ok(()) } diff --git a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs index 9d5d0599df..273f5147d9 100644 --- a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs @@ -1,5 +1,9 @@ use ironrdp_pdu::rdp::autodetect::{AutoDetectRequest, AutoDetectResponse}; -use ironrdp_server::autodetect::{AutoDetectManager, AutoDetectOutcome, BW_BRACKET_MIN_BYTES}; +use ironrdp_server::autodetect::{AutoDetectManager, AutoDetectOutcome}; + +/// The smallest graphics write the server brackets with a bandwidth measurement, 10 KiB. +/// Mirrors a tuning constant that is not public, so retuning it must update this test. +const BW_BRACKET_MIN_BYTES: usize = 10 * 1024; /// Upper bound on ticks to drive while waiting for a bandwidth transaction: /// generous enough to cover the pacing plus the window, tight enough that a @@ -439,16 +443,102 @@ fn bracketed_measurements_are_paced_and_never_overlap() { complete_bracket(&mut mgr, &next, 1_010); } -/// Once the session has large writes to measure across, the tick window (which -/// times idle stretches as well as traffic) stops opening. +/// While large writes keep coming, the tick window (which times idle stretches +/// as well as traffic) stays closed. #[test] -fn the_tick_window_retires_after_a_bracketed_measurement() { +fn the_tick_window_stays_closed_while_brackets_recur() { + let mut mgr = AutoDetectManager::new(); + for round in 0..16u64 { + let now_ms = round * 1_000; + let start = mgr.begin_bandwidth_measure(64 * 1024, now_ms).expect("measured"); + complete_bracket(&mut mgr, &start, now_ms + 10); + for _ in 0..4 { + assert!(mgr.build_bandwidth_measure().is_none()); + } + } +} + +/// A session that stops sending large writes, such as one whose only large +/// frame was the initial render, goes back to the tick window rather than +/// keeping the last bracketed figure forever. +#[test] +fn the_tick_window_resumes_when_brackets_stop() { let mut mgr = AutoDetectManager::new(); let start = mgr.begin_bandwidth_measure(64 * 1024, 0).expect("measured"); complete_bracket(&mut mgr, &start, 10); + assert!( + mgr.build_bandwidth_measure().is_none(), + "not on the tick right after a bracket" + ); + drive_bandwidth_start_and_stop(&mut mgr); +} + +/// A bracketed Start whose Stop is never sent does not block bandwidth +/// measurement for the rest of the session. +#[test] +fn an_unended_bracket_expires() { + let mut mgr = AutoDetectManager::new(); + let _ = mgr.begin_bandwidth_measure(64 * 1024, 0).expect("measured"); + mgr.expire_stale_probes(29_999, 30_000); + assert!( + mgr.begin_bandwidth_measure(64 * 1024, 29_999).is_none(), + "still pending before the maximum age" + ); + mgr.expire_stale_probes(30_000, 30_000); + let next = mgr + .begin_bandwidth_measure(64 * 1024, 30_000) + .expect("a new bracket once the stale one expired"); + complete_bracket(&mut mgr, &next, 30_010); +} + +/// A Stop the client never answers does not block bandwidth measurement for +/// the rest of the session, and results that arrive after it expired are not +/// taken as a measurement. +#[test] +fn an_unanswered_measurement_expires() { + let mut mgr = AutoDetectManager::new(); + let sequence = drive_bandwidth_start_and_stop(&mut mgr); + // The tick has no clock, so the first expiry pass stamps the measurement. + mgr.expire_stale_probes(1_000, 30_000); + mgr.expire_stale_probes(30_999, 30_000); + assert!( + mgr.begin_bandwidth_measure(64 * 1024, 30_999).is_none(), + "still pending before the maximum age" + ); + mgr.expire_stale_probes(31_000, 30_000); + let late = AutoDetectResponse::BandwidthMeasureResults { + sequence_number: sequence, + response_type: ironrdp_pdu::rdp::autodetect::BW_RESULTS_CONTINUOUS, + time_delta_ms: 2, + byte_count: 20_000, + }; + assert_eq!(mgr.handle_response(&late, 31_010), AutoDetectOutcome::Unmatched); + drive_bandwidth_start_and_stop(&mut mgr); +} + +/// A tick window stays open across ticks however far apart they are: only the +/// wait for results after its Stop is timed. +#[test] +fn a_slow_tick_window_is_not_cut_short() { + let mut mgr = AutoDetectManager::new(); + let mut now_ms = 0; + let mut start_sequence = None; for _ in 0..MAX_BANDWIDTH_TICKS { - assert!(mgr.build_bandwidth_measure().is_none()); + now_ms += 20_000; + mgr.expire_stale_probes(now_ms, 30_000); + match (start_sequence, mgr.build_bandwidth_measure()) { + (_, None) => {} + (None, Some(AutoDetectRequest::BandwidthMeasureStart { sequence_number, .. })) => { + start_sequence = Some(sequence_number); + } + (Some(start), Some(AutoDetectRequest::BandwidthMeasureStop { sequence_number, .. })) => { + assert_eq!(sequence_number, start); + return; + } + (_, Some(other)) => panic!("unexpected {other:?}"), + } } + panic!("the tick window never sent its Stop"); } /// Two consecutive measurements landing on the identical kbps figure (plausible on a