From 85be787d9b326b37943780298932af3bfebe8757 Mon Sep 17 00:00:00 2001 From: Greg Lamberson Date: Tue, 29 Sep 2026 22:29:34 -0500 Subject: [PATCH] feat(cliprdr): add a Wayland data-control clipboard client Add the data_control module to ironrdp-cliprdr-native so a Linux CLIPRDR backend can read and set the clipboard without a window. It speaks ext-data-control-v1 and wlr-data-control-unstable-v1 and uses whichever the compositor offers. One thread owns the Wayland connection and DataControl is the handle to call from anywhere. A type advertised without data raises a TransferRequest when something pastes it, which maps onto the CLIPRDR Format Data Response. Dropping the request without answering it closes the paste with no data. A selection the client set itself isn't reported back as a local copy. It is told from another client's copy by a private MIME type that the client advertises next to the real ones. Reads are capped at 100 MiB and time out when the source stalls for 5 seconds. Pastes are bounded as well. At most 16 are served at once, a paste whose data isn't supplied within 5 seconds is closed with no data, and a reader that takes nothing for 5 seconds is dropped. The state tests live in ironrdp-testsuite-extra behind the __test feature, because inline tests aren't built for this crate. The checks that need a compositor are there too, as ignored tests. --- Cargo.lock | 6 + crates/ironrdp-cliprdr-native/Cargo.toml | 12 + crates/ironrdp-cliprdr-native/README.md | 10 +- .../src/data_control/client.rs | 460 ++++++++++ .../src/data_control/dispatch.rs | 176 ++++ .../src/data_control/error.rs | 76 ++ .../src/data_control/mime.rs | 67 ++ .../src/data_control/mod.rs | 72 ++ .../src/data_control/options.rs | 111 +++ .../src/data_control/state.rs | 867 ++++++++++++++++++ .../src/data_control/worker.rs | 264 ++++++ crates/ironrdp-cliprdr-native/src/lib.rs | 3 + crates/ironrdp-testsuite-extra/Cargo.toml | 3 + .../tests/cliprdr_native/data_control.rs | 683 ++++++++++++++ .../tests/cliprdr_native/live.rs | 262 ++++++ .../tests/cliprdr_native/mod.rs | 2 + crates/ironrdp-testsuite-extra/tests/main.rs | 2 + 17 files changed, 3075 insertions(+), 1 deletion(-) create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/client.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/dispatch.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/error.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/mime.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/mod.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/options.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/state.rs create mode 100644 crates/ironrdp-cliprdr-native/src/data_control/worker.rs create mode 100644 crates/ironrdp-testsuite-extra/tests/cliprdr_native/data_control.rs create mode 100644 crates/ironrdp-testsuite-extra/tests/cliprdr_native/live.rs create mode 100644 crates/ironrdp-testsuite-extra/tests/cliprdr_native/mod.rs diff --git a/Cargo.lock b/Cargo.lock index 7e3280ae94..c81363e445 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2816,7 +2816,12 @@ version = "0.7.0" dependencies = [ "ironrdp-cliprdr", "ironrdp-core 0.2.1", + "nix", "tracing", + "visibility", + "wayland-client", + "wayland-protocols", + "wayland-protocols-wlr", "windows", ] @@ -3450,6 +3455,7 @@ dependencies = [ "ironrdp-bulk", "ironrdp-cfg", "ironrdp-client", + "ironrdp-cliprdr-native", "ironrdp-core 0.2.1", "ironrdp-daemon", "ironrdp-dvc", diff --git a/crates/ironrdp-cliprdr-native/Cargo.toml b/crates/ironrdp-cliprdr-native/Cargo.toml index c4c55328ce..c6b1c89bc1 100644 --- a/crates/ironrdp-cliprdr-native/Cargo.toml +++ b/crates/ironrdp-cliprdr-native/Cargo.toml @@ -12,6 +12,11 @@ authors.workspace = true keywords.workspace = true categories.workspace = true +[features] +# Internal (PRIVATE!) feature used to expose internals to the shared integration tests. +# Don't rely on this whatsoever. It may disappear at any time. +__test = ["dep:visibility"] + [lib] doctest = false test = false @@ -21,6 +26,13 @@ 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 = ["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"] } + [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..22212e02b1 100644 --- a/crates/ironrdp-cliprdr-native/README.md +++ b/crates/ironrdp-cliprdr-native/README.md @@ -1,7 +1,15 @@ # IronRDP CLIPRDR native backends -Native CLIPRDR backend implementations. Currently only Windows is supported. +Native CLIPRDR backend implementations. Windows has a full backend; Linux has a Wayland clipboard client that a backend can be built on. This crate is part of the [IronRDP] project. [IronRDP]: https://github.com/Devolutions/IronRDP + +## Linux + +The `data_control` module is a clipboard client for Wayland compositors that offer [`ext-data-control-v1`] or [`wlr-data-control-unstable-v1`]. +It reads and sets the clipboard without a window and supports delayed rendering, so the data for a paste can be produced when the paste happens. + +[`ext-data-control-v1`]: https://gitlab.freedesktop.org/wayland/wayland-protocols/-/blob/main/staging/ext-data-control/ext-data-control-v1.xml +[`wlr-data-control-unstable-v1`]: https://gitlab.freedesktop.org/wlroots/wlr-protocols/-/blob/master/unstable/wlr-data-control-unstable-v1.xml 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..8ca60e0376 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/client.rs @@ -0,0 +1,460 @@ +//! 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, +}; + +use nix::poll::{PollFd, PollFlags, PollTimeout, poll}; + +use super::{ + error::{Error, Result}, + mime::find_mime_match, + options::{Options, Protocol}, + state::{Command, Shared, lock_shared}, + worker, +}; + +/// How long a read waits for the source to deliver more data. +const READ_IDLE_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 the request without answering counts as +/// failing it, so the paste is closed at once and the next paste raises a new +/// request. A paste that waits more than five seconds for an answer is closed +/// with no data, but an answer that comes later is still kept and serves the +/// next paste of that type. +#[derive(Debug)] +#[must_use = "a paste gets no data until the request is answered"] +pub struct TransferRequest { + serial: u32, + mime_type: String, + link: Arc, + answered: bool, +} + +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(mut self, data: impl Into>) -> Result<()> { + self.answered = true; + 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(mut self) -> Result<()> { + self.answered = true; + self.link.send(Command::CompleteTransfer { + serial: self.serial, + data: None, + }) + } +} + +impl Drop for TransferRequest { + fn drop(&mut self) { + if self.answered { + return; + } + // A worker that has already stopped has no paste left to close. + let _ = 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`]. + /// + /// Blocks until the compositor has answered, so call it from a thread that may block. + /// + /// # 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. + /// + /// Blocks until the compositor has answered, so call it from a thread that may block. + /// + /// # 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), + }) + } + + /// The protocol the compositor is speaking. + #[must_use] + pub fn protocol(&self) -> Protocol { + self.protocol + } + + fn shared(&self) -> MutexGuard<'_, Shared> { + lock_shared(&self.shared) + } + + /// 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 + } + + /// Whether this handle owns the clipboard selection. + /// + /// It turns `true` when the worker processes [`set_selection`](Self::set_selection). It turns + /// `false` when another client takes the selection, when the selection is cleared, and when + /// this handle gives it up. Like [`serial`](Self::serial), it reports what the worker has seen + /// so far, so it can trail a request that is still queued. + #[must_use] + pub fn owns_selection(&self) -> bool { + self.shared().own_source_live + } + + /// 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). Clearing the + /// selection is reported as an empty list, even when this handle cleared + /// it with [`clear_selection`](Self::clear_selection). + /// + /// 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), + answered: false, + }); + })); + } + + /// Become the clipboard owner and offer `content`. + /// + /// The selection advertises one private type next to the types of `content`, which marks it as + /// this client's own. Other clients see that type among the types they are offered, and a request + /// for it gets no data. This client leaves it out of the types it reports to the callbacks and + /// lists in [`selection_mime_types`](Self::selection_mime_types). + /// + /// # 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. + /// + /// Blocks until the source has delivered everything, so call it from a thread that may block. + /// The timeout measures silence, not the total time, so a source that keeps delivering data is + /// only stopped by the size limit. + /// + /// # Errors + /// + /// [`Error::SelectionChanged`] if the selection changed before the data + /// arrived, [`Error::TooLarge`] past the size limit, [`Error::Timeout`] if + /// the source stalls for 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, serial) = { + let shared = self.shared(); + (shared.mime_types.clone(), shared.serial) + }; + 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), + })?; + + let mut data = Vec::new(); + let mut chunk = vec![0u8; 64 * 1024]; + loop { + let mut fds = [PollFd::new(reader.as_fd(), PollFlags::POLLIN)]; + let timeout = PollTimeout::try_from(READ_IDLE_TIMEOUT).unwrap_or(PollTimeout::NONE); + 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, unless the selection changed + // under the request, which leaves the pipe empty or holding another selection's data. + Ok(0) if self.serial() != serial => return Err(Error::SelectionChanged), + Ok(0) => return Ok(Some(data)), + Ok(n) => { + if self.max_read_bytes < data.len() + n { + return Err(Error::TooLarge { + size: data.len() + n, + limit: self.max_read_bytes, + }); + } + data.extend_from_slice(&chunk[..n]); + } + Err(error) if error.kind() == std::io::ErrorKind::Interrupted => {} + Err(error) => return Err(error.into()), + } + } + } + + /// Test-only: a handle with no worker thread. The caller receives the commands the handle + /// sends and plays the worker. + /// + /// # Errors + /// + /// [`Error::Io`] if the wake-up socket cannot be created. + #[cfg(feature = "__test")] + #[doc(hidden)] + pub fn without_worker(shared: Arc>) -> Result<(Self, mpsc::Receiver)> { + let (commands, receiver) = mpsc::channel(); + // The other end is dropped, so a wake-up write fails, and `Link::notify` ignores that. + let (wake, _peer) = UnixStream::pair()?; + wake.set_nonblocking(true)?; + let handle = Self { + link: Arc::new(Link { commands, wake }), + shared, + protocol: Protocol::Ext, + max_read_bytes: Options::new().max_read_bytes, + stop: Arc::new(AtomicBool::new(false)), + thread: None, + }; + Ok((handle, receiver)) + } +} + +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(); + } + } +} 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..9aef2b4151 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/dispatch.rs @@ -0,0 +1,176 @@ +//! 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 their handlers are +//! generated from one macro and forward to the shared [`State`]. + +use wayland_client::{ + Connection, Dispatch, 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`. + } +} + +/// Implement the handlers of one data-control protocol. +/// +/// `protocol` names it in log messages, each object is given with the module that holds its +/// events, and `on_data_offer` is the `State` method that takes the protocol's new offer. +macro_rules! dispatch_data_control { + ( + protocol: $protocol:literal, + manager: $manager:ty, + device: $device:ty => $device_events:ident, + source: $source:ty => $source_events:ident, + offer: $offer:ty => $offer_events:ident, + on_data_offer: $on_data_offer:ident $(,)? + ) => { + impl Dispatch<$manager, ()> for Client { + fn event( + _state: &mut Self, + _proxy: &$manager, + _event: <$manager as wayland_client::Proxy>::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The manager has no events. + } + } + + impl Dispatch<$device, ()> for Client { + fn event( + state: &mut Self, + _proxy: &$device, + event: <$device as wayland_client::Proxy>::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + $device_events::Event::DataOffer { id } => { + state.data_control.$on_data_offer(id); + } + $device_events::Event::Selection { id } => { + if id.is_some() { + state.data_control.on_selection(); + } else { + state.data_control.on_selection_cleared(); + } + } + $device_events::Event::Finished => { + state.data_control.on_device_finished(); + } + $device_events::Event::PrimarySelection { .. } => { + // Only the regular clipboard is handled, not the primary selection. + tracing::trace!(protocol = $protocol, "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, $device, [ + $device_events::EVT_DATA_OFFER_OPCODE => ($offer, ()), + ]); + } + + impl Dispatch<$source, ()> for Client { + fn event( + state: &mut Self, + _proxy: &$source, + event: <$source as wayland_client::Proxy>::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + $source_events::Event::Send { mime_type, fd } => { + state.data_control.on_source_send(&mime_type, fd); + } + $source_events::Event::Cancelled => { + state.data_control.on_source_cancelled(); + } + _ => {} + } + } + } + + impl Dispatch<$offer, ()> for Client { + fn event( + state: &mut Self, + _proxy: &$offer, + event: <$offer as wayland_client::Proxy>::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if let $offer_events::Event::Offer { mime_type } = event { + state.data_control.on_offer_mime_type(mime_type); + } + } + } + }; +} + +dispatch_data_control! { + protocol: "ext", + manager: ExtDataControlManagerV1, + device: ExtDataControlDeviceV1 => ext_data_control_device_v1, + source: ExtDataControlSourceV1 => ext_data_control_source_v1, + offer: ExtDataControlOfferV1 => ext_data_control_offer_v1, + on_data_offer: on_data_offer_ext, +} + +dispatch_data_control! { + protocol: "wlr", + manager: ZwlrDataControlManagerV1, + device: ZwlrDataControlDeviceV1 => zwlr_data_control_device_v1, + source: ZwlrDataControlSourceV1 => zwlr_data_control_source_v1, + offer: ZwlrDataControlOfferV1 => zwlr_data_control_offer_v1, + on_data_offer: on_data_offer_wlr, +} 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..48ccfd9d26 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/error.rs @@ -0,0 +1,76 @@ +//! 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, + /// The selection changed while it was being read, so the data is missing or belongs to + /// another selection. + SelectionChanged, + /// 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::SelectionChanged => f.write_str("the clipboard selection changed while it was being read"), + 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..8abbf11e94 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/mime.rs @@ -0,0 +1,67 @@ +//! 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> { + let mut available = available.iter().map(String::as_str); + if let Some(found) = find_charset_variant(requested, available.clone()) { + return Some(found); + } + + if !requested.starts_with("text/") { + return None; + } + let base = requested.split(';').next()?; + + // The last resort: another charset of the same base type. + available.find(|m| m.split(';').next() == Some(base)) +} + +/// Find the entry of `available` that is `requested` apart from its charset parameter. +/// +/// An exact match wins. For `text/` types the charset is then stripped from the request, or +/// `;charset=utf-8` is added to it, and the result is looked up. An entry in another charset never +/// matches, because its bytes are not what the request asks for. [`find_mime_match`] adds that as +/// a last resort for reading, where the caller takes what the source offers. +pub(crate) fn find_charset_variant<'a>( + requested: &str, + mut available: impl Iterator + Clone, +) -> Option<&'a str> { + if let Some(found) = available.clone().find(|m| *m == requested) { + return Some(found); + } + + if !requested.starts_with("text/") { + return None; + } + + if let Some((base, _charset)) = requested.split_once(';') { + // The request carries a charset, so try the bare type. + return available.find(|m| *m == base); + } + + // The request has none, so try the common charset spellings. + [";charset=utf-8", ";charset=UTF-8"].into_iter().find_map(|suffix| { + let with_charset = format!("{requested}{suffix}"); + available.clone().find(|m| *m == with_charset) + }) +} 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..2d922c1f6c --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/mod.rs @@ -0,0 +1,72 @@ +//! 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 [2.2.5.2] Format Data Response PDU 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(()) +//! # } +//! ``` +//! +//! A paste that waits more than five seconds for its data is closed with no +//! data, and an answer that comes later is kept for the next paste. A reader +//! that takes nothing for five seconds is dropped too, and at most 16 pastes +//! are served at once, so a stuck program cannot pile up threads. +//! +//! A selection set here also advertises one private type, so the client can tell it from another +//! client's copy. Other clients see that type among the types they are offered. +//! +//! GNOME's Mutter offers no data-control protocol, so [`DataControl::connect`] +//! returns [`Error::Unsupported`] there. +//! +//! [`ext-data-control-v1`]: https://gitlab.freedesktop.org/wayland/wayland-protocols/-/blob/main/staging/ext-data-control/ext-data-control-v1.xml +//! [`wlr-data-control-unstable-v1`]: https://gitlab.freedesktop.org/wlroots/wlr-protocols/-/blob/master/unstable/wlr-data-control-unstable-v1.xml +//! [2.2.5.2]: https://learn.microsoft.com/en-us/openspecs/windows_protocols/ms-rdpeclip/28c193b8-4cec-413e-a07b-9235e5e15f6b + +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; + +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..7b9c649729 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/options.rs @@ -0,0 +1,111 @@ +//! 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 wayland-protocols staging protocol. + /// + /// [`ext-data-control-v1`]: https://gitlab.freedesktop.org/wayland/wayland-protocols/-/blob/main/staging/ext-data-control/ext-data-control-v1.xml + Ext, + /// [`wlr-data-control-unstable-v1`], the wlroots protocol. + /// + /// [`wlr-data-control-unstable-v1`]: https://gitlab.freedesktop.org/wlroots/wlr-protocols/-/blob/master/unstable/wlr-data-control-unstable-v1.xml + 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..a7f61ad322 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/state.rs @@ -0,0 +1,867 @@ +//! 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). At most +//! `MAX_PASTES` pastes are in flight, and a paste beyond that is closed with +//! no data. A held paste is closed with no data after `PASTE_TIMEOUT`, and +//! a writer gives up when its reader takes nothing for that long. +//! +//! **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 event is the echo of a selection 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::{ + hash::BuildHasher as _, + sync::atomic::{AtomicUsize, Ordering}, + time::Duration, +}; +use std::{ + collections::HashMap, + fs::File, + hash::RandomState, + io::{self, Write as _}, + os::unix::io::{AsFd as _, OwnedFd}, + sync::{Arc, Mutex, MutexGuard, PoisonError}, + time::Instant, +}; + +use nix::{ + errno::Errno, + libc::PIPE_BUF, + poll::{PollFd, PollFlags, PollTimeout, poll}, +}; +use wayland_client::{Dispatch, QueueHandle, 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, +}; + +use super::mime::find_charset_variant; + +// The enums below wrap the 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) { + 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 sent from clipboard backends to the Wayland event loop thread. +#[derive(Debug)] +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) enum Command { + /// 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. +/// +/// The fields are `pub` so that the integration tests can read and set them. The struct is only +/// public under the `__test` feature, so otherwise they reach no further than the crate. +#[derive(Default)] +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) struct Shared { + /// MIME types of the current compositor selection. + pub mime_types: Vec, + /// Serial number, incremented on each selection change. + pub serial: u32, + /// Change notification callback. + /// + /// Called on the event loop thread when selection changes. + pub 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 on_transfer: Option>, + /// Whether our own source is still the compositor's selection source. + /// + /// It is set when a selection is set. It is cleared when another client + /// takes the selection, when the selection is cleared, and when the source + /// is cancelled or given up. [`DataControl::owns_selection`] reads it. + /// + /// [`DataControl::owns_selection`]: super::DataControl::owns_selection + pub own_source_live: bool, +} + +/// Lock the shared state, recovering from a poisoned lock. +/// +/// `Shared` holds plain fields that are each valid on their own, so after a panic while the lock +/// was held the state is still right to read and update. Skipping the update instead would leave +/// the worker and the handle disagreeing about the clipboard. +pub(crate) fn lock_shared(shared: &Mutex) -> MutexGuard<'_, Shared> { + shared.lock().unwrap_or_else(PoisonError::into_inner) +} + +/// The most pastes that can be waiting for data or being written at once. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) const MAX_PASTES: usize = 16; + +/// How long a paste waits for data from `on_transfer`, and how long a reader +/// may take nothing before the writer answering it gives up. +const PASTE_TIMEOUT: Duration = Duration::from_secs(5); + +/// One of the [`MAX_PASTES`] places, given back when the paste holding it ends. +struct PastePlace(Arc); + +impl PastePlace { + /// Take a free place, or `None` when all of them are in use. + fn take(in_flight: &Arc) -> Option { + in_flight + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |count| { + (count < MAX_PASTES).then_some(count + 1) + }) + .ok() + .map(|_| Self(Arc::clone(in_flight))) + } +} + +impl Drop for PastePlace { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::Relaxed); + } +} + +/// A paste in flight: the write end of the pipe its reader waits on, and the +/// place it holds until it ends. +struct Paste { + fd: OwnedFd, + place: PastePlace, +} + +/// A paste waiting on data from `on_transfer`, with every paste of the same +/// MIME type that arrived meanwhile. +struct PendingTransfer { + serial: u32, + /// The waiting pastes, each with the time it stops waiting. + waiters: Vec<(Paste, Instant)>, +} + +/// 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, + /// The private type that marks this client's own selections. + own_marker: OwnMarker, + /// 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, + /// How many pastes hold a place, counting those waiting and those being written. + pastes_in_flight: Arc, + /// Next serial handed to `on_transfer`. + next_transfer_serial: u32, + /// 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(), + own_marker: OwnMarker::random(), + source_data: HashMap::new(), + pending_offer: None, + pending_transfers: HashMap::new(), + pastes_in_flight: Arc::default(), + next_transfer_serial: 1, + 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 type that marks this client's own selections. + #[cfg(feature = "__test")] + pub fn own_marker(&self) -> &str { + self.own_marker.as_str() + } + + /// Test-only: the shared state the callbacks and readers see. + #[cfg(feature = "__test")] + pub fn shared(&self) -> &Arc> { + &self.shared_state + } + + /// 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) { + lock_shared(&self.shared_state).own_source_live = live; + } + + fn set_pending_offer(&mut self, offer: DataControlOffer) { + // An offer that never became the selection still has to be destroyed explicitly. + 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) { + // The offer of the selection being replaced is destroyed explicitly as well. + if let Some(old) = self.current_offer.take() { + old.destroy(); + } + + let mut mime_types = if let Some((offer, pending)) = self.pending_offer.take() { + let types = pending.mime_types; + self.current_offer = Some(offer); + types + } else { + // Without a pending offer there's nothing to promote, so the selection has no types. + Vec::new() + }; + + // The device reports every selection, our own included, and reporting + // our own as a change makes the consumer treat it as another client's + // copy. Our source advertises a private marker type next to the real + // ones, so a selection that carries it is ours and one that doesn't is + // another client's, in whichever order the compositor delivers them and + // whatever types they offer. + let own = self.own_marker.take_from(&mut mime_types); + tracing::debug!( + mime_types = ?mime_types, + own, + "Compositor selection changed" + ); + + // 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 = { + let mut shared = lock_shared(&self.shared_state); + shared.serial = shared.serial.wrapping_add(1); + shared.mime_types.clone_from(&mime_types); + // Another client's copy may arrive before the cancel for our source. Our marker doesn't prove + // that a source of ours is live either, because a clipboard manager can re-offer our content + // after the source is gone. + shared.own_source_live = own && self.current_source.is_some(); + 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(); + } + if let Some((offer, _)) = self.pending_offer.take() { + offer.destroy(); + } + + tracing::debug!("Compositor selection cleared"); + + let callback = { + let mut shared = lock_shared(&self.shared_state); + shared.serial = shared.serial.wrapping_add(1); + shared.mime_types.clear(); + shared.own_source_live = false; + 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`, or until `PASTE_TIMEOUT` is up. + /// With `MAX_PASTES` pastes already in flight the new one is closed with + /// no data. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_source_send(&mut self, mime_type: &str, fd: OwnedFd) { + if OwnMarker::is_marker(mime_type) { + // A client that reads every type it is offered asks for the marker too. It has no data to + // give, and that isn't worth a warning. + return; + } + + let Some(place) = PastePlace::take(&self.pastes_in_flight) else { + tracing::warn!(mime_type, "Too many pastes in flight, closing this one"); + return; + }; + let paste = Paste { fd, place }; + + if let Some(data) = self.cached_data(mime_type) { + write_in_background(paste, data); + return; + } + + let deadline = Instant::now() + PASTE_TIMEOUT; + if let Some(pending) = self.pending_transfers.get_mut(mime_type) { + pending.waiters.push((paste, deadline)); + return; + } + + let on_transfer = lock_shared(&self.shared_state).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, + "Paste closed with no data, the type is not advertised or no transfer callback is set" + ); + // Returning drops the fd, which closes the pipe, so the pasting client sees 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![(paste, deadline)], + }, + ); + tracing::debug!(mime_type, serial, "Paste held until its data arrives"); + on_transfer(serial, mime_type.to_owned()); + } + + /// Data cached for `mime_type`, or for the same type with another spelling of its + /// charset (compositors commonly request `text/plain;charset=utf-8` for `text/plain`, and + /// the reverse). + fn cached_data(&self, mime_type: &str) -> Option> { + let key = find_charset_variant(mime_type, self.source_data.keys().map(String::as_str))?; + self.source_data.get(key).cloned() + } + + /// Process a `CompleteTransfer` command. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn complete_transfer(&mut self, serial: u32, data: Option>) { + 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 data: Arc<[u8]> = data.into(); + for (paste, _) in pending.waiters { + write_in_background(paste, Arc::clone(&data)); + } + self.source_data.insert(mime_type, data); + } + + /// Close the held pastes whose time is up, and return when the next one is + /// due. + /// + /// A closed paste ends with no data. Its transfer stays pending, so a + /// paste that comes meanwhile joins it instead of raising a second + /// request, and an answer that arrives late is still kept for the next + /// paste. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn expire_transfers(&mut self, now: Instant) -> Option { + for pending in self.pending_transfers.values_mut() { + pending.waiters.retain(|(_, deadline)| now < *deadline); + } + self.pending_transfers + .values() + .flat_map(|pending| pending.waiters.iter().map(|(_, deadline)| *deadline)) + .min() + } + + /// 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(); + self.set_own_source_live(false); + } + + /// 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; + } + }; + + for mime_type in mime_types { + new_source.offer(mime_type); + } + // The marker tells the echo of this selection from another client's copy. + new_source.offer(self.own_marker.as_str()); + + 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. + // + // No event for the old source can reach the handlers afterwards. The + // worker drained its queue before running this command, and the backend + // swallows what arrives for a destroyed object. + if let Some(old) = self.current_source.take() { + old.destroy(); + } + self.pending_transfers.clear(); + self.source_data = data + .into_iter() + .map(|(mime_type, bytes)| (mime_type, 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"); + // The selection changed after the caller looked, so there's nothing to receive. + // Dropping the fd gives the caller EOF, and the changed serial tells it why. + } + } + } +} + +/// Type names that start like this belong to some client of this library. +const MARKER_PREFIX: &str = "application/x-ironrdp-source-"; + +/// The private type that marks the selections of one client. +/// +/// A source advertises it next to the real types, so a selection that carries it is content this +/// client put on the clipboard, in whichever order the compositor delivers it relative to other +/// clients' copies. A copy by another client doesn't carry it, whatever types the copy offers. Only a +/// client that re-offers every type of our selection, the marker included, passes it on. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) struct OwnMarker { + /// INVARIANT: starts with `MARKER_PREFIX`, which `take_from` relies on to remove our own marker. + name: String, +} + +impl OwnMarker { + /// A marker that no other client of this library uses, in this process or in another one. + /// + /// The suffix is random and carries no process id, so other clients learn nothing from it, and + /// content that a previous run left on the clipboard doesn't read as our own. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn random() -> Self { + // The standard library seeds `RandomState` from the operating system, which makes it a source of + // random numbers without a dependency. + let suffix = RandomState::new().hash_one(0u8); + Self { + name: format!("{MARKER_PREFIX}{suffix:016x}"), + } + } + + /// The type name to advertise. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn as_str(&self) -> &str { + &self.name + } + + /// Whether `mime_type` is the private marker of some client, ours or another's. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn is_marker(mime_type: &str) -> bool { + mime_type.starts_with(MARKER_PREFIX) + } + + /// Take every private marker out of `mime_types`, and say whether ours was among them. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn take_from(&self, mime_types: &mut Vec) -> bool { + let own = mime_types.contains(&self.name); + mime_types.retain(|mime_type| !Self::is_marker(mime_type)); + own + } +} + +/// 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(paste: Paste, data: Arc<[u8]>) { + let spawned = std::thread::Builder::new() + .name("data-control-send".into()) + .spawn(move || { + let Paste { fd, place } = paste; + let mut pipe = File::from(fd); + if let Err(error) = write_with_idle_timeout(&mut pipe, &data, PASTE_TIMEOUT) { + tracing::debug!(%error, "Paste writer stopped before the end of the data"); + } + // The place goes back before the pipe closes, so a reader that has + // seen the whole paste can paste again at once. + drop(place); + drop(pipe); + }); + if let Err(error) = spawned { + tracing::error!(%error, "Failed to start a paste writer"); + } +} + +/// Write `data` to a pipe, giving up when the reader takes nothing for `idle`. +/// +/// A plain `write_all` waits for as long as the reader does not read. Waiting +/// for the pipe to become writable instead lets the wait end. A pipe is +/// writable once `PIPE_BUF` bytes fit, so writing at most that much at a time +/// cannot block. A slow reader that keeps taking data is never cut off. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) fn write_with_idle_timeout(pipe: &mut File, data: &[u8], idle: Duration) -> io::Result<()> { + let timeout = PollTimeout::try_from(idle).unwrap_or(PollTimeout::MAX); + let mut rest = data; + while !rest.is_empty() { + let mut fds = [PollFd::new(pipe.as_fd(), PollFlags::POLLOUT)]; + match poll(&mut fds, timeout) { + Ok(0) => return Err(io::ErrorKind::TimedOut.into()), + Ok(_) => {} + Err(Errno::EINTR) => continue, + Err(errno) => return Err(errno.into()), + } + + let (chunk, _) = rest.split_at(rest.len().min(PIPE_BUF)); + let written = match pipe.write(chunk) { + Ok(0) => return Err(io::ErrorKind::WriteZero.into()), + Ok(written) => written, + Err(error) if error.kind() == io::ErrorKind::Interrupted => continue, + Err(error) => return Err(error), + }; + rest = rest.split_at(written).1; + } + 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..6cbc1ea2a7 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/worker.rs @@ -0,0 +1,264 @@ +//! 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}, + time::Duration, +}; +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, EventQueue, + backend::WaylandError, + globals::{GlobalList, registry_queue_init}, + protocol::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: &wayland_client::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::debug!( + 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; + } + + // INVARIANT: the event queue is empty when commands run, so no event for an object + // that a command replaces or destroys, such as our previous selection source, is + // still queued. + self.process_commands(); + let next_deadline = 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), + ]; + // Wake for the earliest held paste, so one that nobody answers is closed on time. + let timeout = next_deadline.map_or(PollTimeout::NONE, |deadline| { + // Rounding up keeps the wake-up from landing just before the deadline and spinning. + let wait = deadline.saturating_duration_since(Instant::now()) + Duration::from_millis(1); + PollTimeout::try_from(wait).unwrap_or(PollTimeout::MAX) + }); + match poll(&mut fds, timeout) { + 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::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, + } + } + } +} + +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..4ce1515840 100644 --- a/crates/ironrdp-cliprdr-native/src/lib.rs +++ b/crates/ironrdp-cliprdr-native/src/lib.rs @@ -15,6 +15,9 @@ mod windows; #[cfg(windows)] pub use crate::windows::{HWND, WinClipboard, WinCliprdrError, WinCliprdrResult}; +#[cfg(target_os = "linux")] +pub mod data_control; + mod stub; use std::sync::OnceLock; use std::time::Instant; diff --git a/crates/ironrdp-testsuite-extra/Cargo.toml b/crates/ironrdp-testsuite-extra/Cargo.toml index 3f5a348296..41dfc04762 100644 --- a/crates/ironrdp-testsuite-extra/Cargo.toml +++ b/crates/ironrdp-testsuite-extra/Cargo.toml @@ -66,6 +66,9 @@ tokio = { version = "1", features = ["sync", "time", "net", "rt", "rt-multi-thre tokio-rustls = { version = "0.26", default-features = false } uuid = { version = "1", features = ["v4"] } +[target.'cfg(target_os = "linux")'.dev-dependencies] +ironrdp-cliprdr-native = { path = "../ironrdp-cliprdr-native", features = ["__test"] } + [target.'cfg(windows)'.dev-dependencies] ironrdp-rdpdr-native.path = "../ironrdp-rdpdr-native" diff --git a/crates/ironrdp-testsuite-extra/tests/cliprdr_native/data_control.rs b/crates/ironrdp-testsuite-extra/tests/cliprdr_native/data_control.rs new file mode 100644 index 0000000000..6c22454706 --- /dev/null +++ b/crates/ironrdp-testsuite-extra/tests/cliprdr_native/data_control.rs @@ -0,0 +1,683 @@ +//! Tests for the Wayland data-control clipboard client that need no compositor. + +use core::time::Duration; +use std::{ + collections::HashSet, + io::{Read as _, Write as _}, + os::fd::OwnedFd, + sync::{Arc, Mutex, PoisonError, mpsc}, + time::Instant, +}; + +use ironrdp_cliprdr_native::data_control::state::{ + Command, MAX_PASTES, OwnMarker, Shared, State, write_with_idle_timeout, +}; +use ironrdp_cliprdr_native::data_control::{DataControl, Error, 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 a_selection_that_carries_our_marker_is_ours_and_loses_the_marker() { + let marker = OwnMarker::random(); + let mut types = mimes(&["text/plain;charset=utf-8", marker.as_str()]); + + assert!(marker.take_from(&mut types)); + + assert_eq!(types, mimes(&["text/plain;charset=utf-8"])); +} + +#[test] +fn a_copy_by_another_client_is_foreign_whatever_types_it_offers() { + // The same types as ours, offered in a different order, still don't make it ours. + let marker = OwnMarker::random(); + let mut types = mimes(&["text/plain", "text/plain;charset=utf-8"]); + + assert!(!marker.take_from(&mut types)); + + assert_eq!(types, mimes(&["text/plain", "text/plain;charset=utf-8"])); +} + +#[test] +fn the_marker_of_another_client_is_removed_but_is_not_ours() { + let ours = OwnMarker::random(); + let theirs = OwnMarker::random(); + assert_ne!(ours.as_str(), theirs.as_str()); + let mut types = mimes(&["image/png", theirs.as_str()]); + + assert!(!ours.take_from(&mut types)); + + assert_eq!(types, mimes(&["image/png"])); +} + +#[test] +fn a_clone_of_our_content_by_a_clipboard_manager_is_ours() { + // A manager that re-offers our selection carries every type along, the marker included. + let marker = OwnMarker::random(); + let mut types = mimes(&["text/html", "text/plain", marker.as_str()]); + + assert!(marker.take_from(&mut types)); + + assert_eq!(types, mimes(&["text/html", "text/plain"])); +} + +#[test] +fn markers_made_in_one_process_are_all_different() { + let markers: HashSet = core::iter::repeat_with(|| OwnMarker::random().as_str().to_owned()) + .take(1000) + .collect(); + + assert_eq!(markers.len(), 1000); +} + +#[test] +fn only_private_markers_are_recognised_as_markers() { + assert!(OwnMarker::is_marker(OwnMarker::random().as_str())); + assert!(!OwnMarker::is_marker("text/plain")); + assert!(!OwnMarker::is_marker("application/x-kde-passwordManagerHint")); +} + +#[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().on_change = Some(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 clearing_the_selection_ends_ownership() { + let mut state = State::default(); + state.shared().lock().unwrap().own_source_live = true; + state.on_selection_cleared(); + assert!(!state.shared().lock().unwrap().own_source_live); +} + +#[test] +fn cancelling_our_source_ends_ownership() { + let mut state = State::default(); + state.shared().lock().unwrap().own_source_live = true; + state.on_source_cancelled(); + assert!(!state.shared().lock().unwrap().own_source_live); +} + +#[test] +fn device_finished_ends_ownership() { + let mut state = State::default(); + state.shared().lock().unwrap().own_source_live = true; + state.on_device_finished(); + assert!(!state.shared().lock().unwrap().own_source_live); +} + +#[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"); +} + +#[test] +fn a_paste_of_the_bare_type_is_answered_from_the_utf8_entry() { + let mut state = State::default(); + state.insert_source_data("text/plain;charset=utf-8", b"hello".to_vec()); + let (read, write) = pipe(); + state.on_source_send("text/plain", write); + assert_eq!(read_all(read), b"hello"); +} + +#[test] +fn a_paste_with_a_charset_is_answered_from_the_bare_entry() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello".to_vec()); + let (read, write) = pipe(); + state.on_source_send("text/plain;charset=utf-8", write); + assert_eq!(read_all(read), b"hello"); +} + +#[test] +fn a_paste_is_not_answered_with_bytes_in_another_charset() { + let mut state = State::default(); + state.insert_source_data("text/plain;charset=iso-8859-1", b"\xe9t\xe9".to_vec()); + for requested in ["text/plain;charset=utf-8", "text/plain"] { + let (read, write) = pipe(); + state.on_source_send(requested, write); + assert!(read_all(read).is_empty(), "{requested} got bytes in another charset"); + } +} + +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().on_transfer = Some(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_failed_transfer_lets_the_next_paste_raise_a_new_request() { + let (mut state, requests) = state_with_transfer_callback(); + let (first_read, first_write) = pipe(); + state.on_source_send("image/png", first_write); + + state.complete_transfer(1, None); + assert!(read_all(first_read).is_empty()); + + let (_second_read, second_write) = pipe(); + state.on_source_send("image/png", second_write); + assert_eq!( + *requests.lock().unwrap(), + vec![(1, "image/png".to_owned()), (2, "image/png".to_owned())] + ); +} + +#[test] +fn a_paste_of_the_marker_closes_the_pipe_and_raises_no_transfer() { + let (mut state, requests) = state_with_transfer_callback(); + let marker = state.own_marker().to_owned(); + let (read, write) = pipe(); + + state.on_source_send(&marker, write); + + assert!(read_all(read).is_empty()); + assert!(requests.lock().unwrap().is_empty()); +} + +#[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 pastes_beyond_the_limit_are_closed_without_data() { + let (mut state, requests) = state_with_transfer_callback(); + let mut held = Vec::new(); + for _ in 0..MAX_PASTES { + let (read, write) = pipe(); + state.on_source_send("image/png", write); + held.push(read); + } + // The held pastes share one transfer. + assert_eq!(requests.lock().unwrap().len(), 1); + + // One more is closed at once, so reading it ends without an answer. + let (read, write) = pipe(); + state.on_source_send("image/png", write); + assert!(read_all(read).is_empty()); + + state.complete_transfer(1, Some(b"png bytes".to_vec())); + for read in held { + assert_eq!(read_all(read), b"png bytes"); + } +} + +#[test] +fn a_finished_paste_gives_its_place_back() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello".to_vec()); + // Three times as many pastes as places, one after the other. + for _ in 0..MAX_PASTES * 3 { + let (read, write) = pipe(); + state.on_source_send("text/plain", write); + assert_eq!(read_all(read), b"hello"); + } +} + +#[test] +fn a_held_paste_that_waits_too_long_gets_no_data() { + let (mut state, _requests) = state_with_transfer_callback(); + assert_eq!(state.expire_transfers(Instant::now()), None); + + let (read, write) = pipe(); + state.on_source_send("image/png", write); + let due = state + .expire_transfers(Instant::now()) + .expect("a held paste has a deadline"); + + // Just before the deadline the paste is still held. At the deadline it is closed. + assert_eq!(state.expire_transfers(due - Duration::from_millis(1)), Some(due)); + assert_eq!(state.expire_transfers(due), None); + assert!(read_all(read).is_empty()); +} + +#[test] +fn a_paste_after_the_wait_joins_the_transfer_still_pending() { + let (mut state, requests) = state_with_transfer_callback(); + let (first_read, first_write) = pipe(); + state.on_source_send("image/png", first_write); + let due = state + .expire_transfers(Instant::now()) + .expect("a held paste has a deadline"); + state.expire_transfers(due); + assert!(read_all(first_read).is_empty()); + + // The transfer is still pending, so this paste raises no second request. + let (second_read, second_write) = pipe(); + state.on_source_send("image/png", second_write); + assert_eq!(requests.lock().unwrap().len(), 1); + + state.complete_transfer(1, Some(b"late bytes".to_vec())); + assert_eq!(read_all(second_read), b"late bytes"); +} + +#[test] +fn an_answer_after_the_wait_is_kept_for_the_next_paste() { + let (mut state, requests) = state_with_transfer_callback(); + let (read, write) = pipe(); + state.on_source_send("image/png", write); + let due = state + .expire_transfers(Instant::now()) + .expect("a held paste has a deadline"); + state.expire_transfers(due); + assert!(read_all(read).is_empty()); + + state.complete_transfer(1, Some(b"late bytes".to_vec())); + let (read, write) = pipe(); + state.on_source_send("image/png", write); + assert_eq!(read_all(read), b"late bytes"); + assert_eq!(requests.lock().unwrap().len(), 1); +} + +#[test] +fn a_reader_that_takes_nothing_ends_the_write_after_the_idle_time() { + let (_read, write) = pipe(); + let mut pipe_end = std::fs::File::from(write); + // Far more than a pipe holds, so the write has to wait for the reader. + let data = vec![0x5a; 1 << 20]; + + let error = write_with_idle_timeout(&mut pipe_end, &data, Duration::from_millis(200)).unwrap_err(); + + assert_eq!(error.kind(), std::io::ErrorKind::TimedOut); +} + +#[test] +fn a_slow_reader_that_keeps_reading_is_not_cut_off() { + let (read, write) = pipe(); + let data: Vec = (0..=250u8).cycle().take(256 * 1024).collect(); + let reader = std::thread::spawn(move || { + let mut pipe_end = std::fs::File::from(read); + let mut received = Vec::new(); + let mut buf = [0u8; 16 * 1024]; + loop { + // Every pause is far shorter than the idle time the writer is given. + std::thread::sleep(Duration::from_millis(50)); + match pipe_end.read(&mut buf).unwrap() { + 0 => break received, + n => received.extend_from_slice(&buf[..n]), + } + } + }); + + let mut pipe_end = std::fs::File::from(write); + write_with_idle_timeout(&mut pipe_end, &data, Duration::from_secs(2)).unwrap(); + drop(pipe_end); + + assert_eq!(reader.join().unwrap(), data); +} + +#[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().on_change = Some(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()); +} + +/// Poison the lock of the shared state, as a panic in a thread that holds it does. +fn poison(shared: &Arc>) { + let shared = Arc::clone(shared); + let thread = std::thread::spawn(move || { + let _guard = shared.lock().unwrap(); + unreachable!("poisoning the lock on purpose"); + }); + assert!(thread.join().is_err()); +} + +#[test] +fn a_cleared_selection_is_still_recorded_after_the_lock_was_poisoned() { + let mut state = State::default(); + let called = Arc::new(Mutex::new(false)); + let seen = Arc::clone(&called); + state.shared().lock().unwrap().on_change = Some(Arc::new(move |_types: Vec| { + *seen.lock().unwrap() = true; + })); + poison(state.shared()); + + state.on_selection_cleared(); + + assert!(*called.lock().unwrap()); + assert_eq!(state.shared().lock().unwrap_or_else(PoisonError::into_inner).serial, 1); +} + +#[test] +fn a_paste_still_reaches_the_transfer_callback_after_the_lock_was_poisoned() { + let (mut state, requests) = state_with_transfer_callback(); + poison(state.shared()); + + let (_read, write) = pipe(); + state.on_source_send("image/png", write); + + assert_eq!(*requests.lock().unwrap(), vec![(1, "image/png".to_owned())]); +} + +#[test] +fn cancelling_our_source_still_ends_ownership_after_the_lock_was_poisoned() { + let mut state = State::default(); + state.shared().lock().unwrap().own_source_live = true; + poison(state.shared()); + + state.on_source_cancelled(); + + let shared = state.shared().lock().unwrap_or_else(PoisonError::into_inner); + assert!(!shared.own_source_live); +} + +/// Replace the selection, as a selection event from the compositor does. +fn select(shared: &Mutex, types: &[&str]) { + let mut shared = shared.lock().unwrap(); + shared.serial += 1; + shared.mime_types = mimes(types); +} + +/// A clipboard handle with no worker, over a selection that offers `types`, and the commands it sends. +fn clipboard_offering(types: &[&str]) -> (DataControl, Arc>, mpsc::Receiver) { + let shared: Arc> = Arc::default(); + select(&shared, types); + let (clipboard, commands) = DataControl::without_worker(Arc::clone(&shared)).unwrap(); + (clipboard, shared, commands) +} + +/// The write end of the pipe that a read sent to the worker. +fn receive_request(commands: &mpsc::Receiver) -> OwnedFd { + let Command::ReceiveFromOffer { fd, .. } = commands.recv().unwrap() else { + unreachable!("a read sends ReceiveFromOffer"); + }; + fd +} + +#[test] +fn a_read_returns_what_the_source_wrote() { + let (clipboard, _shared, commands) = clipboard_offering(&["text/plain"]); + let source = std::thread::spawn(move || { + let mut pipe = std::fs::File::from(receive_request(&commands)); + pipe.write_all(b"hello").unwrap(); + }); + + assert_eq!(clipboard.read("text/plain").unwrap(), Some(b"hello".to_vec())); + source.join().unwrap(); +} + +#[test] +fn a_source_that_writes_nothing_gives_an_empty_read() { + let (clipboard, _shared, commands) = clipboard_offering(&["text/plain"]); + let source = std::thread::spawn(move || drop(receive_request(&commands))); + + assert_eq!(clipboard.read("text/plain").unwrap(), Some(Vec::new())); + source.join().unwrap(); +} + +#[test] +fn a_read_of_a_type_the_selection_does_not_offer_is_none() { + let (clipboard, _shared, commands) = clipboard_offering(&["text/plain"]); + + assert_eq!(clipboard.read("image/png").unwrap(), None); + assert!( + commands.try_recv().is_err(), + "nothing is requested for a type that isn't offered" + ); +} + +#[test] +fn a_read_fails_when_the_selection_changed_before_the_worker_answered() { + let (clipboard, shared, commands) = clipboard_offering(&["text/plain"]); + let source = std::thread::spawn(move || { + let pipe = receive_request(&commands); + // The worker handles the selection event ahead of the request, so the offer the request was + // made against is gone and the pipe closes with nothing. + select(&shared, &["image/png"]); + drop(pipe); + }); + + assert!(matches!(clipboard.read("text/plain"), Err(Error::SelectionChanged))); + source.join().unwrap(); +} + +#[test] +fn a_read_fails_when_the_selection_changes_while_the_data_arrives() { + let (clipboard, shared, commands) = clipboard_offering(&["text/plain"]); + let source = std::thread::spawn(move || { + let mut pipe = std::fs::File::from(receive_request(&commands)); + pipe.write_all(b"he").unwrap(); + select(&shared, &["text/plain"]); + pipe.write_all(b"llo").unwrap(); + }); + + assert!(matches!(clipboard.read("text/plain"), Err(Error::SelectionChanged))); + source.join().unwrap(); +} + +/// Raise a transfer for `image/png` through the callback the handle installed, as the worker does for a paste. +fn raise_transfer(shared: &Mutex, serial: u32) { + let on_transfer = shared + .lock() + .unwrap() + .on_transfer + .clone() + .expect("the handle installed a transfer callback"); + on_transfer(serial, "image/png".to_owned()); +} + +/// The answer a transfer request sent to the worker. +fn transfer_answer(commands: &mpsc::Receiver) -> (u32, Option>) { + let Command::CompleteTransfer { serial, data } = commands.try_recv().unwrap() else { + unreachable!("a transfer request answers with CompleteTransfer"); + }; + (serial, data) +} + +#[test] +fn a_dropped_transfer_request_closes_its_paste() { + let (clipboard, shared, commands) = clipboard_offering(&["image/png"]); + clipboard.on_transfer(|_request| {}); + + raise_transfer(&shared, 7); + + assert_eq!(transfer_answer(&commands), (7, None)); + assert!(commands.try_recv().is_err(), "one answer only"); +} + +#[test] +fn a_completed_transfer_request_answers_once() { + let (clipboard, shared, commands) = clipboard_offering(&["image/png"]); + clipboard.on_transfer(|request| request.complete("png bytes").unwrap()); + + raise_transfer(&shared, 7); + + assert_eq!(transfer_answer(&commands), (7, Some(b"png bytes".to_vec()))); + assert!( + commands.try_recv().is_err(), + "dropping the answered request must not answer again" + ); +} + +#[test] +fn a_failed_transfer_request_answers_once() { + let (clipboard, shared, commands) = clipboard_offering(&["image/png"]); + clipboard.on_transfer(|request| request.fail().unwrap()); + + raise_transfer(&shared, 7); + + assert_eq!(transfer_answer(&commands), (7, None)); + assert!( + commands.try_recv().is_err(), + "dropping the failed request must not answer again" + ); +} diff --git a/crates/ironrdp-testsuite-extra/tests/cliprdr_native/live.rs b/crates/ironrdp-testsuite-extra/tests/cliprdr_native/live.rs new file mode 100644 index 0000000000..cadf6333c1 --- /dev/null +++ b/crates/ironrdp-testsuite-extra/tests/cliprdr_native/live.rs @@ -0,0 +1,262 @@ +//! Tests against a running Wayland compositor. +//! +//! They change the clipboard of the compositor they connect to, so they are ignored by default and do +//! nothing unless `IRONRDP_DATA_CONTROL_LIVE=1` is set. Run them inside an isolated compositor, for +//! example `kwin_wayland --virtual --no-lockscreen --socket ironrdp-live`, one at a time because they +//! share its clipboard: +//! +//! ```text +//! WAYLAND_DISPLAY=ironrdp-live IRONRDP_DATA_CONTROL_LIVE=1 \ +//! cargo test -p ironrdp-testsuite-extra -- --ignored --test-threads=1 cliprdr_native::live +//! ``` +//! +//! They use whichever data-control protocol the compositor offers. One of them pastes with `wl-paste` +//! (from wl-clipboard) and fails without it, because the paste has to wait longer than the reader of +//! this library does. + +use core::{ + sync::atomic::{AtomicUsize, Ordering}, + time::Duration, +}; +use std::{ + process::Command, + sync::{Arc, Mutex, mpsc}, + time::Instant, +}; + +use ironrdp_cliprdr_native::data_control::{Content, DataControl, TransferRequest}; + +const TEXT: &str = "text/plain;charset=utf-8"; + +fn enabled() -> bool { + std::env::var("IRONRDP_DATA_CONTROL_LIVE").as_deref() == Ok("1") +} + +fn wait_for_change(rx: &mpsc::Receiver>) -> Vec { + rx.recv_timeout(Duration::from_secs(5)).expect("the selection changed") +} + +/// Wait until the handle has seen a selection newer than `serial`, which is how it learns about one it set +/// itself. +fn wait_for_serial(clipboard: &DataControl, serial: u32) { + let deadline = Instant::now() + Duration::from_secs(5); + while clipboard.serial() == serial { + assert!(Instant::now() < deadline, "no selection event"); + std::thread::sleep(Duration::from_millis(20)); + } +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn set_then_read_back() { + if !enabled() { + return; + } + let clipboard = DataControl::connect().unwrap(); + let before = clipboard.serial(); + + clipboard + .set_selection(Content::new().data(TEXT, "round trip")) + .unwrap(); + wait_for_serial(&clipboard, before); + + assert_eq!(clipboard.read("text/plain").unwrap(), Some(b"round trip".to_vec())); + assert!(clipboard.owns_selection()); + assert_eq!(clipboard.selection_mime_types(), vec![TEXT.to_owned()]); +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn a_second_client_reads_our_selection() { + if !enabled() { + return; + } + let owner = DataControl::connect().unwrap(); + let reader = DataControl::connect().unwrap(); + let (tx, rx) = mpsc::channel(); + reader.on_change(move |types| { + let _ = tx.send(types); + }); + + owner + .set_selection(Content::new().data(TEXT, "seen by another client")) + .unwrap(); + + // The marker of the owner is not among the types the reader is told about. + assert_eq!(wait_for_change(&rx), vec![TEXT.to_owned()]); + assert_eq!(reader.read(TEXT).unwrap(), Some(b"seen by another client".to_vec())); +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn delayed_rendering_produces_data_on_paste() { + if !enabled() { + return; + } + let owner = DataControl::connect().unwrap(); + let reader = DataControl::connect().unwrap(); + let (tx, rx) = mpsc::channel(); + reader.on_change(move |types| { + let _ = tx.send(types); + }); + owner.on_transfer(|request| request.complete("rendered late").unwrap()); + + owner.set_selection(Content::new().advertise(TEXT)).unwrap(); + wait_for_change(&rx); + + assert_eq!(reader.read(TEXT).unwrap(), Some(b"rendered late".to_vec())); +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn a_replaced_selection_is_not_served_from_our_cache() { + if !enabled() { + return; + } + let first = DataControl::connect().unwrap(); + let second = DataControl::connect().unwrap(); + let (tx, rx) = mpsc::channel(); + first.on_change(move |types| { + let _ = tx.send(types); + }); + let before = first.serial(); + + first.set_selection(Content::new().data(TEXT, "old")).unwrap(); + wait_for_serial(&first, before); + second.set_selection(Content::new().data(TEXT, "new")).unwrap(); + wait_for_change(&rx); + + assert_eq!(first.read(TEXT).unwrap(), Some(b"new".to_vec())); + assert!(!first.owns_selection()); +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn our_own_selection_is_not_reported_but_a_copy_with_the_same_types_is() { + if !enabled() { + return; + } + let owner = DataControl::connect().unwrap(); + let other = DataControl::connect().unwrap(); + let (tx, rx) = mpsc::channel(); + owner.on_change(move |types| { + let _ = tx.send(types); + }); + let before = owner.serial(); + + owner.set_selection(Content::new().data(TEXT, "ours")).unwrap(); + wait_for_serial(&owner, before); + // The callback would run right after the serial moves, so a short wait is enough to see it. + assert_eq!( + rx.recv_timeout(Duration::from_millis(500)), + Err(mpsc::RecvTimeoutError::Timeout), + "our own selection was reported" + ); + + other.set_selection(Content::new().data(TEXT, "theirs")).unwrap(); + + // The other client offers the same types as we did, and its marker is not among them. + assert_eq!(wait_for_change(&rx), vec![TEXT.to_owned()]); + assert_eq!(owner.read(TEXT).unwrap(), Some(b"theirs".to_vec())); + assert!(!owner.owns_selection()); +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn our_own_clear_is_reported_as_an_empty_selection() { + if !enabled() { + return; + } + let owner = DataControl::connect().unwrap(); + let (tx, rx) = mpsc::channel(); + owner.on_change(move |types| { + let _ = tx.send(types); + }); + let before = owner.serial(); + owner.set_selection(Content::new().data(TEXT, "ours")).unwrap(); + wait_for_serial(&owner, before); + + owner.clear_selection().unwrap(); + + assert_eq!(wait_for_change(&rx), Vec::::new()); + assert!(!owner.owns_selection()); +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn a_dropped_request_closes_the_paste_at_once_and_the_next_paste_raises_a_new_request() { + if !enabled() { + return; + } + let owner = DataControl::connect().unwrap(); + let reader = DataControl::connect().unwrap(); + let (tx, rx) = mpsc::channel(); + reader.on_change(move |types| { + let _ = tx.send(types); + }); + let raised = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&raised); + owner.on_transfer(move |_request| { + counter.fetch_add(1, Ordering::SeqCst); + }); + owner.set_selection(Content::new().advertise(TEXT)).unwrap(); + wait_for_change(&rx); + + for round in 1..=2 { + let started = Instant::now(); + assert_eq!(reader.read(TEXT).unwrap(), Some(Vec::new())); + assert!( + started.elapsed() < Duration::from_secs(2), + "paste {round} waited too long" + ); + assert_eq!(raised.load(Ordering::SeqCst), round); + } +} + +#[test] +#[ignore = "changes the clipboard of the running compositor, needs wl-paste, set IRONRDP_DATA_CONTROL_LIVE=1 to run"] +fn an_unanswered_paste_is_closed_after_five_seconds_and_a_late_answer_is_kept() { + if !enabled() { + return; + } + let paste = || { + let started = Instant::now(); + let output = Command::new("wl-paste") + .args(["--no-newline", "--type", TEXT]) + .output() + .expect("wl-paste (from wl-clipboard) is needed to paste in this test"); + (output.stdout, started.elapsed()) + }; + + let owner = DataControl::connect().unwrap(); + let held: Arc>> = Arc::default(); + let raised = Arc::new(AtomicUsize::new(0)); + { + let held = Arc::clone(&held); + let raised = Arc::clone(&raised); + owner.on_transfer(move |request| { + raised.fetch_add(1, Ordering::SeqCst); + held.lock().unwrap().push(request); + }); + } + let before = owner.serial(); + owner.set_selection(Content::new().advertise(TEXT)).unwrap(); + wait_for_serial(&owner, before); + + // Nobody answers the first paste, so it is closed with no data after about five seconds. + let (data, waited) = paste(); + assert!(data.is_empty()); + assert!( + Duration::from_millis(4500) < waited && waited < Duration::from_secs(8), + "waited {waited:?}" + ); + assert_eq!(raised.load(Ordering::SeqCst), 1); + + // The answer arrives late and is kept for the next paste, which raises no new request. + let request = held.lock().unwrap().pop().expect("a request was held"); + request.complete("late answer").unwrap(); + let (data, waited) = paste(); + assert_eq!(data, b"late answer"); + assert!(waited < Duration::from_secs(2), "waited {waited:?}"); + assert_eq!(raised.load(Ordering::SeqCst), 1); +} diff --git a/crates/ironrdp-testsuite-extra/tests/cliprdr_native/mod.rs b/crates/ironrdp-testsuite-extra/tests/cliprdr_native/mod.rs new file mode 100644 index 0000000000..8f40d71418 --- /dev/null +++ b/crates/ironrdp-testsuite-extra/tests/cliprdr_native/mod.rs @@ -0,0 +1,2 @@ +mod data_control; +mod live; diff --git a/crates/ironrdp-testsuite-extra/tests/main.rs b/crates/ironrdp-testsuite-extra/tests/main.rs index fb3fedf3ef..cd9ce0044b 100644 --- a/crates/ironrdp-testsuite-extra/tests/main.rs +++ b/crates/ironrdp-testsuite-extra/tests/main.rs @@ -5,6 +5,8 @@ mod agent; mod async_framed; mod capture_helpers; mod client; +#[cfg(target_os = "linux")] +mod cliprdr_native; mod dvc_pipe_proxy; mod e2e; mod gateway_detect;