git.lucas.co / cce-notes
notes editor over the vault (Obsidian-compatible)
git clone https://git.lucas.co/cce-notes.git

commit68d39750fbcc29f7ee075d4432c30bda8940baaf
parente2cb7e6672
authorLucas Galante <lsgalante12@gmail.com>
date2026-10-02 13:18
refactor(instance): single-instance on cce_ui::ipc::instance

This was cce-browser's claim-or-forward copied in, bounded read included;
both live in cce-ui now (40c66d4). instance.rs keeps the notes protocol:
Command, its lines, and `current` answered from the listener.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>

 CLAUDE.md       |   2 +-
 src/instance.rs | 166 +++++++-------------------------------------------------
 2 files changed, 21 insertions(+), 147 deletions(-)

diff --git a/CLAUDE.md b/CLAUDE.md
index 14c58a5..338a8e9 100644
--- a/CLAUDE.md
+++ b/CLAUDE.md
@@ -27,7 +27,7 @@ proposal".
 | `mcp.rs` | MCP tools (`search`, `find_notes`, `read_note`, `backlinks`, `open_note`, `append_daily`, `current_note`) |
 | `images.rs` | Embedded images: link text → vault path → decode thread → upload; the lookup the reading view and the editor share |
 | `paste.rs` | Ctrl+V and drops of a picture or image files: what the clipboard / drop offers, storing into the attachment folder, where the embeds go |
-| `instance.rs` | Single instance on `/tmp/cce-notes-<WAYLAND_DISPLAY>.sock`; the CLI's commands |
+| `instance.rs` | The CLI's commands, on the single-instance socket `/tmp/cce-notes-<WAYLAND_DISPLAY>.sock` (claim, forward and listener are `cce_ui::ipc::instance`) |
 
 ## Behaviour worth knowing before changing it
 
diff --git a/src/instance.rs b/src/instance.rs
index 99b41c3..5b6336e 100644
--- a/src/instance.rs
+++ b/src/instance.rs
@@ -15,20 +15,15 @@
 //!   open) straight from the listener thread; cce-graph's local graph
 //!   polls it to follow the note on screen.
 //!
-//! Connect before binding, as cce-browser does: a refused connect means a
-//! crashed instance left its socket file, which is removed; losing the bind
-//! to a simultaneous launch falls back to one more connect.
+//! The claim, the startup race it closes, the parked listener and the
+//! bounded reads are `cce_ui::ipc::instance`'s; this module is the notes
+//! protocol on top.
 
-use std::io::{BufRead, BufReader, Write};
-use std::os::unix::net::{UnixListener, UnixStream};
 use std::sync::Mutex;
 
 use crate::Message;
 
 const PREFIX: &str = "cce-notes";
-
-static CLAIMED: Mutex<Option<UnixListener>> = Mutex::new(None);
-static OWNED_PATH: Mutex<Option<String>> = Mutex::new(None);
 /// The open note's vault path, for `current`. Set by the app whenever it
 /// changes; read by the listener without a trip through the event loop.
 static CURRENT: Mutex<Option<String>> = Mutex::new(None);
@@ -125,78 +120,34 @@ fn absolute(target: &str) -> String {
 /// Hand `cmd` to a running instance, or claim the socket. True when a
 /// running instance took it and this process should exit.
 pub fn forward_or_claim(cmd: &Command) -> bool {
-    let path = cce_ui::ipc::socket_path(PREFIX);
-    if try_forward(&path, cmd) {
-        return true;
-    }
-    if std::path::Path::new(&path).exists() {
-        let _ = std::fs::remove_file(&path);
-    }
-    match UnixListener::bind(&path) {
-        Ok(listener) => {
-            *CLAIMED.lock().unwrap() = Some(listener);
-            *OWNED_PATH.lock().unwrap() = Some(path);
-            false
-        }
-        Err(_) => try_forward(&path, cmd),
-    }
-}
-
-fn try_forward(path: &str, cmd: &Command) -> bool {
-    let Ok(mut stream) = UnixStream::connect(path) else {
-        return false;
-    };
-    if stream.write_all(format!("{}\n", cmd.to_line()).as_bytes()).is_err() {
-        return false;
-    }
-    // Bounded: an instance whose listener is stuck must not hang this
-    // launch forever; unanswered, the launch is not forwarded.
-    let _ = stream.set_read_timeout(Some(std::time::Duration::from_secs(5)));
-    let mut reply = String::new();
-    BufReader::new(stream).read_line(&mut reply).is_ok()
+    cce_ui::ipc::instance::forward_or_claim(PREFIX, &cmd.to_line())
 }
 
 /// Serve the claimed listener on a thread, feeding the app's loop. False
 /// when this process holds no listener (single-instance handling failed
 /// and it runs standalone).
 pub fn spawn_listener(sender: calloop::channel::Sender<Message>) -> bool {
-    let Some(listener) = CLAIMED.lock().unwrap().take() else {
-        return false;
-    };
-    std::thread::spawn(move || {
-        for conn in listener.incoming() {
-            let Ok(conn) = conn else { continue };
-            let Some(line) = read_request_line(&conn, 64 * 1024, std::time::Duration::from_secs(2)) else {
-                continue;
-            };
-            if line.trim() == "current" {
-                let cur = CURRENT.lock().unwrap_or_else(|e| e.into_inner()).clone();
-                let reply = match cur {
-                    Some(p) => format!("ok {p}\n"),
-                    None => "ok\n".to_string(),
-                };
-                let _ = (&conn).write_all(reply.as_bytes());
-                continue;
+    cce_ui::ipc::instance::serve(move |line| {
+        // Answered from here, without a trip through the event loop.
+        if line == "current" {
+            let cur = CURRENT.lock().unwrap_or_else(|e| e.into_inner()).clone();
+            return Some(match cur {
+                Some(p) => format!("ok {p}"),
+                None => "ok".to_string(),
+            });
+        }
+        match Command::parse(line) {
+            Some(cmd) => {
+                sender.send(Message::Command(cmd)).ok()?;
+                Some("ok".into())
             }
-            let reply: &[u8] = match Command::parse(line.trim()) {
-                Some(cmd) => {
-                    if sender.send(Message::Command(cmd)).is_err() {
-                        return;
-                    }
-                    b"ok\n"
-                }
-                None => b"error unknown command\n",
-            };
-            let _ = (&conn).write_all(reply);
+            None => Some("error unknown command".into()),
         }
-    });
-    true
+    })
 }
 
 pub fn cleanup() {
-    if let Some(path) = OWNED_PATH.lock().unwrap().take() {
-        let _ = std::fs::remove_file(path);
-    }
+    cce_ui::ipc::instance::cleanup();
 }
 
 #[cfg(test)]
@@ -229,80 +180,3 @@ mod tests {
         assert_eq!(Command::parse("open 2026"), Some(Command::Open { target: "2026".into(), line: None }));
     }
 }
-
-/// One request line from a control-socket client, bounded in size and in
-/// TOTAL time. Until 2026-10-02 this was `BufReader::read_line` on a socket
-/// with no timeout at all, so a client that connected and said nothing (or trickled a
-/// byte at a time) held the listener thread for good, and every later
-/// `cce-notes open` forwarded to it waited behind. None on EOF before any byte, timeout,
-/// overflow or a read error.
-fn read_request_line(conn: &std::os::unix::net::UnixStream, limit: usize, deadline: std::time::Duration) -> Option<String> {
-    use std::io::Read;
-    let until = std::time::Instant::now() + deadline;
-    let mut buf: Vec<u8> = Vec::new();
-    let mut chunk = [0u8; 4096];
-    let mut reader = conn;
-    loop {
-        let left = until.checked_duration_since(std::time::Instant::now()).filter(|d| !d.is_zero())?;
-        conn.set_read_timeout(Some(left)).ok()?;
-        let n = reader.read(&mut chunk).ok()?;
-        if n == 0 {
-            if buf.is_empty() {
-                return None;
-            }
-            break;
-        }
-        buf.extend_from_slice(&chunk[..n]);
-        if let Some(end) = buf.iter().position(|&b| b == b'\n') {
-            buf.truncate(end + 1);
-            break;
-        }
-        if buf.len() > limit {
-            return None;
-        }
-    }
-    let _ = conn.set_read_timeout(None);
-    String::from_utf8(buf).ok()
-}
-
-#[cfg(test)]
-mod request_line_tests {
-    use super::read_request_line;
-    use std::io::Write;
-    use std::os::unix::net::UnixStream;
-    use std::time::{Duration, Instant};
-
-    #[test]
-    fn a_request_line_is_bounded_in_time_and_size() {
-        let (mut client, server) = UnixStream::pair().unwrap();
-        client.write_all(b"open note\nmore").unwrap();
-        assert_eq!(read_request_line(&server, 1024, Duration::from_secs(1)).as_deref(), Some("open note\n"));
-
-        // Silent: given up at the deadline, not held forever.
-        let (_quiet, server) = UnixStream::pair().unwrap();
-        let t = Instant::now();
-        assert_eq!(read_request_line(&server, 1024, Duration::from_millis(100)), None);
-        assert!(t.elapsed() < Duration::from_millis(500));
-
-        // Trickling a byte at a time: the TOTAL deadline still ends it.
-        let (mut client, server) = UnixStream::pair().unwrap();
-        let trickle = std::thread::spawn(move || {
-            for _ in 0..40 {
-                if client.write_all(b"x").is_err() {
-                    break;
-                }
-                std::thread::sleep(Duration::from_millis(25));
-            }
-        });
-        let t = Instant::now();
-        assert_eq!(read_request_line(&server, 1024, Duration::from_millis(150)), None);
-        assert!(t.elapsed() < Duration::from_millis(500), "took {:?}", t.elapsed());
-        drop(server);
-        trickle.join().unwrap();
-
-        // Over the size cap.
-        let (mut client, server) = UnixStream::pair().unwrap();
-        client.write_all(&[b'z'; 2000]).unwrap();
-        assert_eq!(read_request_line(&server, 1024, Duration::from_millis(200)), None);
-    }
-}