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 }