Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
157 changes: 141 additions & 16 deletions crates/ironrdp-server/src/autodetect.rs
Original file line number Diff line number Diff line change
@@ -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()`].
//!
Expand Down Expand Up @@ -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(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;

/// Minimum spacing between Network Characteristics Result PDUs.
///
/// The reported figures are averaged over a window and change far more slowly
Expand Down Expand Up @@ -68,6 +79,8 @@ pub struct AutoDetectManager {
bandwidth_kbps: Option<u32>,
/// Counts RTT ticks to pace bandwidth measurements (see [`BW_MEASURE_INTERVAL_TICKS`]).
bw_tick_count: u32,
/// When the last bracketed measurement began, for pacing brackets.
last_bracket_start_ms: Option<u64>,
/// When the last Network Characteristics Result was built, for pacing.
last_netchar_result_ms: Option<u64>,
/// Lowest RTT observed over the whole session, for the wire `baseRTT`.
Expand Down Expand Up @@ -114,8 +127,17 @@ 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.
/// `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<u64> },
}

impl AutoDetectManager {
Expand All @@ -127,6 +149,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,
Expand All @@ -141,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)
}
Expand All @@ -169,20 +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 { .. }) => None,
Some(PendingBandwidth::AwaitingResults { .. } | PendingBandwidth::Bracketing { .. }) => 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,
Expand All @@ -192,6 +216,58 @@ 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 10 KiB, while
/// another measurement is outstanding, and within
/// `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<AutoDetectRequest> {
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);
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))
}

/// 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<AutoDetectRequest> {
let Some(PendingBandwidth::Bracketing { sequence, since_ms }) = self.pending_bw else {
return None;
};
self.pending_bw = Some(PendingBandwidth::AwaitingResults {
sequence,
since_ms: Some(since_ms),
});
Some(AutoDetectRequest::bw_stop_continuous(sequence))
}
Comment thread
glamberson marked this conversation as resolved.

/// Build a Network Characteristics Result reporting the measured network.
///
/// Returns `None` until the network has actually been characterised, which
Expand Down Expand Up @@ -242,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,
Expand Down Expand Up @@ -286,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;
Expand All @@ -299,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 = response.computed_bandwidth_kbps();
self.bandwidth_kbps = measured_bandwidth_kbps(*time_delta_ms, *byte_count);
if self.bandwidth_kbps.is_some() {
self.bandwidth_is_fresh = true;
}
Expand Down Expand Up @@ -357,14 +437,59 @@ 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
}
}

/// 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(time_delta_ms: u32, byte_count: u32) -> Option<u32> {
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 {
Expand Down
79 changes: 69 additions & 10 deletions crates/ironrdp-server/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 => {
Expand Down Expand Up @@ -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.
// 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)?;
writer
Expand Down Expand Up @@ -3482,6 +3489,7 @@ impl RdpServer {
messages: Vec<ironrdp_svc::SvcMessage>,
writer: &mut impl FramedWrite,
user_channel_id: u16,
message_channel_id: Option<u16>,
udp_transport: Option<&multitransport::UdpTransportHandle>,
) -> ServerResult<()> {
let drdynvc_channel_id = self
Expand All @@ -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;
};

Expand Down Expand Up @@ -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")]
Expand All @@ -3566,13 +3586,52 @@ impl RdpServer {
writer: &mut impl FramedWrite,
drdynvc_channel_id: u16,
user_channel_id: u16,
message_channel_id: Option<u16>,
) -> 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.
//
// 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))
.map_err(|e| ServerError::io("write_all", e))?;
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))?;
}
Ok(())
}

/// Writes DRDYNVC output, sending the data of any channel the Soft-Sync
Expand Down
Loading
Loading