From 15e6c989357e25a24c72e9cbc67340972c61fc1a Mon Sep 17 00:00:00 2001 From: uchouT Date: Sun, 27 Sep 2026 13:12:03 +0000 Subject: [PATCH 1/2] refactor(server)!: split per-connection state The connection state now lives in a `ConnectionState` created in `finalize_negotiated`, which both entry paths reach, and passed down to everything that serves the connection. It is dropped when the connection ends. BREAKING CHANGE: `RdpServer::get_svc_processor`, `RdpServer:: get_channel_id_by_type` and `RdpServer::rtt_snapshot` are removed. `RdpServer` can only be borrowed from outside between connections, when the first two found no channels and the last returned the previous connection's figures. Auto-detect results are available through `RdpServer::autodetect_rtt_handle`, `autodetect_baseline_rtt_handle` and `autodetect_bandwidth_handle`. Signed-off-by: uchouT --- crates/ironrdp-server/src/builder.rs | 2 - crates/ironrdp-server/src/server.rs | 547 +++++++++--------- .../tests/server/autodetect.rs | 71 +++ 3 files changed, 338 insertions(+), 282 deletions(-) diff --git a/crates/ironrdp-server/src/builder.rs b/crates/ironrdp-server/src/builder.rs index a511adb8dd..66b624ab60 100644 --- a/crates/ironrdp-server/src/builder.rs +++ b/crates/ironrdp-server/src/builder.rs @@ -424,8 +424,6 @@ impl RdpServerBuilder { 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 diff --git a/crates/ironrdp-server/src/server.rs b/crates/ironrdp-server/src/server.rs index 434cf140df..45e56906a8 100644 --- a/crates/ironrdp-server/src/server.rs +++ b/crates/ironrdp-server/src/server.rs @@ -52,7 +52,7 @@ use tokio::task; use tokio_rustls::TlsAcceptor; use tracing::{debug, error, info, trace, warn}; -use crate::autodetect::{AutoDetectManager, AutoDetectOutcome, RttSnapshot}; +use crate::autodetect::{AutoDetectManager, AutoDetectOutcome}; use crate::clipboard::CliprdrServerFactory; use crate::display::{DisplayUpdate, RdpServerDisplay}; use crate::echo::{EchoDvcBridge, EchoServerHandle, EchoServerMessage, build_echo_request}; @@ -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)] @@ -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 @@ -789,9 +778,10 @@ pub struct RdpServer { /// 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). + /// Increments every time [`Self::autodetect_bandwidth`] is republished: + /// when a Bandwidth Measure transaction completes, whether or not it + /// produced a usable figure (see that field's doc comment on the None + /// case), and when a new connection resets it to `u32::MAX`. /// [`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 @@ -820,34 +810,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 +1142,90 @@ 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, + #[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 { + 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. @@ -1512,7 +1558,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,14 +1570,14 @@ 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, @@ -1558,11 +1603,6 @@ impl RdpServer { 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,29 +1852,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). @@ -1902,8 +1919,6 @@ impl RdpServer { 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 @@ -1929,10 +1944,13 @@ 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`] and read the + /// results through [`Self::autodetect_rtt_handle`], + /// [`Self::autodetect_baseline_rtt_handle`] and + /// [`Self::autodetect_bandwidth_handle`]. 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 +1964,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 +2040,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(), )))); @@ -2229,16 +2239,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 } @@ -2299,17 +2301,39 @@ 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() + }; + // The handles that publish auto-detect measurements to the embedder + // are still owned by the server, so they are reset by hand. Otherwise + // the previous connection's figures, its session-lifetime lowest RTT + // included, would read as this connection's until its first sample. + self.autodetect_rtt.store(u32::MAX, Ordering::Relaxed); + self.autodetect_baseline_rtt.store(u32::MAX, Ordering::Relaxed); + self.autodetect_bandwidth.store(u32::MAX, Ordering::Relaxed); + // An embedder rereads the bandwidth only when the generation moves, so + // the reset has to advance it too, or the previous connection's figure + // would stay cached. Advanced after the store, as for a measurement. + self.autodetect_bandwidth_generation.fetch_add(1, Ordering::Release); 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 +2347,7 @@ impl RdpServer { /// shares. async fn finalize_and_shutdown( &mut self, + conn: &mut ConnectionState, framed: TokioFramed, acceptor: Acceptor, shutdown_label: &str, @@ -2335,7 +2360,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 +2754,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 +2770,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 +2794,7 @@ impl RdpServer { Action::X224 => { if self .handle_x224( + conn, writer, io_channel_id, user_channel_id, @@ -2841,8 +2850,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 +2953,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 +2994,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 +3005,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 +3094,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 +3105,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 +3150,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 +3161,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 +3193,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 +3205,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 +3227,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 +3245,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 +3258,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 +3271,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 +3285,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 +3298,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 +3311,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 +3385,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 +3407,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 +3421,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 +3481,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 +3541,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 +3551,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 +3587,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 +3595,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 +3633,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 +3668,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 +3680,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 +3692,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 +3716,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 +3735,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,7 +3750,7 @@ 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 @@ -3746,7 +3766,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 +3785,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 +3890,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 +3930,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 +4031,7 @@ 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 { + if this.lock().await.1.egfx_on_udp { warn!("UDP transport lost with EGFX on it, ending the connection"); return Err(ServerError::reason( "UDP transport", @@ -4021,8 +4047,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 +4079,7 @@ impl RdpServer { async fn client_accepted( &mut self, + conn: &mut ConnectionState, reader: &mut Framed, writer: &mut Framed, result: AcceptorResult, @@ -4112,6 +4146,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 +4157,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 +4304,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 +4353,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 +4377,7 @@ impl RdpServer { Ok(Action::X224) => { let _ = self .handle_x224( + conn, writer, io_channel_id, user_channel_id, @@ -4445,10 +4487,10 @@ impl RdpServer { Ok(false) } - fn handle_message_channel_data(&mut self, data: SendDataRequest<'_>) { + fn handle_message_channel_data(&mut self, 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); @@ -4507,7 +4549,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 +4557,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 +4568,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, @@ -4549,16 +4596,17 @@ impl RdpServer { } 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 +4614,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 +4678,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 +4760,7 @@ impl RdpServer { match self .client_accepted( + conn, &mut reader, &mut writer, result, @@ -4736,7 +4779,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)?; @@ -5259,6 +5302,7 @@ mod preempt_tests { .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 +5310,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); + server.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; + server.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); + server.handle_message_channel_data(&mut conn, response(&MultitransportResponsePdu::success(1))); + assert!(conn.udp_migration_allowed); } #[tokio::test] @@ -5292,6 +5336,7 @@ mod preempt_tests { .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 +5344,21 @@ mod preempt_tests { }; // No accept pending: a failure response is only logged. - server.handle_message_channel_data(response(&MultitransportResponsePdu::abort(1))); + server.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))); + server.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))); + server.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 +5728,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 +5804,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 +5822,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..e635482bd2 100644 --- a/crates/ironrdp-testsuite-core/tests/server/autodetect.rs +++ b/crates/ironrdp-testsuite-core/tests/server/autodetect.rs @@ -543,6 +543,77 @@ fn with_autodetect_bandwidth_handle_round_trips_the_same_arc() { assert_eq!(server.autodetect_bandwidth_handle().load(Ordering::Relaxed), 42); } +/// The handles report one connection's measurements. A new connection starts +/// back at the sentinel rather than inheriting the figures of the connection +/// before it, the session-lifetime lowest RTT in particular. The bandwidth +/// generation advances with the reset, so an embedder that rereads the +/// bandwidth only on a new generation still sees it. +#[tokio::test] +async fn autodetect_handles_start_at_the_sentinel_for_each_connection() { + use core::net::{Ipv4Addr, SocketAddr}; + use core::sync::atomic::Ordering; + + use ironrdp_core::{WriteBuf, decode, encode_buf}; + use ironrdp_pdu::nego::{self, SecurityProtocol}; + use ironrdp_pdu::x224::X224; + use ironrdp_server::{RdpServer, TransportTls}; + use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _}; + + let mut server = RdpServer::builder() + .with_addr(SocketAddr::from((Ipv4Addr::LOCALHOST, 0))) + .with_no_security() + .with_no_input() + .with_no_display() + .build(); + server.enable_autodetect(); + + let rtt = server.autodetect_rtt_handle(); + let baseline_rtt = server.autodetect_baseline_rtt_handle(); + let bandwidth = server.autodetect_bandwidth_handle(); + let bandwidth_generation = server.autodetect_bandwidth_generation_handle(); + // What an earlier connection would have left behind. + rtt.store(30, Ordering::Relaxed); + baseline_rtt.store(10, Ordering::Relaxed); + bandwidth.store(50_000, Ordering::Relaxed); + bandwidth_generation.store(3, Ordering::Release); + + let (mut client, server_side) = tokio::io::duplex(4096); + + // `RdpServer`'s future is !Send, hence the LocalSet rather than a plain spawn. + let local = tokio::task::LocalSet::new(); + local + .run_until(async move { + let serving = tokio::task::spawn_local(async move { + server.run_connection_with(server_side, TransportTls::AlreadyDone).await + }); + + // Negotiate and hang up: the connection gets as far as finalization + // and ends there, without a single auto-detect exchange. + let request = nego::ConnectionRequest { + nego_data: None, + flags: nego::RequestFlags::empty(), + protocol: SecurityProtocol::empty(), + correlation_info: None, + }; + let mut buf = WriteBuf::new(); + encode_buf(&X224(request), &mut buf).expect("encode connection request"); + client.write_all(buf.filled()).await.expect("send connection request"); + + let mut confirm = [0u8; 128]; + let read = client.read(&mut confirm).await.expect("read connection confirm"); + let _ = decode::>(&confirm[..read]).expect("server answered the negotiation"); + + drop(client); + let _ = serving.await.expect("serving task panicked"); + }) + .await; + + assert_eq!(rtt.load(Ordering::Relaxed), u32::MAX); + assert_eq!(baseline_rtt.load(Ordering::Relaxed), u32::MAX); + assert_eq!(bandwidth_generation.load(Ordering::Acquire), 4); + assert_eq!(bandwidth.load(Ordering::Relaxed), u32::MAX); +} + #[test] fn autodetect_bandwidth_generation_handle_defaults_to_zero() { use core::net::{Ipv4Addr, SocketAddr}; From 32f68d6e946b6089571121fd07ccc64b2b73288c Mon Sep 17 00:00:00 2001 From: uchouT Date: Thu, 1 Oct 2026 15:34:43 +0000 Subject: [PATCH 2/2] removes magic `.1` Signed-off-by: uchouT --- crates/ironrdp-server/src/server.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/ironrdp-server/src/server.rs b/crates/ironrdp-server/src/server.rs index 45e56906a8..0d2748992f 100644 --- a/crates/ironrdp-server/src/server.rs +++ b/crates/ironrdp-server/src/server.rs @@ -4031,7 +4031,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.1.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",