git.lucas.co / cce-core
GUI-free half of the cce toolkit: config, input, IPC, spec parsers
git clone https://git.lucas.co/cce-core.git

src/ipc.rs (13.7K)

  1 //! Helpers for the CCE Unix-socket IPC convention: `/tmp/<prefix>-<WAYLAND_DISPLAY>.sock`.
  2 
  3 use std::io::{Read, Write};
  4 use std::os::unix::net::UnixStream;
  5 
  6 pub mod instance;
  7 
  8 /// Path of a CCE IPC socket for `prefix`, keyed by `$WAYLAND_DISPLAY`.
  9 ///
 10 /// `socket_path("cce")` → `/tmp/cce-<display>.sock` (the compositor control socket);
 11 /// `socket_path("cce-status-interface")` → the status socket. Falls back to
 12 /// `/tmp/<prefix>.sock` when `$WAYLAND_DISPLAY` is unset.
 13 pub fn socket_path(prefix: &str) -> String {
 14     match std::env::var("WAYLAND_DISPLAY") {
 15         Ok(d) if !d.is_empty() => format!("/tmp/{}-{}.sock", prefix, d),
 16         _ => format!("/tmp/{}.sock", prefix),
 17     }
 18 }
 19 
 20 /// Connect to the `prefix` socket, send `command` (newline-terminated), and
 21 /// return the reply text. Errors if the socket can't be reached.
 22 pub fn send_command(prefix: &str, command: &str) -> std::io::Result<String> {
 23     let mut stream = UnixStream::connect(socket_path(prefix))?;
 24     stream.write_all(command.as_bytes())?;
 25     if !command.ends_with('\n') {
 26         stream.write_all(b"\n")?;
 27     }
 28     let mut reply = String::new();
 29     stream.read_to_string(&mut reply)?;
 30     Ok(reply)
 31 }
 32 
 33 /// Ask the compositor to dissolve this client's surfaces out, and return how
 34 /// long it says that will take.
 35 ///
 36 /// The fade is the compositor's, not the app's: it ramps the opacity of the
 37 /// scene subtree, which carries the backdrop blur, the drop shadow and the
 38 /// bevel down with the window. A client fading its own pixels instead leaves
 39 /// its surface fully present, so the blur behind it hangs at full strength
 40 /// over a dissolving window — and any part of its drawing that is not plain
 41 /// vertex alpha (shader-lit plate rims, specular) does not fade at all.
 42 ///
 43 /// The contract is that the caller keeps its surfaces mapped and its process
 44 /// alive for the returned duration and only then exits. `window_runner` does
 45 /// that for every cce-ui `Application`;
 46 /// an app driving its own event loop calls this itself. Zero — no compositor,
 47 /// nothing of ours on screen, or fading configured off — means exit now.
 48 pub fn request_close_fade() -> std::time::Duration {
 49     // This runs on the exit path of every cce-ui app, so it does its own
 50     // socket call rather than `send_command`: that one reads to EOF with no
 51     // deadline, and a compositor wedged mid-frame would hang the quit
 52     // forever. A second is far longer than an IPC round trip and short
 53     // enough that a user who hit Close still sees the window go.
 54     const REPLY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(1);
 55     let Ok(mut stream) = UnixStream::connect(socket_path("cce")) else {
 56         return std::time::Duration::ZERO;
 57     };
 58     let _ = stream.set_write_timeout(Some(REPLY_TIMEOUT));
 59     let _ = stream.set_read_timeout(Some(REPLY_TIMEOUT));
 60     if stream.write_all(b"fade-out\n").is_err() {
 61         return std::time::Duration::ZERO;
 62     }
 63     // The duration is the compositor's to decide (`surface { fade out_ms }`),
 64     // so it is read back rather than assumed: the two sides would otherwise
 65     // drift apart the moment the config changed, and the visible failure —
 66     // the window vanishing partway through its own dissolve — reads as a
 67     // rendering bug rather than a disagreement about a number.
 68     let mut reply = String::new();
 69     if stream.read_to_string(&mut reply).is_err() {
 70         return std::time::Duration::ZERO;
 71     }
 72     let ms: u64 = reply.trim().parse().unwrap_or(0);
 73     // The compositor clamps this already; clamped again here because a
 74     // client must never be made to hang on a number from the other side.
 75     std::time::Duration::from_millis(ms.min(2000))
 76 }
 77 
 78 /// Ask the compositor to bring a window to the user: un-minimize it, focus
 79 /// it, raise it, and pan the camera to it — `ccectl focus-window <query>`.
 80 ///
 81 /// `query` is what the compositor resolves: a numeric window id, else an
 82 /// app_id (an exact match beats a substring one). An app bringing *itself*
 83 /// forward passes its own app_id — which is what a single-instance app does
 84 /// when a later launch forwards to it, or the work lands in a window parked
 85 /// off-camera and the launch looks like it did nothing. (xdg-activation is
 86 /// not the route to this: the compositor deliberately answers it with an
 87 /// attention notification, not focus.)
 88 ///
 89 /// Blocks for one round trip, bounded at a second for the same reason as
 90 /// [`request_close_fade`] — `send_command` reads with no deadline, and a
 91 /// compositor wedged mid-frame must not hang the caller. `Err` when there is
 92 /// no compositor to ask, it does not answer in time, it answers `error: …`
 93 /// (no such window, no seat), or `query` is empty or spans lines.
 94 pub fn focus_window(query: &str) -> std::io::Result<()> {
 95     focus_window_at(&socket_path("cce"), query)
 96 }
 97 
 98 fn focus_window_at(path: &str, query: &str) -> std::io::Result<()> {
 99     use std::io::{Error, ErrorKind};
100     const REPLY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(1);
101     let query = query.trim();
102     // A newline would end this command and start another on the control
103     // socket; refuse rather than pass an injection along.
104     if query.is_empty() || query.contains(['\n', '\r']) {
105         return Err(Error::new(ErrorKind::InvalidInput, format!("not a window query: {query:?}")));
106     }
107     let mut stream = UnixStream::connect(path)?;
108     stream.set_write_timeout(Some(REPLY_TIMEOUT))?;
109     stream.set_read_timeout(Some(REPLY_TIMEOUT))?;
110     stream.write_all(format!("focus-window {query}\n").as_bytes())?;
111     let mut reply = String::new();
112     stream.read_to_string(&mut reply)?;
113     match reply.trim() {
114         "ok" => Ok(()),
115         other => Err(Error::other(format!("focus-window {query}: {other}"))),
116     }
117 }
118 
119 /// One request line from a socket client, bounded in size and in TOTAL time.
120 ///
121 /// For a listener thread that serves one client after another: a plain
122 /// `BufReader::read_line` with no timeout lets a client that connects and
123 /// says nothing (or trickles a byte at a time) hold the thread for good, and
124 /// every later client waits behind it. A per-read timeout alone is not
125 /// enough, since a trickle resets it with every byte; the deadline here is
126 /// for the whole line.
127 ///
128 /// Returns the line with its `\n` (or what arrived before EOF). `None` on
129 /// EOF before any byte, the deadline, more than `limit` bytes without a
130 /// newline, a read error, or invalid UTF-8. The stream's read timeout is
131 /// cleared again on success.
132 ///
133 /// Nothing past the newline is consumed: each chunk is PEEKED first and only
134 /// the bytes up to the newline are read, so whatever the client sent after
135 /// the line is still on the socket for the caller to read. Until 2026-10-06
136 /// the whole chunk was read and the tail truncated away — cce-cloud's
137 /// switcher client sends its request line and then the window list down the
138 /// same connection, and when both had arrived by the time the daemon read,
139 /// the list went with the tail and the switcher opened empty.
140 pub fn read_request_line(conn: &UnixStream, limit: usize, deadline: std::time::Duration) -> Option<String> {
141     let until = std::time::Instant::now() + deadline;
142     let mut buf: Vec<u8> = Vec::new();
143     let mut chunk = [0u8; 4096];
144     let mut reader = conn;
145     loop {
146         let left = until.checked_duration_since(std::time::Instant::now()).filter(|d| !d.is_zero())?;
147         conn.set_read_timeout(Some(left)).ok()?;
148         let n = match peek(conn, &mut chunk) {
149             Ok(n) => n,
150             Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
151             Err(_) => return None,
152         };
153         if n == 0 {
154             if buf.is_empty() {
155                 return None;
156             }
157             break;
158         }
159         // Consume through the newline if the peek holds one, else all of it.
160         let take = chunk[..n].iter().position(|&b| b == b'\n').map_or(n, |end| end + 1);
161         // Already queued, so this returns at once with exactly `take` bytes.
162         reader.read_exact(&mut chunk[..take]).ok()?;
163         buf.extend_from_slice(&chunk[..take]);
164         if buf.last() == Some(&b'\n') {
165             break;
166         }
167         if buf.len() > limit {
168             return None;
169         }
170     }
171     let _ = conn.set_read_timeout(None);
172     String::from_utf8(buf).ok()
173 }
174 
175 /// `recv(MSG_PEEK)`: what is queued on the socket, left there. Blocks (up to
176 /// the read timeout) like `read` when nothing is. `UnixStream::peek` is
177 /// still unstable.
178 fn peek(conn: &UnixStream, buf: &mut [u8]) -> std::io::Result<usize> {
179     use std::os::fd::AsRawFd;
180     // SAFETY: `buf` is a live, writable slice of `buf.len()` bytes.
181     let n = unsafe { libc::recv(conn.as_raw_fd(), buf.as_mut_ptr().cast(), buf.len(), libc::MSG_PEEK) };
182     if n < 0 {
183         Err(std::io::Error::last_os_error())
184     } else {
185         Ok(n as usize)
186     }
187 }
188 
189 #[cfg(test)]
190 mod tests {
191     use super::read_request_line;
192     use std::io::{Read, Write};
193     use std::os::unix::net::UnixStream;
194     use std::time::{Duration, Instant};
195 
196     #[test]
197     fn a_request_line_is_bounded_in_time_and_size() {
198         let (mut client, server) = UnixStream::pair().unwrap();
199         client.write_all(b"open note\nmore").unwrap();
200         assert_eq!(read_request_line(&server, 1024, Duration::from_secs(1)).as_deref(), Some("open note\n"));
201 
202         // Only the line is consumed: what the client sent after it in the
203         // same burst (cce-cloud's switcher list) is still there to read.
204         let (mut client, mut server) = UnixStream::pair().unwrap();
205         client.write_all(b"{\"args\":[]}\nWindow A (a)\nWindow B (b)\n").unwrap();
206         drop(client);
207         assert_eq!(read_request_line(&server, 1024, Duration::from_secs(1)).as_deref(), Some("{\"args\":[]}\n"));
208         let mut rest = String::new();
209         server.read_to_string(&mut rest).unwrap();
210         assert_eq!(rest, "Window A (a)\nWindow B (b)\n");
211 
212         // A line longer than one peek, then a tail: read across chunks,
213         // and the tail is still left.
214         let (mut client, mut server) = UnixStream::pair().unwrap();
215         let long = format!("{}\ntail", "y".repeat(10_000));
216         let writer = std::thread::spawn(move || client.write_all(long.as_bytes()).unwrap());
217         let line = read_request_line(&server, 64 * 1024, Duration::from_secs(1)).unwrap();
218         writer.join().unwrap();
219         assert_eq!(line.len(), 10_001);
220         let mut rest = String::new();
221         server.read_to_string(&mut rest).unwrap();
222         assert_eq!(rest, "tail");
223 
224         // Silent: given up at the deadline, not held forever.
225         let (_quiet, server) = UnixStream::pair().unwrap();
226         let t = Instant::now();
227         assert_eq!(read_request_line(&server, 1024, Duration::from_millis(100)), None);
228         assert!(t.elapsed() < Duration::from_millis(500));
229 
230         // Trickling a byte at a time: the TOTAL deadline still ends it.
231         let (mut client, server) = UnixStream::pair().unwrap();
232         let trickle = std::thread::spawn(move || {
233             for _ in 0..40 {
234                 if client.write_all(b"x").is_err() {
235                     break;
236                 }
237                 std::thread::sleep(Duration::from_millis(25));
238             }
239         });
240         let t = Instant::now();
241         assert_eq!(read_request_line(&server, 1024, Duration::from_millis(150)), None);
242         assert!(t.elapsed() < Duration::from_millis(500), "took {:?}", t.elapsed());
243         drop(server);
244         trickle.join().unwrap();
245 
246         // Over the size cap.
247         let (mut client, server) = UnixStream::pair().unwrap();
248         client.write_all(&[b'z'; 2000]).unwrap();
249         assert_eq!(read_request_line(&server, 1024, Duration::from_millis(200)), None);
250     }
251 
252     /// Against a stand-in compositor on a private path (never the session's
253     /// control socket): the command it sends, `ok` as success, `error: …`
254     /// as an error, a silent compositor bounded, a bad query never sent.
255     #[test]
256     fn focus_window_sends_one_command_and_reads_the_verdict() {
257         use super::focus_window_at;
258         use std::io::Read;
259         use std::os::unix::net::UnixListener;
260 
261         let path = format!("/tmp/cce-ui-focus-test-{}.sock", std::process::id());
262         let _ = std::fs::remove_file(&path);
263         let listener = UnixListener::bind(&path).unwrap();
264         let replies: [&[u8]; 3] = [b"ok\n", b"error: window not found\n", b""];
265         let server = std::thread::spawn(move || {
266             let mut got = Vec::new();
267             for reply in replies {
268                 let (mut conn, _) = listener.accept().unwrap();
269                 let line = read_request_line(&conn, 1024, Duration::from_secs(1)).unwrap();
270                 got.push(line);
271                 if reply.is_empty() {
272                     // Say nothing and keep the connection open: a wedged compositor.
273                     let mut rest = Vec::new();
274                     let _ = conn.read_to_end(&mut rest);
275                 } else {
276                     conn.write_all(reply).unwrap();
277                 }
278             }
279             got
280         });
281 
282         assert!(focus_window_at(&path, "cce-browser").is_ok());
283         let err = focus_window_at(&path, "  nothing-here ").unwrap_err();
284         assert!(err.to_string().contains("window not found"), "{err}");
285         let t = Instant::now();
286         assert!(focus_window_at(&path, "12").is_err(), "a silent compositor is an error");
287         assert!(t.elapsed() < Duration::from_secs(3), "and a bounded one: {:?}", t.elapsed());
288         for bad in ["", "   ", "a\nquit", "a\rb"] {
289             assert_eq!(
290                 focus_window_at(&path, bad).unwrap_err().kind(),
291                 std::io::ErrorKind::InvalidInput,
292                 "{bad:?}"
293             );
294         }
295         let got = server.join().unwrap();
296         assert_eq!(got, ["focus-window cce-browser\n", "focus-window nothing-here\n", "focus-window 12\n"]);
297         let _ = std::fs::remove_file(&path);
298     }
299 }