diff --git a/Cargo.lock b/Cargo.lock index ea162c2dfd..c70d35997e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2812,7 +2812,12 @@ version = "0.7.0" dependencies = [ "ironrdp-cliprdr", "ironrdp-core 0.2.1", + "nix", "tracing", + "visibility", + "wayland-client", + "wayland-protocols", + "wayland-protocols-wlr", "windows", ] @@ -3390,6 +3395,7 @@ dependencies = [ "ironrdp-cfg", "ironrdp-cliprdr", "ironrdp-cliprdr-format", + "ironrdp-cliprdr-native", "ironrdp-connector", "ironrdp-core 0.2.1", "ironrdp-displaycontrol", diff --git a/crates/ironrdp-cliprdr-native/Cargo.toml b/crates/ironrdp-cliprdr-native/Cargo.toml index c4c55328ce..44e78f4bcc 100644 --- a/crates/ironrdp-cliprdr-native/Cargo.toml +++ b/crates/ironrdp-cliprdr-native/Cargo.toml @@ -12,6 +12,10 @@ authors.workspace = true keywords.workspace = true categories.workspace = true +[features] +# Exposes internals to the integration tests in ironrdp-testsuite-core. +__test = ["dep:visibility"] + [lib] doctest = false test = false @@ -21,6 +25,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..6523516266 100644 --- a/crates/ironrdp-cliprdr-native/README.md +++ b/crates/ironrdp-cliprdr-native/README.md @@ -1,7 +1,14 @@ # 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. 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..c0e78366d6 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/client.rs @@ -0,0 +1,384 @@ +//! 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}, + 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 it without answering leaves the paste +/// waiting until the selection is replaced. +#[derive(Debug)] +#[must_use = "a paste stays open until the request is answered"] +pub struct TransferRequest { + serial: u32, + mime_type: String, + link: Arc, +} + +impl core::fmt::Debug for Link { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.debug_struct("Link").finish_non_exhaustive() + } +} + +impl TransferRequest { + /// The MIME type being pasted. + #[must_use] + pub fn mime_type(&self) -> &str { + &self.mime_type + } + + /// Supply the data. It is written to every paste of this type that is + /// waiting, and cached for later pastes. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the clipboard worker has shut down. + pub fn complete(self, data: impl Into>) -> Result<()> { + self.link.send(Command::CompleteTransfer { + serial: self.serial, + data: Some(data.into()), + }) + } + + /// Close the waiting pastes with no data. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the clipboard worker has shut down. + pub fn fail(self) -> Result<()> { + self.link.send(Command::CompleteTransfer { + serial: self.serial, + data: None, + }) + } +} + +/// A connection to the Wayland compositor's clipboard. +/// +/// Connecting spawns one thread that owns the Wayland connection. Dropping the +/// handle stops and joins it. +/// +/// ```no_run +/// use ironrdp_cliprdr_native::data_control::{Content, DataControl}; +/// +/// # fn main() -> ironrdp_cliprdr_native::data_control::Result<()> { +/// let clipboard = DataControl::connect()?; +/// clipboard.set_selection(Content::new().data("text/plain;charset=utf-8", "hello"))?; +/// # Ok(()) +/// # } +/// ``` +pub struct DataControl { + link: Arc, + shared: Arc>, + protocol: Protocol, + max_read_bytes: usize, + stop: Arc, + thread: Option>, +} + +impl core::fmt::Debug for DataControl { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + f.debug_struct("DataControl") + .field("protocol", &self.protocol) + .finish_non_exhaustive() + } +} + +impl DataControl { + /// Connect with the default [`Options`]. + /// + /// # Errors + /// + /// See [`Error`]: no compositor, no data-control protocol, or no seat. + pub fn connect() -> Result { + Self::connect_with(&Options::new()) + } + + /// Connect with explicit options. + /// + /// # Errors + /// + /// See [`Error`]: no compositor, no data-control protocol, or no seat. + pub fn connect_with(options: &Options) -> Result { + let (command_tx, command_rx) = mpsc::channel(); + let (wake_tx, wake_rx) = UnixStream::pair()?; + wake_tx.set_nonblocking(true)?; + wake_rx.set_nonblocking(true)?; + let stop = Arc::new(AtomicBool::new(false)); + let (ready_tx, ready_rx) = mpsc::sync_channel(1); + + let thread_options = options.clone(); + let thread_stop = Arc::clone(&stop); + let thread = std::thread::Builder::new() + .name("ironrdp-cliprdr-data-control".into()) + .spawn( + move || match worker::connect(&thread_options, command_rx, wake_rx, thread_stop) { + Ok((worker, connected)) => { + if ready_tx.send(Ok(connected)).is_ok() { + worker.run(); + } + } + Err(error) => { + let _ = ready_tx.send(Err(error)); + } + }, + )?; + + let connected = match ready_rx.recv() { + Ok(Ok(connected)) => connected, + Ok(Err(error)) => { + let _ = thread.join(); + return Err(error); + } + Err(_) => { + let _ = thread.join(); + return Err(Error::Stopped); + } + }; + + Ok(Self { + link: Arc::new(Link { + commands: command_tx, + wake: wake_tx, + }), + shared: connected.shared, + protocol: connected.protocol, + max_read_bytes: options.max_read_bytes, + stop, + thread: Some(thread), + }) + } + + /// The protocol the compositor is speaking. + #[must_use] + pub fn protocol(&self) -> Protocol { + self.protocol + } + + fn shared(&self) -> MutexGuard<'_, Shared> { + self.shared.lock().unwrap_or_else(std::sync::PoisonError::into_inner) + } + + /// MIME types the current selection offers; empty if there is none. + #[must_use] + pub fn selection_mime_types(&self) -> Vec { + self.shared().mime_types.clone() + } + + /// A counter that increases on every selection change. + #[must_use] + pub fn serial(&self) -> u32 { + self.shared().serial + } + + /// Call `callback` with the offered MIME types whenever another client + /// changes the selection. + /// + /// A selection this handle set itself is not reported; it still advances + /// [`serial`](Self::serial) and updates + /// [`selection_mime_types`](Self::selection_mime_types). + /// + /// The callback runs on the worker thread and must not block. + pub fn on_change(&self, callback: impl Fn(Vec) + Send + Sync + 'static) { + self.shared().on_change = Some(Arc::new(callback)); + } + + /// Call `callback` when something pastes a type that was advertised + /// without data. + /// + /// The callback runs on the worker thread and must not block; hand the + /// request to another thread and answer it there. + pub fn on_transfer(&self, callback: impl Fn(TransferRequest) + Send + Sync + 'static) { + let link = Arc::clone(&self.link); + self.shared().on_transfer = Some(Arc::new(move |serial, mime_type| { + callback(TransferRequest { + serial, + mime_type, + link: Arc::clone(&link), + }); + })); + } + + /// Become the clipboard owner and offer `content`. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the worker has shut down. + pub fn set_selection(&self, content: Content) -> Result<()> { + self.link.send(Command::SetSelection { + mime_types: content.mime_types, + data: content.data, + }) + } + + /// Add data for a type after [`set_selection`](Self::set_selection), + /// without replacing the selection. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the worker has shut down. + pub fn update_data(&self, mime_type: impl Into, data: impl Into>) -> Result<()> { + self.link.send(Command::UpdateSourceData { + mime_type: mime_type.into(), + data: data.into(), + }) + } + + /// Give up the selection if this handle still owns it. + /// + /// # Errors + /// + /// [`Error::Stopped`] if the worker has shut down. + pub fn clear_selection(&self) -> Result<()> { + self.link.send(Command::ClearSelection) + } + + /// Read the current selection as `mime_type`. + /// + /// Charset differences are tolerated, see [`find_mime_match`]. Returns + /// `Ok(None)` when there is no selection or it does not offer the type. + /// + /// # Errors + /// + /// [`Error::TooLarge`] past the size limit, [`Error::Timeout`] if the + /// source 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 = self.selection_mime_types(); + let Some(matched) = find_mime_match(mime_type, &available) else { + return Ok(None); + }; + + let (mut reader, writer) = std::io::pipe()?; + self.link.send(Command::ReceiveFromOffer { + mime_type: matched.to_owned(), + fd: OwnedFd::from(writer), + })?; + + 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. + Ok(0) => return Ok(Some(data)), + Ok(n) => { + if data.len() + n > self.max_read_bytes { + 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()), + } + } + } +} + +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..8a59a78a73 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/dispatch.rs @@ -0,0 +1,231 @@ +//! Wayland event dispatch for the data-control objects. +//! +//! Both protocols (`ext-data-control-v1` and `wlr-data-control-unstable-v1`) +//! have the same object model and the same events, so each handler is written +//! twice and forwards to the shared [`State`]. + +use wayland_client::{ + Connection, Dispatch, QueueHandle, + globals::GlobalListContents, + protocol::{wl_registry::WlRegistry, wl_seat::WlSeat}, +}; +use wayland_protocols::ext::data_control::v1::client::{ + ext_data_control_device_v1::{self, ExtDataControlDeviceV1}, + ext_data_control_manager_v1::ExtDataControlManagerV1, + ext_data_control_offer_v1::{self, ExtDataControlOfferV1}, + ext_data_control_source_v1::{self, ExtDataControlSourceV1}, +}; +use wayland_protocols_wlr::data_control::v1::client::{ + zwlr_data_control_device_v1::{self, ZwlrDataControlDeviceV1}, + zwlr_data_control_manager_v1::ZwlrDataControlManagerV1, + zwlr_data_control_offer_v1::{self, ZwlrDataControlOfferV1}, + zwlr_data_control_source_v1::{self, ZwlrDataControlSourceV1}, +}; + +use super::state::State; + +/// The dispatch target the Wayland event queue calls into. +pub(crate) struct Client { + pub(crate) data_control: State, +} + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &WlRegistry, + _event: ::Event, + _data: &GlobalListContents, + _conn: &Connection, + _qh: &QueueHandle, + ) { + // Globals are looked up once at start-up; later additions are ignored. + } +} + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &WlSeat, + _event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The seat is only an argument to `get_data_device`. + } +} + +// === ext-data-control-v1 === + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &ExtDataControlManagerV1, + _event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The manager has no events. + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ExtDataControlDeviceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + ext_data_control_device_v1::Event::DataOffer { id } => { + state.data_control.on_data_offer_ext(id); + } + ext_data_control_device_v1::Event::Selection { id } => { + if id.is_some() { + state.data_control.on_selection(); + } else { + state.data_control.on_selection_cleared(); + } + } + ext_data_control_device_v1::Event::Finished => { + state.data_control.on_device_finished(); + } + ext_data_control_device_v1::Event::PrimarySelection { .. } => { + // Only the regular clipboard is handled, not the primary selection. + tracing::trace!("ext data control primary selection event (ignored)"); + } + _ => {} + } + } + + // The `data_offer` event creates a child offer object; without this, + // wayland-client's default panics. + wayland_client::event_created_child!(Client, ExtDataControlDeviceV1, [ + ext_data_control_device_v1::EVT_DATA_OFFER_OPCODE => (ExtDataControlOfferV1, ()), + ]); +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ExtDataControlSourceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + ext_data_control_source_v1::Event::Send { mime_type, fd } => { + state.data_control.on_source_send(&mime_type, fd); + } + ext_data_control_source_v1::Event::Cancelled => { + state.data_control.on_source_cancelled(); + } + _ => {} + } + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ExtDataControlOfferV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if let ext_data_control_offer_v1::Event::Offer { mime_type } = event { + state.data_control.on_offer_mime_type(mime_type); + } + } +} + +// === wlr-data-control-unstable-v1 === + +impl Dispatch for Client { + fn event( + _state: &mut Self, + _proxy: &ZwlrDataControlManagerV1, + _event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + // The manager has no events. + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ZwlrDataControlDeviceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + zwlr_data_control_device_v1::Event::DataOffer { id } => { + state.data_control.on_data_offer_wlr(id); + } + zwlr_data_control_device_v1::Event::Selection { id } => { + if id.is_some() { + state.data_control.on_selection(); + } else { + state.data_control.on_selection_cleared(); + } + } + zwlr_data_control_device_v1::Event::Finished => { + state.data_control.on_device_finished(); + } + zwlr_data_control_device_v1::Event::PrimarySelection { .. } => { + tracing::trace!("wlr data control primary selection event (ignored)"); + } + _ => {} + } + } + + wayland_client::event_created_child!(Client, ZwlrDataControlDeviceV1, [ + zwlr_data_control_device_v1::EVT_DATA_OFFER_OPCODE => (ZwlrDataControlOfferV1, ()), + ]); +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ZwlrDataControlSourceV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + match event { + zwlr_data_control_source_v1::Event::Send { mime_type, fd } => { + state.data_control.on_source_send(&mime_type, fd); + } + zwlr_data_control_source_v1::Event::Cancelled => { + state.data_control.on_source_cancelled(); + } + _ => {} + } + } +} + +impl Dispatch for Client { + fn event( + state: &mut Self, + _proxy: &ZwlrDataControlOfferV1, + event: ::Event, + _data: &(), + _conn: &Connection, + _qh: &QueueHandle, + ) { + if let zwlr_data_control_offer_v1::Event::Offer { mime_type } = event { + state.data_control.on_offer_mime_type(mime_type); + } + } +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/error.rs b/crates/ironrdp-cliprdr-native/src/data_control/error.rs new file mode 100644 index 0000000000..1e2c3e6dc4 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/error.rs @@ -0,0 +1,72 @@ +//! Errors returned by the clipboard client. + +use core::fmt; + +/// Errors returned by [`DataControl`](super::DataControl). +#[derive(Debug)] +#[non_exhaustive] +pub enum Error { + /// The connection to the Wayland compositor could not be made. + Connect(String), + /// The compositor offers neither data-control protocol. + Unsupported, + /// The requested protocol is not offered and fallback is disabled. + ProtocolUnavailable, + /// The compositor advertises no seat to attach the clipboard to. + NoSeat, + /// The compositor reported a protocol error or the connection failed. + Wayland(String), + /// The worker thread has stopped, so the request could not be delivered. + Stopped, + /// Clipboard data exceeded the size limit. + TooLarge { + /// Bytes received when the limit was hit. + size: usize, + /// The limit. + limit: usize, + }, + /// The source did not deliver its data in time. + Timeout, + /// An I/O error while creating or reading a pipe. + Io(std::io::Error), +} + +impl fmt::Display for Error { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Connect(reason) => write!(f, "cannot connect to the Wayland compositor: {reason}"), + Self::Unsupported => { + f.write_str("the compositor offers neither ext-data-control-v1 nor wlr-data-control-unstable-v1") + } + Self::ProtocolUnavailable => { + f.write_str("the requested data-control protocol is not offered by the compositor") + } + Self::NoSeat => f.write_str("the compositor advertises no seat"), + Self::Wayland(reason) => write!(f, "Wayland error: {reason}"), + Self::Stopped => f.write_str("the clipboard worker has stopped"), + Self::TooLarge { size, limit } => { + write!(f, "clipboard data of {size} bytes exceeds the {limit} byte limit") + } + Self::Timeout => f.write_str("the clipboard source did not deliver its data in time"), + Self::Io(error) => write!(f, "I/O error: {error}"), + } + } +} + +impl core::error::Error for Error { + fn source(&self) -> Option<&(dyn core::error::Error + 'static)> { + match self { + Self::Io(error) => Some(error), + _ => None, + } + } +} + +impl From for Error { + fn from(error: std::io::Error) -> Self { + Self::Io(error) + } +} + +/// Result alias for this module. +pub type Result = core::result::Result; diff --git a/crates/ironrdp-cliprdr-native/src/data_control/mime.rs b/crates/ironrdp-cliprdr-native/src/data_control/mime.rs new file mode 100644 index 0000000000..978e5f5eb7 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/mime.rs @@ -0,0 +1,53 @@ +//! MIME type matching that tolerates charset parameters. + +/// Find a MIME type match in `available`, tolerating charset differences. +/// +/// A source may offer `text/plain` while a consumer asks for +/// `text/plain;charset=utf-8`, or the other way round. An exact match wins; +/// for `text/` types the charset parameter is then stripped or added, and +/// finally any offered type with the same base type is accepted. +/// +/// Returns the string from `available` that should be used for the request, so +/// that the compositor receives a type the source actually offered. +/// +/// ``` +/// use ironrdp_cliprdr_native::data_control::find_mime_match; +/// +/// let available = vec!["text/plain;charset=utf-8".to_owned()]; +/// assert_eq!( +/// find_mime_match("text/plain", &available), +/// Some("text/plain;charset=utf-8") +/// ); +/// assert_eq!(find_mime_match("image/png", &available), None); +/// ``` +#[must_use] +pub fn find_mime_match<'a>(requested: &str, available: &'a [String]) -> Option<&'a str> { + if let Some(found) = available.iter().find(|m| m.as_str() == requested) { + return Some(found.as_str()); + } + + if !requested.starts_with("text/") { + return None; + } + let base = requested.split(';').next()?; + + if requested.contains(';') { + // The request carries a charset; try the bare type. + if let Some(found) = available.iter().find(|m| m.as_str() == base) { + return Some(found.as_str()); + } + } else { + // The request has none; try the common charset spellings. + for suffix in [";charset=utf-8", ";charset=UTF-8"] { + let with_charset = format!("{requested}{suffix}"); + if let Some(found) = available.iter().find(|m| m.as_str() == with_charset) { + return Some(found.as_str()); + } + } + } + + available + .iter() + .find(|m| m.split(';').next() == Some(base)) + .map(String::as_str) +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/mod.rs b/crates/ironrdp-cliprdr-native/src/data_control/mod.rs new file mode 100644 index 0000000000..b5b91cbdaf --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/mod.rs @@ -0,0 +1,64 @@ +//! A Wayland clipboard client built on the data-control protocols. +//! +//! Data-control lets a client that has no window read and set the clipboard. +//! This module speaks both `ext-data-control-v1` and +//! `wlr-data-control-unstable-v1` and picks whichever the compositor offers. +//! +//! It has no async runtime. One thread owns the Wayland connection, and +//! [`DataControl`] is the handle to call from anywhere. +//! +//! # Reading +//! +//! ```no_run +//! use ironrdp_cliprdr_native::data_control::DataControl; +//! +//! # fn main() -> ironrdp_cliprdr_native::data_control::Result<()> { +//! let clipboard = DataControl::connect()?; +//! if let Some(bytes) = clipboard.read("text/plain;charset=utf-8")? { +//! println!("{}", String::from_utf8_lossy(&bytes)); +//! } +//! # Ok(()) +//! # } +//! ``` +//! +//! # Owning the clipboard, with delayed rendering +//! +//! A type advertised without data raises a [`TransferRequest`] when something +//! pastes it, so the data is produced only if it is wanted. That maps onto +//! CLIPRDR, where the Format Data Response follows the paste. +//! +//! ```no_run +//! use ironrdp_cliprdr_native::data_control::{Content, DataControl}; +//! +//! # fn main() -> ironrdp_cliprdr_native::data_control::Result<()> { +//! let clipboard = DataControl::connect()?; +//! clipboard.on_transfer(|request| { +//! let _ = request.complete("rendered on demand"); +//! }); +//! clipboard.set_selection(Content::new().advertise("text/plain;charset=utf-8"))?; +//! # Ok(()) +//! # } +//! ``` +//! +//! GNOME's Mutter offers no data-control protocol, so [`DataControl::connect`] +//! returns [`Error::Unsupported`] there. + +mod client; +mod dispatch; +mod error; +mod mime; +mod options; +#[cfg(feature = "__test")] +pub mod state; +#[cfg(not(feature = "__test"))] +#[expect( + unreachable_pub, + reason = "the __test feature exposes this module to the shared integration tests" +)] +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..592dacf0e3 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/options.rs @@ -0,0 +1,107 @@ +//! Connection options and the protocol choice. + +use std::fmt; + +/// The data-control protocol in use. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[non_exhaustive] +pub enum Protocol { + /// `ext-data-control-v1`, the standardized protocol. + Ext, + /// `wlr-data-control-unstable-v1`, the wlroots protocol. + Wlr, +} + +impl fmt::Display for Protocol { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Protocol::Ext => f.write_str("ext-data-control-v1"), + Protocol::Wlr => f.write_str("wlr-data-control-unstable-v1"), + } + } +} + +/// Which protocol to prefer when the compositor offers both. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] +#[non_exhaustive] +pub enum Preference { + /// Prefer `ext-data-control-v1`, then `wlr-data-control-unstable-v1`. + #[default] + Auto, + /// Prefer `ext-data-control-v1`. + Ext, + /// Prefer `wlr-data-control-unstable-v1`. + Wlr, +} + +/// Options for [`DataControl::connect_with`](super::DataControl::connect_with). +/// +/// ``` +/// use ironrdp_cliprdr_native::data_control::{Options, Preference}; +/// +/// let options = Options::new() +/// .preference(Preference::Wlr) +/// .allow_fallback(false); +/// ``` +#[derive(Debug, Clone)] +pub struct Options { + pub(crate) preference: Preference, + pub(crate) allow_fallback: bool, + pub(crate) max_read_bytes: usize, +} + +impl Options { + /// The default read size limit: 100 MiB. + pub const DEFAULT_MAX_READ_BYTES: usize = 100 * 1024 * 1024; + + /// Options with automatic protocol selection and fallback enabled. + #[must_use] + pub fn new() -> Self { + Self { + preference: Preference::Auto, + allow_fallback: true, + max_read_bytes: Self::DEFAULT_MAX_READ_BYTES, + } + } + + /// Choose which protocol to prefer. + #[must_use] + pub fn preference(mut self, preference: Preference) -> Self { + self.preference = preference; + self + } + + /// Whether to use the other protocol when the preferred one is missing. + #[must_use] + pub fn allow_fallback(mut self, allow: bool) -> Self { + self.allow_fallback = allow; + self + } + + /// Limit, in bytes, for one [`read`](super::DataControl::read). + #[must_use] + pub fn max_read_bytes(mut self, limit: usize) -> Self { + self.max_read_bytes = limit; + self + } + + /// The protocols to try, in order. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn candidates(&self) -> Vec { + let (first, second) = match self.preference { + Preference::Auto | Preference::Ext => (Protocol::Ext, Protocol::Wlr), + Preference::Wlr => (Protocol::Wlr, Protocol::Ext), + }; + if self.allow_fallback { + vec![first, second] + } else { + vec![first] + } + } +} + +impl Default for Options { + fn default() -> Self { + Self::new() + } +} diff --git a/crates/ironrdp-cliprdr-native/src/data_control/state.rs b/crates/ironrdp-cliprdr-native/src/data_control/state.rs new file mode 100644 index 0000000000..166cb36d93 --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/state.rs @@ -0,0 +1,706 @@ +//! The data-control protocol state machine. +//! +//! [`State`] owns the protocol objects (manager, device, current offer, our +//! source) and turns compositor events into changes of the shared selection +//! state. It is deliberately free of any event loop: the worker in +//! [`super::worker`] feeds it events and commands. +//! +//! # Flows +//! +//! **Set selection (client to compositor):** +//! 1. A [`Command::SetSelection`] arrives with the MIME types and any data. +//! 2. A data-control source is created, the types are advertised on it and it +//! is made the device's selection; the previous source is destroyed only +//! afterwards, so the clipboard is never empty in between. +//! 3. A `send` event for an advertised type writes the cached data to the +//! requested file descriptor on a worker thread. With no cached data and an +//! `on_transfer` callback set, the paste is held until +//! [`Command::CompleteTransfer`] answers it (delayed rendering). +//! +//! **Selection changed (compositor to client):** +//! 1. `data_offer`, then one `offer` per MIME type, then `selection`. +//! 2. The types are stored in [`Shared`] and the change callback is called, +//! unless the selection is the one this client set itself. +//! +//! **Read (client reads the compositor's selection):** +//! [`Command::ReceiveFromOffer`] calls `offer.receive(mime, fd)`; the caller +//! reads the other end of the pipe. + +use std::{ + collections::HashMap, + os::unix::io::OwnedFd, + sync::{Arc, Mutex}, +}; + +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, +}; + +// === Protocol object enums === +// These wrap both ext and wlr variants so State can work +// with either protocol transparently. + +/// Data control manager (either ext or wlr). +pub(crate) enum DataControlManager { + /// ext-data-control-v1 manager. + Ext(ExtDataControlManagerV1), + /// wlr-data-control-unstable-v1 manager. + Wlr(ZwlrDataControlManagerV1), +} + +/// Data control device (either ext or wlr). +pub(crate) enum DataControlDevice { + /// ext-data-control-v1 device. + Ext(ExtDataControlDeviceV1), + /// wlr-data-control-unstable-v1 device. + Wlr(ZwlrDataControlDeviceV1), +} + +/// Data control offer (either ext or wlr). +pub(crate) enum DataControlOffer { + /// ext-data-control-v1 offer. + Ext(ExtDataControlOfferV1), + /// wlr-data-control-unstable-v1 offer. + Wlr(ZwlrDataControlOfferV1), +} + +impl DataControlOffer { + /// Request data transfer for a MIME type. + /// + /// Tells the source client to write data to the provided fd. + pub(crate) fn receive(&self, mime_type: &str, fd: &OwnedFd) { + use std::os::unix::io::AsFd as _; + match self { + DataControlOffer::Ext(offer) => offer.receive(mime_type.to_owned(), fd.as_fd()), + DataControlOffer::Wlr(offer) => offer.receive(mime_type.to_owned(), fd.as_fd()), + } + } + + /// Destroy this offer. + pub(crate) fn destroy(&self) { + match self { + DataControlOffer::Ext(offer) => offer.destroy(), + DataControlOffer::Wlr(offer) => offer.destroy(), + } + } +} + +/// Data control source (either ext or wlr). +enum DataControlSource { + /// ext-data-control-v1 source. + Ext(ExtDataControlSourceV1), + /// wlr-data-control-unstable-v1 source. + Wlr(ZwlrDataControlSourceV1), +} + +impl DataControlSource { + /// Advertise a MIME type on this source. + fn offer(&self, mime_type: &str) { + match self { + DataControlSource::Ext(source) => source.offer(mime_type.to_owned()), + DataControlSource::Wlr(source) => source.offer(mime_type.to_owned()), + } + } + + /// Destroy this source. + fn destroy(&self) { + match self { + DataControlSource::Ext(source) => source.destroy(), + DataControlSource::Wlr(source) => source.destroy(), + } + } +} + +// === Commands and shared state === + +/// Commands sent from clipboard backends to the Wayland event loop thread. +#[derive(Debug)] +pub(crate) enum Command { + /// Set the clipboard selection on the compositor. + /// + /// Creates a data control source with the offered MIME types and + /// stores the data for responding to `send` events. + SetSelection { + /// MIME types to advertise. + mime_types: Vec, + /// Data for each MIME type (written to fd on `send` event). + data: HashMap>, + }, + /// Update source data for a MIME type without re-creating the source. + /// + /// Used when data wasn't available at `SetSelection` time (eager fetch + /// from a remote clipboard). The Wayland data source stays unchanged; + /// only the cached data map is updated so the next `send` event can + /// serve the requested MIME type. + UpdateSourceData { + /// MIME type key for the data. + mime_type: String, + /// Data bytes to cache. + data: Vec, + }, + /// Receive clipboard data from the current compositor offer. + /// + /// Calls `offer.receive(mime_type, fd)` on the event loop thread. + /// The caller reads from the other end of the pipe. + ReceiveFromOffer { + /// MIME type to request. + mime_type: String, + /// Write end of the pipe (compositor writes data here). + fd: OwnedFd, + }, + /// Answer a transfer raised through `Shared::on_transfer`. + /// + /// Writes `data` to every paste waiting on that transfer and caches it + /// for later pastes of the same MIME type. `None` closes the waiting + /// pastes with no data. + CompleteTransfer { + /// Serial passed to the `on_transfer` callback. + serial: u32, + /// The data, or `None` if the remote could not supply it. + data: Option>, + }, + /// Give up our selection, if we still own it. + ClearSelection, +} + +/// Shared clipboard state readable from any thread. +/// +/// Updated by the event loop thread when the compositor's selection changes. +/// Read by clipboard backends to report current state. +#[derive(Default)] +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) struct Shared { + /// MIME types of the current compositor selection. + pub(crate) mime_types: Vec, + /// Serial number, incremented on each selection change. + pub(crate) serial: u32, + /// Change notification callback. + /// + /// Called on the event loop thread when selection changes. + /// Typically captures a tokio channel sender for async notification. + pub(crate) on_change: Option) + Send + Sync>>, + /// Paste callback for data not supplied up front. + /// + /// When set, a paste of an advertised MIME type with no cached data is + /// held open and this is called with a serial and the MIME type; the + /// answer comes back as `Command::CompleteTransfer`. Without + /// it such a paste gets no data. + pub(crate) on_transfer: Option>, + /// Whether our own source is still the compositor's selection source. + /// + /// While it is, data we cached ourselves can be read back without a + /// round trip; once another client takes the selection it must not be. + pub(crate) own_source_live: bool, +} + +#[cfg(feature = "__test")] +impl Shared { + /// Test-only: MIME types of the current selection. + pub fn mime_types(&self) -> &[String] { + &self.mime_types + } + + /// Test-only: the selection change counter. + pub fn serial(&self) -> u32 { + self.serial + } + + /// Test-only: install the change callback. + pub fn set_on_change(&mut self, callback: Arc) + Send + Sync>) { + self.on_change = Some(callback); + } + + /// Test-only: install the transfer callback. + pub fn set_on_transfer(&mut self, callback: Arc) { + self.on_transfer = Some(callback); + } +} + +// === Data control state === + +/// A paste waiting on data from `on_transfer`, with every paste of the same +/// MIME type that arrived meanwhile. +struct PendingTransfer { + serial: u32, + waiters: Vec, +} + +/// Accumulated MIME types for a pending data offer. +/// +/// Between `data_offer` and `selection` events, the compositor sends +/// `offer` events with MIME types. We collect them here. +#[derive(Default)] +struct PendingOffer { + /// MIME types accumulated from `offer` events. + mime_types: Vec, +} + +/// Central data control state. +/// +/// Manages the lifecycle of data control protocol objects and routes +/// events to the shared clipboard state. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) struct State { + /// The data control manager global. + pub(crate) manager: Option, + /// The data control device (per-seat). + pub(crate) device: Option, + /// The current selection offer from the compositor. + current_offer: Option, + /// The data source we created for `SetSelection` (if any). + current_source: Option, + /// MIME types advertised on `current_source`. + pub(crate) current_source_mime_types: Vec, + /// Data cached for our source's `send` events. + pub(crate) source_data: HashMap>, + /// Pending offer being built up (between `data_offer` and `selection` events). + pending_offer: Option<(DataControlOffer, PendingOffer)>, + /// Pastes waiting on `on_transfer`, by requested MIME type. + pending_transfers: HashMap, + /// Next serial handed to `on_transfer`. + next_transfer_serial: u32, + /// Shared clipboard state for cross-thread access. + #[expect(clippy::struct_field_names)] + pub(crate) shared_state: Arc>, +} + +impl Default for State { + fn default() -> Self { + Self { + manager: None, + device: None, + current_offer: None, + current_source: None, + current_source_mime_types: Vec::new(), + source_data: HashMap::new(), + pending_offer: None, + pending_transfers: HashMap::new(), + next_transfer_serial: 1, + 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); + } + + /// Test-only: whether any source data is cached. + #[cfg(feature = "__test")] + pub fn has_source_data(&self) -> bool { + !self.source_data.is_empty() + } + + /// Test-only: pretend our source advertised these MIME types. + #[cfg(feature = "__test")] + pub fn set_advertised(&mut self, mime_types: Vec) { + self.current_source_mime_types = mime_types; + } + + /// Test-only: the shared state the callbacks and readers see. + #[cfg(feature = "__test")] + pub fn shared(&self) -> &Arc> { + &self.shared_state + } + + /// Create a data control device from the manager and seat. + /// + /// Must be called after both the manager and seat are bound. + pub(crate) fn create_device(&mut self, seat: &WlSeat, qh: &QueueHandle) + where + D: Dispatch + Dispatch + 'static, + { + let device = match &self.manager { + Some(DataControlManager::Ext(mgr)) => DataControlDevice::Ext(mgr.get_data_device(seat, qh, ())), + Some(DataControlManager::Wlr(mgr)) => DataControlDevice::Wlr(mgr.get_data_device(seat, qh, ())), + None => { + tracing::error!("Cannot create data control device: manager not bound"); + return; + } + }; + + tracing::debug!("Created data control device"); + self.device = Some(device); + } + + /// Handle a `data_offer` event from the device. + /// + /// A new offer is being introduced. Store it and start collecting + /// MIME types from subsequent `offer` events. + pub(crate) fn on_data_offer_ext(&mut self, offer: ExtDataControlOfferV1) { + self.set_pending_offer(DataControlOffer::Ext(offer)); + } + + /// Handle a `data_offer` event from the device (wlr variant). + pub(crate) fn on_data_offer_wlr(&mut self, offer: ZwlrDataControlOfferV1) { + self.set_pending_offer(DataControlOffer::Wlr(offer)); + } + + fn set_own_source_live(&self, live: bool) { + if let Ok(mut shared) = self.shared_state.lock() { + shared.own_source_live = live; + } + } + + fn set_pending_offer(&mut self, offer: DataControlOffer) { + // Destroy any previous pending offer that wasn't used + if let Some((old_offer, _)) = self.pending_offer.take() { + old_offer.destroy(); + } + self.pending_offer = Some((offer, PendingOffer::default())); + } + + /// Handle an `offer` event on a data offer (MIME type offered). + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_offer_mime_type(&mut self, mime_type: String) { + if let Some((_, ref mut pending)) = self.pending_offer { + pending.mime_types.push(mime_type); + } + } + + /// Handle the `selection` event from the device. + /// + /// The compositor's selection has changed. The pending offer + /// (with accumulated MIME types) becomes the current offer. + pub(crate) fn on_selection(&mut self) { + // Destroy the old current offer + if let Some(old) = self.current_offer.take() { + old.destroy(); + } + + // Promote the pending offer to current + let mime_types = if let Some((offer, pending)) = self.pending_offer.take() { + let types = pending.mime_types; + self.current_offer = Some(offer); + types + } else { + // NULL selection (clipboard cleared) + Vec::new() + }; + + // The device reports every selection, our own included. Reporting + // our own back as a change makes the consumer treat it as another + // client's copy. + let own = is_own_selection( + self.current_source.is_some(), + &self.current_source_mime_types, + &mime_types, + ); + tracing::debug!( + mime_types = ?mime_types, + own, + "Compositor selection changed" + ); + + // The callback runs after the lock is released: it may call back into + // the handle (`serial`, `selection_mime_types`), which locks the same state. + let callback = self.shared_state.lock().ok().and_then(|mut shared| { + shared.serial += 1; + shared.mime_types.clone_from(&mime_types); + if own { None } else { shared.on_change.clone() } + }); + if let Some(callback) = callback { + callback(mime_types); + } + } + + /// Handle the `selection` event with a NULL offer (selection cleared). + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_selection_cleared(&mut self) { + if let Some(old) = self.current_offer.take() { + old.destroy(); + } + // Clear any pending offer too + if let Some((offer, _)) = self.pending_offer.take() { + offer.destroy(); + } + + tracing::debug!("Compositor selection cleared"); + + let callback = self.shared_state.lock().ok().and_then(|mut shared| { + shared.serial += 1; + shared.mime_types.clear(); + shared.on_change.clone() + }); + if let Some(callback) = callback { + callback(Vec::new()); + } + } + + /// Update cached source data for a MIME type. + /// + /// Inserts or replaces data in the `source_data` map without + /// re-creating the Wayland data source. Used when data arrives + /// after `set_selection` was called with an empty data map + /// (eager fetch from a remote clipboard). + pub(crate) fn update_source_data(&mut self, mime_type: String, data: Vec) { + tracing::debug!( + mime_type = %mime_type, + bytes = data.len(), + "Source data updated (post-announcement)" + ); + self.source_data.insert(mime_type, data); + } + + /// Handle a `send` event on our data source. + /// + /// The compositor (or another client pasting) wants data in the + /// specified MIME type. Cached data is written at once; otherwise, with + /// an `on_transfer` callback set, the paste is held until the data + /// arrives through `CompleteTransfer`. + #[cfg_attr(feature = "__test", visibility::make(pub))] + pub(crate) fn on_source_send(&mut self, mime_type: &str, fd: OwnedFd) { + if let Some(data) = self.cached_data(mime_type) { + write_in_background(fd, data); + return; + } + + if let Some(pending) = self.pending_transfers.get_mut(mime_type) { + pending.waiters.push(fd); + return; + } + + let on_transfer = self + .shared_state + .lock() + .ok() + .and_then(|shared| shared.on_transfer.clone()); + let advertised = self.current_source_mime_types.iter().any(|m| m == mime_type); + let (Some(on_transfer), true) = (on_transfer, advertised) else { + tracing::warn!(mime_type, "Source send event for unknown MIME type"); + // fd is dropped/closed here, signaling no data + return; + }; + + let serial = self.next_transfer_serial; + self.next_transfer_serial = self.next_transfer_serial.wrapping_add(1); + self.pending_transfers.insert( + mime_type.to_owned(), + PendingTransfer { + serial, + waiters: vec![fd], + }, + ); + tracing::debug!(mime_type, serial, "Paste held until its data arrives"); + on_transfer(serial, mime_type.to_owned()); + } + + /// Data cached for `mime_type`, tolerating a charset parameter + /// (compositors commonly request `text/plain;charset=utf-8` for + /// `text/plain`). + fn cached_data(&self, mime_type: &str) -> Option> { + self.source_data + .get(mime_type) + .or_else(|| { + let base = mime_type.split(';').next()?.trim(); + self.source_data.get(base) + }) + .map(|data| Arc::from(data.as_slice())) + } + + /// 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 shared: Arc<[u8]> = Arc::from(data.as_slice()); + for fd in pending.waiters { + write_in_background(fd, Arc::clone(&shared)); + } + self.source_data.insert(mime_type, data); + } + + /// 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; + } + }; + + // Advertise MIME types + for mime_type in mime_types { + new_source.offer(mime_type); + } + + // Set as selection on the device + match (&self.device, &new_source) { + (Some(DataControlDevice::Ext(dev)), DataControlSource::Ext(src)) => { + dev.set_selection(Some(src)); + } + (Some(DataControlDevice::Wlr(dev)), DataControlSource::Wlr(src)) => { + dev.set_selection(Some(src)); + } + _ => { + tracing::error!("Cannot set selection: device/source protocol mismatch or no device"); + new_source.destroy(); + return; + } + } + + tracing::debug!( + mime_types = ?mime_types, + "Set clipboard selection on compositor" + ); + + // Replace, then destroy: destroying the current source first would + // leave the clipboard empty for a moment, which clipboard managers + // (Klipper's "prevent empty clipboard") answer by re-offering their + // last history item. + if let Some(old) = self.current_source.take() { + old.destroy(); + } + self.pending_transfers.clear(); + self.source_data = data; + self.current_source = Some(new_source); + self.current_source_mime_types = mime_types.to_vec(); + self.set_own_source_live(true); + } + + /// Process a `ReceiveFromOffer` command. + /// + /// Calls `offer.receive(mime_type, fd)` to request data from the + /// compositor. The caller reads from the other end of the pipe. + #[expect( + clippy::needless_pass_by_value, + reason = "OwnedFd must be owned so it is dropped after the receive call" + )] + pub(crate) fn receive_from_offer(&self, mime_type: &str, fd: OwnedFd) { + match &self.current_offer { + Some(offer) => { + offer.receive(mime_type, &fd); + tracing::debug!(mime_type, "Requested clipboard data from compositor offer"); + } + None => { + tracing::warn!(mime_type, "ReceiveFromOffer but no current offer"); + // fd is dropped, closing the pipe — caller will get EOF + } + } + } +} + +/// Whether a selection event reports the selection this client set itself. +/// +/// Another client taking the selection cancels our source first, so while +/// our source is live the selection is still ours. The offered MIME set must +/// also match what we advertised: that keeps a foreign copy from being +/// swallowed on a compositor that delivers the selection before the cancel. +#[cfg_attr(feature = "__test", visibility::make(pub))] +pub(crate) fn is_own_selection(source_live: bool, advertised: &[String], offered: &[String]) -> bool { + if !source_live || advertised.len() != offered.len() { + return false; + } + offered.iter().all(|mime| advertised.contains(mime)) +} + +/// Answer a paste on a worker thread: a large image written into a slow +/// reader would otherwise block every other Wayland event. +fn write_in_background(fd: OwnedFd, data: Arc<[u8]>) { + use std::io::Write as _; + + let spawned = std::thread::Builder::new() + .name("data-control-send".into()) + .spawn(move || { + let mut file = std::fs::File::from(fd); + if let Err(e) = file.write_all(&data) { + tracing::debug!(error = %e, "Paste reader went away"); + } + }); + if let Err(e) = spawned { + tracing::error!(error = %e, "Failed to start a paste writer"); + } +} 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..4ff046e2fe --- /dev/null +++ b/crates/ironrdp-cliprdr-native/src/data_control/worker.rs @@ -0,0 +1,250 @@ +//! The worker thread that owns the Wayland connection. +//! +//! Wayland objects live on this one thread. Callers talk to it through a +//! command channel plus a wake socket, so a request is served immediately +//! instead of at the next poll timeout. + +use core::sync::atomic::{AtomicBool, Ordering}; +use std::{ + io::Read as _, + os::unix::{io::AsFd as _, net::UnixStream}, + sync::{Arc, Mutex, mpsc}, +}; + +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::info!( + preferred = %candidates[0], + used = %protocol, + "Preferred data-control protocol unavailable, using the other" + ); + } + return Ok((manager, *protocol)); + } + } + + if options.allow_fallback { + Err(Error::Unsupported) + } else { + Err(Error::ProtocolUnavailable) + } +} + +impl Worker { + /// Serve events and commands until asked to stop or the connection fails. + pub(crate) fn run(mut self) { + tracing::debug!("Data-control worker started"); + + while !self.stop.load(Ordering::Relaxed) { + if let Err(error) = self.queue.dispatch_pending(&mut self.client) { + tracing::error!(%error, "Wayland dispatch failed"); + break; + } + + self.process_commands(); + + let flush_blocked = match self.connection.flush() { + Ok(()) => false, + Err(WaylandError::Io(ref error)) if error.kind() == std::io::ErrorKind::WouldBlock => true, + Err(error) => { + tracing::error!(%error, "Wayland flush failed"); + break; + } + }; + + // `None` means events arrived after the dispatch above; loop to + // handle them before waiting. + let Some(guard) = self.queue.prepare_read() else { + continue; + }; + + let wayland_events = if flush_blocked { + PollFlags::POLLIN | PollFlags::POLLOUT + } else { + PollFlags::POLLIN + }; + let mut fds = [ + PollFd::new(guard.connection_fd(), wayland_events), + PollFd::new(self.wake.as_fd(), PollFlags::POLLIN), + ]; + match poll(&mut fds, PollTimeout::NONE) { + 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-core/Cargo.toml b/crates/ironrdp-testsuite-core/Cargo.toml index 2654492a8b..2d394583ac 100644 --- a/crates/ironrdp-testsuite-core/Cargo.toml +++ b/crates/ironrdp-testsuite-core/Cargo.toml @@ -39,6 +39,7 @@ expect-test.workspace = true hex = "0.4" ironrdp-cliprdr-format.path = "../ironrdp-cliprdr-format" ironrdp-cliprdr = { path = "../ironrdp-cliprdr", features = ["__test"] } +ironrdp-cliprdr-native = { path = "../ironrdp-cliprdr-native", features = ["__test"] } ironrdp-acceptor.path = "../ironrdp-acceptor" ironrdp-async.path = "../ironrdp-async" ironrdp-tokio.path = "../ironrdp-tokio" diff --git a/crates/ironrdp-testsuite-core/tests/cliprdr_native/data_control.rs b/crates/ironrdp-testsuite-core/tests/cliprdr_native/data_control.rs new file mode 100644 index 0000000000..3458692168 --- /dev/null +++ b/crates/ironrdp-testsuite-core/tests/cliprdr_native/data_control.rs @@ -0,0 +1,261 @@ +//! Tests for the Wayland data-control clipboard client that need no compositor. + +use std::{ + io::Read as _, + os::fd::OwnedFd, + sync::{Arc, Mutex}, +}; + +use ironrdp_cliprdr_native::data_control::state::{State, is_own_selection}; +use ironrdp_cliprdr_native::data_control::{Options, Preference, Protocol, find_mime_match}; + +fn mimes(types: &[&str]) -> Vec { + types.iter().map(|t| (*t).to_owned()).collect() +} + +fn read_all(fd: OwnedFd) -> Vec { + let mut buf = Vec::new(); + std::fs::File::from(fd).read_to_end(&mut buf).unwrap(); + buf +} + +fn pipe() -> (OwnedFd, OwnedFd) { + let (reader, writer) = std::io::pipe().unwrap(); + (OwnedFd::from(reader), OwnedFd::from(writer)) +} + +#[test] +fn mime_exact_match_is_returned() { + let available = mimes(&["text/plain", "text/html"]); + assert_eq!(find_mime_match("text/html", &available), Some("text/html")); +} + +#[test] +fn mime_requested_charset_is_stripped() { + let available = mimes(&["text/plain"]); + assert_eq!( + find_mime_match("text/plain;charset=utf-8", &available), + Some("text/plain") + ); +} + +#[test] +fn mime_missing_charset_is_added() { + let available = mimes(&["text/plain;charset=utf-8"]); + assert_eq!( + find_mime_match("text/plain", &available), + Some("text/plain;charset=utf-8") + ); +} + +#[test] +fn mime_non_text_types_get_no_charset_fallback() { + let available = mimes(&["image/png;foo=bar"]); + assert_eq!(find_mime_match("image/png", &available), None); +} + +#[test] +fn mime_exact_match_beats_a_charset_variant() { + let available = mimes(&["text/plain;charset=utf-8", "text/plain"]); + assert_eq!(find_mime_match("text/plain", &available), Some("text/plain")); +} + +#[test] +fn mime_any_charset_of_the_same_base_type_is_accepted() { + let available = mimes(&["text/plain;charset=iso-8859-1"]); + assert_eq!( + find_mime_match("text/plain", &available), + Some("text/plain;charset=iso-8859-1") + ); +} + +#[test] +fn auto_preference_tries_ext_then_wlr() { + assert_eq!(Options::new().candidates(), [Protocol::Ext, Protocol::Wlr]); +} + +#[test] +fn wlr_preference_tries_wlr_first() { + let options = Options::new().preference(Preference::Wlr); + assert_eq!(options.candidates(), [Protocol::Wlr, Protocol::Ext]); +} + +#[test] +fn disabling_fallback_leaves_only_the_preferred_protocol() { + let options = Options::new().preference(Preference::Wlr).allow_fallback(false); + assert_eq!(options.candidates(), [Protocol::Wlr]); +} + +#[test] +fn own_selection_is_recognised_while_our_source_is_live() { + let ours = mimes(&["text/plain;charset=utf-8"]); + assert!(is_own_selection(true, &ours, &mimes(&["text/plain;charset=utf-8"]))); +} + +#[test] +fn a_selection_after_our_source_was_cancelled_is_foreign() { + let ours = mimes(&["text/plain;charset=utf-8"]); + assert!(!is_own_selection(false, &ours, &mimes(&["text/plain;charset=utf-8"]))); +} + +#[test] +fn a_different_offer_while_our_source_is_live_is_foreign() { + let ours = mimes(&["text/plain;charset=utf-8"]); + let other = mimes(&[ + "text/plain;charset=utf-8", + "text/plain", + "TEXT", + "STRING", + "UTF8_STRING", + ]); + assert!(!is_own_selection(true, &ours, &other)); +} + +#[test] +fn own_selection_ignores_offer_order() { + let ours = mimes(&["image/png", "text/html"]); + assert!(is_own_selection(true, &ours, &mimes(&["text/html", "image/png"]))); +} + +#[test] +fn clearing_the_selection_notifies_and_advances_the_serial() { + let mut state = State::default(); + let called = Arc::new(Mutex::new(false)); + let seen = Arc::clone(&called); + state + .shared() + .lock() + .unwrap() + .set_on_change(Arc::new(move |types: Vec| { + assert!(types.is_empty()); + *seen.lock().unwrap() = true; + })); + + state.on_selection_cleared(); + + assert!(*called.lock().unwrap()); + let shared = state.shared().lock().unwrap(); + assert!(shared.mime_types().is_empty()); + assert_eq!(shared.serial(), 1); +} + +#[test] +fn cancelling_our_source_drops_its_data() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello".to_vec()); + state.on_source_cancelled(); + assert!(!state.has_source_data()); +} + +#[test] +fn device_finished_drops_source_data() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello".to_vec()); + state.on_device_finished(); + assert!(!state.has_source_data()); +} + +#[test] +fn a_mime_type_without_a_pending_offer_is_ignored() { + let mut state = State::default(); + state.on_offer_mime_type("text/plain".to_owned()); +} + +#[test] +fn a_paste_of_an_unknown_type_closes_the_pipe_with_no_data() { + let mut state = State::default(); + let (read, write) = pipe(); + state.on_source_send("text/unknown", write); + assert!(read_all(read).is_empty()); +} + +#[test] +fn a_paste_is_answered_from_cached_data() { + let mut state = State::default(); + state.insert_source_data("text/plain", b"hello world".to_vec()); + let (read, write) = pipe(); + state.on_source_send("text/plain", write); + assert_eq!(read_all(read), b"hello world"); +} + +type TransferLog = Arc>>; + +/// State advertising `image/png` with a transfer callback that records what it was asked for. +fn state_with_transfer_callback() -> (State, TransferLog) { + let mut state = State::default(); + state.set_advertised(mimes(&["image/png"])); + let requests: TransferLog = Arc::default(); + let seen = Arc::clone(&requests); + state + .shared() + .lock() + .unwrap() + .set_on_transfer(Arc::new(move |serial, mime| { + seen.lock().unwrap().push((serial, mime)); + })); + (state, requests) +} + +#[test] +fn a_paste_without_data_is_held_until_the_transfer_completes() { + let (mut state, requests) = state_with_transfer_callback(); + let (first_read, first_write) = pipe(); + let (second_read, second_write) = pipe(); + + state.on_source_send("image/png", first_write); + state.on_source_send("image/png", second_write); + // Both pastes share one transfer. + assert_eq!(*requests.lock().unwrap(), vec![(1, "image/png".to_owned())]); + + state.complete_transfer(1, Some(b"png bytes".to_vec())); + assert_eq!(read_all(first_read), b"png bytes"); + assert_eq!(read_all(second_read), b"png bytes"); + + // A later paste is served from the cache without a new transfer. + let (third_read, third_write) = pipe(); + state.on_source_send("image/png", third_write); + assert_eq!(read_all(third_read), b"png bytes"); + assert_eq!(requests.lock().unwrap().len(), 1); +} + +#[test] +fn a_held_paste_gets_nothing_when_the_selection_is_lost() { + let (mut state, _requests) = state_with_transfer_callback(); + let (read, write) = pipe(); + state.on_source_send("image/png", write); + + state.on_source_cancelled(); + assert!(read_all(read).is_empty()); + // An answer for the lost selection is ignored. + state.complete_transfer(1, Some(b"late".to_vec())); + assert!(!state.has_source_data()); +} + +#[test] +fn a_paste_of_an_unadvertised_type_raises_no_transfer() { + let (mut state, requests) = state_with_transfer_callback(); + let (read, write) = pipe(); + state.on_source_send("text/html", write); + assert!(read_all(read).is_empty()); + assert!(requests.lock().unwrap().is_empty()); +} + +#[test] +fn the_change_callback_may_lock_the_shared_state() { + let mut state = State::default(); + let shared = Arc::clone(state.shared()); + let relocked = Arc::new(Mutex::new(false)); + let seen = Arc::clone(&relocked); + state + .shared() + .lock() + .unwrap() + .set_on_change(Arc::new(move |_types: Vec| { + // A callback that reads the handle locks this state; it must not be held. + *seen.lock().unwrap() = shared.try_lock().is_ok(); + })); + + state.on_selection_cleared(); + + assert!(*relocked.lock().unwrap()); +} diff --git a/crates/ironrdp-testsuite-core/tests/cliprdr_native/mod.rs b/crates/ironrdp-testsuite-core/tests/cliprdr_native/mod.rs new file mode 100644 index 0000000000..4a68c78fc9 --- /dev/null +++ b/crates/ironrdp-testsuite-core/tests/cliprdr_native/mod.rs @@ -0,0 +1 @@ +mod data_control; diff --git a/crates/ironrdp-testsuite-core/tests/main.rs b/crates/ironrdp-testsuite-core/tests/main.rs index 161c01bed0..107932172d 100644 --- a/crates/ironrdp-testsuite-core/tests/main.rs +++ b/crates/ironrdp-testsuite-core/tests/main.rs @@ -14,6 +14,8 @@ mod cfg; mod clipboard; +#[cfg(target_os = "linux")] +mod cliprdr_native; mod connector; mod displaycontrol; mod dvc;