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/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 }