diff --git a/crates/ironrdp-server/src/autodetect.rs b/crates/ironrdp-server/src/autodetect.rs index fadf157d52..6b4fa61f2b 100644 --- a/crates/ironrdp-server/src/autodetect.rs +++ b/crates/ironrdp-server/src/autodetect.rs @@ -9,7 +9,9 @@ //! //! [MS-RDPBCGR 2.2.14]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpbcgr/dc672839-4f4e-40b1-a71c-cd6a959baa38 +use core::sync::atomic::AtomicU32; use std::collections::VecDeque; +use std::sync::Arc; use ironrdp_pdu::rdp::autodetect::{AutoDetectRequest, AutoDetectResponse}; @@ -388,3 +390,59 @@ pub struct RttSnapshot { /// Number of samples in the current window. pub sample_count: usize, } + +/// The latest auto-detect measurements of one connection, written by the server +/// as they arrive. +/// +/// [`Default`] creates handles at their initial values: no measurement yet, +/// generation 0. The server creates them that way for each connection, so a +/// new connection never sees the figures of the one before it. They stay at +/// those values while auto-detect is disabled (see +/// [`RdpServer::enable_autodetect`](crate::RdpServer::enable_autodetect)). +#[derive(Debug, Clone)] +#[non_exhaustive] +pub struct AutoDetectHandles { + /// Latest round-trip time in milliseconds, or `u32::MAX` until the first + /// measurement. Updated on each RTT Measure Response, so a reader gets a + /// fresh, frame-traffic-independent network RTT for flow control. + pub rtt: Arc, + + /// Lowest RTT of the connection in milliseconds (`baseRTT` per + /// [MS-RDPBCGR] 2.2.14.1.5), or `u32::MAX` until the first measurement. + /// Updated at the same point as [`Self::rtt`], but unlike it this figure + /// never rises: pair it with that figure to derive queueing delay + /// (`averageRTT - baseRTT`), which `rtt` alone cannot give since it is a + /// sliding-window value that rises as low samples age out. + pub baseline_rtt: Arc, + + /// Latest measured bandwidth in kilobits per second, or `u32::MAX` until + /// the first measurement completes. Updated whenever a Bandwidth Measure + /// Results response completes a measurement, with the figure the server + /// also reports to the client on the wire; a measurement without a usable + /// figure sets it back to `u32::MAX`. + pub bandwidth: Arc, + + /// Pairs with [`Self::bandwidth`]: increments every time that figure is + /// republished, since the figure itself repeats too often to be its own + /// freshness signal. Load this with `Ordering::Acquire` to detect a fresh + /// measurement window, then read the bandwidth: the server increments this + /// with `Ordering::Release` after storing the value, so the bandwidth read + /// is at least as new as the generation observed. It is not an exact pair: + /// if the next window closes between the two reads, the value can already + /// belong to that later window. + /// + /// Starts at 0 with each connection, so compare it only with generations + /// read from this same handle. + pub bandwidth_generation: Arc, +} + +impl Default for AutoDetectHandles { + fn default() -> Self { + Self { + rtt: Arc::new(AtomicU32::new(u32::MAX)), + baseline_rtt: Arc::new(AtomicU32::new(u32::MAX)), + bandwidth: Arc::new(AtomicU32::new(u32::MAX)), + bandwidth_generation: Arc::new(AtomicU32::new(0)), + } + } +} diff --git a/crates/ironrdp-server/src/builder.rs b/crates/ironrdp-server/src/builder.rs index a511adb8dd..6d361ec7e3 100644 --- a/crates/ironrdp-server/src/builder.rs +++ b/crates/ironrdp-server/src/builder.rs @@ -1,5 +1,4 @@ use core::net::SocketAddr; -use core::sync::atomic::{AtomicBool, AtomicU32}; use std::sync::Arc; use ironrdp_pdu::codecs::rfx::Quant; @@ -8,7 +7,7 @@ use ironrdp_pdu::rdp::session_info::ServerAutoReconnect; use tokio_rustls::TlsAcceptor; use super::clipboard::CliprdrServerFactory; -use super::display::{DesktopSize, RdpServerDisplay}; +use super::display::{DesktopSize, DisplayContext, RdpServerDisplay}; #[cfg(feature = "egfx")] use super::gfx::GfxServerFactory; use super::handler::{KeyboardEvent, MouseEvent, RdpServerInputHandler}; @@ -56,11 +55,6 @@ pub struct BuilderDone { gfx_factory: Option>, #[cfg(feature = "usb")] usb_factory: Option>, - display_suppressed: Option>, - autodetect_rtt: Option>, - autodetect_baseline_rtt: Option>, - autodetect_bandwidth: Option>, - autodetect_bandwidth_generation: Option>, honor_client_desktop_size: Option, auto_reconnect_cookie: Option, connection_policy: ConnectionPolicy, @@ -171,11 +165,6 @@ impl RdpServerBuilder { gfx_factory: None, #[cfg(feature = "usb")] usb_factory: None, - display_suppressed: None, - autodetect_rtt: None, - autodetect_baseline_rtt: None, - autodetect_bandwidth: None, - autodetect_bandwidth_generation: None, honor_client_desktop_size: None, connection_policy: ConnectionPolicy::default(), auto_reconnect_cookie: None, @@ -207,11 +196,6 @@ impl RdpServerBuilder { gfx_factory: None, #[cfg(feature = "usb")] usb_factory: None, - display_suppressed: None, - autodetect_rtt: None, - autodetect_baseline_rtt: None, - autodetect_bandwidth: None, - autodetect_bandwidth_generation: None, honor_client_desktop_size: None, connection_policy: ConnectionPolicy::default(), auto_reconnect_cookie: None, @@ -293,26 +277,6 @@ impl RdpServerBuilder { self } - /// Share the server's "display suppressed" flag with the display - /// backend before construction. - /// - /// The flag is `true` while the connected client has sent - /// `SuppressOutput { desktop_rect: None }` (e.g., mstsc minimized). - /// Display backends that want to skip frame emission while the - /// client is minimized create one `Arc` in the - /// application, hand a clone to the display, and pass the same - /// `Arc` here so the server's per-connection PDU handler writes to - /// the same instance the backend reads. - /// - /// When this is not called, the server allocates its own internal - /// flag (still readable via [`RdpServer::display_suppressed_handle`]) - /// — useful when the backend can call `display_suppressed_handle()` - /// after construction to obtain a handle, rather than sharing one in. - pub fn with_display_suppressed_handle(mut self, handle: Arc) -> Self { - self.state.display_suppressed = Some(handle); - self - } - /// Negotiate each session at the desktop size the client requests in its /// Client Core Data, rather than the size reported by the display handler. /// @@ -387,55 +351,6 @@ impl RdpServerBuilder { self } - /// Inject a shared NetworkAutoDetect RTT handle (milliseconds, `u32::MAX` - /// until the first measurement). The server writes the latest measured RTT - /// to the same instance the backend reads. When not called, the server - /// allocates its own (still readable via - /// [`RdpServer::autodetect_rtt_handle`]). The value stays `u32::MAX` unless - /// auto-detect is enabled via [`RdpServer::enable_autodetect`]. - pub fn with_autodetect_rtt_handle(mut self, handle: Arc) -> Self { - self.state.autodetect_rtt = Some(handle); - self - } - - /// Inject a shared session-lifetime baseline RTT handle (milliseconds, - /// `u32::MAX` until the first measurement; see - /// [`RdpServer::autodetect_baseline_rtt_handle`] for what distinguishes - /// this from [`Self::with_autodetect_rtt_handle`]). The server writes the - /// latest baseline to the same instance the backend reads. When not - /// called, the server allocates its own (still readable via - /// [`RdpServer::autodetect_baseline_rtt_handle`]). The value stays - /// `u32::MAX` unless auto-detect is enabled via - /// [`RdpServer::enable_autodetect`]. - pub fn with_autodetect_baseline_rtt_handle(mut self, handle: Arc) -> Self { - self.state.autodetect_baseline_rtt = Some(handle); - self - } - - /// Inject a shared NetworkAutoDetect bandwidth handle (kilobits per - /// second, `u32::MAX` until the first measurement completes). The server - /// writes the latest measured bandwidth to the same instance the backend - /// reads. When not called, the server allocates its own (still readable - /// via [`RdpServer::autodetect_bandwidth_handle`]). The value stays - /// `u32::MAX` unless auto-detect is enabled via - /// [`RdpServer::enable_autodetect`]. - pub fn with_autodetect_bandwidth_handle(mut self, handle: Arc) -> Self { - self.state.autodetect_bandwidth = Some(handle); - self - } - - /// Inject a shared handle that increments every time a Bandwidth Measure - /// transaction completes, whether or not it produced a usable figure. - /// Pairs with [`Self::with_autodetect_bandwidth_handle`]: the bandwidth - /// figure alone repeats too often to tell a fresh measurement window - /// apart from a stale one. When not called, the server allocates its own - /// (still readable via - /// [`RdpServer::autodetect_bandwidth_generation_handle`]). - pub fn with_autodetect_bandwidth_generation_handle(mut self, handle: Arc) -> Self { - self.state.autodetect_bandwidth_generation = Some(handle); - self - } - /// Provision the Server Auto-Reconnect Cookie (MS-RDPBCGR 2.2.4.2 /// `ARC_SC_PRIVATE_PACKET`) handed to the client during logon. /// @@ -538,13 +453,8 @@ impl RdpServerBuilder { self.state.connection_handler, #[cfg(feature = "egfx")] self.state.gfx_factory, - self.state.display_suppressed, #[cfg(feature = "usb")] self.state.usb_factory, - self.state.autodetect_rtt, - self.state.autodetect_baseline_rtt, - self.state.autodetect_bandwidth, - self.state.autodetect_bandwidth_generation, ); server.set_credential_validator(self.state.credential_validator); server.set_auto_reconnect_cookie(self.state.auto_reconnect_cookie); @@ -577,7 +487,7 @@ impl RdpServerDisplay for NoopDisplay { DesktopSize { width: 0, height: 0 } } - async fn updates(&mut self) -> ServerResult> { + async fn updates(&mut self, _: DisplayContext) -> ServerResult> { Ok(Box::new(NoopDisplayUpdates {})) } } diff --git a/crates/ironrdp-server/src/display.rs b/crates/ironrdp-server/src/display.rs index 3d232cf37c..085a5c5056 100644 --- a/crates/ironrdp-server/src/display.rs +++ b/crates/ironrdp-server/src/display.rs @@ -1,4 +1,6 @@ use core::num::{NonZeroU16, NonZeroUsize}; +use core::sync::atomic::AtomicBool; +use std::sync::Arc; use bytes::{Bytes, BytesMut}; use ironrdp_displaycontrol::pdu::DisplayControlMonitorLayout; @@ -6,6 +8,7 @@ use ironrdp_graphics::diff; use ironrdp_pdu::pointer::PointerPositionAttribute; use tracing::{debug, warn}; +use crate::autodetect::AutoDetectHandles; use crate::error::ServerResult; #[rustfmt::skip] @@ -281,12 +284,50 @@ pub trait RdpServerDisplayUpdates { async fn next_update(&mut self) -> ServerResult>; } +/// What a connection publishes to its display backend, handed to +/// [`RdpServerDisplay::updates`]. +/// +/// The server writes these values while the connection runs; the backend keeps +/// the handles it needs and reads them. Each connection has its own handles, +/// starting from the initial values documented on each field, so nothing the +/// previous connection set or measured carries over. A +/// Deactivation-Reactivation Sequence keeps the connection, so `updates` is +/// then called again with a context holding the same handles. +#[derive(Debug)] +#[non_exhaustive] +pub struct DisplayContext { + /// `true` while the client has sent `SuppressOutput { desktop_rect: None }` + /// (e.g., mstsc minimized), `false` at the start of the connection. Cleared + /// on `SuppressOutput { Some(rect) }` or `RefreshRectangle`. + /// + /// A backend can skip frame emission while it's set, so the client doesn't + /// accumulate a backlog of frames it can't present until refocus. + /// + /// **Caveat:** some clients (notably mstsc) send + /// `SuppressOutput { desktop_rect: None }` during their connect + /// handshake *before* their display surface is fully initialized; a + /// backend that honors the flag blindly will block that first frame + /// and leave the client with a half-initialized surface that doesn't + /// recover on un-suppress (visible as a frozen desktop on first + /// connect). Backends are advised to defer acting on the flag until + /// after the first frame has been delivered to the client, and to + /// debounce transient flaps (some clients pulse this PDU under wire + /// pressure on heavy CPU/IO loads) — e.g., only engage the gate once + /// the flag has been steady-`true` for ~1 s. + pub display_suppressed: Arc, + + /// The connection's auto-detect measurements, for flow control. + pub autodetect: AutoDetectHandles, +} + /// Display for an RDP server /// /// # Example /// /// ``` -/// use ironrdp_server::{DesktopSize, DisplayUpdate, RdpServerDisplay, RdpServerDisplayUpdates, ServerResult}; +/// use ironrdp_server::{ +/// DesktopSize, DisplayContext, DisplayUpdate, RdpServerDisplay, RdpServerDisplayUpdates, ServerResult, +/// }; /// /// pub struct DisplayUpdates { /// receiver: tokio::sync::mpsc::Receiver, @@ -310,7 +351,7 @@ pub trait RdpServerDisplayUpdates { /// DesktopSize { width: self.width, height: self.height } /// } /// -/// async fn updates(&mut self) -> ServerResult> { +/// async fn updates(&mut self, _ctx: DisplayContext) -> ServerResult> { /// Ok(Box::new(DisplayUpdates { receiver: todo!() })) /// } /// } @@ -329,7 +370,11 @@ pub trait RdpServerDisplay: Send { } /// Return a display updates receiver - async fn updates(&mut self) -> ServerResult>; + /// + /// Called when a connection's session starts, and again after each + /// Deactivation-Reactivation Sequence of that connection. `ctx` carries + /// what the connection publishes to the display; see [`DisplayContext`]. + async fn updates(&mut self, ctx: DisplayContext) -> ServerResult>; /// Request a new size for the display fn request_layout(&mut self, layout: DisplayControlMonitorLayout) { diff --git a/crates/ironrdp-server/src/lib.rs b/crates/ironrdp-server/src/lib.rs index b1f9554080..473c108765 100644 --- a/crates/ironrdp-server/src/lib.rs +++ b/crates/ironrdp-server/src/lib.rs @@ -29,8 +29,8 @@ mod urbdrc; pub use clipboard::CliprdrServerFactory; pub use display::{ - BitmapUpdate, ColorPointer, DesktopSize, DisplayUpdate, Framebuffer, LargePointer, PixelFormat, RGBAPointer, - RdpServerDisplay, RdpServerDisplayUpdates, + BitmapUpdate, ColorPointer, DesktopSize, DisplayContext, DisplayUpdate, Framebuffer, LargePointer, PixelFormat, + RGBAPointer, RdpServerDisplay, RdpServerDisplayUpdates, }; pub use echo::{EchoDvcBridge, EchoRoundTripMeasurement, EchoServerHandle, EchoServerMessage}; pub use error::{ServerError, ServerErrorExt, ServerErrorKind, ServerResult, ServerResultExt}; diff --git a/crates/ironrdp-server/src/server.rs b/crates/ironrdp-server/src/server.rs index 434cf140df..2ba69b8758 100644 --- a/crates/ironrdp-server/src/server.rs +++ b/crates/ironrdp-server/src/server.rs @@ -1,7 +1,7 @@ use core::cell::RefCell; use core::fmt; use core::net::{IpAddr, SocketAddr}; -use core::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering}; +use core::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use core::time::Duration; #[cfg(feature = "usb")] use std::collections::HashMap; @@ -52,9 +52,9 @@ use tokio::task; use tokio_rustls::TlsAcceptor; use tracing::{debug, error, info, trace, warn}; -use crate::autodetect::{AutoDetectManager, AutoDetectOutcome, RttSnapshot}; +use crate::autodetect::{AutoDetectHandles, AutoDetectManager, AutoDetectOutcome}; use crate::clipboard::CliprdrServerFactory; -use crate::display::{DisplayUpdate, RdpServerDisplay}; +use crate::display::{DisplayContext, DisplayUpdate, RdpServerDisplay}; use crate::echo::{EchoDvcBridge, EchoServerHandle, EchoServerMessage, build_echo_request}; use crate::encoder::{UpdateEncoder, UpdateEncoderCodecs}; use crate::error::{ServerError, ServerErrorExt as _, ServerErrorKind, ServerResult}; @@ -614,23 +614,12 @@ impl DisplayControlHandler for DisplayControlBackend { } #[cfg(feature = "usb")] +#[derive(Default)] struct ServerUsbManager { - factory: Box, comp_iface_alloc: InterfaceAlloc, router: HashMap>, } -#[cfg(feature = "usb")] -impl ServerUsbManager { - fn new(inner: Box) -> Self { - Self { - factory: inner, - comp_iface_alloc: InterfaceAlloc::default(), - router: HashMap::new(), - } - } -} - /// Selects who performs the TLS handshake for a connection accepted via /// [`RdpServer::run_connection_with`]. #[derive(Debug, Clone, Copy)] @@ -657,7 +646,7 @@ pub enum TransportTls { /// ``` /// use ironrdp_server::{RdpServer, RdpServerInputHandler, RdpServerDisplay, RdpServerDisplayUpdates}; /// -///# use ironrdp_server::{DisplayUpdate, DesktopSize, KeyboardEvent, MouseEvent, ServerResult}; +///# use ironrdp_server::{DisplayContext, DisplayUpdate, DesktopSize, KeyboardEvent, MouseEvent, ServerResult}; ///# use tokio_rustls::TlsAcceptor; ///# struct NoopInputHandler; ///# impl RdpServerInputHandler for NoopInputHandler { @@ -670,7 +659,7 @@ pub enum TransportTls { ///# async fn size(&mut self) -> DesktopSize { ///# todo!() ///# } -///# async fn updates(&mut self) -> ServerResult> { +///# async fn updates(&mut self, _: DisplayContext) -> ServerResult> { ///# todo!() ///# } ///# } @@ -710,7 +699,6 @@ pub struct RdpServer { // FIXME: replace with a channel and poll/process the handler? handler: Arc>>, display: Arc>>, - static_channels: StaticChannelSet, static_channel_factories: Vec>, sound_factory: Option>, cliprdr_factory: Option>, @@ -723,7 +711,7 @@ pub struct RdpServer { #[cfg(feature = "egfx")] gfx_handle: Option, #[cfg(feature = "usb")] - usb_man: Option, + usb_factory: Option>, ev_sender: mpsc::UnboundedSender, ev_receiver: Arc>>, creds: Option, @@ -732,7 +720,8 @@ pub struct RdpServer { /// The local address the current client reached. See /// [`Self::set_connection_local_addr`]. connection_local_addr: Option, - autodetect: Option, + /// Whether each connection runs auto-detect; see [`Self::enable_autodetect`]. + autodetect_enabled: bool, heartbeat: Option, connection_handler: Option>, /// Anti-storm net for [`ConnectionPolicy::Preempt`]: the @@ -749,57 +738,6 @@ pub struct RdpServer { /// IP, since the source port changes on every reconnect. Cleared once a /// session ends on its own terms rather than being replaced. recently_evicted: Option, - /// True while the client has sent `SuppressOutput { desktop_rect: None }` - /// — the standard RDP "I don't need display updates right now" signal - /// (mstsc raises it on window minimize). Cleared on - /// `SuppressOutput { Some(rect) }` or `RefreshRectangle` (sent on - /// refocus). Exposed via [`Self::display_suppressed_handle`] so display - /// backends can hold a clone and skip frame emission while it's set — - /// without this, a server keeps streaming high-bitrate - /// EGFX/H.264 frames into a minimized client, which accumulates them - /// and locks up its input dispatch for seconds on refocus while it - /// chews through the backlog. - display_suppressed: Arc, - - /// Latest NetworkAutoDetect round-trip time in milliseconds, or `u32::MAX` - /// until the first measurement (and while auto-detect is disabled). Updated - /// on each RTT Measure Response when auto-detect is enabled (see - /// [`Self::enable_autodetect`]). Exposed via [`Self::autodetect_rtt_handle`] - /// so display backends can read a fresh, frame-traffic-independent network - /// RTT for flow control. - autodetect_rtt: Arc, - - /// Session-lifetime lowest RTT in milliseconds (`baseRTT` per MS-RDPBCGR - /// 2.2.14.1.5), or `u32::MAX` until the first measurement. Unlike - /// [`Self::autodetect_rtt`], this never rises: it is the floor over the - /// whole session, not a sliding-window figure, which is what makes - /// `averageRTT - baseRTT` a queueing-delay signal rather than two - /// unrelated latency numbers. Updated at the same point as - /// [`Self::autodetect_rtt`]. Exposed via - /// [`Self::autodetect_baseline_rtt_handle`]. - autodetect_baseline_rtt: Arc, - - /// Latest NetworkAutoDetect measured bandwidth in kilobits per second, or - /// `u32::MAX` until the first measurement completes (and while auto-detect - /// is disabled). Updated whenever a Bandwidth Measure Results response is - /// processed, same trigger point as [`Self::autodetect_rtt`]. Exposed via - /// [`Self::autodetect_bandwidth_handle`]: without it, the server can tell - /// the *client* its measured bandwidth over the wire but has no way to - /// tell the embedder, which the connect-time figure carried to the client - /// alone does not fix. - autodetect_bandwidth: Arc, - - /// Increments every time a Bandwidth Measure transaction completes, - /// whether or not it produced a usable figure (see - /// [`Self::autodetect_bandwidth`]'s doc comment on the None case). - /// [`Self::autodetect_bandwidth`] alone cannot tell an embedder "a new - /// window just closed" apart from "the value happens to repeat": that - /// value repeats often (a quiet link reads the same low figure for - /// several consecutive windows), so diffing it is not a valid freshness - /// signal. Incremented with `Release` after the bandwidth value is - /// stored, so an `Acquire` load of it makes that value visible. - /// Exposed via [`Self::autodetect_bandwidth_generation_handle`]. - autodetect_bandwidth_generation: Arc, /// Optional Server Auto-Reconnect Cookie (MS-RDPBCGR 2.2.4.2 /// `ARC_SC_PRIVATE_PACKET`). When `Some`, the server validates a returning @@ -820,34 +758,6 @@ pub struct RdpServer { /// Tracks whether the current cookie has reached a client. Subsequent /// connections and hourly updates replace it with a new random. auto_reconnect_sent: bool, - - /// Abort handle of the current connection's pending UDP multitransport - /// accept, if one is running. A client that could not establish the - /// sideband transport answers with a failure Initiate Multitransport - /// Response, often after finalization has completed; the message-channel - /// handler uses this to stop the accept instead of letting it hold its - /// socket until `multitransport::UDP_ACCEPT_TIMEOUT`. - pending_udp_accept_abort: Option, - /// Whether the current connection negotiated SOFT_SYNC_TCP_TO_UDP, one of - /// the two conditions for migrating a channel onto the sideband - /// transport (see [`Self::udp_migration_allowed`]). - soft_sync_negotiated: bool, - /// Whether the current connection may migrate EGFX onto the sideband - /// transport. MS-RDPEDYC 3.1.5.3/3.3.5.3.1: Soft-Sync MUST NOT be used - /// unless both peers negotiated SOFT_SYNC_TCP_TO_UDP and a successful - /// Initiate Multitransport Response was received. Finalization sets it - /// when the response came during finalization; the message-channel - /// handler sets it when the response comes later, which is the usual - /// case with mstsc, whose UDP bootstrap outlasts the TCP finalization. - udp_migration_allowed: bool, - /// Whether the current connection's EGFX data has started going over the - /// sideband transport, so the switch is logged once. - egfx_on_udp: bool, - /// Tunnel payloads that arrived while a Soft-Sync request was waiting for - /// its response. The client writes on the tunnel right after sending the - /// response over TCP, so its first tunnel data can overtake it; this holds - /// that data until the response is in instead of dropping it. - early_tunnel_payloads: VecDeque>, } /// Cloneable handle for updating the Server Auto-Reconnect Cookie while @@ -1180,6 +1090,110 @@ struct NegotiatedConnection { local_addr: Option, } +/// State that belongs to a single client connection. +/// +/// Created when a [`NegotiatedConnection`] enters finalization and dropped +/// when that connection ends, on every path out of it. Nothing stored here can +/// outlive its connection or be observed by the next one, so it needs no +/// manual reset between connections. +#[derive(Default)] +struct ConnectionState { + /// The static channels this connection negotiated. + /// + /// Empty until the acceptor hands them back after finalization. Handed to + /// the acceptor again for each deactivation-reactivation, and returned + /// once that finalization completes. + static_channels: StaticChannelSet, + /// Whether the client advertised `SUPPORT_HEART_BEAT_PDU` in its GCC + /// Client Core Data. + client_supports_heartbeat: bool, + /// Auto-detect state, present when the server has auto-detect enabled. + /// + /// Probes in flight, RTT samples and the session-lifetime lowest RTT all + /// describe this connection's network path, so they start over with each + /// connection. + autodetect: Option, + /// The latest auto-detect measurements, published to the display backend + /// through [`DisplayContext::autodetect`]. + autodetect_handles: AutoDetectHandles, + /// Whether the client asked the server to stop sending display updates + /// (`SuppressOutput { desktop_rect: None }`), published to the display + /// backend through [`DisplayContext::display_suppressed`]. Without it, a + /// server keeps streaming high-bitrate EGFX/H.264 frames into a minimized + /// client, which accumulates them and locks up its input dispatch for + /// seconds on refocus while it chews through the backlog. + display_suppressed: Arc, + #[cfg(feature = "usb")] + usb_man: ServerUsbManager, + /// Abort handle of this connection's pending UDP multitransport accept, + /// if one is running. A client that could not establish the sideband + /// transport answers with a failure Initiate Multitransport Response, + /// often after finalization has completed; the message-channel handler + /// uses this to stop the accept instead of letting it hold its socket + /// until `multitransport::UDP_ACCEPT_TIMEOUT`. + pending_udp_accept_abort: Option, + /// Whether this connection negotiated SOFT_SYNC_TCP_TO_UDP, one of the two + /// conditions for migrating a channel onto the sideband transport (see + /// [`Self::udp_migration_allowed`]). + soft_sync_negotiated: bool, + /// Whether this connection may migrate EGFX onto the sideband transport. + /// MS-RDPEDYC 3.1.5.3/3.3.5.3.1: Soft-Sync MUST NOT be used unless both + /// peers negotiated SOFT_SYNC_TCP_TO_UDP and a successful Initiate + /// Multitransport Response was received. Finalization sets it when the + /// response came during finalization; the message-channel handler sets it + /// when the response comes later, which is the usual case with mstsc, + /// whose UDP bootstrap outlasts the TCP finalization. + udp_migration_allowed: bool, + /// Whether this connection's EGFX data has started going over the + /// sideband transport, so the switch is logged once. + egfx_on_udp: bool, + /// Tunnel payloads that arrived while a Soft-Sync request was waiting for + /// its response. The client writes on the tunnel right after sending the + /// response over TCP, so its first tunnel data can overtake it; this holds + /// that data until the response is in instead of dropping it. + early_tunnel_payloads: VecDeque>, +} + +impl ConnectionState { + /// The handles this connection hands to the display backend, cloned at + /// each activation so every call to [`RdpServerDisplay::updates`] on the + /// same connection gets the same ones. + fn display_context(&self) -> DisplayContext { + DisplayContext { + display_suppressed: Arc::clone(&self.display_suppressed), + autodetect: self.autodetect_handles.clone(), + } + } + + fn get_svc_processor(&mut self) -> Option<&mut T> { + self.static_channels + .get_by_type_mut::() + .and_then(|svc| svc.channel_processor_downcast_mut()) + } + + fn get_channel_id_by_type(&self) -> Option { + self.static_channels.get_channel_id_by_type::() + } + + #[cfg(feature = "usb")] + fn remove_usb_device(&mut self, dvc_id: DynamicChannelId) { + let Some(device) = self.usb_man.router.remove(&dvc_id) else { + trace!(dvc_id, "Closed USB device is absent from request router"); + return; + }; + + // Set the terminal state before failing waiters: a woken PendingRequest + // must not enqueue CANCEL_REQUEST for a removed DVC. The pending map is + // shared, so dropping the router entry no longer drops it. + device.mark_closed(); + let pending_requests = device.drain_pending(); + debug!( + dvc_id, + pending_requests, "Removed closed USB device from request router" + ); + } +} + /// Advance a stream that is now past the security upgrade: mark the acceptor /// accordingly and, under [`RdpServerSecurity::Hybrid`], run the CredSSP /// exchange. @@ -1480,12 +1494,7 @@ impl RdpServer { mut rdpeai_factory: Option>, connection_handler: Option>, #[cfg(feature = "egfx")] mut gfx_factory: Option>, - display_suppressed: Option>, #[cfg(feature = "usb")] usb_factory: Option>, - autodetect_rtt: Option>, - autodetect_baseline_rtt: Option>, - autodetect_bandwidth: Option>, - autodetect_bandwidth_generation: Option>, ) -> Self { let (ev_sender, ev_receiver) = ServerEvent::create_channel(); if let Some(cliprdr) = cliprdr_factory.as_mut() { @@ -1512,7 +1521,6 @@ impl RdpServer { opts, handler: Arc::new(Mutex::new(handler)), display: Arc::new(Mutex::new(display)), - static_channels: StaticChannelSet::new(), static_channel_factories, sound_factory, cliprdr_factory, @@ -1525,44 +1533,20 @@ impl RdpServer { #[cfg(feature = "egfx")] gfx_handle: None, #[cfg(feature = "usb")] - usb_man: usb_factory.map(ServerUsbManager::new), + usb_factory, ev_sender, ev_receiver: Arc::new(Mutex::new(ev_receiver)), creds: None, credential_validator: None, local_addr: None, connection_local_addr: None, - autodetect: None, + autodetect_enabled: false, heartbeat: None, connection_handler, recently_evicted: None, - display_suppressed: display_suppressed.unwrap_or_else(|| Arc::new(AtomicBool::new(false))), - autodetect_rtt: { - // Reset to the sentinel: an injected handle must not expose a stale value before the first measurement. - let handle = autodetect_rtt.unwrap_or_else(|| Arc::new(AtomicU32::new(u32::MAX))); - handle.store(u32::MAX, Ordering::Relaxed); - handle - }, - autodetect_baseline_rtt: { - let handle = autodetect_baseline_rtt.unwrap_or_else(|| Arc::new(AtomicU32::new(u32::MAX))); - handle.store(u32::MAX, Ordering::Relaxed); - handle - }, - autodetect_bandwidth: { - let handle = autodetect_bandwidth.unwrap_or_else(|| Arc::new(AtomicU32::new(u32::MAX))); - handle.store(u32::MAX, Ordering::Relaxed); - handle - }, - autodetect_bandwidth_generation: autodetect_bandwidth_generation - .unwrap_or_else(|| Arc::new(AtomicU32::new(0))), auto_reconnect_cookie: None, previous_auto_reconnect_cookie: None, auto_reconnect_sent: false, - pending_udp_accept_abort: None, - soft_sync_negotiated: false, - udp_migration_allowed: false, - egfx_on_udp: false, - early_tunnel_payloads: VecDeque::new(), } } @@ -1812,112 +1796,6 @@ impl RdpServer { &self.ev_sender } - #[cfg(feature = "usb")] - fn remove_usb_device(&mut self, dvc_id: DynamicChannelId) { - let Some(usb_man) = self.usb_man.as_mut() else { - warn!("Missing USB device factory"); - return; - }; - - let Some(device) = usb_man.router.remove(&dvc_id) else { - trace!(dvc_id, "Closed USB device is absent from request router"); - return; - }; - - // Set the terminal state before failing waiters: a woken PendingRequest - // must not enqueue CANCEL_REQUEST for a removed DVC. The pending map is - // shared, so dropping the router entry no longer drops it. - device.mark_closed(); - let pending_requests = device.drain_pending(); - debug!( - dvc_id, - pending_requests, "Removed closed USB device from request router" - ); - } - - /// Returns the shared "display suppressed" flag — `true` while the - /// connected client has sent `SuppressOutput { desktop_rect: None }` - /// (e.g., mstsc minimized). - /// - /// Display backends should hold a clone of this `Arc` and skip frame - /// emission while it's set, so the client doesn't accumulate a backlog - /// of frames it can't present until refocus. Cleared by the per- - /// connection PDU handler on `SuppressOutput { Some(rect) }` or - /// `RefreshRectangle`. - /// - /// **Caveat:** some clients (notably mstsc) send - /// `SuppressOutput { desktop_rect: None }` during their connect - /// handshake *before* their display surface is fully initialized; a - /// backend that honors the flag blindly will block that first frame - /// and leave the client with a half-initialized surface that doesn't - /// recover on un-suppress (visible as a frozen desktop on first - /// connect). Backends are advised to defer acting on the flag until - /// after the first frame has been delivered to the client, and to - /// debounce transient flaps (some clients pulse this PDU under wire - /// pressure on heavy CPU/IO loads) — e.g., only engage the gate once - /// the flag has been steady-`true` for ~1 s. - /// - /// The display backend typically needs to share this flag with the - /// server before any client connects (so the same `Arc` is read by - /// the backend's polling thread and written by the per-connection - /// PDU handler). To inject the shared instance at construction time, - /// use [`RdpServerBuilder::with_display_suppressed_handle`](crate::RdpServerBuilder::with_display_suppressed_handle). - /// - /// [crate::RdpServerBuilder]: crate::RdpServerBuilder - pub fn display_suppressed_handle(&self) -> Arc { - Arc::clone(&self.display_suppressed) - } - - /// Returns a handle to the latest NetworkAutoDetect RTT in milliseconds - /// (`u32::MAX` until the first measurement, and while auto-detect is - /// disabled). The server updates it on each RTT Measure Response; backends - /// clone the handle to read a fresh network RTT for flow control. Inject a - /// shared instance at construction with - /// [`RdpServerBuilder::with_autodetect_rtt_handle`](crate::RdpServerBuilder::with_autodetect_rtt_handle). - pub fn autodetect_rtt_handle(&self) -> Arc { - Arc::clone(&self.autodetect_rtt) - } - - /// Returns a handle to the session-lifetime lowest RTT in milliseconds - /// (`baseRTT` per MS-RDPBCGR 2.2.14.1.5; `u32::MAX` until the first - /// measurement, and while auto-detect is disabled). Unlike - /// [`Self::autodetect_rtt_handle`], this figure never rises: pair it with - /// that handle's average to derive queueing delay - /// (`averageRTT - baseRTT`), which `autodetect_rtt_handle` alone cannot - /// give since its figure is a sliding-window value that rises as low - /// samples age out. Inject a shared instance at construction with - /// [`RdpServerBuilder::with_autodetect_baseline_rtt_handle`](crate::RdpServerBuilder::with_autodetect_baseline_rtt_handle). - pub fn autodetect_baseline_rtt_handle(&self) -> Arc { - Arc::clone(&self.autodetect_baseline_rtt) - } - - /// Returns a handle to the latest NetworkAutoDetect measured bandwidth in - /// kilobits per second (`u32::MAX` until the first measurement completes, - /// and while auto-detect is disabled). The server updates it whenever a - /// Bandwidth Measure Results response completes a measurement; backends - /// clone the handle to read the figure the server also reports to the - /// client on the wire. Inject a shared instance at construction with - /// [`RdpServerBuilder::with_autodetect_bandwidth_handle`](crate::RdpServerBuilder::with_autodetect_bandwidth_handle). - pub fn autodetect_bandwidth_handle(&self) -> Arc { - Arc::clone(&self.autodetect_bandwidth) - } - - /// Returns a handle that increments every time a Bandwidth Measure - /// transaction completes, whether or not it produced a usable figure. - /// Pairs with [`Self::autodetect_bandwidth_handle`]: load this with - /// `Ordering::Acquire` to detect a fresh measurement window (the - /// bandwidth figure itself repeats too often to be its own freshness - /// signal), then read the bandwidth handle for the value. The server - /// increments this with `Ordering::Release` after storing the value, so - /// the bandwidth read is at least as new as the generation observed. It - /// is not an exact pair: If the next window closes between the two reads, - /// the value can already belong to that later window. Inject a shared - /// instance at construction with - /// [`RdpServerBuilder::with_autodetect_bandwidth_generation_handle`](crate::RdpServerBuilder::with_autodetect_bandwidth_generation_handle). - pub fn autodetect_bandwidth_generation_handle(&self) -> Arc { - Arc::clone(&self.autodetect_bandwidth_generation) - } - /// Returns the shared ECHO server handle for runtime probe requests and RTT measurements. pub fn echo_handle(&self) -> &EchoServerHandle { &self.echo_handle @@ -1929,10 +1807,11 @@ impl RdpServer { /// separate from the ECHO DVC. It supports bandwidth measurement /// in addition to RTT and works even when DVC is unavailable. /// - /// Send probes via [`ServerEvent::AutoDetectRttRequest`] and - /// query results with [`rtt_snapshot()`](Self::rtt_snapshot). + /// Send probes via [`ServerEvent::AutoDetectRttRequest`]. The display + /// backend reads the results through [`DisplayContext::autodetect`]. Each connection measures its own network + /// path, starting from scratch. pub fn enable_autodetect(&mut self) { - self.autodetect = Some(AutoDetectManager::new()); + self.autodetect_enabled = true; } /// Enable periodic Server Heartbeat PDUs (MS-RDPBCGR 2.2.16.1). @@ -1946,14 +1825,6 @@ impl RdpServer { self.heartbeat = Some(config); } - /// Get the latest auto-detect RTT snapshot. - /// - /// Returns `None` if auto-detect is not enabled or no measurements - /// have been received yet. - pub fn rtt_snapshot(&self) -> Option { - self.autodetect.as_ref().and_then(|ad| ad.snapshot()) - } - /// Returns the shared EGFX server handle for proactive frame submission. /// /// Available after `build_server_with_handle()` returns `Some` during @@ -2030,7 +1901,7 @@ impl RdpServer { #[cfg(feature = "usb")] let dvc = { let mut dvc = dvc; - if self.usb_man.is_some() { + if self.usb_factory.is_some() { dvc = dvc.with_dynamic_channel(UrbdrcControlServer::new(Box::new(UsbControlHandle::new( self.ev_sender.clone(), )))); @@ -2119,8 +1990,6 @@ impl RdpServer { /// on, a preemption winner is indistinguishable from a normally-accepted /// connection. async fn serve_negotiated(&mut self, candidate: Box) -> ServerResult<()> { - self.display_suppressed.store(false, Ordering::Relaxed); - let mut candidate = candidate; // Only NOW build the channel backends: this connection has // authenticated and is about to be served, so the factories run @@ -2229,16 +2098,8 @@ impl RdpServer { S: AsyncRead + AsyncWrite + Send + Sync + Unpin, { let result = self.run_connection_inner(stream, tls).await; - - // The static channels belong to the connection that negotiated them, - // and their backends own real resources: an rdpsnd handler is stopped - // through `Drop`, so an audio backend keeps capturing until the set is - // replaced. `run` cleared the set itself, which left embedders driving - // connections through this method with the previous session's backends - // still live until the next client attached new ones. - self.static_channels = StaticChannelSet::new(); + // Given for this connection only; see `set_connection_local_addr`. self.connection_local_addr = None; - result } @@ -2246,17 +2107,6 @@ impl RdpServer { where S: AsyncRead + AsyncWrite + Send + Sync + Unpin, { - // Per-connection state must start fresh: if the previous client - // disconnected while it had sent `SuppressOutput { None }` (e.g., - // closed the mstsc window while minimized so the matching resume - // PDU never arrived), the flag would still read `true` here and the - // display backend would silently drop frames for the entire new - // session until/unless the new client happens to send a - // `RefreshRectangle` or `SuppressOutput { Some(rect) }`. Resetting - // here also covers backends that share an externally-created Arc via - // `set_display_suppressed_handle()`. - self.display_suppressed.store(false, Ordering::Relaxed); - let size = self.display.lock().await.size().await; let monitor_count = self.display.lock().await.monitor_count().await; let capabilities = capabilities::capabilities(&self.opts, size); @@ -2299,17 +2149,28 @@ impl RdpServer { if local_addr.is_some() { self.connection_local_addr = local_addr; } + // Both entry paths, `run_connection_inner` and the preemption winner's + // `serve_negotiated`, reach this point, so dropping the state when this + // function returns releases it on every way out of a connection. In + // particular the static channel backends own real resources (an rdpsnd + // handler is stopped through `Drop`), which must not stay live until + // the next client attaches new ones. + let mut conn = ConnectionState { + autodetect: self.autodetect_enabled.then(AutoDetectManager::new), + ..ConnectionState::default() + }; match transport { // No security upgrade happened, so there is no TLS session to shut // down — matches the pre-existing `BeginResult::Continue` arm. NegotiatedTransport::Continued(framed) => { - self.accept_finalize(framed, acceptor).await?; + self.accept_finalize(&mut conn, framed, acceptor).await?; } NegotiatedTransport::Tls(framed) => { - self.finalize_and_shutdown(*framed, acceptor, "TLS connection").await?; + self.finalize_and_shutdown(&mut conn, *framed, acceptor, "TLS connection") + .await?; } NegotiatedTransport::Offloaded(framed) => { - self.finalize_and_shutdown(framed, acceptor, "TLS-offloaded stream") + self.finalize_and_shutdown(&mut conn, framed, acceptor, "TLS-offloaded stream") .await?; } } @@ -2323,6 +2184,7 @@ impl RdpServer { /// shares. async fn finalize_and_shutdown( &mut self, + conn: &mut ConnectionState, framed: TokioFramed, acceptor: Acceptor, shutdown_label: &str, @@ -2335,7 +2197,7 @@ impl RdpServer { // `PendingConnection::negotiate_and_authenticate` before this function // ever runs -- upstream's un-refactored equivalent still does that // work at this point, since it has no separate negotiation step. - let framed = self.accept_finalize(framed, acceptor).await?; + let framed = self.accept_finalize(conn, framed, acceptor).await?; debug!("Shutting down {}", shutdown_label); let (mut inner, _) = framed.into_inner(); if let Err(e) = inner.shutdown().await { @@ -2729,16 +2591,8 @@ impl RdpServer { error!(?error, "Connection error"); } - // NOT redundant with `run_connection_with`'s own reset (added - // upstream, #1721) despite resetting the same field: a preemption - // winner reaches this point via `serve_negotiated`, which never - // calls `run_connection`/`run_connection_with` at all -- so this - // is the only reset that path gets. Removing this because - // `run_connection_with` "already handles it" would silently - // reintroduce #1721's leak (channel backends, e.g. rdpsnd's audio - // capture, held open until the next client) for every preemption - // takeover. - self.static_channels = StaticChannelSet::new(); + // Neither path above goes through `run_connection_with`, which + // clears this for embedders. self.connection_local_addr = None; if let Some(ref mut handler) = self.connection_handler { @@ -2753,22 +2607,13 @@ impl RdpServer { Ok(()) } - pub fn get_svc_processor(&mut self) -> Option<&mut T> { - self.static_channels - .get_by_type_mut::() - .and_then(|svc| svc.channel_processor_downcast_mut()) - } - - pub fn get_channel_id_by_type(&self) -> Option { - self.static_channels.get_channel_id_by_type::() - } - #[expect( clippy::too_many_arguments, reason = "private per-connection dispatch; the parameters are the connection's negotiated identifiers and transports" )] async fn dispatch_pdu( &mut self, + conn: &mut ConnectionState, action: Action, bytes: bytes::BytesMut, writer: &mut impl FramedWrite, @@ -2786,6 +2631,7 @@ impl RdpServer { Action::X224 => { if self .handle_x224( + conn, writer, io_channel_id, user_channel_id, @@ -2841,8 +2687,13 @@ impl RdpServer { Ok((RunState::Continue, encoder)) } + #[expect( + clippy::too_many_arguments, + reason = "private per-connection dispatch; the parameters are the connection's negotiated identifiers and transports" + )] async fn dispatch_server_events( &mut self, + conn: &mut ConnectionState, events: &mut Vec, writer: &mut impl FramedWrite, io_channel_id: u16, @@ -2939,7 +2790,7 @@ impl RdpServer { .await?; } ServerEvent::Rdpsnd(s) => { - let Some(rdpsnd) = self.get_svc_processor::() else { + let Some(rdpsnd) = conn.get_svc_processor::() else { warn!("No rdpsnd channel, dropping event"); continue; }; @@ -2980,7 +2831,7 @@ impl RdpServer { } } .map_err_kind("failed to send rdpsnd event", ServerErrorKind::Pdu)?; - let channel_id = self + let channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("SVC channel not found"))?; let data = server_encode_svc_messages(msgs.into(), channel_id, user_channel_id) @@ -2991,7 +2842,7 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } ServerEvent::Rdpdr(msg) => { - let Some(rdpdr) = self.get_svc_processor::() else { + let Some(rdpdr) = conn.get_svc_processor::() else { warn!("No rdpdr channel, dropping event"); continue; }; @@ -3080,7 +2931,7 @@ impl RdpServer { ), } .map_err_kind("failed to send rdpdr event", ServerErrorKind::Pdu)?; - let channel_id = self + let channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("SVC channel not found"))?; let data = @@ -3091,7 +2942,7 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } ServerEvent::Rdpeai(msg) => { - let Some(drdynvc) = self.get_svc_processor::() else { + let Some(drdynvc) = conn.get_svc_processor::() else { warn!("No drdynvc channel, dropping AUDIO_INPUT event"); continue; }; @@ -3136,7 +2987,7 @@ impl RdpServer { }; let dvc_messages = dvc::encode_dvc_messages(channel_id, msgs, ChannelFlags::SHOW_PROTOCOL) .map_err(ServerError::encode)?; - let drdynvc_channel_id = self + let drdynvc_channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; let data = server_encode_svc_messages(dvc_messages, drdynvc_channel_id, user_channel_id) @@ -3147,7 +2998,7 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } ServerEvent::Clipboard(c) => { - let Some(cliprdr) = self.get_svc_processor::() else { + let Some(cliprdr) = conn.get_svc_processor::() else { warn!("No clipboard channel, dropping event"); continue; }; @@ -3179,7 +3030,7 @@ impl RdpServer { } }; - let channel_id = self + let channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("SVC channel not found"))?; let data = server_encode_svc_messages(msgs.into(), channel_id, user_channel_id) @@ -3191,7 +3042,7 @@ impl RdpServer { } ServerEvent::Echo(msg) => match msg { EchoServerMessage::SendRequest { payload } => { - let Some(drdynvc) = self.get_svc_processor::() else { + let Some(drdynvc) = conn.get_svc_processor::() else { warn!("No drdynvc channel, dropping ECHO request"); continue; }; @@ -3213,7 +3064,7 @@ impl RdpServer { dvc::encode_dvc_messages(echo_channel_id, vec![request], ChannelFlags::SHOW_PROTOCOL) .map_err(ServerError::encode)?; - let drdynvc_channel_id = self + let drdynvc_channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; @@ -3231,11 +3082,11 @@ impl RdpServer { let create_dvc_msg = { use crate::urbdrc::UsbRedirServer; - let Some(usb_man) = self.usb_man.as_mut() else { + let Some(usb_factory) = self.usb_factory.as_mut() else { warn!("Missing USB device factory"); continue; }; - let Some(drdynvc) = self + let Some(drdynvc) = conn .static_channels .get_by_type_mut::() .and_then(|svc| svc.channel_processor_downcast_mut::()) @@ -3244,12 +3095,12 @@ impl RdpServer { continue; }; - let Some(comp_iface) = usb_man.comp_iface_alloc.alloc() else { + let Some(comp_iface) = conn.usb_man.comp_iface_alloc.alloc() else { warn!("Run out of URBDRC interface IDs"); continue; }; - let Some(device_backend) = usb_man.factory.create_device() else { + let Some(device_backend) = usb_factory.create_device() else { warn!("Failed to create USB device backend"); continue; }; @@ -3257,7 +3108,7 @@ impl RdpServer { drdynvc .create_channel_with(|dvc_id| { let handle = UsbDeviceHandle::new(self.ev_sender.clone(), dvc_id); - if usb_man.router.insert(dvc_id, handle.device()).is_some() { + if conn.usb_man.router.insert(dvc_id, handle.device()).is_some() { warn!(dvc_id = dvc_id, "Replacing USB device pending-request map"); } Ok::<_, PduError>( @@ -3271,7 +3122,7 @@ impl RdpServer { .map_err_kind("create URBDRC device channel", ServerErrorKind::Pdu)? }; - let drdynvc_channel_id = self + let drdynvc_channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; let data = @@ -3284,12 +3135,7 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } UrbdrcServerMessage::Device { dvc_id, dev_msg } => { - let Some(device) = self - .usb_man - .as_ref() - .and_then(|usb_man| usb_man.router.get(&dvc_id)) - .map(Arc::clone) - else { + let Some(device) = conn.usb_man.router.get(&dvc_id).map(Arc::clone) else { warn!(dvc_id, "Missing USB device state"); continue; }; @@ -3302,7 +3148,7 @@ impl RdpServer { continue; } - let Some(drdynvc) = self.get_svc_processor::() else { + let Some(drdynvc) = conn.get_svc_processor::() else { warn!("No drdynvc channel, dropping URBDRC request"); continue; }; @@ -3376,17 +3222,17 @@ impl RdpServer { .map_err(ServerError::encode)?; if close_dev { - let close_message = self + let close_message = conn .get_svc_processor::() .and_then(|drdynvc| drdynvc.close_channel(dvc_id)) .ok_or_else(|| { ServerError::channel("URBDRC dynamic channel disappeared before close") })?; - self.remove_usb_device(dvc_id); + conn.remove_usb_device(dvc_id); messages.push(close_message); } - let drdynvc_channel_id = self + let drdynvc_channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; @@ -3398,13 +3244,13 @@ impl RdpServer { .map_err(|e| ServerError::io("write_all", e))?; } UrbdrcServerMessage::DeviceClosed { dvc_id } => { - self.remove_usb_device(dvc_id); + conn.remove_usb_device(dvc_id); } }, #[cfg(feature = "egfx")] ServerEvent::Egfx(msg) => match msg { EgfxServerMessage::SendMessages { messages } => { - self.dispatch_egfx_messages(messages, writer, user_channel_id, udp_transport) + self.dispatch_egfx_messages(conn, messages, writer, user_channel_id, udp_transport) .await?; } }, @@ -3412,7 +3258,7 @@ impl RdpServer { // Auto-detect requests ride the MCS message channel // ([MS-RDPBCGR] 2.2.14.3). With none negotiated (the client // did not request it), there is nowhere to send them. - if let (Some(ad), Some(message_channel_id)) = (self.autodetect.as_mut(), message_channel_id) { + if let (Some(ad), Some(message_channel_id)) = (conn.autodetect.as_mut(), message_channel_id) { let now_ms = monotonic_now_ms(); ad.expire_stale_probes(now_ms, crate::autodetect::RTT_PROBE_MAX_AGE_MS); let request = ad.send_rtt_request(now_ms); @@ -3472,27 +3318,28 @@ impl RdpServer { #[cfg(feature = "egfx")] async fn dispatch_egfx_messages( &mut self, + conn: &mut ConnectionState, messages: Vec, writer: &mut impl FramedWrite, user_channel_id: u16, udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult<()> { - let drdynvc_channel_id = self + let drdynvc_channel_id = conn .get_channel_id_by_type::() .ok_or_else(|| ServerError::channel("DRDYNVC channel not found"))?; let mut route_over_udp = false; - // `self.udp_migration_allowed` gates the Soft-Sync Request itself + // `conn.udp_migration_allowed` gates the Soft-Sync Request itself // (`request_reliable_udp` below): MS-RDPEDYC 3.1.5.3/3.3.5.3.1 forbid // it unless both peers negotiated SOFT_SYNC_TCP_TO_UDP and a // successful Initiate Multitransport Response was actually received, // neither of which the sideband transport's own handshake succeeding // (`udp_transport.is_some()`) establishes on its own. let mut newly_on_udp = None; - if self.udp_migration_allowed + if conn.udp_migration_allowed && let Some(udp_transport) = udp_transport - && let Some(drdynvc) = self.get_svc_processor::() + && let Some(drdynvc) = conn.get_svc_processor::() { let Some(egfx_dvc_id) = crate::gfx::egfx_channel_id(drdynvc) else { trace!("EGFX channel not open yet, staying on TCP"); @@ -3531,7 +3378,7 @@ impl RdpServer { drdynvc.tunnel_for_outgoing_channel(egfx_dvc_id) == Some(dvc::pdu::SoftSyncTunnelType::RELIABLE_UDP); if route_over_udp { - if !self.egfx_on_udp { + if !conn.egfx_on_udp { newly_on_udp = Some(egfx_dvc_id); } for message in &messages { @@ -3541,7 +3388,7 @@ impl RdpServer { } } if let Some(egfx_dvc_id) = newly_on_udp { - self.egfx_on_udp = true; + conn.egfx_on_udp = true; debug!(egfx_dvc_id, "EGFX is now sent over the UDP transport"); } @@ -3577,6 +3424,7 @@ impl RdpServer { /// Without a live tunnel everything goes over TCP. async fn write_drdynvc_output( &mut self, + conn: &mut ConnectionState, messages: Vec, writer: &mut impl FramedWrite, drdynvc_channel_id: u16, @@ -3584,7 +3432,7 @@ impl RdpServer { udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult<()> { let mut over_tcp = Vec::with_capacity(messages.len()); - match (udp_transport, self.get_svc_processor::()) { + match (udp_transport, conn.get_svc_processor::()) { (Some(udp_transport), Some(drdynvc)) => { for message in messages { let payload = message.encode_unframed_pdu().map_err(ServerError::encode)?; @@ -3622,23 +3470,24 @@ impl RdpServer { /// waiting on has been processed. async fn replay_early_tunnel_payloads( &mut self, + conn: &mut ConnectionState, writer: &mut impl FramedWrite, user_channel_id: u16, udp_transport: Option<&multitransport::UdpTransportHandle>, ) -> ServerResult<()> { - if self.early_tunnel_payloads.is_empty() - || !self + if conn.early_tunnel_payloads.is_empty() + || !conn .get_svc_processor::() .is_some_and(|drdynvc| drdynvc.soft_sync_response_received()) { return Ok(()); } debug!( - count = self.early_tunnel_payloads.len(), + count = conn.early_tunnel_payloads.len(), "Soft-Sync response received, processing the tunnel payloads held for it" ); - while let Some(payload) = self.early_tunnel_payloads.pop_front() { - self.dispatch_udp_tunnel_payload(&payload, writer, user_channel_id, udp_transport) + while let Some(payload) = conn.early_tunnel_payloads.pop_front() { + self.dispatch_udp_tunnel_payload(conn, &payload, writer, user_channel_id, udp_transport) .await?; } Ok(()) @@ -3656,6 +3505,7 @@ impl RdpServer { /// session depends on the sideband transport's correctness. async fn dispatch_udp_tunnel_payload( &mut self, + conn: &mut ConnectionState, payload: &[u8], writer: &mut impl FramedWrite, user_channel_id: u16, @@ -3667,8 +3517,8 @@ impl RdpServer { // is inconsistent, not a case this function can meaningfully // distinguish from "no drdynvc channel". let (Some(drdynvc_channel_id), Some(drdynvc)) = ( - self.get_channel_id_by_type::(), - self.get_svc_processor::(), + conn.get_channel_id_by_type::(), + conn.get_svc_processor::(), ) else { warn!("No drdynvc channel, dropping UDP tunnel payload"); return Ok(RunState::Continue); @@ -3679,12 +3529,12 @@ impl RdpServer { // it; `replay_early_tunnel_payloads` feeds it through once the // response is in. if drdynvc.soft_sync_awaiting_response() { - if self.early_tunnel_payloads.len() < MAX_EARLY_TUNNEL_PAYLOADS { + if conn.early_tunnel_payloads.len() < MAX_EARLY_TUNNEL_PAYLOADS { trace!( len = payload.len(), "Holding a tunnel payload until the Soft-Sync response arrives" ); - self.early_tunnel_payloads.push_back(payload.to_vec()); + conn.early_tunnel_payloads.push_back(payload.to_vec()); } else { warn!("Too many tunnel payloads ahead of the Soft-Sync response, dropping one"); } @@ -3703,8 +3553,15 @@ impl RdpServer { return Ok(RunState::Continue); } - self.write_drdynvc_output(messages, writer, drdynvc_channel_id, user_channel_id, udp_transport) - .await?; + self.write_drdynvc_output( + conn, + messages, + writer, + drdynvc_channel_id, + user_channel_id, + udp_transport, + ) + .await?; Ok(RunState::Continue) } @@ -3715,12 +3572,12 @@ impl RdpServer { )] async fn client_loop( &mut self, + conn: &mut ConnectionState, reader: &mut Framed, writer: &mut Framed, io_channel_id: u16, user_channel_id: u16, message_channel_id: Option, - client_supports_heartbeat: bool, mut encoder: UpdateEncoder, udp_transport: Rc>>, pending_udp_accept: Option>>, @@ -3730,12 +3587,12 @@ impl RdpServer { W: FramedWrite, { debug!("Starting client loop"); - let heartbeat = if client_supports_heartbeat { + let heartbeat = if conn.client_supports_heartbeat { self.heartbeat } else { None }; - let mut display_updates = self.display.lock().await.updates().await?; + let mut display_updates = self.display.lock().await.updates(conn.display_context()).await?; let mut writer = SharedWriter::new(writer); let mut display_writer = writer.clone(); let mut event_writer = writer.clone(); @@ -3746,7 +3603,7 @@ impl RdpServer { let udp_transport_for_pdus = Rc::clone(&udp_transport); let write_counter = writer.write_counter(); let ev_receiver = Arc::clone(&self.ev_receiver); - let s = Rc::new(Mutex::new(self)); + let s = Rc::new(Mutex::new((self, conn))); let this = Rc::clone(&s); let dispatch_pdu = async move { @@ -3765,8 +3622,10 @@ impl RdpServer { let dispatch_start = Instant::now(); let current_udp_transport = udp_transport_for_pdus.borrow().clone(); - let result = this + let (server, conn) = &mut *this; + let result = server .dispatch_pdu( + conn, action, bytes, &mut writer, @@ -3868,8 +3727,10 @@ impl RdpServer { // an await point. `UdpTransportHandle` is cheap to clone (see // its own doc comment). let current_udp_transport = udp_transport_for_events.borrow().clone(); - let result = this + let (server, conn) = &mut *this; + let result = server .dispatch_server_events( + conn, &mut events, &mut event_writer, io_channel_id, @@ -3906,7 +3767,9 @@ impl RdpServer { loop { interval.tick().await; let mut this = this.lock().await; - this.rotate_auto_reconnect_cookie(&mut auto_reconnect_writer, io_channel_id, user_channel_id) + let (server, _) = &mut *this; + server + .rotate_auto_reconnect_cookie(&mut auto_reconnect_writer, io_channel_id, user_channel_id) .await?; } }; @@ -4005,7 +3868,8 @@ impl RdpServer { // (MS-RDPEMT 1.3.3). A client whose channels moved has // nowhere left to read them, so end the connection and let // it reconnect rather than keep a session it cannot draw. - if this.lock().await.egfx_on_udp { + let (_, conn) = &mut *this.lock().await; + if conn.egfx_on_udp { warn!("UDP transport lost with EGFX on it, ending the connection"); return Err(ServerError::reason( "UDP transport", @@ -4021,8 +3885,15 @@ impl RdpServer { return core::future::pending::>().await; }; let mut this = this.lock().await; - let result = this - .dispatch_udp_tunnel_payload(&payload, &mut udp_tunnel_writer, user_channel_id, Some(&transport)) + let (server, conn) = &mut *this; + let result = server + .dispatch_udp_tunnel_payload( + conn, + &payload, + &mut udp_tunnel_writer, + user_channel_id, + Some(&transport), + ) .await?; match result { RunState::Continue => continue, @@ -4046,6 +3917,7 @@ impl RdpServer { async fn client_accepted( &mut self, + conn: &mut ConnectionState, reader: &mut Framed, writer: &mut Framed, result: AcceptorResult, @@ -4112,6 +3984,7 @@ impl RdpServer { // Set on a reactivation pass, where EGFX may already be on the tunnel. let current_udp_transport = udp_transport.borrow().clone(); self.handle_input_backlog( + conn, writer, result.io_channel_id, result.user_channel_id, @@ -4122,9 +3995,12 @@ impl RdpServer { .await?; } - self.static_channels = result.static_channels; + conn.static_channels = result.static_channels; + conn.client_supports_heartbeat = result + .client_early_capability_flags + .contains(ironrdp_pdu::gcc::ClientEarlyCapabilityFlags::SUPPORT_HEART_BEAT_PDU); if !result.reactivation { - for (_channel_key, channel, channel_id) in self.static_channels.iter_by_key_mut() { + for (_channel_key, channel, channel_id) in conn.static_channels.iter_by_key_mut() { debug!(?channel, ?channel_id, "Start"); let Some(channel_id) = channel_id else { continue; @@ -4266,28 +4142,26 @@ impl RdpServer { let pending_udp_accept = Self::drop_declined_udp_accept(pending_udp_accept, result.multitransport_response_success); - self.pending_udp_accept_abort = pending_udp_accept.as_ref().map(task::JoinHandle::abort_handle); + conn.pending_udp_accept_abort = pending_udp_accept.as_ref().map(task::JoinHandle::abort_handle); // See `udp_migration_allowed`: a successful response that arrives // after this point enables migration from the message-channel handler. // Only ever raised here, never lowered: a deactivation-reactivation // pass lands here again without repeating the multitransport exchange, // and must not take back a migration the session may already be using. - self.soft_sync_negotiated |= result + conn.soft_sync_negotiated |= result .multitransport_flags .contains(ironrdp_pdu::gcc::MultiTransportFlags::SOFT_SYNC_TCP_TO_UDP); - self.udp_migration_allowed |= self.soft_sync_negotiated && result.multitransport_response_success == Some(true); + conn.udp_migration_allowed |= conn.soft_sync_negotiated && result.multitransport_response_success == Some(true); let state = self .client_loop( + conn, reader, writer, result.io_channel_id, result.user_channel_id, result.message_channel_id, - result - .client_early_capability_flags - .contains(ironrdp_pdu::gcc::ClientEarlyCapabilityFlags::SUPPORT_HEART_BEAT_PDU), encoder, udp_transport, pending_udp_accept, @@ -4317,8 +4191,13 @@ impl RdpServer { } } + #[expect( + clippy::too_many_arguments, + reason = "private per-connection dispatch; the parameters are the connection's negotiated identifiers and transports" + )] async fn handle_input_backlog( &mut self, + conn: &mut ConnectionState, writer: &mut impl FramedWrite, io_channel_id: u16, user_channel_id: u16, @@ -4336,6 +4215,7 @@ impl RdpServer { Ok(Action::X224) => { let _ = self .handle_x224( + conn, writer, io_channel_id, user_channel_id, @@ -4390,7 +4270,11 @@ impl RdpServer { } } - async fn handle_io_channel_data(&mut self, data: SendDataRequest<'_>) -> ServerResult { + async fn handle_io_channel_data( + &mut self, + conn: &mut ConnectionState, + data: SendDataRequest<'_>, + ) -> ServerResult { let control: rdp::headers::ShareControlHeader = decode(data.user_data.as_ref()).map_err(ServerError::decode)?; match control.share_control_pdu { @@ -4415,7 +4299,7 @@ impl RdpServer { // set. rdp::headers::ShareDataPdu::SuppressOutput(pdu) => { let suppress = pdu.desktop_rect.is_none(); - self.display_suppressed.store(suppress, Ordering::Relaxed); + conn.display_suppressed.store(suppress, Ordering::Relaxed); debug!(suppress, "client suppress-output state changed"); } @@ -4427,7 +4311,7 @@ impl RdpServer { // this; clearing here is belt-and-braces against clients // that send only one of the two.) rdp::headers::ShareDataPdu::RefreshRectangle(_) => { - if self.display_suppressed.swap(false, Ordering::Relaxed) { + if conn.display_suppressed.swap(false, Ordering::Relaxed) { debug!("client RefreshRectangle cleared suppress-output state"); } } @@ -4445,20 +4329,22 @@ impl RdpServer { Ok(false) } - fn handle_message_channel_data(&mut self, data: SendDataRequest<'_>) { + fn handle_message_channel_data(conn: &mut ConnectionState, data: SendDataRequest<'_>) { match decode::(data.user_data.as_ref()) { Ok(rdp::message_channel::ClientMessageChannelPdu::AutoDetectResponse(pdu)) => { - if let Some(ref mut ad) = self.autodetect { + if let Some(ref mut ad) = conn.autodetect { match ad.handle_response(&pdu.response, monotonic_now_ms()) { AutoDetectOutcome::Rtt(rtt_ms) => { - self.autodetect_rtt.store(rtt_ms, Ordering::Relaxed); + conn.autodetect_handles.rtt.store(rtt_ms, Ordering::Relaxed); // A matched RTT sample always updates the session-lifetime low in the // same call (see `handle_response`'s RttResponse arm), so it is available // unconditionally here, not just on a new low. let baseline_rtt_ms = ad .baseline_rtt_ms() .expect("handle_response just recorded a sample above"); - self.autodetect_baseline_rtt.store(baseline_rtt_ms, Ordering::Relaxed); + conn.autodetect_handles + .baseline_rtt + .store(baseline_rtt_ms, Ordering::Relaxed); debug!( rtt_ms, baseline_rtt_ms, @@ -4467,8 +4353,12 @@ impl RdpServer { ); } AutoDetectOutcome::Bandwidth(Some(bandwidth_kbps)) => { - self.autodetect_bandwidth.store(bandwidth_kbps, Ordering::Relaxed); - self.autodetect_bandwidth_generation.fetch_add(1, Ordering::Release); + conn.autodetect_handles + .bandwidth + .store(bandwidth_kbps, Ordering::Relaxed); + conn.autodetect_handles + .bandwidth_generation + .fetch_add(1, Ordering::Release); // Logging the whole response, not just the computed figure: a // damage-driven video source makes any single measurement // window's byte count wildly bimodal (near-idle vs. a real @@ -4481,8 +4371,10 @@ impl RdpServer { // The manager just cleared its own figure rather than keep // reporting a stale one (see `handle_response`'s doc comment); // mirror that here so the exposed handle does not disagree. - self.autodetect_bandwidth.store(u32::MAX, Ordering::Relaxed); - self.autodetect_bandwidth_generation.fetch_add(1, Ordering::Release); + conn.autodetect_handles.bandwidth.store(u32::MAX, Ordering::Relaxed); + conn.autodetect_handles + .bandwidth_generation + .fetch_add(1, Ordering::Release); trace!( seq = pdu.response.sequence_number(), "Bandwidth measurement completed without a usable figure" @@ -4507,7 +4399,7 @@ impl RdpServer { // E_ABORT: the client could not establish the multitransport // connection (MS-RDPBCGR 2.2.15.2), so no UDP handshake is coming. if !pdu.is_success() - && let Some(abort) = self.pending_udp_accept_abort.take() + && let Some(abort) = conn.pending_udp_accept_abort.take() { abort.abort(); debug!("Client could not establish the UDP multitransport connection, continuing TCP-only"); @@ -4515,8 +4407,8 @@ impl RdpServer { // A success after finalization, the usual order with mstsc, // completes the condition finalization could not see (see // `udp_migration_allowed`), so EGFX can migrate from here on. - if pdu.is_success() && self.soft_sync_negotiated && !self.udp_migration_allowed { - self.udp_migration_allowed = true; + if pdu.is_success() && conn.soft_sync_negotiated && !conn.udp_migration_allowed { + conn.udp_migration_allowed = true; debug!("Multitransport confirmed after finalization, EGFX may migrate to UDP"); } } @@ -4526,8 +4418,13 @@ impl RdpServer { } } + #[expect( + clippy::too_many_arguments, + reason = "private per-connection dispatch; the parameters are the connection's negotiated identifiers and transports" + )] async fn handle_x224( &mut self, + conn: &mut ConnectionState, writer: &mut impl FramedWrite, io_channel_id: u16, user_channel_id: u16, @@ -4545,20 +4442,21 @@ impl RdpServer { "McsMessage::SendDataRequest" ); if data.channel_id == io_channel_id { - return self.handle_io_channel_data(data).await; + return self.handle_io_channel_data(conn, data).await; } if message_channel_id == Some(data.channel_id) { - self.handle_message_channel_data(data); + Self::handle_message_channel_data(conn, data); return Ok(false); } - if let Some(svc) = self.static_channels.get_by_channel_id_mut(data.channel_id) { + if let Some(svc) = conn.static_channels.get_by_channel_id_mut(data.channel_id) { let response_pdus = svc .process(&data.user_data) .map_err_kind("svc process", ServerErrorKind::Pdu)?; - if self.get_channel_id_by_type::() == Some(data.channel_id) { + if conn.get_channel_id_by_type::() == Some(data.channel_id) { self.write_drdynvc_output( + conn, response_pdus, writer, data.channel_id, @@ -4566,7 +4464,7 @@ impl RdpServer { udp_transport, ) .await?; - self.replay_early_tunnel_payloads(writer, user_channel_id, udp_transport) + self.replay_early_tunnel_payloads(conn, writer, user_channel_id, udp_transport) .await?; } else { let response = server_encode_svc_messages(response_pdus, data.channel_id, user_channel_id) @@ -4630,19 +4528,13 @@ impl RdpServer { async fn accept_finalize( &mut self, + conn: &mut ConnectionState, mut framed: TokioFramed, mut acceptor: Acceptor, ) -> ServerResult> where S: AsyncRead + AsyncWrite + Sync + Send + Unpin, { - // Per-connection: set again once this connection's finalization knows - // its own negotiated flags and response. - self.soft_sync_negotiated = false; - self.udp_migration_allowed = false; - self.egfx_on_udp = false; - self.early_tunnel_payloads.clear(); - let udp_bind_addr = self .opts .udp_bind_addr @@ -4718,6 +4610,7 @@ impl RdpServer { match self .client_accepted( + conn, &mut reader, &mut writer, result, @@ -4736,7 +4629,7 @@ impl RdpServer { // various state issues during client resize. acceptor = Acceptor::new_deactivation_reactivation( acceptor, - core::mem::take(&mut self.static_channels), + core::mem::take(&mut conn.static_channels), desktop_size, ) .map_err_kind("deactivation-reactivation acceptor", ServerErrorKind::Connector)?; @@ -5008,7 +4901,7 @@ mod preempt_tests { height: 768, } } - async fn updates(&mut self) -> ServerResult> { + async fn updates(&mut self, _: DisplayContext) -> ServerResult> { unreachable!("negotiation never asks for updates") } } @@ -5253,12 +5146,7 @@ mod preempt_tests { fn a_late_multitransport_success_enables_migration_only_with_soft_sync() { use ironrdp_pdu::rdp::multitransport::MultitransportResponsePdu; - let mut server = RdpServer::builder() - .with_addr((Ipv4Addr::LOCALHOST, 0)) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); + let mut conn = ConnectionState::default(); let response = |pdu: &MultitransportResponsePdu| SendDataRequest { initiator_id: 1007, channel_id: 1008, @@ -5266,17 +5154,17 @@ mod preempt_tests { }; // Soft-Sync not negotiated: MS-RDPEDYC forbids migration whatever the response. - server.handle_message_channel_data(response(&MultitransportResponsePdu::success(1))); - assert!(!server.udp_migration_allowed); + RdpServer::handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::success(1))); + assert!(!conn.udp_migration_allowed); // Negotiated, but the client could not bring UDP up. - server.soft_sync_negotiated = true; - server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); - assert!(!server.udp_migration_allowed); + conn.soft_sync_negotiated = true; + RdpServer::handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::abort(1))); + assert!(!conn.udp_migration_allowed); // Negotiated, and the success arrives after finalization. - server.handle_message_channel_data(response(&MultitransportResponsePdu::success(1))); - assert!(server.udp_migration_allowed); + RdpServer::handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::success(1))); + assert!(conn.udp_migration_allowed); } #[tokio::test] @@ -5286,12 +5174,7 @@ mod preempt_tests { let local = task::LocalSet::new(); local .run_until(async { - let mut server = RdpServer::builder() - .with_addr((Ipv4Addr::LOCALHOST, 0)) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); + let mut conn = ConnectionState::default(); let response = |pdu: &MultitransportResponsePdu| SendDataRequest { initiator_id: 1007, channel_id: 1008, @@ -5299,21 +5182,21 @@ mod preempt_tests { }; // No accept pending: a failure response is only logged. - server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); + RdpServer::handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::abort(1))); let accept = task::spawn_local(core::future::pending::<()>()); - server.pending_udp_accept_abort = Some(accept.abort_handle()); + conn.pending_udp_accept_abort = Some(accept.abort_handle()); // Success leaves the accept running. - server.handle_message_channel_data(response(&MultitransportResponsePdu::success(1))); + RdpServer::handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::success(1))); task::yield_now().await; assert!(!accept.is_finished()); - server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); + RdpServer::handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::abort(1))); task::yield_now().await; assert!(accept.is_finished()); assert!(accept.await.expect_err("accept was aborted").is_cancelled()); - assert!(server.pending_udp_accept_abort.is_none()); + assert!(conn.pending_udp_accept_abort.is_none()); }) .await; } @@ -5683,64 +5566,6 @@ mod preempt_tests { } } -#[cfg(test)] -mod tests { - use ironrdp_core::impl_as_any; - use ironrdp_pdu::gcc::ChannelName; - use ironrdp_svc::{SvcMessage, SvcServerProcessor}; - - use super::*; - - /// A channel backend that owns a resource, released on drop the way - /// `RdpsndServer` stops its handler. - #[derive(Debug)] - struct ResourceChannel(Arc); - - impl Drop for ResourceChannel { - fn drop(&mut self) { - self.0.store(true, Ordering::Relaxed); - } - } - - impl_as_any!(ResourceChannel); - - impl SvcProcessor for ResourceChannel { - fn channel_name(&self) -> ChannelName { - ChannelName::from_static(b"testchan") - } - - fn process(&mut self, _payload: &[u8]) -> PduResult> { - Ok(Vec::new()) - } - } - - impl SvcServerProcessor for ResourceChannel {} - - #[tokio::test] - async fn run_connection_releases_the_static_channels() { - let mut server = RdpServer::builder() - .with_addr(([127, 0, 0, 1], 0)) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); - - let released = Arc::new(AtomicBool::new(false)); - server.static_channels.insert(ResourceChannel(Arc::clone(&released))); - - // A stream that is already at EOF: the connection ends early, which - // is the path an embedder's accept loop sees when a client vanishes. - let (client, server_side) = tokio::io::duplex(64); - drop(client); - let _ = server.run_connection(server_side).await; - - assert!( - released.load(Ordering::Relaxed), - "the channel backends of a finished connection must be released, not held until the next client" - ); - } -} - /// A failing clipboard event must not disconnect the session. /// /// `dispatch_server_events` used to `?` the result of every `CliprdrServer` @@ -5817,9 +5642,9 @@ mod cliprdr_error_tests { // Left in its initial state, so `require_ready` refuses the request // below -- the cheapest reproduction of "the channel said no". let cliprdr: CliprdrServer = Cliprdr::new(Box::new(SilentBackend)); - server.static_channels.insert(cliprdr); - server - .static_channels + let mut conn = ConnectionState::default(); + conn.static_channels.insert(cliprdr); + conn.static_channels .attach_channel_id(TypeId::of::(), 1004); let mut events = vec![ServerEvent::Clipboard(ClipboardMessage::SendFileContentsRequest( @@ -5835,7 +5660,7 @@ mod cliprdr_error_tests { let mut writer = CapturingWriter::default(); let state = server - .dispatch_server_events(&mut events, &mut writer, 1003, 1002, None, None) + .dispatch_server_events(&mut conn, &mut events, &mut writer, 1003, 1002, None, None) .await .expect("a refused clipboard message must not surface as a session error"); diff --git a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs index 8b1e6a6391..5f8f87943a 100644 --- a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs @@ -411,190 +411,6 @@ fn sequence_number_wraps_at_u16_max() { assert_eq!(req2.sequence_number(), 0, "should wrap around"); } -#[test] -fn autodetect_rtt_handle_defaults_to_sentinel() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::Ordering; - - use ironrdp_server::RdpServer; - - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); - - assert_eq!(server.autodetect_rtt_handle().load(Ordering::Relaxed), u32::MAX); -} - -#[test] -fn with_autodetect_rtt_handle_round_trips_the_same_arc() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::{AtomicU32, Ordering}; - use std::sync::Arc; - - use ironrdp_server::RdpServer; - - let handle = Arc::new(AtomicU32::new(42)); - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .with_autodetect_rtt_handle(Arc::clone(&handle)) - .build(); - - assert!(Arc::ptr_eq(&handle, &server.autodetect_rtt_handle())); - // The server resets an injected handle to the sentinel at construction. - assert_eq!(server.autodetect_rtt_handle().load(Ordering::Relaxed), u32::MAX); - // The Arc is shared: mutating the original is visible through the server's handle. - handle.store(42, Ordering::Relaxed); - assert_eq!(server.autodetect_rtt_handle().load(Ordering::Relaxed), 42); -} - -#[test] -fn autodetect_baseline_rtt_handle_defaults_to_sentinel() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::Ordering; - - use ironrdp_server::RdpServer; - - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); - - assert_eq!( - server.autodetect_baseline_rtt_handle().load(Ordering::Relaxed), - u32::MAX - ); -} - -#[test] -fn with_autodetect_baseline_rtt_handle_round_trips_the_same_arc() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::{AtomicU32, Ordering}; - use std::sync::Arc; - - use ironrdp_server::RdpServer; - - let handle = Arc::new(AtomicU32::new(42)); - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .with_autodetect_baseline_rtt_handle(Arc::clone(&handle)) - .build(); - - assert!(Arc::ptr_eq(&handle, &server.autodetect_baseline_rtt_handle())); - // The server resets an injected handle to the sentinel at construction. - assert_eq!( - server.autodetect_baseline_rtt_handle().load(Ordering::Relaxed), - u32::MAX - ); - // The Arc is shared: mutating the original is visible through the server's handle. - handle.store(42, Ordering::Relaxed); - assert_eq!(server.autodetect_baseline_rtt_handle().load(Ordering::Relaxed), 42); -} - -#[test] -fn autodetect_bandwidth_handle_defaults_to_sentinel() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::Ordering; - - use ironrdp_server::RdpServer; - - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); - - assert_eq!(server.autodetect_bandwidth_handle().load(Ordering::Relaxed), u32::MAX); -} - -#[test] -fn with_autodetect_bandwidth_handle_round_trips_the_same_arc() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::{AtomicU32, Ordering}; - use std::sync::Arc; - - use ironrdp_server::RdpServer; - - let handle = Arc::new(AtomicU32::new(42)); - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .with_autodetect_bandwidth_handle(Arc::clone(&handle)) - .build(); - - assert!(Arc::ptr_eq(&handle, &server.autodetect_bandwidth_handle())); - // The server resets an injected handle to the sentinel at construction. - assert_eq!(server.autodetect_bandwidth_handle().load(Ordering::Relaxed), u32::MAX); - // The Arc is shared: mutating the original is visible through the server's handle. - handle.store(42, Ordering::Relaxed); - assert_eq!(server.autodetect_bandwidth_handle().load(Ordering::Relaxed), 42); -} - -#[test] -fn autodetect_bandwidth_generation_handle_defaults_to_zero() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::Ordering; - - use ironrdp_server::RdpServer; - - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .build(); - - assert_eq!( - server.autodetect_bandwidth_generation_handle().load(Ordering::Acquire), - 0 - ); -} - -#[test] -fn with_autodetect_bandwidth_generation_handle_round_trips_the_same_arc() { - use core::net::{Ipv4Addr, SocketAddr}; - use core::sync::atomic::{AtomicU32, Ordering}; - use std::sync::Arc; - - use ironrdp_server::RdpServer; - - let handle = Arc::new(AtomicU32::new(7)); - let server = RdpServer::builder() - .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) - .with_no_security() - .with_no_input() - .with_no_display() - .with_autodetect_bandwidth_generation_handle(Arc::clone(&handle)) - .build(); - - assert!(Arc::ptr_eq(&handle, &server.autodetect_bandwidth_generation_handle())); - // Unlike the bandwidth value, an injected generation counter is not reset at - // construction: an embedder sharing one counter across servers keeps counting. - assert_eq!( - server.autodetect_bandwidth_generation_handle().load(Ordering::Acquire), - 7 - ); - // The Arc is shared: advancing the original is visible through the server's handle. - handle.fetch_add(1, Ordering::Release); - assert_eq!( - server.autodetect_bandwidth_generation_handle().load(Ordering::Acquire), - 8 - ); -} - #[test] fn stale_probe_expiry() { let mut mgr = AutoDetectManager::new(); diff --git a/crates/ironrdp-testsuite-extra/tests/e2e.rs b/crates/ironrdp-testsuite-extra/tests/e2e.rs index bf0db21a3e..f31aab89d9 100644 --- a/crates/ironrdp-testsuite-extra/tests/e2e.rs +++ b/crates/ironrdp-testsuite-extra/tests/e2e.rs @@ -19,8 +19,9 @@ use ironrdp::pdu::rdp::client_info::CompressionType as PduCompressionType; use ironrdp::pdu::rdp::headers::CompressionFlags; use ironrdp::pdu::{self, gcc}; use ironrdp::server::{ - self, Acceptor, DesktopSize, DisplayUpdate, KeyboardEvent, MouseEvent, PixelFormat, RdpServer, RdpServerDisplay, - RdpServerDisplayUpdates, RdpServerInputHandler, ServerEvent, ServerResult, StaticChannelFactory, TlsIdentityCtx, + self, Acceptor, DesktopSize, DisplayContext, DisplayUpdate, KeyboardEvent, MouseEvent, PixelFormat, RdpServer, + RdpServerDisplay, RdpServerDisplayUpdates, RdpServerInputHandler, ServerEvent, ServerResult, StaticChannelFactory, + TlsIdentityCtx, }; use ironrdp::session::image::DecodedImage; use ironrdp::session::{self, ActiveStage, ActiveStageBuilder, ActiveStageOutput}; @@ -496,7 +497,7 @@ impl RdpServerDisplay for TestDisplay { } } - async fn updates(&mut self) -> ServerResult> { + async fn updates(&mut self, _: DisplayContext) -> ServerResult> { Ok(Box::new(TestDisplayUpdates { rx: Arc::clone(&self.rx), })) diff --git a/crates/ironrdp/examples/server.rs b/crates/ironrdp/examples/server.rs index 7e6c9c9e3f..5054f19e6f 100644 --- a/crates/ironrdp/examples/server.rs +++ b/crates/ironrdp/examples/server.rs @@ -15,9 +15,9 @@ use ironrdp::connector::DesktopSize; use ironrdp::rdpsnd::pdu::{AudioFormat, WaveFormat}; use ironrdp::rdpsnd::server::{NegotiatedFormat, RdpsndError, RdpsndServerHandler, RdpsndServerMessage}; use ironrdp::server::{ - BitmapUpdate, CliprdrServerFactory, Credentials, DisplayUpdate, KeyboardEvent, MouseEvent, PixelFormat, RdpServer, - RdpServerDisplay, RdpServerDisplayUpdates, RdpServerInputHandler, ServerEvent, ServerEventSender, - SoundServerFactory, TlsIdentityCtx, + BitmapUpdate, CliprdrServerFactory, Credentials, DisplayContext, DisplayUpdate, KeyboardEvent, MouseEvent, + PixelFormat, RdpServer, RdpServerDisplay, RdpServerDisplayUpdates, RdpServerInputHandler, ServerEvent, + ServerEventSender, SoundServerFactory, TlsIdentityCtx, }; use ironrdp_cliprdr_native::StubCliprdrBackend; use rand::prelude::*; @@ -210,7 +210,7 @@ impl RdpServerDisplay for Handler { } } - async fn updates(&mut self) -> ironrdp::server::ServerResult> { + async fn updates(&mut self, _: DisplayContext) -> ironrdp::server::ServerResult> { Ok(Box::new(DisplayUpdates {})) } }