diff --git a/Cargo.lock b/Cargo.lock index ea162c2dfd..f93db0054d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2811,9 +2811,16 @@ name = "ironrdp-cliprdr-native" version = "0.7.0" dependencies = [ "ironrdp-cliprdr", + "ironrdp-cliprdr-format", "ironrdp-core 0.2.1", + "nix", "tracing", + "visibility", + "wayland-client", + "wayland-protocols", + "wayland-protocols-wlr", "windows", + "x11rb", ] [[package]] @@ -3390,6 +3397,7 @@ dependencies = [ "ironrdp-cfg", "ironrdp-cliprdr", "ironrdp-cliprdr-format", + "ironrdp-cliprdr-native", "ironrdp-connector", "ironrdp-core 0.2.1", "ironrdp-displaycontrol", diff --git a/crates/ironrdp-client/src/lib.rs b/crates/ironrdp-client/src/lib.rs index fe4b112898..34cf5be4e3 100644 --- a/crates/ironrdp-client/src/lib.rs +++ b/crates/ironrdp-client/src/lib.rs @@ -12,7 +12,7 @@ pub mod output_channel; pub mod rail; pub mod rdp; -#[cfg(all(windows, feature = "clipboard"))] +#[cfg(all(any(windows, target_os = "linux"), feature = "clipboard"))] mod clipboard; mod ws; diff --git a/crates/ironrdp-client/src/rdp.rs b/crates/ironrdp-client/src/rdp.rs index d184f13b11..b85520a38a 100644 --- a/crates/ironrdp-client/src/rdp.rs +++ b/crates/ironrdp-client/src/rdp.rs @@ -790,14 +790,21 @@ impl RdpClient { // ── Clipboard initialisation (compile-time gated) ───────────────────── // // On Windows the WinClipboard object must outlive the entire connection loop, so we - // keep it alive via `_win_clipboard`. On non-Windows a StubClipboard backend is used - // and its ownership can be released immediately after the factory is extracted. + // keep it alive via `_win_clipboard`; the same goes for LinuxClipboard and + // `_linux_clipboard`. Elsewhere a StubClipboard backend is used and its ownership can be + // released immediately after the factory is extracted. #[cfg(all(windows, feature = "clipboard"))] #[expect( clippy::collection_is_never_read, reason = "binding owns the Windows clipboard so it stays alive for the connection's lifetime" )] let _win_clipboard; + #[cfg(all(target_os = "linux", feature = "clipboard"))] + #[expect( + clippy::collection_is_never_read, + reason = "binding owns the Linux clipboard so it stays alive for the connection's lifetime" + )] + let mut _linux_clipboard = None; #[cfg(feature = "clipboard")] let cliprdr_factory: Option>; @@ -852,7 +859,25 @@ impl RdpClient { } } - #[cfg(not(windows))] + #[cfg(target_os = "linux")] + { + use crate::clipboard::ClientClipboardMessageProxy; + use ironrdp_cliprdr_native::{LinuxClipboard, StubClipboard}; + match LinuxClipboard::new(ClientClipboardMessageProxy::new(self.input_event_sender.clone())) { + Ok(clipboard) => { + cliprdr_factory = Some(clipboard.backend_factory()); + _linux_clipboard = Some(clipboard); + } + Err(error) => { + // Without a desktop clipboard there is nothing to bridge, and the + // session is still useful, so this is not a connection failure. + warn!(%error, "OS clipboard unavailable; clipboard redirection is off for this session"); + cliprdr_factory = Some(StubClipboard::new().backend_factory()); + } + } + } + + #[cfg(not(any(windows, target_os = "linux")))] { use ironrdp_cliprdr_native::StubClipboard; let stub = StubClipboard::new(); diff --git a/crates/ironrdp-cliprdr-native/Cargo.toml b/crates/ironrdp-cliprdr-native/Cargo.toml index c4c55328ce..af61fea7e5 100644 --- a/crates/ironrdp-cliprdr-native/Cargo.toml +++ b/crates/ironrdp-cliprdr-native/Cargo.toml @@ -12,6 +12,10 @@ authors.workspace = true keywords.workspace = true categories.workspace = true +[features] +# Exposes internals to the integration tests in ironrdp-testsuite-core. +__test = ["dep:visibility"] + [lib] doctest = false test = false @@ -21,6 +25,15 @@ ironrdp-cliprdr = { path = "../ironrdp-cliprdr", version = "0.7" } # public ironrdp-core = { path = "../ironrdp-core", version = "0.2" } tracing = { version = "0.1", features = ["log"] } +[target.'cfg(target_os = "linux")'.dependencies] +nix = { version = "0.31", features = ["fs", "poll"] } +visibility = { version = "0.1", optional = true } +wayland-client = "0.31" +wayland-protocols = { version = "0.32", features = ["client", "staging"] } +wayland-protocols-wlr = { version = "0.3", features = ["client"] } +x11rb = { version = "0.13", features = ["xfixes"] } +ironrdp-cliprdr-format = { path = "../ironrdp-cliprdr-format", version = "0.2" } + [target.'cfg(windows)'.dependencies] windows = { version = "0.62", features = [ "Win32_Foundation", diff --git a/crates/ironrdp-cliprdr-native/README.md b/crates/ironrdp-cliprdr-native/README.md index 42c396b7a9..5425841ab8 100644 --- a/crates/ironrdp-cliprdr-native/README.md +++ b/crates/ironrdp-cliprdr-native/README.md @@ -1,7 +1,35 @@ # IronRDP CLIPRDR native backends -Native CLIPRDR backend implementations. Currently only Windows is supported. +Native CLIPRDR backend implementations: Windows (`WinClipboard`) and Linux desktops over X11 and Wayland (`LinuxClipboard`, text and images). This crate is part of the [IronRDP] project. [IronRDP]: https://github.com/Devolutions/IronRDP + +## Linux + +`LinuxClipboard` uses the shared `data_control` client for Wayland compositors offering `ext-data-control-v1` or `wlr-data-control-unstable-v1`. +When data-control is unavailable, including on GNOME, it falls back to X11/XWayland through XFixes and selection ownership. +A desktop without either clipboard service uses the client's stub backend. + +Both backends react to selection changes and advertise remote formats without fetching their contents. +Text (`CF_UNICODETEXT`) and images (`CF_DIB`/`CF_DIBV5`, converted to/from PNG by `ironrdp-cliprdr-format`) are requested only when an application pastes. +Local clipboard content is read only when the peer requests it. +File clipboard transfer and HTML are not supported. + +Only one Format Data Request is outstanding at a time, because responses do not identify their request. +A paste times out after five seconds, but its outstanding request is retained until the late response is drained. +If the peer never replies, new remote pastes fail until the clipboard channel is reinitialized; this prevents stale content from being mistaken for a newer copy. +Selection changes cancel older paste waiters, and the existing `ironrdp_cliprdr::loop_detector` detects content echoed by clipboard managers. + +The clipboard workers bound data transfers and pending pastes. +X11 transfers use INCR for large content and separate requestor windows so late selection responses cannot complete newer reads. + +State-machine tests run in the normal CI suite with `cargo test -p ironrdp-testsuite-core cliprdr_native`. +Three additional integration tests require a dedicated X11 display; the interoperability test also requires `xclip`: + +```sh +DISPLAY=:N cargo test -p ironrdp-testsuite-core cliprdr_native::linux::x11_ -- --ignored --test-threads=1 +``` + +Use a private display such as Xvfb for these tests, because they replace its clipboard selection. diff --git a/crates/ironrdp-cliprdr-native/src/data_control/client.rs b/crates/ironrdp-cliprdr-native/src/data_control/client.rs new file mode 100644 index 0000000000..f94ab0e3dd --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/client.rs @@ -0,0 +1,411 @@ +//! The public clipboard handle. + +use core::sync::atomic::{AtomicBool, Ordering}; +use core::time::Duration; +use std::{ + collections::HashMap, + io::{Read as _, Write as _}, + os::{ + fd::OwnedFd, + unix::{io::AsFd as _, net::UnixStream}, + }, + sync::{Arc, Mutex, MutexGuard, mpsc}, + thread::JoinHandle, + time::Instant, +}; + +use nix::poll::{PollFd, PollFlags, PollTimeout, poll}; + +use super::{ + error::{Error, Result}, + mime::find_mime_match, + options::{Options, Protocol}, + state::{Command, Shared}, + worker, +}; + +/// Maximum total time for a clipboard read, including a slowly streaming source. +const READ_TIMEOUT: Duration = Duration::from_secs(5); + +/// What to offer as the clipboard selection. +/// +/// Each MIME type is either given data up front with [`data`](Self::data), or +/// only advertised with [`advertise`](Self::advertise). A paste of an +/// advertised type with no data raises a [`TransferRequest`], which the owner +/// answers later. That is delayed rendering: the data is produced only if +/// something pastes it. +#[derive(Debug, Clone, Default)] +pub struct Content { + mime_types: Vec, + data: HashMap>, +} + +impl Content { + /// Empty content. + #[must_use] + pub fn new() -> Self { + Self::default() + } + + /// Offer `mime_type` with its data. + #[must_use] + pub fn data(mut self, mime_type: impl Into, data: impl Into>) -> Self { + let mime_type = mime_type.into(); + self.push_type(&mime_type); + self.data.insert(mime_type, data.into()); + self + } + + /// Offer `mime_type` without data; pastes of it raise a [`TransferRequest`]. + #[must_use] + pub fn advertise(mut self, mime_type: impl Into) -> Self { + self.push_type(&mime_type.into()); + self + } + + fn push_type(&mut self, mime_type: &str) { + if !self.mime_types.iter().any(|m| m == mime_type) { + self.mime_types.push(mime_type.to_owned()); + } + } +} + +/// Channel to the worker thread: a command queue and a socket that wakes its +/// poll loop. +struct Link { + commands: mpsc::Sender, + wake: UnixStream, +} + +impl Link { + fn send(&self, command: Command) -> Result<()> { + self.commands.send(command).map_err(|_| Error::Stopped)?; + self.notify(); + Ok(()) + } + + fn notify(&self) { + // A full socket buffer means a wake-up is already pending. + let _ = (&self.wake).write(&[1]); + } +} + +/// A paste waiting for data that was advertised without any. +/// +/// Answer it with [`complete`](Self::complete), or [`fail`](Self::fail) if the +/// data cannot be produced. Dropping it without answering leaves the paste +/// waiting until the selection is replaced or its five-second deadline expires. +#[derive(Debug)] +#[must_use = "a paste stays open until the request is answered"] +pub struct TransferRequest { + serial: u32, + mime_type: String, + link: Arc, +} + +impl core::fmt::Debug for Link { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.debug_struct("Link").finish_non_exhaustive() + } +} + +impl TransferRequest { + /// The MIME type being pasted. + #[must_use] + pub fn mime_type(&self) -> &str { + &self.mime_type + } + + /// Supply the data. It is written to every paste of this type that is + /// waiting, and cached for later pastes. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the clipboard worker has shut down. + pub fn complete(self, data: impl Into>) -> Result<()> { + self.link.send(Command::CompleteTransfer { + serial: self.serial, + data: Some(data.into()), + }) + } + + /// Close the waiting pastes with no data. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the clipboard worker has shut down. + pub fn fail(self) -> Result<()> { + self.link.send(Command::CompleteTransfer { + serial: self.serial, + data: None, + }) + } +} + +/// A connection to the Wayland compositor's clipboard. +/// +/// Connecting spawns one thread that owns the Wayland connection. Dropping the +/// handle stops and joins it. +/// +/// ```no_run +/// use ironrdp_cliprdr_native::data_control::{Content, DataControl}; +/// +/// # fn main() -> ironrdp_cliprdr_native::data_control::Result<()> { +/// let clipboard = DataControl::connect()?; +/// clipboard.set_selection(Content::new().data("text/plain;charset=utf-8", "hello"))?; +/// # Ok(()) +/// # } +/// ``` +pub struct DataControl { + link: Arc, + shared: Arc>, + protocol: Protocol, + max_read_bytes: usize, + stop: Arc, + thread: Option>, +} + +impl core::fmt::Debug for DataControl { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.debug_struct("DataControl") + .field("protocol", &self.protocol) + .finish_non_exhaustive() + } +} + +impl DataControl { + /// Connect with the default [`Options`]. + /// + /// # Errors + /// + /// See [`Error`]: no compositor, no data-control protocol, or no seat. + pub fn connect() -> Result { + Self::connect_with(&Options::new()) + } + + /// Connect with explicit options. + /// + /// # Errors + /// + /// See [`Error`]: no compositor, no data-control protocol, or no seat. + pub fn connect_with(options: &Options) -> Result { + let (command_tx, command_rx) = mpsc::channel(); + let (wake_tx, wake_rx) = UnixStream::pair()?; + wake_tx.set_nonblocking(true)?; + wake_rx.set_nonblocking(true)?; + let stop = Arc::new(AtomicBool::new(false)); + let (ready_tx, ready_rx) = mpsc::sync_channel(1); + + let thread_options = options.clone(); + let thread_stop = Arc::clone(&stop); + let thread = std::thread::Builder::new() + .name("ironrdp-cliprdr-data-control".into()) + .spawn( + move || match worker::connect(&thread_options, command_rx, wake_rx, thread_stop) { + Ok((worker, connected)) => { + if ready_tx.send(Ok(connected)).is_ok() { + worker.run(); + } + } + Err(error) => { + let _ = ready_tx.send(Err(error)); + } + }, + )?; + + let connected = match ready_rx.recv() { + Ok(Ok(connected)) => connected, + Ok(Err(error)) => { + let _ = thread.join(); + return Err(error); + } + Err(_) => { + let _ = thread.join(); + return Err(Error::Stopped); + } + }; + + Ok(Self { + link: Arc::new(Link { + commands: command_tx, + wake: wake_tx, + }), + shared: connected.shared, + protocol: connected.protocol, + max_read_bytes: options.max_read_bytes, + stop, + thread: Some(thread), + }) + } + + /// Wait for the compositor to process earlier commands and selection events + /// before inspecting ownership in the native bridge. + pub(crate) fn synchronize(&self) -> Result<()> { + let (sender, receiver) = mpsc::channel(); + self.link.send(Command::Synchronize(sender))?; + receiver.recv_timeout(READ_TIMEOUT).map_err(|error| match error { + mpsc::RecvTimeoutError::Timeout => Error::Timeout, + mpsc::RecvTimeoutError::Disconnected => Error::Stopped, + }) + } + + pub(crate) fn owns_selection(&self) -> bool { + self.shared().own_source_live + } + + /// The protocol the compositor is speaking. + #[must_use] + pub fn protocol(&self) -> Protocol { + self.protocol + } + + fn shared(&self) -> MutexGuard<'_, Shared> { + self.shared.lock().unwrap_or_else(std::sync::PoisonError::into_inner) + } + + /// MIME types the current selection offers; empty if there is none. + #[must_use] + pub fn selection_mime_types(&self) -> Vec { + self.shared().mime_types.clone() + } + + /// A counter that increases on every selection change. + #[must_use] + pub fn serial(&self) -> u32 { + self.shared().serial + } + + /// Call `callback` with the offered MIME types whenever another client + /// changes the selection. + /// + /// A selection this handle set itself is not reported; it still advances + /// [`serial`](Self::serial) and updates + /// [`selection_mime_types`](Self::selection_mime_types). + /// + /// The callback runs on the worker thread and must not block. + pub fn on_change(&self, callback: impl Fn(Vec) + Send + Sync + 'static) { + self.shared().on_change = Some(Arc::new(callback)); + } + + /// Call `callback` when something pastes a type that was advertised + /// without data. + /// + /// The callback runs on the worker thread and must not block; hand the + /// request to another thread and answer it there. + pub fn on_transfer(&self, callback: impl Fn(TransferRequest) + Send + Sync + 'static) { + let link = Arc::clone(&self.link); + self.shared().on_transfer = Some(Arc::new(move |serial, mime_type| { + callback(TransferRequest { + serial, + mime_type, + link: Arc::clone(&link), + }); + })); + } + + /// Become the clipboard owner and offer `content`. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the worker has shut down. + pub fn set_selection(&self, content: Content) -> Result<()> { + self.link.send(Command::SetSelection { + mime_types: content.mime_types, + data: content.data, + }) + } + + /// Add data for a type after [`set_selection`](Self::set_selection), + /// without replacing the selection. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the worker has shut down. + pub fn update_data(&self, mime_type: impl Into, data: impl Into>) -> Result<()> { + self.link.send(Command::UpdateSourceData { + mime_type: mime_type.into(), + data: data.into(), + }) + } + + /// Give up the selection if this handle still owns it. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the worker has shut down. + pub fn clear_selection(&self) -> Result<()> { + self.link.send(Command::ClearSelection) + } + + /// Read the current selection as `mime_type`. + /// + /// Charset differences are tolerated, see [`find_mime_match`]. Returns + /// `Ok(None)` when there is no selection or it does not offer the type. + /// + /// # Errors + /// + /// [`Error::TooLarge`] past the size limit, [`Error::Timeout`] if the + /// source does not complete within five seconds, [`Error::Stopped`] if the worker has + /// shut down, [`Error::Io`] for pipe failures. + pub fn read(&self, mime_type: &str) -> Result>> { + let available = self.selection_mime_types(); + let Some(matched) = find_mime_match(mime_type, &available) else { + return Ok(None); + }; + + let (mut reader, writer) = std::io::pipe()?; + self.link.send(Command::ReceiveFromOffer { + mime_type: matched.to_owned(), + fd: OwnedFd::from(writer), + })?; + + read_pipe(&mut reader, self.max_read_bytes, READ_TIMEOUT).map(Some) + } +} + +impl Drop for DataControl { + fn drop(&mut self) { + self.stop.store(true, Ordering::Relaxed); + self.link.notify(); + if let Some(thread) = self.thread.take() { + let _ = thread.join(); + } + } +} + +/// Read a pipe with both a total deadline and an allocation limit. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) fn read_pipe(reader: &mut std::io::PipeReader, max_read_bytes: usize, timeout: Duration) -> Result> { + let deadline = Instant::now() + timeout; + let mut data = Vec::new(); + let mut chunk = vec![0u8; 64 * 1024]; + loop { + let mut fds = [PollFd::new(reader.as_fd(), PollFlags::POLLIN)]; + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Err(Error::Timeout); + } + let timeout = PollTimeout::try_from(remaining).unwrap_or(PollTimeout::MAX); + match poll(&mut fds, timeout) { + Ok(0) => return Err(Error::Timeout), + Ok(_) => {} + Err(nix::errno::Errno::EINTR) => continue, + Err(errno) => return Err(std::io::Error::from(errno).into()), + } + match reader.read(&mut chunk) { + // The writer closed: the source has sent everything. + Ok(0) => return Ok(data), + Ok(n) => { + if data.len() + n > max_read_bytes { + return Err(Error::TooLarge { + size: data.len() + n, + limit: max_read_bytes, + }); + } + data.extend_from_slice(&chunk[..n]); + } + Err(error) if error.kind() == std::io::ErrorKind::Interrupted => {} + Err(error) => return Err(error.into()), + } + } +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/dispatch.rs b/crates/ironrdp-cliprdr-native/src/data_control/dispatch.rs new file mode 100644 index 0000000000..43345bbe24 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/dispatch.rs @@ -0,0 +1,237 @@ +//! Wayland event dispatch for the data-control objects. +//! +//! Both protocols (`ext-data-control-v1` and `wlr-data-control-unstable-v1`) +//! have the same object model and the same events, so each handler is written +//! twice and forwards to the shared [`State`]. + +use wayland_client::{ + Connection, Dispatch, Proxy as _, QueueHandle, + globals::GlobalListContents, + protocol::{wl_registry::WlRegistry, wl_seat::WlSeat}, +}; +use wayland_protocols::ext::data_control::v1::client::{ + ext_data_control_device_v1::{self, ExtDataControlDeviceV1}, + ext_data_control_manager_v1::ExtDataControlManagerV1, + ext_data_control_offer_v1::{self, ExtDataControlOfferV1}, + ext_data_control_source_v1::{self, ExtDataControlSourceV1}, +}; +use wayland_protocols_wlr::data_control::v1::client::{ + zwlr_data_control_device_v1::{self, ZwlrDataControlDeviceV1}, + zwlr_data_control_manager_v1::ZwlrDataControlManagerV1, + zwlr_data_control_offer_v1::{self, ZwlrDataControlOfferV1}, + zwlr_data_control_source_v1::{self, ZwlrDataControlSourceV1}, +}; + +use super::state::State; + +/// The dispatch target the Wayland event queue calls into. +pub(crate) struct Client { + pub(crate) data_control: State, +} + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &WlRegistry, + _event: ::Event, + _data: &GlobalListContents, + _conn: &Connection, + _qh: &QueueHandle, + ) { + // Globals are looked up once at start-up; later additions are ignored. + } +} + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &WlSeat, + _event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The seat is only an argument to `get_data_device`. + } +} + +// === ext-data-control-v1 === + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &ExtDataControlManagerV1, + _event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The manager has no events. + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ExtDataControlDeviceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + ext_data_control_device_v1::Event::DataOffer { id } => { + state.data_control.on_data_offer_ext(id); + } + ext_data_control_device_v1::Event::Selection { id } => { + if id.is_some() { + state.data_control.on_selection(); + } else { + state.data_control.on_selection_cleared(); + } + } + ext_data_control_device_v1::Event::Finished => { + state.data_control.on_device_finished(); + } + ext_data_control_device_v1::Event::PrimarySelection { .. } => { + // Only the regular clipboard is handled, not the primary selection. + tracing::trace!("ext data control primary selection event (ignored)"); + } + _ => {} + } + } + + // The `data_offer` event creates a child offer object; without this, + // wayland-client's default panics. + wayland_client::event_created_child!(Client, ExtDataControlDeviceV1, [ + ext_data_control_device_v1::EVT_DATA_OFFER_OPCODE => (ExtDataControlOfferV1, ()), + ]); +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + proxy: &ExtDataControlSourceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if !state.data_control.is_current_source(&proxy.id()) { + return; + } + match event { + ext_data_control_source_v1::Event::Send { mime_type, fd } => { + state.data_control.on_source_send(&mime_type, fd); + } + ext_data_control_source_v1::Event::Cancelled => { + state.data_control.on_source_cancelled(); + } + _ => {} + } + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ExtDataControlOfferV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if let ext_data_control_offer_v1::Event::Offer { mime_type } = event { + state.data_control.on_offer_mime_type(mime_type); + } + } +} + +// === wlr-data-control-unstable-v1 === + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &ZwlrDataControlManagerV1, + _event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The manager has no events. + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ZwlrDataControlDeviceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + zwlr_data_control_device_v1::Event::DataOffer { id } => { + state.data_control.on_data_offer_wlr(id); + } + zwlr_data_control_device_v1::Event::Selection { id } => { + if id.is_some() { + state.data_control.on_selection(); + } else { + state.data_control.on_selection_cleared(); + } + } + zwlr_data_control_device_v1::Event::Finished => { + state.data_control.on_device_finished(); + } + zwlr_data_control_device_v1::Event::PrimarySelection { .. } => { + tracing::trace!("wlr data control primary selection event (ignored)"); + } + _ => {} + } + } + + wayland_client::event_created_child!(Client, ZwlrDataControlDeviceV1, [ + zwlr_data_control_device_v1::EVT_DATA_OFFER_OPCODE => (ZwlrDataControlOfferV1, ()), + ]); +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + proxy: &ZwlrDataControlSourceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if !state.data_control.is_current_source(&proxy.id()) { + return; + } + match event { + zwlr_data_control_source_v1::Event::Send { mime_type, fd } => { + state.data_control.on_source_send(&mime_type, fd); + } + zwlr_data_control_source_v1::Event::Cancelled => { + state.data_control.on_source_cancelled(); + } + _ => {} + } + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ZwlrDataControlOfferV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if let zwlr_data_control_offer_v1::Event::Offer { mime_type } = event { + state.data_control.on_offer_mime_type(mime_type); + } + } +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/error.rs b/crates/ironrdp-cliprdr-native/src/data_control/error.rs new file mode 100644 index 0000000000..1e2c3e6dc4 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/error.rs @@ -0,0 +1,72 @@ +//! Errors returned by the clipboard client. + +use core::fmt; + +/// Errors returned by [`DataControl`](super::DataControl). +#[derive(Debug)] +#[non_exhaustive] +pub enum Error { + /// The connection to the Wayland compositor could not be made. + Connect(String), + /// The compositor offers neither data-control protocol. + Unsupported, + /// The requested protocol is not offered and fallback is disabled. + ProtocolUnavailable, + /// The compositor advertises no seat to attach the clipboard to. + NoSeat, + /// The compositor reported a protocol error or the connection failed. + Wayland(String), + /// The worker thread has stopped, so the request could not be delivered. + Stopped, + /// Clipboard data exceeded the size limit. + TooLarge { + /// Bytes received when the limit was hit. + size: usize, + /// The limit. + limit: usize, + }, + /// The source did not deliver its data in time. + Timeout, + /// An I/O error while creating or reading a pipe. + Io(std::io::Error), +} + +impl fmt::Display for Error { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Connect(reason) => write!(f, "cannot connect to the Wayland compositor: {reason}"), + Self::Unsupported => { + f.write_str("the compositor offers neither ext-data-control-v1 nor wlr-data-control-unstable-v1") + } + Self::ProtocolUnavailable => { + f.write_str("the requested data-control protocol is not offered by the compositor") + } + Self::NoSeat => f.write_str("the compositor advertises no seat"), + Self::Wayland(reason) => write!(f, "Wayland error: {reason}"), + Self::Stopped => f.write_str("the clipboard worker has stopped"), + Self::TooLarge { size, limit } => { + write!(f, "clipboard data of {size} bytes exceeds the {limit} byte limit") + } + Self::Timeout => f.write_str("the clipboard source did not deliver its data in time"), + Self::Io(error) => write!(f, "I/O error: {error}"), + } + } +} + +impl core::error::Error for Error { + fn source(&self) -> Option<&(dyn core::error::Error + 'static)> { + match self { + Self::Io(error) => Some(error), + _ => None, + } + } +} + +impl From for Error { + fn from(error: std::io::Error) -> Self { + Self::Io(error) + } +} + +/// Result alias for this module. +pub type Result = core::result::Result; diff --git a/crates/ironrdp-cliprdr-native/src/data_control/mime.rs b/crates/ironrdp-cliprdr-native/src/data_control/mime.rs new file mode 100644 index 0000000000..978e5f5eb7 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/mime.rs @@ -0,0 +1,53 @@ +//! MIME type matching that tolerates charset parameters. + +/// Find a MIME type match in `available`, tolerating charset differences. +/// +/// A source may offer `text/plain` while a consumer asks for +/// `text/plain;charset=utf-8`, or the other way round. An exact match wins; +/// for `text/` types the charset parameter is then stripped or added, and +/// finally any offered type with the same base type is accepted. +/// +/// Returns the string from `available` that should be used for the request, so +/// that the compositor receives a type the source actually offered. +/// +/// ``` +/// use ironrdp_cliprdr_native::data_control::find_mime_match; +/// +/// let available = vec!["text/plain;charset=utf-8".to_owned()]; +/// assert_eq!( +/// find_mime_match("text/plain", &available), +/// Some("text/plain;charset=utf-8") +/// ); +/// assert_eq!(find_mime_match("image/png", &available), None); +/// ``` +#[must_use] +pub fn find_mime_match<'a>(requested: &str, available: &'a [String]) -> Option<&'a str> { + if let Some(found) = available.iter().find(|m| m.as_str() == requested) { + return Some(found.as_str()); + } + + if !requested.starts_with("text/") { + return None; + } + let base = requested.split(';').next()?; + + if requested.contains(';') { + // The request carries a charset; try the bare type. + if let Some(found) = available.iter().find(|m| m.as_str() == base) { + return Some(found.as_str()); + } + } else { + // The request has none; try the common charset spellings. + for suffix in [";charset=utf-8", ";charset=UTF-8"] { + let with_charset = format!("{requested}{suffix}"); + if let Some(found) = available.iter().find(|m| m.as_str() == with_charset) { + return Some(found.as_str()); + } + } + } + + available + .iter() + .find(|m| m.split(';').next() == Some(base)) + .map(String::as_str) +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/mod.rs b/crates/ironrdp-cliprdr-native/src/data_control/mod.rs new file mode 100644 index 0000000000..3ff9afcb94 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/mod.rs @@ -0,0 +1,62 @@ +//! A Wayland clipboard client built on the data-control protocols. +//! +//! Data-control lets a client that has no window read and set the clipboard. +//! This module speaks both `ext-data-control-v1` and +//! `wlr-data-control-unstable-v1` and picks whichever the compositor offers. +//! +//! It has no async runtime. One thread owns the Wayland connection, and +//! [`DataControl`] is the handle to call from anywhere. +//! +//! # Reading +//! +//! ```no_run +//! use ironrdp_cliprdr_native::data_control::DataControl; +//! +//! # fn main() -> ironrdp_cliprdr_native::data_control::Result<()> { +//! let clipboard = DataControl::connect()?; +//! if let Some(bytes) = clipboard.read("text/plain;charset=utf-8")? { +//! println!("{}", String::from_utf8_lossy(&bytes)); +//! } +//! # Ok(()) +//! # } +//! ``` +//! +//! # Owning the clipboard, with delayed rendering +//! +//! A type advertised without data raises a [`TransferRequest`] when something +//! pastes it, so the data is produced only if it is wanted. That maps onto +//! CLIPRDR, where the Format Data Response follows the paste. +//! +//! ```no_run +//! use ironrdp_cliprdr_native::data_control::{Content, DataControl}; +//! +//! # fn main() -> ironrdp_cliprdr_native::data_control::Result<()> { +//! let clipboard = DataControl::connect()?; +//! clipboard.on_transfer(|request| { +//! let _ = request.complete("rendered on demand"); +//! }); +//! clipboard.set_selection(Content::new().advertise("text/plain;charset=utf-8"))?; +//! # Ok(()) +//! # } +//! ``` +//! +//! GNOME's Mutter offers no data-control protocol, so [`DataControl::connect`] +//! returns [`Error::Unsupported`] there. + +mod client; +mod dispatch; +mod error; +mod mime; +mod options; +#[cfg(feature = "__test")] +pub mod state; +#[cfg(not(feature = "__test"))] +pub(crate) mod state; +mod worker; + +#[cfg(feature = "__test")] +pub use self::client::read_pipe; +pub use self::client::{Content, DataControl, TransferRequest}; +pub use self::error::{Error, Result}; +pub use self::mime::find_mime_match; +pub use self::options::{Options, Preference, Protocol}; diff --git a/crates/ironrdp-cliprdr-native/src/data_control/options.rs b/crates/ironrdp-cliprdr-native/src/data_control/options.rs new file mode 100644 index 0000000000..592dacf0e3 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/options.rs @@ -0,0 +1,107 @@ +//! Connection options and the protocol choice. + +use std::fmt; + +/// The data-control protocol in use. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[non_exhaustive] +pub enum Protocol { + /// `ext-data-control-v1`, the standardized protocol. + Ext, + /// `wlr-data-control-unstable-v1`, the wlroots protocol. + Wlr, +} + +impl fmt::Display for Protocol { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Protocol::Ext => f.write_str("ext-data-control-v1"), + Protocol::Wlr => f.write_str("wlr-data-control-unstable-v1"), + } + } +} + +/// Which protocol to prefer when the compositor offers both. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +#[non_exhaustive] +pub enum Preference { + /// Prefer `ext-data-control-v1`, then `wlr-data-control-unstable-v1`. + #[default] + Auto, + /// Prefer `ext-data-control-v1`. + Ext, + /// Prefer `wlr-data-control-unstable-v1`. + Wlr, +} + +/// Options for [`DataControl::connect_with`](super::DataControl::connect_with). +/// +/// ``` +/// use ironrdp_cliprdr_native::data_control::{Options, Preference}; +/// +/// let options = Options::new() +/// .preference(Preference::Wlr) +/// .allow_fallback(false); +/// ``` +#[derive(Debug, Clone)] +pub struct Options { + pub(crate) preference: Preference, + pub(crate) allow_fallback: bool, + pub(crate) max_read_bytes: usize, +} + +impl Options { + /// The default read size limit: 100 MiB. + pub const DEFAULT_MAX_READ_BYTES: usize = 100 * 1024 * 1024; + + /// Options with automatic protocol selection and fallback enabled. + #[must_use] + pub fn new() -> Self { + Self { + preference: Preference::Auto, + allow_fallback: true, + max_read_bytes: Self::DEFAULT_MAX_READ_BYTES, + } + } + + /// Choose which protocol to prefer. + #[must_use] + pub fn preference(mut self, preference: Preference) -> Self { + self.preference = preference; + self + } + + /// Whether to use the other protocol when the preferred one is missing. + #[must_use] + pub fn allow_fallback(mut self, allow: bool) -> Self { + self.allow_fallback = allow; + self + } + + /// Limit, in bytes, for one [`read`](super::DataControl::read). + #[must_use] + pub fn max_read_bytes(mut self, limit: usize) -> Self { + self.max_read_bytes = limit; + self + } + + /// The protocols to try, in order. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn candidates(&self) -> Vec { + let (first, second) = match self.preference { + Preference::Auto | Preference::Ext => (Protocol::Ext, Protocol::Wlr), + Preference::Wlr => (Protocol::Wlr, Protocol::Ext), + }; + if self.allow_fallback { + vec![first, second] + } else { + vec![first] + } + } +} + +impl Default for Options { + fn default() -> Self { + Self::new() + } +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/state.rs b/crates/ironrdp-cliprdr-native/src/data_control/state.rs new file mode 100644 index 0000000000..0a756b9027 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/state.rs @@ -0,0 +1,834 @@ +//! The data-control protocol state machine. +//! +//! [`State`] owns the protocol objects (manager, device, current offer, our +//! source) and turns compositor events into changes of the shared selection +//! state. It is deliberately free of any event loop: the worker in +//! [`super::worker`] feeds it events and commands. +//! +//! # Flows +//! +//! **Set selection (client to compositor):** +//! 1. A [`Command::SetSelection`] arrives with the MIME types and any data. +//! 2. A data-control source is created, the types are advertised on it and it +//! is made the device's selection; the previous source is destroyed only +//! afterwards, so the clipboard is never empty in between. +//! 3. A `send` event for an advertised type writes the cached data to the +//! requested file descriptor on a worker thread. With no cached data and an +//! `on_transfer` callback set, the paste is held until +//! [`Command::CompleteTransfer`] answers it (delayed rendering). +//! +//! **Selection changed (compositor to client):** +//! 1. `data_offer`, then one `offer` per MIME type, then `selection`. +//! 2. The types are stored in [`Shared`] and the change callback is called, +//! unless the selection is the one this client set itself. +//! +//! **Read (client reads the compositor's selection):** +//! [`Command::ReceiveFromOffer`] calls `offer.receive(mime, fd)`; the caller +//! reads the other end of the pipe. + +use core::sync::atomic::{AtomicUsize, Ordering}; +use core::time::Duration; +use std::{ + collections::HashMap, + os::unix::io::{AsFd as _, OwnedFd}, + sync::{Arc, Mutex}, + time::Instant, +}; + +use nix::fcntl::{FcntlArg, OFlag, fcntl}; +use nix::poll::{PollFd, PollFlags, PollTimeout, poll}; + +use wayland_client::{Dispatch, Proxy as _, QueueHandle, backend::ObjectId, protocol::wl_seat::WlSeat}; +use wayland_protocols::ext::data_control::v1::client::{ + ext_data_control_device_v1::ExtDataControlDeviceV1, ext_data_control_manager_v1::ExtDataControlManagerV1, + ext_data_control_offer_v1::ExtDataControlOfferV1, ext_data_control_source_v1::ExtDataControlSourceV1, +}; +use wayland_protocols_wlr::data_control::v1::client::{ + zwlr_data_control_device_v1::ZwlrDataControlDeviceV1, zwlr_data_control_manager_v1::ZwlrDataControlManagerV1, + zwlr_data_control_offer_v1::ZwlrDataControlOfferV1, zwlr_data_control_source_v1::ZwlrDataControlSourceV1, +}; + +const MAX_PASTES: usize = 16; +const TRANSFER_TIMEOUT: Duration = Duration::from_secs(5); + +// === Protocol object enums === +// These wrap both ext and wlr variants so State can work +// with either protocol transparently. + +/// Data control manager (either ext or wlr). +pub(crate) enum DataControlManager { + /// ext-data-control-v1 manager. + Ext(ExtDataControlManagerV1), + /// wlr-data-control-unstable-v1 manager. + Wlr(ZwlrDataControlManagerV1), +} + +/// Data control device (either ext or wlr). +pub(crate) enum DataControlDevice { + /// ext-data-control-v1 device. + Ext(ExtDataControlDeviceV1), + /// wlr-data-control-unstable-v1 device. + Wlr(ZwlrDataControlDeviceV1), +} + +/// Data control offer (either ext or wlr). +pub(crate) enum DataControlOffer { + /// ext-data-control-v1 offer. + Ext(ExtDataControlOfferV1), + /// wlr-data-control-unstable-v1 offer. + Wlr(ZwlrDataControlOfferV1), +} + +impl DataControlOffer { + /// Request data transfer for a MIME type. + /// + /// Tells the source client to write data to the provided fd. + pub(crate) fn receive(&self, mime_type: &str, fd: &OwnedFd) { + use std::os::unix::io::AsFd as _; + match self { + DataControlOffer::Ext(offer) => offer.receive(mime_type.to_owned(), fd.as_fd()), + DataControlOffer::Wlr(offer) => offer.receive(mime_type.to_owned(), fd.as_fd()), + } + } + + /// Destroy this offer. + pub(crate) fn destroy(&self) { + match self { + DataControlOffer::Ext(offer) => offer.destroy(), + DataControlOffer::Wlr(offer) => offer.destroy(), + } + } +} + +/// Data control source (either ext or wlr). +enum DataControlSource { + /// ext-data-control-v1 source. + Ext(ExtDataControlSourceV1), + /// wlr-data-control-unstable-v1 source. + Wlr(ZwlrDataControlSourceV1), +} + +impl DataControlSource { + /// Advertise a MIME type on this source. + fn offer(&self, mime_type: &str) { + match self { + DataControlSource::Ext(source) => source.offer(mime_type.to_owned()), + DataControlSource::Wlr(source) => source.offer(mime_type.to_owned()), + } + } + + /// Destroy this source. + fn destroy(&self) { + match self { + DataControlSource::Ext(source) => source.destroy(), + DataControlSource::Wlr(source) => source.destroy(), + } + } +} + +// === Commands and shared state === + +/// Commands sent from clipboard backends to the Wayland event loop thread. +#[derive(Debug)] +pub(crate) enum Command { + /// Acknowledge earlier commands and selection events after a compositor sync. + Synchronize(std::sync::mpsc::Sender<()>), + /// Set the clipboard selection on the compositor. + /// + /// Creates a data control source with the offered MIME types and + /// stores the data for responding to `send` events. + SetSelection { + /// MIME types to advertise. + mime_types: Vec, + /// Data for each MIME type (written to fd on `send` event). + data: HashMap>, + }, + /// Update source data for a MIME type without re-creating the source. + /// + /// Used when data wasn't available at `SetSelection` time (eager fetch + /// from a remote clipboard). The Wayland data source stays unchanged; + /// only the cached data map is updated so the next `send` event can + /// serve the requested MIME type. + UpdateSourceData { + /// MIME type key for the data. + mime_type: String, + /// Data bytes to cache. + data: Vec, + }, + /// Receive clipboard data from the current compositor offer. + /// + /// Calls `offer.receive(mime_type, fd)` on the event loop thread. + /// The caller reads from the other end of the pipe. + ReceiveFromOffer { + /// MIME type to request. + mime_type: String, + /// Write end of the pipe (compositor writes data here). + fd: OwnedFd, + }, + /// Answer a transfer raised through `Shared::on_transfer`. + /// + /// Writes `data` to every paste waiting on that transfer and caches it + /// for later pastes of the same MIME type. `None` closes the waiting + /// pastes with no data. + CompleteTransfer { + /// Serial passed to the `on_transfer` callback. + serial: u32, + /// The data, or `None` if the remote could not supply it. + data: Option>, + }, + /// Give up our selection, if we still own it. + ClearSelection, +} + +/// Shared clipboard state readable from any thread. +/// +/// Updated by the event loop thread when the compositor's selection changes. +/// Read by clipboard backends to report current state. +#[derive(Default)] +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) struct Shared { + /// MIME types of the current compositor selection. + pub(crate) mime_types: Vec, + /// Serial number, incremented on each selection change. + pub(crate) serial: u32, + /// Change notification callback. + /// + /// Called on the event loop thread when selection changes. + /// Typically captures a tokio channel sender for async notification. + pub(crate) on_change: Option) + Send + Sync>>, + /// Paste callback for data not supplied up front. + /// + /// When set, a paste of an advertised MIME type with no cached data is + /// held open and this is called with a serial and the MIME type; the + /// answer comes back as `Command::CompleteTransfer`. Without + /// it such a paste gets no data. + pub(crate) on_transfer: Option>, + /// Whether our own source is still the compositor's selection source. + /// + /// While it is, data we cached ourselves can be read back without a + /// round trip; once another client takes the selection it must not be. + pub(crate) own_source_live: bool, +} + +#[cfg(feature = "__test")] +impl Shared { + /// Test-only: MIME types of the current selection. + pub fn mime_types(&self) -> &[String] { + &self.mime_types + } + + /// Test-only: the selection change counter. + pub fn serial(&self) -> u32 { + self.serial + } + + /// Test-only: install the change callback. + pub fn set_on_change(&mut self, callback: Arc) + Send + Sync>) { + self.on_change = Some(callback); + } + + /// Test-only: install the transfer callback. + pub fn set_on_transfer(&mut self, callback: Arc) { + self.on_transfer = Some(callback); + } +} + +// === Data control state === + +/// A paste waiting on data from `on_transfer`, with every paste of the same +/// MIME type that arrived meanwhile. +struct PendingTransfer { + serial: u32, + waiters: Vec, +} + +/// A slot covers both a held descriptor and the thread writing its response. +struct PastePermit(Arc); + +impl Drop for PastePermit { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::Relaxed); + } +} + +struct PasteWriter { + fd: OwnedFd, + deadline: Instant, + _permit: PastePermit, +} + +/// Accumulated MIME types for a pending data offer. +/// +/// Between `data_offer` and `selection` events, the compositor sends +/// `offer` events with MIME types. We collect them here. +#[derive(Default)] +struct PendingOffer { + /// MIME types accumulated from `offer` events. + mime_types: Vec, +} + +/// Central data control state. +/// +/// Manages the lifecycle of data control protocol objects and routes +/// events to the shared clipboard state. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) struct State { + /// The data control manager global. + pub(crate) manager: Option, + /// The data control device (per-seat). + pub(crate) device: Option, + /// The current selection offer from the compositor. + current_offer: Option, + /// The data source we created for `SetSelection` (if any). + current_source: Option, + /// MIME types advertised on `current_source`. + pub(crate) current_source_mime_types: Vec, + /// Data cached for our source's `send` events. + pub(crate) source_data: HashMap>, + /// Pending offer being built up (between `data_offer` and `selection` events). + pending_offer: Option<(DataControlOffer, PendingOffer)>, + /// Pastes waiting on `on_transfer`, by requested MIME type. + pending_transfers: HashMap, + /// Next serial handed to `on_transfer`. + next_transfer_serial: u32, + paste_count: Arc, + transfer_timeout: Duration, + /// Shared clipboard state for cross-thread access. + #[expect(clippy::struct_field_names)] + pub(crate) shared_state: Arc>, +} + +impl Default for State { + fn default() -> Self { + Self { + manager: None, + device: None, + current_offer: None, + current_source: None, + current_source_mime_types: Vec::new(), + source_data: HashMap::new(), + pending_offer: None, + pending_transfers: HashMap::new(), + next_transfer_serial: 1, + paste_count: Arc::default(), + transfer_timeout: TRANSFER_TIMEOUT, + shared_state: Arc::new(Mutex::new(Shared::default())), + } + } +} + +impl State { + /// Test-only: cache data for a MIME type as if it had been set. + #[cfg(feature = "__test")] + pub fn insert_source_data(&mut self, mime_type: &str, data: Vec) { + self.source_data.insert(mime_type.to_owned(), data.into()); + } + + /// Test-only: whether any source data is cached. + #[cfg(feature = "__test")] + pub fn has_source_data(&self) -> bool { + !self.source_data.is_empty() + } + + /// Test-only: pretend our source advertised these MIME types. + #[cfg(feature = "__test")] + pub fn set_advertised(&mut self, mime_types: Vec) { + self.current_source_mime_types = mime_types; + } + + /// Test-only: the shared state the callbacks and readers see. + #[cfg(feature = "__test")] + pub fn shared(&self) -> &Arc> { + &self.shared_state + } + + /// Test-only: number of held descriptors and active writer threads. + #[cfg(feature = "__test")] + pub fn paste_count(&self) -> usize { + self.paste_count.load(Ordering::Relaxed) + } + + /// Test-only: shorten transfer deadlines without sleeping for five seconds. + #[cfg(feature = "__test")] + pub fn set_transfer_timeout(&mut self, timeout: Duration) { + self.transfer_timeout = timeout; + } + + /// Ignore queued events from a source that a newer copy already replaced. + pub(crate) fn is_current_source(&self, id: &ObjectId) -> bool { + match &self.current_source { + Some(DataControlSource::Ext(source)) => source.id() == *id, + Some(DataControlSource::Wlr(source)) => source.id() == *id, + None => false, + } + } + + /// Release abandoned delayed pastes even when the callback never replies. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn expire_transfers(&mut self, now: Instant) { + self.pending_transfers.retain(|_, pending| { + pending.waiters.retain(|waiter| now < waiter.deadline); + !pending.waiters.is_empty() + }); + } + + /// Create a data control device from the manager and seat. + /// + /// Must be called after both the manager and seat are bound. + pub(crate) fn create_device(&mut self, seat: &WlSeat, qh: &QueueHandle) + where + D: Dispatch + Dispatch + 'static, + { + let device = match &self.manager { + Some(DataControlManager::Ext(mgr)) => DataControlDevice::Ext(mgr.get_data_device(seat, qh, ())), + Some(DataControlManager::Wlr(mgr)) => DataControlDevice::Wlr(mgr.get_data_device(seat, qh, ())), + None => { + tracing::error!("Cannot create data control device: manager not bound"); + return; + } + }; + + tracing::debug!("Created data control device"); + self.device = Some(device); + } + + /// Handle a `data_offer` event from the device. + /// + /// A new offer is being introduced. Store it and start collecting + /// MIME types from subsequent `offer` events. + pub(crate) fn on_data_offer_ext(&mut self, offer: ExtDataControlOfferV1) { + self.set_pending_offer(DataControlOffer::Ext(offer)); + } + + /// Handle a `data_offer` event from the device (wlr variant). + pub(crate) fn on_data_offer_wlr(&mut self, offer: ZwlrDataControlOfferV1) { + self.set_pending_offer(DataControlOffer::Wlr(offer)); + } + + fn set_own_source_live(&self, live: bool) { + if let Ok(mut shared) = self.shared_state.lock() { + shared.own_source_live = live; + } + } + + fn set_pending_offer(&mut self, offer: DataControlOffer) { + // Destroy any previous pending offer that wasn't used + if let Some((old_offer, _)) = self.pending_offer.take() { + old_offer.destroy(); + } + self.pending_offer = Some((offer, PendingOffer::default())); + } + + /// Handle an `offer` event on a data offer (MIME type offered). + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_offer_mime_type(&mut self, mime_type: String) { + if let Some((_, ref mut pending)) = self.pending_offer { + pending.mime_types.push(mime_type); + } + } + + /// Handle the `selection` event from the device. + /// + /// The compositor's selection has changed. The pending offer + /// (with accumulated MIME types) becomes the current offer. + pub(crate) fn on_selection(&mut self) { + // Destroy the old current offer + if let Some(old) = self.current_offer.take() { + old.destroy(); + } + + // Promote the pending offer to current + let mime_types = if let Some((offer, pending)) = self.pending_offer.take() { + let types = pending.mime_types; + self.current_offer = Some(offer); + types + } else { + // NULL selection (clipboard cleared) + Vec::new() + }; + + // The device reports every selection, our own included. Reporting + // our own back as a change makes the consumer treat it as another + // client's copy. + let own = is_own_selection( + self.current_source.is_some(), + &self.current_source_mime_types, + &mime_types, + ); + tracing::debug!( + mime_types = ?mime_types, + own, + "Compositor selection changed" + ); + + self.publish_selection(mime_types, own); + } + + /// Publish a selection after its ownership and offered formats are known. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn publish_selection(&self, mime_types: Vec, own: bool) { + // The callback runs after the lock is released: it may call back into + // the handle (`serial`, `selection_mime_types`), which locks the same state. + let callback = self.shared_state.lock().ok().and_then(|mut shared| { + shared.serial = shared.serial.wrapping_add(1); + shared.own_source_live = own; + shared.mime_types.clone_from(&mime_types); + if own { None } else { shared.on_change.clone() } + }); + if let Some(callback) = callback { + callback(mime_types); + } + } + + /// Handle the `selection` event with a NULL offer (selection cleared). + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_selection_cleared(&mut self) { + if let Some(old) = self.current_offer.take() { + old.destroy(); + } + // Clear any pending offer too + if let Some((offer, _)) = self.pending_offer.take() { + offer.destroy(); + } + + tracing::debug!("Compositor selection cleared"); + + let callback = self.shared_state.lock().ok().and_then(|mut shared| { + shared.serial = shared.serial.wrapping_add(1); + shared.own_source_live = false; + shared.mime_types.clear(); + shared.on_change.clone() + }); + if let Some(callback) = callback { + callback(Vec::new()); + } + } + + /// Update cached source data for a MIME type. + /// + /// Inserts or replaces data in the `source_data` map without + /// re-creating the Wayland data source. Used when data arrives + /// after `set_selection` was called with an empty data map + /// (eager fetch from a remote clipboard). + pub(crate) fn update_source_data(&mut self, mime_type: String, data: Vec) { + tracing::debug!( + mime_type = %mime_type, + bytes = data.len(), + "Source data updated (post-announcement)" + ); + self.source_data.insert(mime_type, data.into()); + } + + /// Handle a `send` event on our data source. + /// + /// The compositor (or another client pasting) wants data in the + /// specified MIME type. Cached data is written at once; otherwise, with + /// an `on_transfer` callback set, the paste is held until the data + /// arrives through `CompleteTransfer`. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_source_send(&mut self, mime_type: &str, fd: OwnedFd) { + self.expire_transfers(Instant::now()); + if self + .paste_count + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |count| { + (count < MAX_PASTES).then_some(count + 1) + }) + .is_err() + { + tracing::debug!("Too many clipboard pastes; closing the new request"); + return; + } + let writer = PasteWriter { + fd, + deadline: Instant::now() + self.transfer_timeout, + _permit: PastePermit(Arc::clone(&self.paste_count)), + }; + if let Some(data) = self.cached_data(mime_type) { + write_in_background(writer, data); + return; + } + + if let Some(pending) = self.pending_transfers.get_mut(mime_type) { + pending.waiters.push(writer); + return; + } + + let on_transfer = self + .shared_state + .lock() + .ok() + .and_then(|shared| shared.on_transfer.clone()); + let advertised = self.current_source_mime_types.iter().any(|m| m == mime_type); + let (Some(on_transfer), true) = (on_transfer, advertised) else { + tracing::warn!(mime_type, "Source send event for unknown MIME type"); + // fd is dropped/closed here, signaling no data + return; + }; + + let serial = self.next_transfer_serial; + self.next_transfer_serial = self.next_transfer_serial.wrapping_add(1); + self.pending_transfers.insert( + mime_type.to_owned(), + PendingTransfer { + serial, + waiters: vec![writer], + }, + ); + tracing::debug!(mime_type, serial, "Paste held until its data arrives"); + on_transfer(serial, mime_type.to_owned()); + } + + /// Data cached for `mime_type`, tolerating a charset parameter + /// (compositors commonly request `text/plain;charset=utf-8` for + /// `text/plain`). + fn cached_data(&self, mime_type: &str) -> Option> { + self.source_data + .get(mime_type) + .or_else(|| { + let base = mime_type.split(';').next()?.trim(); + self.source_data.get(base) + }) + .cloned() + } + + /// Process a `CompleteTransfer` command. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn complete_transfer(&mut self, serial: u32, data: Option>) { + self.expire_transfers(Instant::now()); + let Some(mime_type) = self + .pending_transfers + .iter() + .find(|(_, pending)| pending.serial == serial) + .map(|(mime, _)| mime.clone()) + else { + tracing::debug!(serial, "Transfer answered after its selection was replaced"); + return; + }; + let Some(pending) = self.pending_transfers.remove(&mime_type) else { + return; + }; + let Some(data) = data.filter(|d| !d.is_empty()) else { + tracing::debug!(mime_type, waiters = pending.waiters.len(), "Transfer returned no data"); + return; + }; + let shared: Arc<[u8]> = data.into(); + for writer in pending.waiters { + write_in_background(writer, Arc::clone(&shared)); + } + self.source_data.insert(mime_type, shared); + } + + /// Process a `ClearSelection` command: give up our selection if the + /// compositor still has it. + pub(crate) fn clear_selection(&mut self) { + let Some(source) = self.current_source.take() else { + return; + }; + match &self.device { + Some(DataControlDevice::Ext(dev)) => dev.set_selection(None), + Some(DataControlDevice::Wlr(dev)) => dev.set_selection(None), + None => {} + } + source.destroy(); + self.current_source_mime_types.clear(); + self.source_data.clear(); + self.pending_transfers.clear(); + self.set_own_source_live(false); + tracing::debug!("Released our clipboard selection"); + } + + /// Handle the `cancelled` event on our data source. + /// + /// Our source has been replaced by another. Clean up. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_source_cancelled(&mut self) { + tracing::debug!("Data control source cancelled"); + if let Some(source) = self.current_source.take() { + source.destroy(); + } + self.current_source_mime_types.clear(); + self.source_data.clear(); + // Waiting pastes see end-of-file. + self.pending_transfers.clear(); + // A same-format foreign selection can arrive before this cancellation + // and be mistaken for our own echo. In that ordering its callback was + // suppressed, so ownership loss must publish the current MIME snapshot. + let changed = self.shared_state.lock().ok().and_then(|mut shared| { + let was_own = shared.own_source_live; + shared.own_source_live = false; + if was_own { + shared.serial = shared.serial.wrapping_add(1); + shared + .on_change + .clone() + .map(|callback| (callback, shared.mime_types.clone())) + } else { + None + } + }); + if let Some((callback, mime_types)) = changed { + callback(mime_types); + } + } + + /// Handle the `finished` event on the device. + /// + /// The data control device is no longer valid. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_device_finished(&mut self) { + tracing::debug!("Data control device finished"); + self.device = None; + + if let Some(offer) = self.current_offer.take() { + offer.destroy(); + } + if let Some((offer, _)) = self.pending_offer.take() { + offer.destroy(); + } + if let Some(source) = self.current_source.take() { + source.destroy(); + } + self.current_source_mime_types.clear(); + self.source_data.clear(); + self.pending_transfers.clear(); + self.set_own_source_live(false); + } + + /// Process a `SetSelection` command. + /// + /// Creates a new data control source, advertises MIME types, and + /// sets it as the selection on the device. + pub(crate) fn set_selection( + &mut self, + mime_types: &[String], + data: HashMap>, + qh: &QueueHandle, + ) where + D: Dispatch + Dispatch + 'static, + { + let new_source = match &self.manager { + Some(DataControlManager::Ext(mgr)) => DataControlSource::Ext(mgr.create_data_source(qh, ())), + Some(DataControlManager::Wlr(mgr)) => DataControlSource::Wlr(mgr.create_data_source(qh, ())), + None => { + tracing::error!("Cannot set selection: manager not bound"); + return; + } + }; + + // Advertise MIME types + for mime_type in mime_types { + new_source.offer(mime_type); + } + + // Set as selection on the device + match (&self.device, &new_source) { + (Some(DataControlDevice::Ext(dev)), DataControlSource::Ext(src)) => { + dev.set_selection(Some(src)); + } + (Some(DataControlDevice::Wlr(dev)), DataControlSource::Wlr(src)) => { + dev.set_selection(Some(src)); + } + _ => { + tracing::error!("Cannot set selection: device/source protocol mismatch or no device"); + new_source.destroy(); + return; + } + } + + tracing::debug!( + mime_types = ?mime_types, + "Set clipboard selection on compositor" + ); + + // Replace, then destroy: destroying the current source first would + // leave the clipboard empty for a moment, which clipboard managers + // (Klipper's "prevent empty clipboard") answer by re-offering their + // last history item. + if let Some(old) = self.current_source.take() { + old.destroy(); + } + self.pending_transfers.clear(); + self.source_data = data.into_iter().map(|(mime, bytes)| (mime, bytes.into())).collect(); + self.current_source = Some(new_source); + self.current_source_mime_types = mime_types.to_vec(); + self.set_own_source_live(true); + } + + /// Process a `ReceiveFromOffer` command. + /// + /// Calls `offer.receive(mime_type, fd)` to request data from the + /// compositor. The caller reads from the other end of the pipe. + #[expect( + clippy::needless_pass_by_value, + reason = "OwnedFd must be owned so it is dropped after the receive call" + )] + pub(crate) fn receive_from_offer(&self, mime_type: &str, fd: OwnedFd) { + match &self.current_offer { + Some(offer) => { + offer.receive(mime_type, &fd); + tracing::debug!(mime_type, "Requested clipboard data from compositor offer"); + } + None => { + tracing::warn!(mime_type, "ReceiveFromOffer but no current offer"); + // fd is dropped, closing the pipe — caller will get EOF + } + } + } +} + +/// Whether a selection event reports the selection this client set itself. +/// +/// Another client taking the selection cancels our source first, so while +/// our source is live the selection is still ours. The offered MIME set must +/// also match what we advertised: that keeps a foreign copy from being +/// swallowed on a compositor that delivers the selection before the cancel. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) fn is_own_selection(source_live: bool, advertised: &[String], offered: &[String]) -> bool { + if !source_live || advertised.len() != offered.len() { + return false; + } + offered.iter().all(|mime| advertised.contains(mime)) +} + +/// Answer a paste on a worker thread: a large image written into a slow +/// reader would otherwise block every other Wayland event. +fn write_in_background(writer: PasteWriter, data: Arc<[u8]>) { + let spawned = std::thread::Builder::new() + .name("data-control-send".into()) + .spawn(move || { + // Keep the permit until the descriptor closes, including failed writes. + let _permit = writer._permit; + if let Err(error) = write_with_deadline(&writer.fd, &data, writer.deadline) { + tracing::debug!(%error, "Could not finish a clipboard paste"); + } + drop(writer.fd); + }); + if let Err(error) = spawned { + tracing::error!(%error, "Failed to start a paste writer"); + } +} + +fn write_with_deadline(fd: &OwnedFd, mut data: &[u8], deadline: Instant) -> std::io::Result<()> { + let flags = OFlag::from_bits_truncate(fcntl(fd, FcntlArg::F_GETFL)?); + fcntl(fd, FcntlArg::F_SETFL(flags | OFlag::O_NONBLOCK))?; + while !data.is_empty() { + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Err(std::io::ErrorKind::TimedOut.into()); + } + let mut fds = [PollFd::new(fd.as_fd(), PollFlags::POLLOUT)]; + let timeout = PollTimeout::try_from(remaining).unwrap_or(PollTimeout::MAX); + match poll(&mut fds, timeout) { + Ok(0) | Err(nix::errno::Errno::EINTR) => continue, + Ok(_) => {} + Err(error) => return Err(error.into()), + } + match nix::unistd::write(fd, data) { + Ok(0) => return Err(std::io::ErrorKind::WriteZero.into()), + Ok(count) => data = &data[count..], + Err(nix::errno::Errno::EINTR | nix::errno::Errno::EAGAIN) => {} + Err(error) => return Err(error.into()), + } + } + Ok(()) +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/worker.rs b/crates/ironrdp-cliprdr-native/src/data_control/worker.rs new file mode 100644 index 0000000000..c86a2cd5e9 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/worker.rs @@ -0,0 +1,271 @@ +//! The worker thread that owns the Wayland connection. +//! +//! Wayland objects live on this one thread. Callers talk to it through a +//! command channel plus a wake socket, so a request is served immediately +//! instead of at the next poll timeout. + +use core::sync::atomic::{AtomicBool, Ordering}; +use std::{ + io::Read as _, + os::unix::{io::AsFd as _, net::UnixStream}, + sync::{Arc, Mutex, mpsc}, + time::Instant, +}; + +use nix::{ + errno::Errno, + poll::{PollFd, PollFlags, PollTimeout, poll}, +}; +use wayland_client::{ + Connection, Dispatch, EventQueue, QueueHandle, + backend::WaylandError, + globals::{GlobalList, registry_queue_init}, + protocol::{wl_callback::WlCallback, wl_seat::WlSeat}, +}; +use wayland_protocols::ext::data_control::v1::client::ext_data_control_manager_v1::ExtDataControlManagerV1; +use wayland_protocols_wlr::data_control::v1::client::zwlr_data_control_manager_v1::ZwlrDataControlManagerV1; + +use super::{ + dispatch::Client, + error::{Error, Result}, + options::{Options, Protocol}, + state::{Command, DataControlManager, Shared, State}, +}; + +/// A connected worker, ready to run. +pub(crate) struct Worker { + connection: Connection, + queue: EventQueue, + client: Client, + /// Kept alive for the device's lifetime. + _seat: WlSeat, + commands: mpsc::Receiver, + wake: UnixStream, + stop: Arc, +} + +/// What the caller learns once the worker is connected. +pub(crate) struct Connected { + pub(crate) protocol: Protocol, + pub(crate) shared: Arc>, +} + +/// Connect to the compositor, bind a data-control manager and the seat, and +/// create the device. +/// +/// Runs on the worker thread, so any Wayland object created here stays there. +pub(crate) fn connect( + options: &Options, + commands: mpsc::Receiver, + wake: UnixStream, + stop: Arc, +) -> Result<(Worker, Connected)> { + let connection = Connection::connect_to_env().map_err(|error| Error::Connect(error.to_string()))?; + let (globals, mut queue) = + registry_queue_init::(&connection).map_err(|error| Error::Wayland(error.to_string()))?; + let queue_handle = queue.handle(); + + let mut client = Client { + data_control: State::default(), + }; + + let (manager, protocol) = bind_manager(&globals, &queue_handle, options)?; + client.data_control.manager = Some(manager); + + let seat = globals + .bind::(&queue_handle, 1..=1, ()) + .map_err(|_| Error::NoSeat)?; + client.data_control.create_device(&seat, &queue_handle); + + // The compositor sends the current selection right after the device is + // created; wait for it so `selection_mime_types` is right from the start. + queue + .roundtrip(&mut client) + .map_err(|error| Error::Wayland(error.to_string()))?; + + let shared = Arc::clone(&client.data_control.shared_state); + tracing::info!(%protocol, "Data-control clipboard connected"); + + Ok(( + Worker { + connection, + queue, + client, + _seat: seat, + commands, + wake, + stop, + }, + Connected { protocol, shared }, + )) +} + +/// Bind the first protocol in the option's candidate order that the +/// compositor offers. +fn bind_manager( + globals: &GlobalList, + queue_handle: &QueueHandle, + options: &Options, +) -> Result<(DataControlManager, Protocol)> { + let candidates = options.candidates(); + for (index, protocol) in candidates.iter().enumerate() { + let bound = match protocol { + Protocol::Ext => globals + .bind::(queue_handle, 1..=1, ()) + .ok() + .map(DataControlManager::Ext), + Protocol::Wlr => globals + .bind::(queue_handle, 1..=2, ()) + .ok() + .map(DataControlManager::Wlr), + }; + if let Some(manager) = bound { + if index > 0 { + tracing::info!( + preferred = %candidates[0], + used = %protocol, + "Preferred data-control protocol unavailable, using the other" + ); + } + return Ok((manager, *protocol)); + } + } + + if options.allow_fallback { + Err(Error::Unsupported) + } else { + Err(Error::ProtocolUnavailable) + } +} + +impl Worker { + /// Serve events and commands until asked to stop or the connection fails. + pub(crate) fn run(mut self) { + tracing::debug!("Data-control worker started"); + + while !self.stop.load(Ordering::Relaxed) { + if let Err(error) = self.queue.dispatch_pending(&mut self.client) { + tracing::error!(%error, "Wayland dispatch failed"); + break; + } + + self.process_commands(); + self.client.data_control.expire_transfers(Instant::now()); + + let flush_blocked = match self.connection.flush() { + Ok(()) => false, + Err(WaylandError::Io(ref error)) if error.kind() == std::io::ErrorKind::WouldBlock => true, + Err(error) => { + tracing::error!(%error, "Wayland flush failed"); + break; + } + }; + + // `None` means events arrived after the dispatch above; loop to + // handle them before waiting. + let Some(guard) = self.queue.prepare_read() else { + continue; + }; + + let wayland_events = if flush_blocked { + PollFlags::POLLIN | PollFlags::POLLOUT + } else { + PollFlags::POLLIN + }; + let mut fds = [ + PollFd::new(guard.connection_fd(), wayland_events), + PollFd::new(self.wake.as_fd(), PollFlags::POLLIN), + ]; + match poll(&mut fds, PollTimeout::from(100u16)) { + Ok(_) => {} + Err(Errno::EINTR) => continue, + Err(error) => { + tracing::error!(%error, "Wayland poll failed"); + break; + } + } + let wayland_readable = readable(&fds[0]); + let woken = readable(&fds[1]); + + if wayland_readable { + match guard.read() { + Ok(_) => {} + // Another reader took the data between poll and read. + Err(WaylandError::Io(ref error)) if error.kind() == std::io::ErrorKind::WouldBlock => {} + Err(error) => { + tracing::error!(%error, "Wayland read failed"); + break; + } + } + } else { + drop(guard); + } + + if woken { + self.drain_wake(); + } + } + + tracing::debug!("Data-control worker stopped"); + } + + fn process_commands(&mut self) { + let queue_handle = self.queue.handle(); + while let Ok(command) = self.commands.try_recv() { + match command { + Command::Synchronize(sender) => { + // The callback is dispatched after the compositor has + // processed prior selection requests and sent their events. + // Keep servicing the event loop while waiting for it. + self.connection.display().sync(&queue_handle, sender); + } + Command::SetSelection { mime_types, data } => { + self.client.data_control.set_selection(&mime_types, data, &queue_handle); + } + Command::UpdateSourceData { mime_type, data } => { + self.client.data_control.update_source_data(mime_type, data); + } + Command::ReceiveFromOffer { mime_type, fd } => { + self.client.data_control.receive_from_offer(&mime_type, fd); + } + Command::CompleteTransfer { serial, data } => { + self.client.data_control.complete_transfer(serial, data); + } + Command::ClearSelection => { + self.client.data_control.clear_selection(); + } + } + } + } + + /// Empty the wake socket so the next poll blocks again. + fn drain_wake(&mut self) { + let mut buffer = [0u8; 64]; + loop { + match self.wake.read(&mut buffer) { + Ok(0) => break, + Ok(_) => {} + Err(error) if error.kind() == std::io::ErrorKind::Interrupted => {} + Err(_) => break, + } + } + } +} + +impl Dispatch> for Client { + fn event( + _state: &mut Self, + _proxy: &WlCallback, + _event: ::Event, + data: &mpsc::Sender<()>, + _connection: &Connection, + _queue_handle: &QueueHandle, + ) { + let _ = data.send(()); + } +} + +fn readable(fd: &PollFd<'_>) -> bool { + fd.revents() + .is_some_and(|events| events.intersects(PollFlags::POLLIN | PollFlags::POLLHUP | PollFlags::POLLERR)) +} diff --git a/crates/ironrdp-cliprdr-native/src/lib.rs b/crates/ironrdp-cliprdr-native/src/lib.rs index b278ce487f..331cf3c38a 100644 --- a/crates/ironrdp-cliprdr-native/src/lib.rs +++ b/crates/ironrdp-cliprdr-native/src/lib.rs @@ -15,6 +15,19 @@ mod windows; #[cfg(windows)] pub use crate::windows::{HWND, WinClipboard, WinCliprdrError, WinCliprdrResult}; +#[cfg(target_os = "linux")] +pub mod data_control; +#[cfg(all(target_os = "linux", feature = "__test"))] +pub mod linux; +#[cfg(all(target_os = "linux", not(feature = "__test")))] +#[expect( + unreachable_pub, + reason = "the __test feature exposes internals to the shared test suite" +)] +pub(crate) mod linux; +#[cfg(target_os = "linux")] +pub use crate::linux::{LinuxClipboard, LinuxCliprdrBackend, LinuxCliprdrError}; + mod stub; use std::sync::OnceLock; use std::time::Instant; diff --git a/crates/ironrdp-cliprdr-native/src/linux/mod.rs b/crates/ironrdp-cliprdr-native/src/linux/mod.rs new file mode 100644 index 0000000000..5ce75ecca6 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/linux/mod.rs @@ -0,0 +1,206 @@ +//! Linux CLIPRDR backend with clipboard change events and delayed rendering. +//! +//! Wayland uses the shared [`crate::data_control`] client. X11/XWayland uses +//! XFixes selection notifications and selection ownership. Both advertise +//! remote formats immediately and request data only when an application pastes. +//! Plain text (`CF_UNICODETEXT`) and PNG images (`CF_DIB`/`CF_DIBV5`) travel in +//! both directions; file clipboard transfer and HTML are not supported. + +pub mod os; +pub mod worker; +pub mod x11; + +use core::fmt; +use std::sync::mpsc::{self, Sender}; +use std::thread::JoinHandle; + +use ironrdp_cliprdr::backend::{ClipboardMessageProxy, CliprdrBackend, CliprdrBackendFactory}; +use ironrdp_cliprdr::pdu::{ + ClipboardFormat, ClipboardGeneralCapabilityFlags, FileContentsRequest, FileContentsResponse, FormatDataRequest, + FormatDataResponse, LockDataId, +}; +use ironrdp_core::impl_as_any; +use tracing::{debug, warn}; + +use self::worker::{Command, Worker}; + +#[derive(Debug)] +pub enum LinuxCliprdrError { + /// No usable clipboard: neither a Wayland data-control clipboard nor an X11 + /// display could be opened. + Clipboard(String), + Spawn(std::io::Error), + /// The clipboard thread stopped before reporting whether it could start. + WorkerStopped, +} + +impl fmt::Display for LinuxCliprdrError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Clipboard(error) => write!(f, "cannot open the OS clipboard: {error}"), + Self::Spawn(error) => write!(f, "cannot start the clipboard thread: {error}"), + Self::WorkerStopped => f.write_str("the clipboard thread stopped unexpectedly"), + } + } +} + +impl core::error::Error for LinuxCliprdrError { + fn source(&self) -> Option<&(dyn core::error::Error + 'static)> { + match self { + Self::Clipboard(_) => None, + Self::Spawn(error) => Some(error), + Self::WorkerStopped => None, + } + } +} + +/// The Linux clipboard bridge. Owns the thread that talks to the OS clipboard; +/// keep it alive for as long as the connection it serves, and build one +/// [`CliprdrBackend`] per channel initialization through +/// [`backend_factory`](Self::backend_factory). +pub struct LinuxClipboard { + commands: Sender, + thread: Option>, +} + +impl fmt::Debug for LinuxClipboard { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("LinuxClipboard").finish_non_exhaustive() + } +} + +impl LinuxClipboard { + /// Opens the OS clipboard and starts the clipboard thread. + /// + /// `message_proxy` delivers the resulting `CLIPRDR` messages to the RDP + /// session; it is called from the clipboard thread. + pub fn new(message_proxy: impl ClipboardMessageProxy + 'static) -> Result { + let (commands, receiver) = mpsc::channel(); + let (ready_sender, ready) = mpsc::channel(); + + let events = commands.clone(); + let thread = std::thread::Builder::new() + .name("cliprdr-linux".to_owned()) + .spawn(move || { + let os = match os::NativeClipboard::open(events) { + Ok(os) => os, + Err(error) => { + let _ = ready_sender.send(Err(error)); + return; + } + }; + let _ = ready_sender.send(Ok(())); + Worker::new(os, message_proxy).run(&receiver); + debug!("Clipboard thread stopped"); + }) + .map_err(LinuxCliprdrError::Spawn)?; + + match ready.recv() { + Ok(Ok(())) => Ok(Self { + commands, + thread: Some(thread), + }), + Ok(Err(error)) => Err(LinuxCliprdrError::Clipboard(error)), + Err(_) => Err(LinuxCliprdrError::WorkerStopped), + } + } + + pub fn backend_factory(&self) -> Box { + Box::new(LinuxCliprdrBackendFactory { + commands: self.commands.clone(), + }) + } +} + +impl Drop for LinuxClipboard { + fn drop(&mut self) { + let _ = self.commands.send(Command::Shutdown); + if let Some(thread) = self.thread.take() { + let _ = thread.join(); + } + } +} + +struct LinuxCliprdrBackendFactory { + commands: Sender, +} + +impl CliprdrBackendFactory for LinuxCliprdrBackendFactory { + fn build_cliprdr_backend(&self) -> Box { + let _ = self.commands.send(Command::Reset); + Box::new(LinuxCliprdrBackend { + commands: self.commands.clone(), + }) + } +} + +/// Forwards the `CLIPRDR` callbacks to the clipboard thread. +#[derive(Debug)] +pub struct LinuxCliprdrBackend { + commands: Sender, +} + +impl_as_any!(LinuxCliprdrBackend); + +impl LinuxCliprdrBackend { + fn send(&self, command: Command) { + if self.commands.send(command).is_err() { + warn!("The clipboard thread is gone; dropping a CLIPRDR event"); + } + } +} + +impl CliprdrBackend for LinuxCliprdrBackend { + fn temporary_directory(&self) -> &str { + ".cliprdr" + } + + fn client_capabilities(&self) -> ClipboardGeneralCapabilityFlags { + // No file transfer: no STREAM_FILECLIP_ENABLED, no CAN_LOCK_CLIPDATA. + ClipboardGeneralCapabilityFlags::empty() + } + + fn on_ready(&mut self) {} + + fn on_request_format_list(&mut self) { + self.send(Command::AdvertiseLocal); + } + + fn on_process_negotiated_capabilities(&mut self, capabilities: ClipboardGeneralCapabilityFlags) { + debug!(?capabilities, "CLIPRDR capabilities negotiated"); + } + + fn on_remote_copy(&mut self, available_formats: &[ClipboardFormat]) { + self.send(Command::RemoteCopy(available_formats.to_vec())); + } + + fn on_format_data_request(&mut self, request: FormatDataRequest) { + self.send(Command::RenderLocal(request.format)); + } + + fn on_format_data_response(&mut self, response: FormatDataResponse<'_>) { + let data = (!response.is_error() && response.data().len() <= worker::MAX_TRANSFER_BYTES) + .then(|| response.data().to_vec()); + self.send(Command::RemoteData(data)); + } + + fn on_file_contents_request(&mut self, request: FileContentsRequest) { + debug!(?request, "File transfer is not supported"); + } + + fn on_file_contents_response(&mut self, response: FileContentsResponse<'_>) { + debug!(?response, "File transfer is not supported"); + } + + fn on_lock(&mut self, data_id: LockDataId) { + debug!(?data_id); + } + + fn on_unlock(&mut self, data_id: LockDataId) { + debug!(?data_id); + } + + fn now_ms(&self) -> u64 { + crate::native_now_ms() + } +} diff --git a/crates/ironrdp-cliprdr-native/src/linux/os.rs b/crates/ironrdp-cliprdr-native/src/linux/os.rs new file mode 100644 index 0000000000..567df10bfd --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/linux/os.rs @@ -0,0 +1,163 @@ +//! The clipboard interface used by the protocol worker. + +use std::sync::mpsc::Sender; + +use super::worker::Command; +use super::x11::X11Clipboard; +use crate::data_control::{Content, DataControl}; + +pub const TEXT: &str = "text/plain;charset=utf-8"; +pub const PNG: &str = "image/png"; + +type CompletePaste = Box>) + Send>; + +/// A delayed OS paste. Every completion, including failure, closes its waiter. +pub struct Transfer { + mime: String, + generation: u64, + complete: Option, +} + +impl core::fmt::Debug for Transfer { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.debug_struct("Transfer") + .field("mime", &self.mime) + .field("generation", &self.generation) + .finish_non_exhaustive() + } +} + +impl Transfer { + pub fn new(mime: String, generation: u64, complete: impl FnOnce(Option>) + Send + 'static) -> Self { + Self { + mime, + generation, + complete: Some(Box::new(complete)), + } + } + + pub fn mime_type(&self) -> &str { + &self.mime + } + + pub fn generation(&self) -> u64 { + self.generation + } + + pub fn finish(mut self, data: Option>) { + if let Some(complete) = self.complete.take() { + complete(data); + } + } +} + +impl Drop for Transfer { + fn drop(&mut self) { + if let Some(complete) = self.complete.take() { + complete(None); + } + } +} + +pub trait OsClipboard { + fn is_owner(&self) -> bool; + fn mime_types(&self) -> Vec; + fn read(&mut self, mime: &str) -> Result>, String>; + fn offer(&mut self, mimes: &[String], generation: u64) -> Result<(), String>; + fn clear(&mut self) -> Result<(), String>; +} + +pub enum NativeClipboard { + Wayland { + clipboard: DataControl, + events: Sender, + }, + X11(X11Clipboard), +} + +impl NativeClipboard { + pub fn open(events: Sender) -> Result { + match DataControl::connect() { + Ok(clipboard) => { + let changes = events.clone(); + clipboard.on_change(move |_| { + let _ = changes.send(Command::LocalChanged); + }); + Ok(Self::Wayland { clipboard, events }) + } + Err(wayland_error) => { + tracing::debug!(%wayland_error, "Wayland data-control unavailable; trying X11"); + X11Clipboard::open(events) + .map(Self::X11) + .map_err(|x11_error| format!("Wayland: {wayland_error}; X11: {x11_error}")) + } + } + } +} + +impl OsClipboard for NativeClipboard { + fn is_owner(&self) -> bool { + match self { + Self::Wayland { clipboard, .. } => clipboard.owns_selection(), + Self::X11(clipboard) => clipboard.is_owner(), + } + } + + fn mime_types(&self) -> Vec { + match self { + Self::Wayland { clipboard, .. } => clipboard.selection_mime_types(), + Self::X11(clipboard) => clipboard.mime_types(), + } + } + + fn read(&mut self, mime: &str) -> Result>, String> { + match self { + Self::Wayland { clipboard, .. } => { + let serial = clipboard.serial(); + let data = clipboard.read(mime).map_err(|error| error.to_string())?; + if serial != clipboard.serial() { + return Err("clipboard selection changed while reading".into()); + } + Ok(data) + } + Self::X11(clipboard) => clipboard.read(mime), + } + } + + fn offer(&mut self, mimes: &[String], generation: u64) -> Result<(), String> { + match self { + Self::Wayland { clipboard, events } => { + let events = events.clone(); + clipboard.on_transfer(move |request| { + let transfer = Transfer::new(request.mime_type().to_owned(), generation, move |data| { + let result = match data { + Some(data) => request.complete(data), + None => request.fail(), + }; + if let Err(error) = result { + tracing::debug!(%error, "Could not finish a Wayland paste"); + } + }); + let _ = events.send(Command::Paste(transfer)); + }); + let content = mimes + .iter() + .fold(Content::new(), |content, mime| content.advertise(mime)); + clipboard.set_selection(content).map_err(|error| error.to_string())?; + clipboard.synchronize().map_err(|error| error.to_string())?; + if !clipboard.owns_selection() { + return Err("clipboard ownership was not acquired".into()); + } + Ok(()) + } + Self::X11(clipboard) => clipboard.offer(mimes, generation), + } + } + + fn clear(&mut self) -> Result<(), String> { + match self { + Self::Wayland { clipboard, .. } => clipboard.clear_selection().map_err(|error| error.to_string()), + Self::X11(clipboard) => clipboard.clear(), + } + } +} diff --git a/crates/ironrdp-cliprdr-native/src/linux/worker.rs b/crates/ironrdp-cliprdr-native/src/linux/worker.rs new file mode 100644 index 0000000000..66905f8b55 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/linux/worker.rs @@ -0,0 +1,340 @@ +//! Serializes the CLIPRDR request/response stream and OS clipboard events. + +use core::time::Duration; +use std::collections::{HashMap, VecDeque}; +use std::sync::mpsc::{Receiver, RecvTimeoutError}; +use std::time::Instant; + +use ironrdp_cliprdr::backend::{ClipboardMessage, ClipboardMessageProxy}; +use ironrdp_cliprdr::loop_detector::{ClipboardSource, LoopDetector}; +use ironrdp_cliprdr::pdu::{ClipboardFormat, ClipboardFormatId, FormatDataResponse}; +use ironrdp_cliprdr_format::bitmap; +use ironrdp_core::IntoOwned as _; +use tracing::{debug, warn}; + +use super::os::{OsClipboard, PNG, TEXT, Transfer}; +use crate::data_control::find_mime_match; + +pub const PASTE_TIMEOUT: Duration = Duration::from_secs(5); +pub const MAX_TRANSFER_BYTES: usize = 100 * 1024 * 1024; +const MAX_WAITERS: usize = 16; + +#[derive(Debug)] +pub enum Command { + AdvertiseLocal, + LocalChanged, + RemoteCopy(Vec), + RemoteData(Option>), + RenderLocal(ClipboardFormatId), + Paste(Transfer), + Reset, + Shutdown, +} + +struct Pending { + format: ClipboardFormatId, + generation: u64, + deadline: Instant, + expired: bool, + waiters: Vec, +} + +struct Queued { + deadline: Instant, + transfer: Transfer, +} + +pub struct Worker { + os: O, + proxy: P, + ready: bool, + local_mimes: Vec, + remote: Vec, + owns_remote: bool, + generation: u64, + pending: Option, + queue: VecDeque, + cache: HashMap>, + loops: LoopDetector, +} + +impl Worker { + pub fn new(os: O, proxy: P) -> Self { + let local_mimes = os.mime_types(); + Self { + os, + proxy, + local_mimes, + ready: false, + remote: Vec::new(), + owns_remote: false, + generation: 0, + pending: None, + queue: VecDeque::new(), + cache: HashMap::new(), + loops: LoopDetector::new(), + } + } + + pub fn run(mut self, receiver: &Receiver) { + loop { + match receiver.recv_timeout(Duration::from_millis(100)) { + Ok(Command::Shutdown) | Err(RecvTimeoutError::Disconnected) => break, + Ok(command) => self.handle(command), + Err(RecvTimeoutError::Timeout) => {} + } + self.expire(Instant::now()); + } + let _ = self.os.clear(); + } + + pub fn handle(&mut self, command: Command) { + match command { + Command::AdvertiseLocal => { + self.ready = true; + if !self.owns_remote { + self.local_mimes = self.os.mime_types(); + self.advertise_local(); + } + } + Command::LocalChanged => { + // An OS event queued before offer() acknowledged our ownership + // must not invalidate the selection that was just installed. + if self.os.is_owner() { + return; + } + self.invalidate(); + self.owns_remote = false; + self.remote.clear(); + self.local_mimes = self.os.mime_types(); + if self.ready { + self.advertise_local(); + } + } + Command::RemoteCopy(formats) => { + self.invalidate(); + self.remote = formats; + self.owns_remote = true; + let mut mimes = Vec::new(); + if self.remote.iter().any(|f| f.id() == ClipboardFormatId::CF_UNICODETEXT) { + mimes.extend([TEXT.to_owned(), "text/plain".to_owned()]); + } + if self + .remote + .iter() + .any(|f| matches!(f.id(), ClipboardFormatId::CF_DIB | ClipboardFormatId::CF_DIBV5)) + { + mimes.push(PNG.to_owned()); + } + if let Err(error) = self.os.offer(&mimes, self.generation) { + warn!(%error, "Could not advertise the remote clipboard"); + let _ = self.os.clear(); + self.owns_remote = false; + self.remote.clear(); + self.local_mimes = self.os.mime_types(); + if self.ready { + self.advertise_local(); + } + } + } + Command::Paste(transfer) => self.paste(transfer), + Command::RemoteData(data) => self.remote_data(data), + Command::RenderLocal(format) => self.render_local(format), + Command::Reset => { + self.invalidate(); + self.pending = None; + self.ready = false; + self.remote.clear(); + if self.owns_remote { + let _ = self.os.clear(); + } + self.owns_remote = false; + self.loops.clear(); + } + Command::Shutdown => {} + } + } + + fn invalidate(&mut self) { + self.generation = self.generation.wrapping_add(1); + self.queue.clear(); + self.cache.clear(); + if let Some(pending) = &mut self.pending { + pending.waiters.clear(); + } + } + + fn advertise_local(&self) { + let mut formats = Vec::new(); + if find_mime_match(TEXT, &self.local_mimes).is_some() { + formats.push(ClipboardFormat::new(ClipboardFormatId::CF_UNICODETEXT)); + } + if find_mime_match(PNG, &self.local_mimes).is_some() { + formats.extend([ + ClipboardFormat::new(ClipboardFormatId::CF_DIB), + ClipboardFormat::new(ClipboardFormatId::CF_DIBV5), + ]); + } + self.proxy + .send_clipboard_message(ClipboardMessage::SendInitiateCopy(formats)); + } + + fn paste(&mut self, transfer: Transfer) { + if !self.owns_remote || transfer.generation() != self.generation { + return; + } + let has = |id| self.remote.iter().any(|f| f.id() == id); + let format = + if transfer.mime_type().split(';').next() == Some("text/plain") && has(ClipboardFormatId::CF_UNICODETEXT) { + ClipboardFormatId::CF_UNICODETEXT + } else if transfer.mime_type() == PNG && has(ClipboardFormatId::CF_DIBV5) { + ClipboardFormatId::CF_DIBV5 + } else if transfer.mime_type() == PNG && has(ClipboardFormatId::CF_DIB) { + ClipboardFormatId::CF_DIB + } else { + return; + }; + if let Some(data) = self.cache.get(&format) { + transfer.finish(Some(data.clone())); + return; + } + if let Some(pending) = &mut self.pending { + if pending.expired || MAX_WAITERS <= pending.waiters.len() + self.queue.len() { + return; + } + if pending.generation == self.generation && pending.format == format { + pending.waiters.push(transfer); + } else { + self.queue.push_back(Queued { + deadline: Instant::now() + PASTE_TIMEOUT, + transfer, + }); + } + } else { + self.start(format, transfer); + } + } + + fn start(&mut self, format: ClipboardFormatId, transfer: Transfer) { + self.pending = Some(Pending { + format, + generation: self.generation, + deadline: Instant::now() + PASTE_TIMEOUT, + expired: false, + waiters: vec![transfer], + }); + self.proxy + .send_clipboard_message(ClipboardMessage::SendInitiatePaste(format)); + } + + /// Expire OS waiters, but keep the outstanding wire request until its reply arrives. + /// [MS-RDPECLIP 2.2.5.2] has no response identifier, so starting another request + /// after a timeout would misinterpret the late reply as the new clipboard content. + /// + /// [MS-RDPECLIP 2.2.5.2]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpeclip/28c193b8-4cec-413e-a07b-9235e5e15f6b + pub fn expire(&mut self, now: Instant) { + if let Some(pending) = &mut self.pending { + if pending.deadline <= now { + pending.expired = true; + pending.waiters.clear(); + } + } + self.queue.retain(|queued| now < queued.deadline); + } + + fn remote_data(&mut self, data: Option>) { + let Some(pending) = self.pending.take() else { + debug!("Discarding an unsolicited clipboard response"); + return; + }; + if !pending.expired + && Instant::now() < pending.deadline + && pending.generation == self.generation + && self.owns_remote + { + let decoded = data.filter(|data| data.len() <= MAX_TRANSFER_BYTES).and_then(|data| { + let result = match pending.format { + ClipboardFormatId::CF_UNICODETEXT => FormatDataResponse::new_data(data.as_slice()) + .to_unicode_string() + .map(String::into_bytes) + .map_err(|error| error.to_string()), + ClipboardFormatId::CF_DIB => bitmap::dib_to_png(&data).map_err(|error| error.to_string()), + ClipboardFormatId::CF_DIBV5 => bitmap::dibv5_to_png(&data).map_err(|error| error.to_string()), + _ => return None, + }; + match result { + Ok(data) => Some(data), + Err(error) => { + warn!(%error, "Could not decode the remote clipboard"); + None + } + } + }); + if let Some(data) = &decoded { + self.loops + .record_content(data, ClipboardSource::Remote, crate::native_now_ms()); + self.cache.insert(pending.format, data.clone()); + } + for transfer in pending.waiters { + transfer.finish(decoded.clone()); + } + } + while self.pending.is_none() { + let Some(queued) = self.queue.pop_front() else { + break; + }; + if Instant::now() < queued.deadline { + self.paste(queued.transfer); + } + } + } + + fn render_local(&mut self, format: ClipboardFormatId) { + let result = (|| -> Result, String> { + if self.owns_remote { + return Err("local clipboard selection was replaced".into()); + } + let mime = match format { + ClipboardFormatId::CF_UNICODETEXT => TEXT, + ClipboardFormatId::CF_DIB | ClipboardFormatId::CF_DIBV5 => PNG, + _ => return Err("unsupported clipboard format".into()), + }; + let data = self.os.read(mime)?.ok_or("clipboard format is no longer available")?; + if MAX_TRANSFER_BYTES < data.len() { + return Err("clipboard data exceeds the size limit".into()); + } + // Native APIs suppress our own selection events. This also stops a + // clipboard manager that republishes recently received content. + if self + .loops + .would_cause_content_loop(&data, ClipboardSource::Local, crate::native_now_ms()) + { + return Err("clipboard content would echo a recent remote paste".into()); + } + self.loops + .record_content(&data, ClipboardSource::Local, crate::native_now_ms()); + match format { + ClipboardFormatId::CF_UNICODETEXT => { + let text = core::str::from_utf8(&data).map_err(|error| error.to_string())?; + let units = text.encode_utf16().count().saturating_add(1); + if MAX_TRANSFER_BYTES / 2 < units { + return Err("encoded clipboard text exceeds the size limit".into()); + } + Ok(FormatDataResponse::new_unicode_string(text).into_data().into_owned()) + } + ClipboardFormatId::CF_DIB => bitmap::png_to_cf_dib(&data).map_err(|error| error.to_string()), + ClipboardFormatId::CF_DIBV5 => bitmap::png_to_cf_dibv5(&data).map_err(|error| error.to_string()), + _ => Err("unsupported clipboard format".into()), + } + })(); + let response = match result { + Ok(data) => FormatDataResponse::new_data(data).into_owned(), + Err(error) => { + debug!(%error, "Could not render the local clipboard"); + FormatDataResponse::new_error() + } + }; + self.proxy + .send_clipboard_message(ClipboardMessage::SendFormatData(response)); + } +} diff --git a/crates/ironrdp-cliprdr-native/src/linux/x11.rs b/crates/ironrdp-cliprdr-native/src/linux/x11.rs new file mode 100644 index 0000000000..908df457db --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/linux/x11.rs @@ -0,0 +1,763 @@ +//! XFixes selection notifications and ICCCM selection transfers. +//! +//! Every conversion uses its own requestor window, so a timed-out owner's late +//! reply cannot satisfy a newer request. Large selections use the INCR protocol. +//! +//! See [ICCCM section 2.7.2, INCR Properties]. +//! +//! [ICCCM section 2.7.2, INCR Properties]: https://www.x.org/releases/X11R7.7/doc/xorg-docs/icccm/icccm.html + +use core::time::Duration; +use std::collections::{HashMap, VecDeque}; +use std::sync::{Arc, Mutex, mpsc}; +use std::thread::JoinHandle; +use std::time::Instant; + +use x11rb::connection::Connection as _; +use x11rb::protocol::Event; +use x11rb::protocol::xfixes::{ConnectionExt as _, SelectionEventMask}; +use x11rb::protocol::xproto::{ + Atom, AtomEnum, ChangeWindowAttributesAux, ConnectionExt as _, CreateWindowAux, EventMask, PropMode, Property, + SelectionNotifyEvent, SelectionRequestEvent, Window, WindowClass, +}; +use x11rb::rust_connection::RustConnection; +use x11rb::wrapper::ConnectionExt as _; +use x11rb::{COPY_DEPTH_FROM_PARENT, CURRENT_TIME, NONE}; + +use super::os::{OsClipboard, PNG, TEXT, Transfer}; +use super::worker::{Command, MAX_TRANSFER_BYTES, PASTE_TIMEOUT}; + +// Fits the X11 core request size even without BIG-REQUESTS. +const CHUNK_BYTES: usize = 64 * 1024; +const MAX_TRANSFERS: usize = 16; +type ReadReply = mpsc::Sender>, String>>; +type XResult = Result>; + +x11rb::atom_manager! { + Atoms: AtomsCookie { + CLIPBOARD, TARGETS, TIMESTAMP, INCR, UTF8_STRING, + TEXT_UTF8: b"text/plain;charset=utf-8", + TEXT_PLAIN: b"text/plain", + PNG: b"image/png", + TRANSFER: b"IRONRDP_CLIPBOARD_TRANSFER", + } +} + +enum Request { + Read(String, ReadReply), + Offer(Vec, u64, mpsc::Sender>), + IsOwner(mpsc::Sender), + Clear, + Complete(SelectionRequestEvent, u64, Option>), + Shutdown, +} + +pub struct X11Clipboard { + commands: mpsc::Sender, + mime_types: Arc>>, + thread: Option>, +} + +impl X11Clipboard { + pub fn open(events: mpsc::Sender) -> Result { + let (commands, receiver) = mpsc::channel(); + let mime_types = Arc::new(Mutex::new(Vec::new())); + let state = + State::connect(events, commands.clone(), Arc::clone(&mime_types)).map_err(|error| error.to_string())?; + let thread = std::thread::Builder::new() + .name("cliprdr-x11".into()) + .spawn(move || { + if let Err(error) = state.run(receiver) { + tracing::warn!(%error, "X11 clipboard worker stopped"); + } + }) + .map_err(|error| error.to_string())?; + Ok(Self { + commands, + mime_types, + thread: Some(thread), + }) + } +} + +impl OsClipboard for X11Clipboard { + fn is_owner(&self) -> bool { + let (sender, receiver) = mpsc::channel(); + if self.commands.send(Request::IsOwner(sender)).is_err() { + return false; + } + receiver.recv_timeout(PASTE_TIMEOUT).unwrap_or(false) + } + + fn mime_types(&self) -> Vec { + self.mime_types + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() + } + + fn read(&mut self, mime: &str) -> Result>, String> { + let (sender, receiver) = mpsc::channel(); + self.commands + .send(Request::Read(mime.to_owned(), sender)) + .map_err(|error| error.to_string())?; + receiver + .recv_timeout(PASTE_TIMEOUT + Duration::from_secs(1)) + .map_err(|error| error.to_string())? + } + + fn offer(&mut self, mimes: &[String], generation: u64) -> Result<(), String> { + let (sender, receiver) = mpsc::channel(); + self.commands + .send(Request::Offer(mimes.to_vec(), generation, sender)) + .map_err(|error| error.to_string())?; + receiver + .recv_timeout(PASTE_TIMEOUT) + .map_err(|error| error.to_string())? + } + + fn clear(&mut self) -> Result<(), String> { + self.commands.send(Request::Clear).map_err(|error| error.to_string()) + } +} + +impl Drop for X11Clipboard { + fn drop(&mut self) { + let _ = self.commands.send(Request::Shutdown); + if let Some(thread) = self.thread.take() { + let _ = thread.join(); + } + } +} + +enum ReadKind { + Targets, + Data(ReadReply), +} + +struct Incoming { + target: Atom, + kind: ReadKind, + bytes: Vec, + incremental: bool, + deadline: Instant, +} + +struct PendingPaste { + request: SelectionRequestEvent, + deadline: Instant, + expired: bool, +} + +struct Outgoing { + target: Atom, + data: Vec, + offset: usize, + deadline: Instant, +} + +struct State { + connection: RustConnection, + window: Window, + root: Window, + atoms: Atoms, + events: mpsc::Sender, + commands: mpsc::Sender, + mime_types: Arc>>, + targets: Vec, + owner: Window, + owns: bool, + generation: u64, + timestamp: u32, + incoming: HashMap, + outgoing: HashMap<(Window, Atom), Outgoing>, + pending_pastes: HashMap<(Window, Atom), PendingPaste>, + deferred: VecDeque, +} + +impl State { + fn connect( + events: mpsc::Sender, + commands: mpsc::Sender, + mime_types: Arc>>, + ) -> XResult { + let (connection, screen) = x11rb::connect(None)?; + connection.xfixes_query_version(5, 0)?.reply()?; + let root = connection.setup().roots[screen].root; + let window = connection.generate_id()?; + connection + .create_window( + COPY_DEPTH_FROM_PARENT, + window, + root, + 0, + 0, + 1, + 1, + 0, + WindowClass::INPUT_OUTPUT, + 0, + &CreateWindowAux::new().event_mask(EventMask::PROPERTY_CHANGE), + )? + .check()?; + let atoms = Atoms::new(&connection)?.reply()?; + connection + .xfixes_select_selection_input( + window, + atoms.CLIPBOARD, + SelectionEventMask::SET_SELECTION_OWNER + | SelectionEventMask::SELECTION_WINDOW_DESTROY + | SelectionEventMask::SELECTION_CLIENT_CLOSE, + )? + .check()?; + let owner = connection.get_selection_owner(atoms.CLIPBOARD)?.reply()?.owner; + let mut state = Self { + connection, + window, + root, + atoms, + events, + commands, + mime_types, + targets: Vec::new(), + owner, + owns: false, + generation: 0, + timestamp: CURRENT_TIME, + incoming: HashMap::new(), + outgoing: HashMap::new(), + pending_pastes: HashMap::new(), + deferred: VecDeque::new(), + }; + state.timestamp = state.server_timestamp()?; + if owner != NONE { + state.begin_read(state.atoms.TARGETS, ReadKind::Targets)?; + } + state.connection.flush()?; + Ok(state) + } + + fn run(mut self, receiver: mpsc::Receiver) -> XResult<()> { + loop { + for _ in 0..64 { + let event = match self.deferred.pop_front() { + Some(event) => Some(event), + None => self.connection.poll_for_event()?, + }; + let Some(event) = event else { + break; + }; + // BadWindow from an application that closed mid-paste is local to + // that transfer; it must not stop clipboard synchronization. + if let Err(error) = self.event(event) { + tracing::debug!(%error, "X11 clipboard event failed"); + } + } + match receiver.recv_timeout(Duration::from_millis(10)) { + Ok(Request::Shutdown) | Err(mpsc::RecvTimeoutError::Disconnected) => break, + Ok(request) => { + if let Err(error) = self.request(request) { + tracing::debug!(%error, "X11 clipboard request failed"); + } + } + Err(mpsc::RecvTimeoutError::Timeout) => {} + } + let now = Instant::now(); + let expired: Vec<_> = self + .incoming + .iter() + .filter(|(_, read)| read.deadline <= now) + .map(|(window, _)| *window) + .collect(); + for window in expired { + self.finish_read(window, Err("X11 clipboard read timed out".into()))?; + } + self.outgoing.retain(|_, write| now < write.deadline); + let mut refused = Vec::new(); + for paste in self.pending_pastes.values_mut() { + if !paste.expired && paste.deadline <= now { + paste.expired = true; + refused.push(paste.request); + } + } + for request in refused { + let _ = self.notify(request, NONE); + } + // Keep expired pending slots until Complete is consumed, so cached + // responses already queued on the channel remain within the budget. + self.connection.flush()?; + } + self.clear()?; + Ok(()) + } + + fn request(&mut self, request: Request) -> XResult<()> { + match request { + Request::Read(mime, sender) => { + let target = if mime == PNG { + self.targets.contains(&self.atoms.PNG).then_some(self.atoms.PNG) + } else { + [ + self.atoms.UTF8_STRING, + self.atoms.TEXT_UTF8, + self.atoms.TEXT_PLAIN, + AtomEnum::STRING.into(), + ] + .into_iter() + .find(|target| self.targets.contains(target)) + }; + if let Some(target) = target.filter(|_| !self.owns && self.owner != NONE) { + self.begin_read(target, ReadKind::Data(sender))?; + } else { + let _ = sender.send(Ok(None)); + } + } + Request::IsOwner(sender) => { + let owner = self + .connection + .get_selection_owner(self.atoms.CLIPBOARD)? + .reply()? + .owner; + let _ = sender.send(owner == self.window); + } + Request::Offer(mimes, generation, sender) => { + let result = (|| -> XResult<()> { + self.generation = generation; + self.targets.clear(); + for mime in mimes { + if mime == PNG { + self.targets.push(self.atoms.PNG); + } + if mime == TEXT { + self.targets + .extend([self.atoms.UTF8_STRING, self.atoms.TEXT_UTF8, self.atoms.TEXT_PLAIN]); + } + } + self.timestamp = self.server_timestamp()?; + self.connection + .set_selection_owner(self.window, self.atoms.CLIPBOARD, self.timestamp)? + .check()?; + self.owner = self + .connection + .get_selection_owner(self.atoms.CLIPBOARD)? + .reply()? + .owner; + self.owns = self.owner == self.window; + let readers: Vec<_> = self.incoming.keys().copied().collect(); + for window in readers { + self.finish_read(window, Err("clipboard selection changed while reading".into()))?; + } + if !self.owns { + return Err("clipboard ownership was not acquired".into()); + } + Ok(()) + })(); + let _ = sender.send(result.map_err(|error| error.to_string())); + } + Request::Clear => self.clear()?, + Request::Complete(request, generation, data) => { + let property = if request.property == NONE { + request.target + } else { + request.property + }; + let Some(pending) = self.pending_pastes.remove(&(request.requestor, property)) else { + return Ok(()); + }; + if pending.expired { + return Ok(()); + } + if generation != self.generation || !self.owns { + self.notify(request, NONE)?; + } else if let Some(data) = data.filter(|d| d.len() <= MAX_TRANSFER_BYTES) { + let property = if request.property == NONE { + request.target + } else { + request.property + }; + if data.len() <= CHUNK_BYTES { + self.connection + .change_property8(PropMode::REPLACE, request.requestor, property, request.target, &data)? + .check()?; + self.notify(request, property)?; + } else if self.outgoing.len() < MAX_TRANSFERS { + self.connection + .change_window_attributes( + request.requestor, + &ChangeWindowAttributesAux::new().event_mask(EventMask::PROPERTY_CHANGE), + )? + .check()?; + let length = u32::try_from(data.len())?; + self.connection + .change_property32( + PropMode::REPLACE, + request.requestor, + property, + self.atoms.INCR, + &[length], + )? + .check()?; + self.outgoing.insert( + (request.requestor, property), + Outgoing { + target: request.target, + data, + offset: 0, + deadline: Instant::now() + PASTE_TIMEOUT, + }, + ); + self.notify(request, property)?; + } else { + self.notify(request, NONE)?; + } + } else { + self.notify(request, NONE)?; + } + } + Request::Shutdown => {} + } + Ok(()) + } + + // ICCCM uses the event timestamp that triggered ownership. RDP commands + // have no X timestamp, so obtain one from PropertyNotify and retain all + // other events for normal dispatch. + fn server_timestamp(&mut self) -> XResult { + self.connection + .change_property8( + PropMode::APPEND, + self.window, + self.atoms.TIMESTAMP, + AtomEnum::INTEGER, + &[], + )? + .check()?; + self.connection.flush()?; + loop { + let event = self.connection.wait_for_event()?; + if let Event::PropertyNotify(property) = &event { + if property.window == self.window && property.atom == self.atoms.TIMESTAMP { + return Ok(property.time); + } + } + self.deferred.push_back(event); + } + } + + fn clear(&mut self) -> XResult<()> { + if self + .connection + .get_selection_owner(self.atoms.CLIPBOARD)? + .reply()? + .owner + == self.window + { + self.connection + .set_selection_owner(NONE, self.atoms.CLIPBOARD, self.timestamp)? + .check()?; + } + self.owns = false; + Ok(()) + } + + fn begin_read(&mut self, target: Atom, kind: ReadKind) -> XResult<()> { + let window = self.connection.generate_id()?; + self.connection + .create_window( + COPY_DEPTH_FROM_PARENT, + window, + self.root, + 0, + 0, + 1, + 1, + 0, + WindowClass::INPUT_OUTPUT, + 0, + &CreateWindowAux::new().event_mask(EventMask::PROPERTY_CHANGE), + )? + .check()?; + self.connection + .convert_selection( + window, + self.atoms.CLIPBOARD, + target, + self.atoms.TRANSFER, + self.timestamp, + )? + .check()?; + self.incoming.insert( + window, + Incoming { + target, + kind, + bytes: Vec::new(), + incremental: false, + deadline: Instant::now() + PASTE_TIMEOUT, + }, + ); + Ok(()) + } + + fn event(&mut self, event: Event) -> XResult<()> { + match event { + Event::XfixesSelectionNotify(event) if event.selection == self.atoms.CLIPBOARD => { + if self + .connection + .get_selection_owner(self.atoms.CLIPBOARD)? + .reply()? + .owner + != event.owner + { + return Ok(()); + } + self.owner = event.owner; + if event.owner == self.window { + self.owns = true; + return Ok(()); + } + self.timestamp = event.selection_timestamp; + self.owns = false; + self.targets.clear(); + let readers: Vec<_> = self.incoming.keys().copied().collect(); + for window in readers { + self.finish_read(window, Err("clipboard selection changed while reading".into()))?; + } + if event.owner != NONE { + self.begin_read(self.atoms.TARGETS, ReadKind::Targets)?; + } else { + self.mime_types + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clear(); + let _ = self.events.send(Command::LocalChanged); + } + } + Event::SelectionNotify(event) + if event.selection == self.atoms.CLIPBOARD + && self + .incoming + .get(&event.requestor) + .is_some_and(|read| read.target == event.target) => + { + if event.property != self.atoms.TRANSFER { + self.finish_read(event.requestor, Err("clipboard conversion refused".into()))?; + } else { + self.read_property(event.requestor)?; + } + } + Event::PropertyNotify(event) => { + if event.state == Property::NEW_VALUE && self.incoming.get(&event.window).is_some_and(|r| r.incremental) + { + self.read_property(event.window)?; + } + if event.state == Property::DELETE { + if let Some(mut write) = self.outgoing.remove(&(event.window, event.atom)) { + let end = write.data.len().min(write.offset.saturating_add(CHUNK_BYTES)); + let chunk = &write.data[write.offset..end]; + self.connection + .change_property8(PropMode::REPLACE, event.window, event.atom, write.target, chunk)? + .check()?; + if !chunk.is_empty() { + write.offset = end; + write.deadline = Instant::now() + PASTE_TIMEOUT; + self.outgoing.insert((event.window, event.atom), write); + } + } + } + } + Event::SelectionRequest(request) => self.selection_request(request)?, + Event::SelectionClear(event) if event.selection == self.atoms.CLIPBOARD => { + // A clear from an earlier owner can be queued while Offer is + // obtaining its timestamp. Consult the server before applying + // it to a selection that has since been reacquired. + self.owns = self + .connection + .get_selection_owner(self.atoms.CLIPBOARD)? + .reply()? + .owner + == self.window; + } + _ => {} + } + Ok(()) + } + + fn read_property(&mut self, window: Window) -> XResult<()> { + let reply = self + .connection + .get_property( + true, + window, + self.atoms.TRANSFER, + AtomEnum::ANY, + 0, + u32::try_from(MAX_TRANSFER_BYTES / 4 + 1)?, + )? + .reply()?; + let Some(read) = self.incoming.get_mut(&window) else { + return Ok(()); + }; + if reply.type_ == self.atoms.INCR { + let size = reply.value32().and_then(|mut values| values.next()).unwrap_or(u32::MAX); + if read.incremental || MAX_TRANSFER_BYTES < usize::try_from(size)? { + self.finish_read(window, Err("X11 clipboard data exceeds the size limit".into()))?; + } else { + read.incremental = true; + } + return Ok(()); + } + if reply.bytes_after != 0 || MAX_TRANSFER_BYTES < read.bytes.len().saturating_add(reply.value.len()) { + self.finish_read(window, Err("X11 clipboard data exceeds the size limit".into()))?; + return Ok(()); + } + let expected_type = if read.target == self.atoms.TARGETS { + AtomEnum::ATOM.into() + } else { + read.target + }; + let expected_format = if read.target == self.atoms.TARGETS { 32 } else { 8 }; + if reply.type_ != expected_type || reply.format != expected_format { + self.finish_read(window, Err("invalid X11 clipboard property type".into()))?; + return Ok(()); + } + let finished = !read.incremental || reply.value.is_empty(); + read.bytes.extend_from_slice(&reply.value); + if finished { + let bytes = core::mem::take(&mut read.bytes); + self.finish_read(window, Ok(bytes))?; + } + Ok(()) + } + + fn finish_read(&mut self, window: Window, result: Result, String>) -> XResult<()> { + let Some(read) = self.incoming.remove(&window) else { + return Ok(()); + }; + self.connection.destroy_window(window)?.check()?; + match read.kind { + ReadKind::Data(sender) => { + let result = result.map(|bytes| { + if read.target == AtomEnum::STRING.into() { + bytes.into_iter().map(char::from).collect::().into_bytes() + } else { + bytes + } + }); + let _ = sender.send(result.map(Some)); + } + ReadKind::Targets => { + if self.owns { + return Ok(()); + } + self.targets = result + .unwrap_or_default() + .chunks_exact(4) + .map(|chunk| u32::from_ne_bytes([chunk[0], chunk[1], chunk[2], chunk[3]])) + .collect(); + let mut mimes = Vec::new(); + if [ + self.atoms.UTF8_STRING, + self.atoms.TEXT_UTF8, + self.atoms.TEXT_PLAIN, + AtomEnum::STRING.into(), + ] + .iter() + .any(|target| self.targets.contains(target)) + { + mimes.push(TEXT.to_owned()); + } + if self.targets.contains(&self.atoms.PNG) { + mimes.push(PNG.to_owned()); + } + *self + .mime_types + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = mimes; + let _ = self.events.send(Command::LocalChanged); + } + } + Ok(()) + } + + fn selection_request(&mut self, request: SelectionRequestEvent) -> XResult<()> { + if request.selection != self.atoms.CLIPBOARD || !self.owns { + return self.notify(request, NONE); + } + // X timestamps wrap at 32 bits. We cannot render a selection from + // before our current ownership interval, including queued old requests. + if request.time != CURRENT_TIME && 0x7fff_ffff < request.time.wrapping_sub(self.timestamp) { + return self.notify(request, NONE); + } + let property = if request.property == NONE { + request.target + } else { + request.property + }; + if request.target == self.atoms.TARGETS { + let mut targets = self.targets.clone(); + targets.extend([self.atoms.TARGETS, self.atoms.TIMESTAMP]); + self.connection + .change_property32(PropMode::REPLACE, request.requestor, property, AtomEnum::ATOM, &targets)? + .check()?; + return self.notify(request, property); + } + if request.target == self.atoms.TIMESTAMP { + self.connection + .change_property32( + PropMode::REPLACE, + request.requestor, + property, + AtomEnum::INTEGER, + &[self.timestamp], + )? + .check()?; + return self.notify(request, property); + } + if !self.targets.contains(&request.target) { + return self.notify(request, NONE); + } + // Reserve before notifying the protocol worker: cached replies can + // otherwise enqueue unbounded full clipboard buffers before we consume + // their Complete commands. A slot covers waiting, queued and INCR data. + let key = (request.requestor, property); + if MAX_TRANSFERS <= self.pending_pastes.len() + self.outgoing.len() + || self.pending_pastes.contains_key(&key) + || self.outgoing.contains_key(&key) + { + return self.notify(request, NONE); + } + self.pending_pastes.insert( + key, + PendingPaste { + request, + deadline: Instant::now() + PASTE_TIMEOUT, + expired: false, + }, + ); + let mime = if request.target == self.atoms.PNG { PNG } else { TEXT }; + let commands = self.commands.clone(); + let generation = self.generation; + let transfer = Transfer::new(mime.to_owned(), generation, move |data| { + let _ = commands.send(Request::Complete(request, generation, data)); + }); + let _ = self.events.send(Command::Paste(transfer)); + Ok(()) + } + + fn notify(&self, request: SelectionRequestEvent, property: Atom) -> XResult<()> { + let notify = SelectionNotifyEvent { + response_type: x11rb::protocol::xproto::SELECTION_NOTIFY_EVENT, + sequence: 0, + time: request.time, + requestor: request.requestor, + selection: request.selection, + target: request.target, + property, + }; + self.connection + .send_event(false, request.requestor, EventMask::NO_EVENT, notify)? + .check()?; + Ok(()) + } +} diff --git a/crates/ironrdp-testsuite-core/Cargo.toml b/crates/ironrdp-testsuite-core/Cargo.toml index 2654492a8b..2d394583ac 100644 --- a/crates/ironrdp-testsuite-core/Cargo.toml +++ b/crates/ironrdp-testsuite-core/Cargo.toml @@ -39,6 +39,7 @@ expect-test.workspace = true hex = "0.4" ironrdp-cliprdr-format.path = "../ironrdp-cliprdr-format" ironrdp-cliprdr = { path = "../ironrdp-cliprdr", features = ["__test"] } +ironrdp-cliprdr-native = { path = "../ironrdp-cliprdr-native", features = ["__test"] } ironrdp-acceptor.path = "../ironrdp-acceptor" ironrdp-async.path = "../ironrdp-async" ironrdp-tokio.path = "../ironrdp-tokio" diff --git a/crates/ironrdp-testsuite-core/tests/cliprdr_native/data_control.rs b/crates/ironrdp-testsuite-core/tests/cliprdr_native/data_control.rs new file mode 100644 index 0000000000..371380de21 --- /dev/null +++ b/crates/ironrdp-testsuite-core/tests/cliprdr_native/data_control.rs @@ -0,0 +1,450 @@ +//! Tests for the Wayland data-control clipboard client that need no compositor. + +use std::{ + io::Read as _, + os::fd::OwnedFd, + sync::{Arc, Mutex}, +}; + +use ironrdp_cliprdr_native::data_control::state::{State, is_own_selection}; +use ironrdp_cliprdr_native::data_control::{Options, Preference, Protocol, find_mime_match}; + +fn mimes(types: &[&str]) -> Vec { + types.iter().map(|t| (*t).to_owned()).collect() +} + +fn read_all(fd: OwnedFd) -> Vec { + let mut buf = Vec::new(); + std::fs::File::from(fd).read_to_end(&mut buf).unwrap(); + buf +} + +fn pipe() -> (OwnedFd, OwnedFd) { + let (reader, writer) = std::io::pipe().unwrap(); + (OwnedFd::from(reader), OwnedFd::from(writer)) +} + +#[test] +fn mime_exact_match_is_returned() { + let available = mimes(&["text/plain", "text/html"]); + assert_eq!(find_mime_match("text/html", &available), Some("text/html")); +} + +#[test] +fn mime_requested_charset_is_stripped() { + let available = mimes(&["text/plain"]); + assert_eq!( + find_mime_match("text/plain;charset=utf-8", &available), + Some("text/plain") + ); +} + +#[test] +fn mime_missing_charset_is_added() { + let available = mimes(&["text/plain;charset=utf-8"]); + assert_eq!( + find_mime_match("text/plain", &available), + Some("text/plain;charset=utf-8") + ); +} + +#[test] +fn mime_non_text_types_get_no_charset_fallback() { + let available = mimes(&["image/png;foo=bar"]); + assert_eq!(find_mime_match("image/png", &available), None); +} + +#[test] +fn mime_exact_match_beats_a_charset_variant() { + let available = mimes(&["text/plain;charset=utf-8", "text/plain"]); + assert_eq!(find_mime_match("text/plain", &available), Some("text/plain")); +} + +#[test] +fn mime_any_charset_of_the_same_base_type_is_accepted() { + let available = mimes(&["text/plain;charset=iso-8859-1"]); + assert_eq!( + find_mime_match("text/plain", &available), + Some("text/plain;charset=iso-8859-1") + ); +} + +#[test] +fn auto_preference_tries_ext_then_wlr() { + assert_eq!(Options::new().candidates(), [Protocol::Ext, Protocol::Wlr]); +} + +#[test] +fn wlr_preference_tries_wlr_first() { + let options = Options::new().preference(Preference::Wlr); + assert_eq!(options.candidates(), [Protocol::Wlr, Protocol::Ext]); +} + +#[test] +fn disabling_fallback_leaves_only_the_preferred_protocol() { + let options = Options::new().preference(Preference::Wlr).allow_fallback(false); + assert_eq!(options.candidates(), [Protocol::Wlr]); +} + +#[test] +fn own_selection_is_recognised_while_our_source_is_live() { + let ours = mimes(&["text/plain;charset=utf-8"]); + assert!(is_own_selection(true, &ours, &mimes(&["text/plain;charset=utf-8"]))); +} + +#[test] +fn a_selection_after_our_source_was_cancelled_is_foreign() { + let ours = mimes(&["text/plain;charset=utf-8"]); + assert!(!is_own_selection(false, &ours, &mimes(&["text/plain;charset=utf-8"]))); +} + +#[test] +fn a_different_offer_while_our_source_is_live_is_foreign() { + let ours = mimes(&["text/plain;charset=utf-8"]); + let other = mimes(&[ + "text/plain;charset=utf-8", + "text/plain", + "TEXT", + "STRING", + "UTF8_STRING", + ]); + assert!(!is_own_selection(true, &ours, &other)); +} + +#[test] +fn own_selection_ignores_offer_order() { + let ours = mimes(&["image/png", "text/html"]); + assert!(is_own_selection(true, &ours, &mimes(&["text/html", "image/png"]))); +} + +#[test] +fn clearing_the_selection_notifies_and_advances_the_serial() { + let mut state = State::default(); + let called = Arc::new(Mutex::new(false)); + let seen = Arc::clone(&called); + state + .shared() + .lock() + .unwrap() + .set_on_change(Arc::new(move |types: Vec| { + assert!(types.is_empty()); + *seen.lock().unwrap() = true; + })); + + state.on_selection_cleared(); + + assert!(*called.lock().unwrap()); + let shared = state.shared().lock().unwrap(); + assert!(shared.mime_types().is_empty()); + assert_eq!(shared.serial(), 1); +} + +#[test] +fn cancelling_our_source_drops_its_data() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello".to_vec()); + state.on_source_cancelled(); + assert!(!state.has_source_data()); +} + +#[test] +fn device_finished_drops_source_data() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello".to_vec()); + state.on_device_finished(); + assert!(!state.has_source_data()); +} + +#[test] +fn a_mime_type_without_a_pending_offer_is_ignored() { + let mut state = State::default(); + state.on_offer_mime_type("text/plain".to_owned()); +} + +#[test] +fn a_paste_of_an_unknown_type_closes_the_pipe_with_no_data() { + let mut state = State::default(); + let (read, write) = pipe(); + state.on_source_send("text/unknown", write); + assert!(read_all(read).is_empty()); +} + +#[test] +fn a_paste_is_answered_from_cached_data() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello world".to_vec()); + let (read, write) = pipe(); + state.on_source_send("text/plain", write); + assert_eq!(read_all(read), b"hello world"); +} + +type TransferLog = Arc>>; + +/// State advertising `image/png` with a transfer callback that records what it was asked for. +fn state_with_transfer_callback() -> (State, TransferLog) { + let mut state = State::default(); + state.set_advertised(mimes(&["image/png"])); + let requests: TransferLog = Arc::default(); + let seen = Arc::clone(&requests); + state + .shared() + .lock() + .unwrap() + .set_on_transfer(Arc::new(move |serial, mime| { + seen.lock().unwrap().push((serial, mime)); + })); + (state, requests) +} + +#[test] +fn a_paste_without_data_is_held_until_the_transfer_completes() { + let (mut state, requests) = state_with_transfer_callback(); + let (first_read, first_write) = pipe(); + let (second_read, second_write) = pipe(); + + state.on_source_send("image/png", first_write); + state.on_source_send("image/png", second_write); + // Both pastes share one transfer. + assert_eq!(*requests.lock().unwrap(), vec![(1, "image/png".to_owned())]); + + state.complete_transfer(1, Some(b"png bytes".to_vec())); + assert_eq!(read_all(first_read), b"png bytes"); + assert_eq!(read_all(second_read), b"png bytes"); + + // A later paste is served from the cache without a new transfer. + let (third_read, third_write) = pipe(); + state.on_source_send("image/png", third_write); + assert_eq!(read_all(third_read), b"png bytes"); + assert_eq!(requests.lock().unwrap().len(), 1); +} + +#[test] +fn a_held_paste_gets_nothing_when_the_selection_is_lost() { + let (mut state, _requests) = state_with_transfer_callback(); + let (read, write) = pipe(); + state.on_source_send("image/png", write); + + state.on_source_cancelled(); + assert!(read_all(read).is_empty()); + // An answer for the lost selection is ignored. + state.complete_transfer(1, Some(b"late".to_vec())); + assert!(!state.has_source_data()); +} + +#[test] +fn a_paste_of_an_unadvertised_type_raises_no_transfer() { + let (mut state, requests) = state_with_transfer_callback(); + let (read, write) = pipe(); + state.on_source_send("text/html", write); + assert!(read_all(read).is_empty()); + assert!(requests.lock().unwrap().is_empty()); +} + +#[test] +fn the_change_callback_may_lock_the_shared_state() { + let mut state = State::default(); + let shared = Arc::clone(state.shared()); + let relocked = Arc::new(Mutex::new(false)); + let seen = Arc::clone(&relocked); + state + .shared() + .lock() + .unwrap() + .set_on_change(Arc::new(move |_types: Vec| { + // A callback that reads the handle locks this state; it must not be held. + *seen.lock().unwrap() = shared.try_lock().is_ok(); + })); + + state.on_selection_cleared(); + + assert!(*relocked.lock().unwrap()); +} + +#[test] +fn held_paste_descriptors_are_bounded_and_released_on_cancellation() { + let (mut state, requests) = state_with_transfer_callback(); + let mut readers = Vec::new(); + for _ in 0..16 { + let (reader, writer) = pipe(); + state.on_source_send("image/png", writer); + readers.push(reader); + } + assert_eq!(state.paste_count(), 16); + let (rejected, writer) = pipe(); + state.on_source_send("image/png", writer); + assert!(read_all(rejected).is_empty()); + assert_eq!(requests.lock().unwrap().len(), 1); + + state.on_source_cancelled(); + assert_eq!(state.paste_count(), 0); + for reader in readers { + assert!(read_all(reader).is_empty()); + } + state.set_advertised(mimes(&["image/png"])); + let (reader, writer) = pipe(); + state.on_source_send("image/png", writer); + // The previous response cannot answer the next selection's paste. + state.complete_transfer(1, Some(b"stale".to_vec())); + assert_eq!(state.paste_count(), 1); + state.complete_transfer(2, Some(b"current".to_vec())); + assert_eq!(read_all(reader), b"current"); +} + +#[test] +fn unanswered_pastes_expire_without_replacing_the_selection() { + use core::time::Duration; + use std::time::Instant; + + let (mut state, requests) = state_with_transfer_callback(); + let (reader, writer) = pipe(); + state.on_source_send("image/png", writer); + state.expire_transfers(Instant::now() + Duration::from_secs(6)); + assert!(read_all(reader).is_empty()); + assert_eq!(state.paste_count(), 0); + state.complete_transfer(1, Some(b"late".to_vec())); + assert!(!state.has_source_data()); + + let (reader, writer) = pipe(); + state.on_source_send("image/png", writer); + assert_eq!(requests.lock().unwrap().len(), 2); + state.complete_transfer(2, Some(b"new".to_vec())); + assert_eq!(read_all(reader), b"new"); +} + +#[test] +fn cached_pastes_share_the_descriptor_and_writer_limit() { + let mut state = State::default(); + state.insert_source_data("image/png", vec![0x55; 1024 * 1024]); + let mut readers = Vec::new(); + for _ in 0..16 { + let (reader, writer) = pipe(); + state.on_source_send("image/png", writer); + readers.push(reader); + } + assert_eq!(state.paste_count(), 16); + let (rejected, writer) = pipe(); + state.on_source_send("image/png", writer); + assert!(read_all(rejected).is_empty()); + // Breaking each blocked pipe releases its thread and its permit. + drop(readers); + wait_for_pastes_to_finish(&state); + let (reader, writer) = pipe(); + state.insert_source_data("image/png", b"available again".to_vec()); + state.on_source_send("image/png", writer); + assert_eq!(read_all(reader), b"available again"); +} + +#[test] +fn a_blocked_paste_writer_times_out_and_closes_its_descriptor() { + use core::time::Duration; + + let mut state = State::default(); + state.set_transfer_timeout(Duration::from_millis(30)); + state.insert_source_data("image/png", vec![0x55; 1024 * 1024]); + let (reader, writer) = pipe(); + state.on_source_send("image/png", writer); + // Keep the reader open without draining its pipe: write_all used to block forever. + wait_for_pastes_to_finish(&state); + assert!(read_all(reader).len() < 1024 * 1024); +} + +fn wait_for_pastes_to_finish(state: &State) { + use core::time::Duration; + use std::time::Instant; + + let deadline = Instant::now() + Duration::from_secs(2); + while state.paste_count() != 0 { + assert!(Instant::now() < deadline, "clipboard writer did not release its permit"); + std::thread::sleep(Duration::from_millis(1)); + } +} + +#[test] +fn a_read_deadline_is_not_extended_by_slow_progress() { + use core::time::Duration; + use std::io::Write as _; + use std::sync::mpsc; + + use ironrdp_cliprdr_native::data_control::{Error, read_pipe}; + + let (mut reader, mut writer) = std::io::pipe().unwrap(); + let (started, ready) = mpsc::sync_channel(0); + let thread = std::thread::spawn(move || { + writer.write_all(b"a").unwrap(); + started.send(()).unwrap(); + for _ in 0..1000 { + std::thread::sleep(Duration::from_millis(1)); + if writer.write_all(b"a").is_err() { + break; + } + } + }); + ready.recv().unwrap(); + let result = read_pipe(&mut reader, 1024 * 1024, Duration::from_millis(30)); + drop(reader); + thread.join().unwrap(); + assert!(matches!(result, Err(Error::Timeout))); +} + +#[test] +fn same_format_selection_before_cancellation_still_notifies_ownership_loss() { + let mut state = State::default(); + let formats = mimes(&["text/plain"]); + let observed = Arc::new(Mutex::new(Vec::new())); + let seen = Arc::clone(&observed); + let shared = Arc::clone(state.shared()); + state.shared().lock().unwrap().set_on_change(Arc::new(move |mimes| { + assert!( + shared.try_lock().is_ok(), + "callback must run after releasing shared state" + ); + seen.lock().unwrap().push(mimes); + })); + state.set_advertised(formats.clone()); + state.publish_selection(formats.clone(), true); + // Without the cancellation yet, the existing ownership heuristic cannot + // distinguish this external selection from our own same-format echo. + let assumed_own = is_own_selection(true, &formats, &formats); + state.publish_selection(formats.clone(), assumed_own); + assert!(observed.lock().unwrap().is_empty()); + + state.on_source_cancelled(); + assert_eq!(*observed.lock().unwrap(), vec![formats]); + assert_eq!(state.shared().lock().unwrap().serial(), 3); +} + +#[test] +fn cancellation_before_external_selection_converges_on_the_new_formats() { + let mut state = State::default(); + let observed = Arc::new(Mutex::new(Vec::new())); + let seen = Arc::clone(&observed); + state + .shared() + .lock() + .unwrap() + .set_on_change(Arc::new(move |mimes| seen.lock().unwrap().push(mimes))); + state.publish_selection(mimes(&["text/plain"]), true); + state.on_source_cancelled(); + state.publish_selection(mimes(&["image/png"]), false); + assert_eq!( + *observed.lock().unwrap(), + vec![mimes(&["text/plain"]), mimes(&["image/png"])] + ); + assert_eq!(state.shared().lock().unwrap().mime_types(), mimes(&["image/png"])); +} + +#[test] +fn cancellation_after_a_foreign_selection_does_not_repeat_its_notification() { + let mut state = State::default(); + let observed = Arc::new(Mutex::new(Vec::new())); + let seen = Arc::clone(&observed); + state + .shared() + .lock() + .unwrap() + .set_on_change(Arc::new(move |mimes| seen.lock().unwrap().push(mimes))); + state.publish_selection(mimes(&["text/plain"]), true); + state.publish_selection(mimes(&["image/png"]), false); + state.on_source_cancelled(); + assert_eq!(*observed.lock().unwrap(), vec![mimes(&["image/png"])]); +} diff --git a/crates/ironrdp-testsuite-core/tests/cliprdr_native/linux.rs b/crates/ironrdp-testsuite-core/tests/cliprdr_native/linux.rs new file mode 100644 index 0000000000..e0042d1cc0 --- /dev/null +++ b/crates/ironrdp-testsuite-core/tests/cliprdr_native/linux.rs @@ -0,0 +1,559 @@ +use std::collections::HashMap; +use std::sync::{Arc, Mutex, mpsc}; +use std::time::Instant; + +use ironrdp_cliprdr::backend::{ClipboardMessage, ClipboardMessageProxy}; +use ironrdp_cliprdr::pdu::{ClipboardFormat, ClipboardFormatId, FormatDataResponse}; +use ironrdp_cliprdr_format::bitmap; +use ironrdp_cliprdr_native::linux::os::{OsClipboard, PNG, TEXT, Transfer}; +use ironrdp_cliprdr_native::linux::worker::{Command, PASTE_TIMEOUT, Worker}; + +#[derive(Default)] +struct FakeOs { + data: HashMap>, + offers: Vec<(Vec, u64)>, + reads: usize, + clears: usize, + owns: bool, +} + +struct FakeClipboard(Arc>); + +impl OsClipboard for FakeClipboard { + fn is_owner(&self) -> bool { + self.0.lock().unwrap().owns + } + fn mime_types(&self) -> Vec { + self.0.lock().unwrap().data.keys().cloned().collect() + } + fn read(&mut self, mime: &str) -> Result>, String> { + let mut os = self.0.lock().unwrap(); + os.reads += 1; + Ok(os.data.get(mime).cloned()) + } + fn offer(&mut self, mimes: &[String], generation: u64) -> Result<(), String> { + self.0.lock().unwrap().offers.push((mimes.to_vec(), generation)); + Ok(()) + } + fn clear(&mut self) -> Result<(), String> { + self.0.lock().unwrap().clears += 1; + Ok(()) + } +} + +#[derive(Default, Debug, Clone)] +struct Recorder(Arc>>); + +impl ClipboardMessageProxy for Recorder { + fn send_clipboard_message(&self, message: ClipboardMessage) { + self.0.lock().unwrap().push(message); + } +} + +impl Recorder { + fn take(&self) -> Vec { + core::mem::take(&mut *self.0.lock().unwrap()) + } +} + +type TestWorker = Worker; +fn worker() -> (TestWorker, Arc>, Recorder) { + let os = Arc::new(Mutex::new(FakeOs::default())); + let recorder = Recorder::default(); + ( + Worker::new(FakeClipboard(Arc::clone(&os)), recorder.clone()), + os, + recorder, + ) +} + +fn remote(worker: &mut TestWorker, formats: &[ClipboardFormatId]) { + worker.handle(Command::RemoteCopy( + formats.iter().copied().map(ClipboardFormat::new).collect(), + )); +} + +fn paste(worker: &mut TestWorker, os: &Arc>, mime: &str) -> mpsc::Receiver>> { + let generation = os.lock().unwrap().offers.last().unwrap().1; + let (sender, receiver) = mpsc::channel(); + worker.handle(Command::Paste(Transfer::new( + mime.to_owned(), + generation, + move |data| { + sender.send(data).unwrap(); + }, + ))); + receiver +} + +fn text_response(worker: &mut TestWorker, text: &str) { + worker.handle(Command::RemoteData(Some( + FormatDataResponse::new_unicode_string(text).data().to_vec(), + ))); +} + +fn png() -> Vec { + let mut bytes = Vec::new(); + { + let mut encoder = png::Encoder::new(&mut bytes, 1, 1); + encoder.set_color(png::ColorType::Rgba); + encoder.set_depth(png::BitDepth::Eight); + encoder + .write_header() + .unwrap() + .write_image_data(&[0x22, 0x44, 0x66, 0xff]) + .unwrap(); + } + bytes +} + +#[test] +fn selection_events_wait_for_monitor_ready_and_do_not_read_content() { + let (mut worker, os, recorder) = worker(); + os.lock().unwrap().data.insert(TEXT.into(), b"hello".to_vec()); + worker.handle(Command::LocalChanged); + assert!(recorder.take().is_empty()); + worker.handle(Command::AdvertiseLocal); + assert!( + matches!(recorder.take().as_slice(), [ClipboardMessage::SendInitiateCopy(formats)] if formats.len() == 1 && formats[0].id() == ClipboardFormatId::CF_UNICODETEXT) + ); + assert_eq!(os.lock().unwrap().reads, 0); +} + +#[test] +fn a_new_local_copy_with_the_same_formats_is_announced() { + let (mut worker, _, recorder) = worker(); + worker.handle(Command::AdvertiseLocal); + recorder.take(); + worker.handle(Command::LocalChanged); + worker.handle(Command::LocalChanged); + assert_eq!(recorder.take().len(), 2); +} + +#[test] +fn remote_formats_are_advertised_without_fetching_data() { + let (mut worker, os, recorder) = worker(); + remote( + &mut worker, + &[ClipboardFormatId::CF_UNICODETEXT, ClipboardFormatId::CF_DIBV5], + ); + assert!(recorder.take().is_empty()); + assert_eq!(os.lock().unwrap().offers[0].0, [TEXT, "text/plain", PNG]); + let response = paste(&mut worker, &os, TEXT); + assert!(matches!( + recorder.take().as_slice(), + [ClipboardMessage::SendInitiatePaste(ClipboardFormatId::CF_UNICODETEXT)] + )); + text_response(&mut worker, "only when pasted"); + assert_eq!(response.recv().unwrap(), Some(b"only when pasted".to_vec())); +} + +#[test] +fn same_format_waiters_share_one_wire_request_and_cache() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let first = paste(&mut worker, &os, TEXT); + let second = paste(&mut worker, &os, "text/plain"); + assert_eq!(recorder.take().len(), 1); + text_response(&mut worker, "shared"); + assert_eq!(first.recv().unwrap(), Some(b"shared".to_vec())); + assert_eq!(second.recv().unwrap(), Some(b"shared".to_vec())); + let cached = paste(&mut worker, &os, TEXT); + assert_eq!(cached.recv().unwrap(), Some(b"shared".to_vec())); + assert!(recorder.take().is_empty()); +} + +#[test] +fn different_formats_are_requested_one_at_a_time() { + let (mut worker, os, recorder) = worker(); + remote( + &mut worker, + &[ + ClipboardFormatId::CF_UNICODETEXT, + ClipboardFormatId::CF_DIB, + ClipboardFormatId::CF_DIBV5, + ], + ); + let text = paste(&mut worker, &os, TEXT); + let image = paste(&mut worker, &os, PNG); + assert_eq!(recorder.take().len(), 1); + text_response(&mut worker, "text"); + assert_eq!(text.recv().unwrap(), Some(b"text".to_vec())); + assert!(matches!( + recorder.take().as_slice(), + [ClipboardMessage::SendInitiatePaste(ClipboardFormatId::CF_DIBV5)] + )); + worker.handle(Command::RemoteData(Some(bitmap::png_to_cf_dibv5(&png()).unwrap()))); + let decoded = image.recv().unwrap().unwrap(); + assert_eq!( + bitmap::png_to_cf_dibv5(&decoded).unwrap(), + bitmap::png_to_cf_dibv5(&png()).unwrap() + ); +} + +#[test] +fn a_timed_out_request_is_drained_before_another_can_start() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let old = paste(&mut worker, &os, TEXT); + recorder.take(); + worker.expire(Instant::now() + PASTE_TIMEOUT); + assert_eq!(old.recv().unwrap(), None); + remote(&mut worker, &[ClipboardFormatId::CF_DIB]); + let blocked = paste(&mut worker, &os, PNG); + assert_eq!(blocked.recv().unwrap(), None); + assert!( + recorder.take().is_empty(), + "timeout must not start an ambiguous second request" + ); + text_response(&mut worker, "late old response"); + let current = paste(&mut worker, &os, PNG); + assert!(matches!( + recorder.take().as_slice(), + [ClipboardMessage::SendInitiatePaste(ClipboardFormatId::CF_DIB)] + )); + worker.handle(Command::RemoteData(Some(bitmap::png_to_cf_dib(&png()).unwrap()))); + assert!(current.recv().unwrap().is_some()); +} + +#[test] +fn replaced_remote_selection_discards_the_old_response() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let old = paste(&mut worker, &os, TEXT); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + assert_eq!(old.recv().unwrap(), None); + let new = paste(&mut worker, &os, TEXT); + assert_eq!(recorder.take().len(), 1); + text_response(&mut worker, "old"); + assert_eq!( + recorder.take().len(), + 1, + "new request starts only after old reply is drained" + ); + text_response(&mut worker, "new"); + assert_eq!(new.recv().unwrap(), Some(b"new".to_vec())); +} + +#[test] +fn local_selection_change_cancels_waiters_and_discards_remote_data() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let response = paste(&mut worker, &os, TEXT); + worker.handle(Command::LocalChanged); + assert_eq!(response.recv().unwrap(), None); + recorder.take(); + text_response(&mut worker, "stale"); + assert!(recorder.take().is_empty()); + assert_eq!( + os.lock().unwrap().offers.len(), + 1, + "a response cannot reclaim the clipboard" + ); +} + +#[test] +fn stale_os_paste_and_unadvertised_formats_are_refused() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let generation = os.lock().unwrap().offers[0].1; + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let (sender, receiver) = mpsc::channel(); + worker.handle(Command::Paste(Transfer::new(TEXT.into(), generation, move |data| { + sender.send(data).unwrap(); + }))); + assert_eq!(receiver.recv().unwrap(), None); + assert_eq!(paste(&mut worker, &os, PNG).recv().unwrap(), None); + assert!(recorder.take().is_empty()); +} + +#[test] +fn failed_remote_response_closes_the_os_paste() { + let (mut worker, os, _) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let response = paste(&mut worker, &os, TEXT); + worker.handle(Command::RemoteData(None)); + assert_eq!(response.recv().unwrap(), None); +} + +#[test] +fn pending_os_pastes_are_bounded() { + let (mut worker, os, _) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let waiters: Vec<_> = core::iter::repeat_with(|| paste(&mut worker, &os, TEXT)) + .take(16) + .collect(); + assert_eq!(paste(&mut worker, &os, TEXT).recv().unwrap(), None); + text_response(&mut worker, "bounded"); + for waiter in waiters { + assert_eq!(waiter.recv().unwrap(), Some(b"bounded".to_vec())); + } +} + +#[test] +fn image_dimensions_are_checked_by_the_shared_bitmap_converter() { + let (mut worker, os, _) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_DIB]); + let response = paste(&mut worker, &os, PNG); + let mut header = vec![0u8; 40]; + header[0..4].copy_from_slice(&40u32.to_le_bytes()); + header[4..8].copy_from_slice(&100_000i32.to_le_bytes()); + header[8..12].copy_from_slice(&100_000i32.to_le_bytes()); + header[12..14].copy_from_slice(&1u16.to_le_bytes()); + header[14..16].copy_from_slice(&32u16.to_le_bytes()); + worker.handle(Command::RemoteData(Some(header))); + assert_eq!(response.recv().unwrap(), None); +} + +#[test] +fn local_content_is_read_only_when_the_peer_requests_it() { + let (mut worker, os, recorder) = worker(); + os.lock() + .unwrap() + .data + .insert(TEXT.into(), "local \u{03bb}".as_bytes().to_vec()); + os.lock().unwrap().data.insert(PNG.into(), png()); + worker.handle(Command::AdvertiseLocal); + recorder.take(); + assert_eq!(os.lock().unwrap().reads, 0); + worker.handle(Command::RenderLocal(ClipboardFormatId::CF_UNICODETEXT)); + assert!( + matches!(recorder.take().as_slice(), [ClipboardMessage::SendFormatData(response)] if response.to_unicode_string().unwrap() == "local \u{03bb}") + ); + worker.handle(Command::RenderLocal(ClipboardFormatId::CF_DIBV5)); + assert!( + matches!(recorder.take().as_slice(), [ClipboardMessage::SendFormatData(response)] if bitmap::validate_dibv5(response.data()).is_ok()) + ); + assert_eq!(os.lock().unwrap().reads, 2); +} + +#[test] +fn shared_loop_detector_rejects_a_clipboard_manager_echo() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let response = paste(&mut worker, &os, TEXT); + text_response(&mut worker, "echoed"); + response.recv().unwrap(); + os.lock().unwrap().data.insert(TEXT.into(), b"echoed".to_vec()); + worker.handle(Command::LocalChanged); + recorder.take(); + worker.handle(Command::RenderLocal(ClipboardFormatId::CF_UNICODETEXT)); + assert!(matches!(recorder.take().as_slice(), [ClipboardMessage::SendFormatData(response)] if response.is_error())); + os.lock().unwrap().data.insert(TEXT.into(), b"different".to_vec()); + worker.handle(Command::RenderLocal(ClipboardFormatId::CF_UNICODETEXT)); + assert!(matches!(recorder.take().as_slice(), [ClipboardMessage::SendFormatData(response)] if !response.is_error())); +} + +#[test] +fn channel_reset_cancels_pending_pastes_and_releases_ownership() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + let response = paste(&mut worker, &os, TEXT); + worker.handle(Command::Reset); + assert_eq!(response.recv().unwrap(), None); + assert_eq!(os.lock().unwrap().clears, 1); + recorder.take(); + worker.handle(Command::LocalChanged); + assert!(recorder.take().is_empty()); +} + +/// Run against a private X server: DISPLAY=:N cargo test -p ironrdp-testsuite-core +/// x11_delayed_pastes -- --ignored --test-threads=1. +#[test] +#[ignore = "requires a dedicated X11 display"] +fn x11_delayed_pastes_use_notifications_and_incremental_transfers() { + use core::time::Duration; + + use ironrdp_cliprdr_native::linux::x11::X11Clipboard; + + let (owner_events, owner_receiver) = mpsc::channel(); + let mut owner = X11Clipboard::open(owner_events).unwrap(); + let (reader_events, reader_receiver) = mpsc::channel(); + let reader = X11Clipboard::open(reader_events).unwrap(); + owner.offer(&[TEXT.into(), PNG.into()], 42).unwrap(); + loop { + if let Command::LocalChanged = reader_receiver.recv_timeout(Duration::from_secs(5)).unwrap() { + let mimes = reader.mime_types(); + if mimes.contains(&TEXT.to_owned()) && mimes.contains(&PNG.to_owned()) { + break; + } + } + } + assert!( + !owner_receiver + .try_iter() + .any(|event| matches!(event, Command::Paste(_))) + ); + let thread = std::thread::spawn(move || { + let mut reader = reader; + let result = reader.read(TEXT); + (reader, result) + }); + let transfer = loop { + if let Command::Paste(transfer) = owner_receiver.recv_timeout(Duration::from_secs(5)).unwrap() { + break transfer; + } + }; + assert_eq!(transfer.generation(), 42); + // Exceeds the core X11 request length and must use INCR in both directions. + let bytes = "clipboard \u{03bb}\n".repeat(30_000).into_bytes(); + transfer.finish(Some(bytes.clone())); + let (mut reader, received) = thread.join().unwrap(); + assert_eq!(received.unwrap(), Some(bytes)); + + // Becoming the owner must not report our own selection back as a local copy. + assert!( + !owner_receiver + .try_iter() + .any(|event| matches!(event, Command::LocalChanged)) + ); + reader.offer(&[TEXT.into()], 43).unwrap(); + loop { + if let Command::LocalChanged = owner_receiver.recv_timeout(Duration::from_secs(5)).unwrap() { + if owner.mime_types().contains(&TEXT.to_owned()) { + break; + } + } + } + assert_eq!(owner.mime_types(), [TEXT]); +} + +#[test] +#[ignore = "requires a dedicated X11 display and xclip"] +fn x11_clipboard_interoperates_with_xclip_in_both_directions() { + use core::time::Duration; + use std::io::Write as _; + use std::process::{Command as ProcessCommand, Stdio}; + + use ironrdp_cliprdr_native::linux::x11::X11Clipboard; + + let (events, receiver) = mpsc::channel(); + let mut clipboard = X11Clipboard::open(events).unwrap(); + clipboard.offer(&[TEXT.into()], 1).unwrap(); + // The read is processed after the ownership command on the X11 worker. + // Our own selection is not read back, but this waits for the claim to finish. + assert!(clipboard.read(TEXT).unwrap().is_none()); + let reader = std::thread::spawn(|| { + ProcessCommand::new("xclip") + .args(["-selection", "clipboard", "-out"]) + .output() + .unwrap() + }); + let transfer = loop { + if let Command::Paste(transfer) = receiver.recv_timeout(Duration::from_secs(5)).unwrap() { + break transfer; + } + }; + let data = "remote unicode \u{03bb}\n".repeat(20_000).into_bytes(); + transfer.finish(Some(data.clone())); + // ICCCM requires an already-started INCR transfer to finish even after + // the owner relinquishes the selection. + clipboard.clear().unwrap(); + let output = reader.join().unwrap(); + assert!(output.status.success(), "{}", String::from_utf8_lossy(&output.stderr)); + assert_eq!(output.stdout, data); + + let mut writer = ProcessCommand::new("xclip") + .args(["-selection", "clipboard", "-in", "-quiet"]) + .stdin(Stdio::piped()) + .stderr(Stdio::null()) + .spawn() + .unwrap(); + let mut stdin = writer.stdin.take().unwrap(); + stdin.write_all(&data).unwrap(); + drop(stdin); + loop { + if let Command::LocalChanged = receiver.recv_timeout(Duration::from_secs(5)).unwrap() { + if clipboard.mime_types().contains(&TEXT.to_owned()) { + break; + } + } + } + let result = clipboard.read(TEXT); + let _ = writer.kill(); + let _ = writer.wait(); + assert_eq!(result.unwrap(), Some(data)); +} + +#[test] +fn queued_local_event_does_not_cancel_an_acknowledged_remote_selection() { + let (mut worker, os, recorder) = worker(); + remote(&mut worker, &[ClipboardFormatId::CF_UNICODETEXT]); + os.lock().unwrap().owns = true; + worker.handle(Command::LocalChanged); + let response = paste(&mut worker, &os, TEXT); + assert!(matches!( + recorder.take().as_slice(), + [ClipboardMessage::SendInitiatePaste(ClipboardFormatId::CF_UNICODETEXT)] + )); + text_response(&mut worker, "current remote"); + assert_eq!(response.recv().unwrap(), Some(b"current remote".to_vec())); +} + +#[test] +fn local_notifications_refresh_the_current_selection_formats() { + let (mut worker, os, recorder) = worker(); + worker.handle(Command::AdvertiseLocal); + recorder.take(); + os.lock() + .unwrap() + .data + .insert(TEXT.into(), b"newest selection".to_vec()); + worker.handle(Command::LocalChanged); + assert!( + matches!(recorder.take().as_slice(), [ClipboardMessage::SendInitiateCopy(formats)] if formats.len() == 1 && formats[0].id() == ClipboardFormatId::CF_UNICODETEXT) + ); +} + +#[test] +#[ignore = "requires a dedicated X11 display"] +fn x11_bounds_requests_before_queuing_clipboard_data() { + use core::time::Duration; + use ironrdp_cliprdr_native::linux::x11::X11Clipboard; + + let (owner_events, owner_receiver) = mpsc::channel(); + let mut owner = X11Clipboard::open(owner_events).unwrap(); + owner.offer(&[TEXT.into()], 1).unwrap(); + let (results, received) = mpsc::channel(); + let readers: Vec<_> = core::iter::repeat_with(|| { + let results = results.clone(); + std::thread::spawn(move || { + let (events, notifications) = mpsc::channel(); + let mut reader = X11Clipboard::open(events).unwrap(); + loop { + if let Command::LocalChanged = notifications.recv_timeout(Duration::from_secs(5)).unwrap() { + if reader.mime_types().contains(&TEXT.to_owned()) { + break; + } + } + } + results.send(reader.read(TEXT)).unwrap(); + }) + }) + .take(17) + .collect(); + let mut held = Vec::new(); + while held.len() < 16 { + if let Command::Paste(transfer) = owner_receiver.recv_timeout(Duration::from_secs(5)).unwrap() { + held.push(transfer); + } + } + assert!( + received.recv_timeout(Duration::from_secs(2)).unwrap().is_err(), + "the seventeenth read is refused before a Paste is queued" + ); + assert!( + !owner_receiver + .try_iter() + .any(|event| matches!(event, Command::Paste(_))) + ); + for transfer in held { + transfer.finish(None); + } + for _ in 0..16 { + assert!(received.recv_timeout(Duration::from_secs(5)).unwrap().is_err()); + } + for reader in readers { + reader.join().unwrap(); + } +} diff --git a/crates/ironrdp-testsuite-core/tests/cliprdr_native/mod.rs b/crates/ironrdp-testsuite-core/tests/cliprdr_native/mod.rs new file mode 100644 index 0000000000..89e6075b5b --- /dev/null +++ b/crates/ironrdp-testsuite-core/tests/cliprdr_native/mod.rs @@ -0,0 +1,2 @@ +mod data_control; +mod linux; diff --git a/crates/ironrdp-testsuite-core/tests/main.rs b/crates/ironrdp-testsuite-core/tests/main.rs index 161c01bed0..107932172d 100644 --- a/crates/ironrdp-testsuite-core/tests/main.rs +++ b/crates/ironrdp-testsuite-core/tests/main.rs @@ -14,6 +14,8 @@ mod cfg; mod clipboard; +#[cfg(target_os = "linux")] +mod cliprdr_native; mod connector; mod displaycontrol; mod dvc;