GUI-free half of the cce toolkit: config, input, IPC, spec parsers
git clone https://git.lucas.co/cce-core.git
src/ipc/instance.rs (8.5K)
1 //! Single-instance apps over the CCE socket convention: claim or forward.
2 //!
3 //! An app that should run once per session (a link opened from elsewhere
4 //! becomes a tab, `cce-notes open X` shows X in the running window) listens
5 //! on `/tmp/<prefix>-<WAYLAND_DISPLAY>.sock` — keyed by display, so shadow
6 //! sessions stay apart for free — and every later launch hands it one line
7 //! and exits before any Wayland or engine work.
8 //!
9 //! ```ignore
10 //! fn main() {
11 //! if cce_ui::ipc::instance::forward_or_claim("cce-foo", &launch_line()) {
12 //! return; // the running instance took it
13 //! }
14 //! cce_ui::engine::run::<Foo>();
15 //! cce_ui::ipc::instance::cleanup();
16 //! }
17 //!
18 //! // in Application::new, once the loop's sender exists:
19 //! cce_ui::ipc::instance::serve(move |line| {
20 //! let msg = parse(line)?; // None: hang up without a reply
21 //! sender.send(msg).ok()?;
22 //! Some("ok".into())
23 //! });
24 //! ```
25 //!
26 //! The order in [`forward_or_claim`] is what closes the startup race: try to
27 //! connect, and only bind after a connect has failed. A refused connection
28 //! means the socket file outlived a crashed instance and is removed before
29 //! binding; losing the bind to a simultaneous launch falls back to one more
30 //! connect. If that also fails the launch proceeds un-listened rather than
31 //! not at all.
32 //!
33 //! The claimed listener has to survive from `main()` (before the engine
34 //! starts) to `Application::new` (where the app's loop sender first exists),
35 //! and `engine::run` takes no arguments, so it parks in this module until
36 //! [`serve`] adopts it. One claim per process.
37 //!
38 //! The wire protocol is the app's: one request line in, one reply line out.
39 //! Anything a forwarded argument needs to mean the same thing in the
40 //! instance — a relative path made absolute against the *sender's* cwd —
41 //! is the caller's to do before building the line.
42
43 use std::io::{BufRead, BufReader, Write};
44 use std::os::unix::net::{UnixListener, UnixStream};
45 use std::sync::Mutex;
46 use std::time::Duration;
47
48 /// The listener claimed by [`forward_or_claim`], waiting for [`serve`].
49 static CLAIMED: Mutex<Option<UnixListener>> = Mutex::new(None);
50 /// The socket path this process bound (and must unlink on exit), if any.
51 static OWNED_PATH: Mutex<Option<String>> = Mutex::new(None);
52
53 /// How long a launch waits for the running instance to answer. Bounded so an
54 /// instance whose listener is stuck cannot hang every later launch forever;
55 /// unanswered, the launch is not forwarded.
56 const FORWARD_TIMEOUT: Duration = Duration::from_secs(5);
57
58 /// How long, and how much, the listener reads from one client before giving
59 /// up on it and serving the next (see [`super::read_request_line`]).
60 const REQUEST_TIMEOUT: Duration = Duration::from_secs(2);
61 const REQUEST_LIMIT: usize = 64 * 1024;
62
63 /// Send `line` to the instance listening on `prefix`'s socket and return its
64 /// reply line, trimmed.
65 ///
66 /// `None` when nothing is listening, the write fails, or no reply arrives
67 /// within the timeout. `Some("")` when the instance read the line and hung
68 /// up without answering — it still received it. The client half alone, for
69 /// an app that talks to another app's instance.
70 pub fn forward(prefix: &str, line: &str) -> Option<String> {
71 forward_to(&super::socket_path(prefix), line)
72 }
73
74 fn forward_to(path: &str, line: &str) -> Option<String> {
75 let mut stream = UnixStream::connect(path).ok()?;
76 let mut msg = line.trim_end_matches('\n').to_string();
77 msg.push('\n');
78 stream.write_all(msg.as_bytes()).ok()?;
79 // Wait for the answer: returning (and exiting) on the write alone races
80 // the instance actually reading the line.
81 stream.set_read_timeout(Some(FORWARD_TIMEOUT)).ok()?;
82 let mut reply = String::new();
83 BufReader::new(stream).read_line(&mut reply).ok()?;
84 Some(reply.trim().to_string())
85 }
86
87 /// Hand `line` to a running instance, or claim the instance socket.
88 ///
89 /// `true` when a running instance took the launch: the caller should exit
90 /// without starting. `false` when this process is now the instance — the
91 /// listener parked for [`serve`] — or when single-instance handling failed
92 /// entirely and the launch should proceed standalone.
93 pub fn forward_or_claim(prefix: &str, line: &str) -> bool {
94 let path = super::socket_path(prefix);
95 if forward_to(&path, line).is_some() {
96 return true;
97 }
98 // Nothing answered. A socket file that still exists is a leftover from a
99 // crashed instance; binding needs it gone.
100 if std::path::Path::new(&path).exists() {
101 let _ = std::fs::remove_file(&path);
102 }
103 match UnixListener::bind(&path) {
104 Ok(listener) => {
105 *CLAIMED.lock().unwrap_or_else(|e| e.into_inner()) = Some(listener);
106 *OWNED_PATH.lock().unwrap_or_else(|e| e.into_inner()) = Some(path);
107 false
108 }
109 // Lost the bind race to a simultaneous launch: it is the instance.
110 Err(_) => forward_to(&path, line).is_some(),
111 }
112 }
113
114 /// Serve the listener [`forward_or_claim`] claimed, on a thread of its own.
115 ///
116 /// `handle` gets each request line, trimmed, and returns the reply line to
117 /// write back (a trailing newline is added), or `None` to hang up without
118 /// one. It runs on the listener thread, so it should hand work to the app's
119 /// loop (a calloop channel) rather than do it; answering straight from
120 /// shared state is fine for a query the loop need not see.
121 ///
122 /// Each client is read under a total deadline and size cap, so one that
123 /// connects and says nothing cannot wedge the listener for every launch
124 /// after it. `false` when this process claimed nothing (it runs standalone,
125 /// or `serve` already took the listener).
126 pub fn serve<F>(mut handle: F) -> bool
127 where
128 F: FnMut(&str) -> Option<String> + Send + 'static,
129 {
130 let Some(listener) = CLAIMED.lock().unwrap_or_else(|e| e.into_inner()).take() else {
131 return false;
132 };
133 std::thread::spawn(move || {
134 for conn in listener.incoming() {
135 let Ok(conn) = conn else { continue };
136 let Some(line) = super::read_request_line(&conn, REQUEST_LIMIT, REQUEST_TIMEOUT) else {
137 continue;
138 };
139 if let Some(mut reply) = handle(line.trim()) {
140 reply.push('\n');
141 let _ = (&conn).write_all(reply.as_bytes());
142 }
143 }
144 });
145 true
146 }
147
148 /// Unlink the socket if this process bound it. Call it after the engine loop
149 /// returns; a crash skips it, which is what the stale-socket removal in
150 /// [`forward_or_claim`] exists for.
151 pub fn cleanup() {
152 if let Some(path) = OWNED_PATH.lock().unwrap_or_else(|e| e.into_inner()).take() {
153 let _ = std::fs::remove_file(path);
154 }
155 }
156
157 #[cfg(test)]
158 mod tests {
159 use super::*;
160
161 /// The whole round trip in one process: the first call claims, `serve`
162 /// answers, a second launch is forwarded and sees the reply, and
163 /// `cleanup` removes the socket. One test, because the claim is
164 /// process-wide.
165 #[test]
166 fn claim_serve_forward_cleanup() {
167 let prefix = format!("cce-ui-instance-test-{}", std::process::id());
168 let path = crate::ipc::socket_path(&prefix);
169 // A stale file from a "crashed" instance is replaced, not fatal.
170 std::fs::write(&path, b"").unwrap();
171
172 assert!(!forward_or_claim(&prefix, "first"), "nothing was running: this launch is the instance");
173 let (tx, rx) = std::sync::mpsc::channel();
174 assert!(serve(move |line| {
175 tx.send(line.to_string()).ok()?;
176 match line {
177 "silent" => None,
178 l => Some(format!("ok {l}")),
179 }
180 }));
181 assert!(!serve(|_| None), "the listener is served once");
182
183 assert_eq!(forward(&prefix, "open x\n").as_deref(), Some("ok open x"));
184 assert!(forward_or_claim(&prefix, "second"), "a running instance takes the launch");
185 // Read and hung up on: still received.
186 assert_eq!(forward(&prefix, "silent").as_deref(), Some(""));
187 let got: Vec<String> = rx.try_iter().collect();
188 assert_eq!(got, ["open x", "second", "silent"]);
189
190 // A client that connects and says nothing does not wedge the listener.
191 let _quiet = UnixStream::connect(&path).unwrap();
192 let t = std::time::Instant::now();
193 assert_eq!(forward(&prefix, "after").as_deref(), Some("ok after"));
194 assert!(t.elapsed() < FORWARD_TIMEOUT, "took {:?}", t.elapsed());
195
196 cleanup();
197 assert!(!std::path::Path::new(&path).exists());
198 assert_eq!(forward(&prefix, "gone"), None);
199 }
200 }