diff --git a/Cargo.lock b/Cargo.lock index 9673431..da014f7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -184,7 +184,7 @@ dependencies = [ [[package]] name = "opencode-pty" -version = "0.1.13" +version = "0.2.0" dependencies = [ "anyhow", "base64", diff --git a/Cargo.toml b/Cargo.toml index 638cf20..bc6e1e1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "opencode-pty" -version = "0.1.13" +version = "0.2.0" edition = "2024" rust-version = "1.90" description = "Persistent PTY service for OpenCode" diff --git a/README.md b/README.md index 46f1730..c53b865 100644 --- a/README.md +++ b/README.md @@ -12,7 +12,7 @@ that connection until exit; exiting the playground stops the daemon and all its terminals. Observer commands connect to an existing daemon without taking ownership. The service uses a private authenticated Unix socket and atomic registration file. -Integrations launch `opencode-pty daemon` (protocol 7). +Integrations launch `opencode-pty daemon --name NAME [--runtime-dir DIR]` (protocol 7). The server must claim the daemon within 5 seconds by sending the authenticated framed envelope `{"token":"...","request":{"op":"own","instance_id":"..."}}`. The response is `{"type":"owned"}`; that connection stays open as the sole @@ -33,6 +33,17 @@ cancels the handoff. An ordinary authenticated `shutdown` request, or one from t current owner, stops the daemon even during handoff. No ownership or handoff state is persisted. +Every command requires `--name NAME`, which selects the runtime directory +`DIR/NAME`. `--runtime-dir DIR` is optional and defaults to OpenCode's state +directory, `${XDG_STATE_HOME:-~/.local/state}/opencode/pty`. Names are a single +path component of letters, digits, `.`, `_`, or `-`. The registration +(`service.json`) and lock (`service.lock`) live in that directory. They only +default to a temporary directory when no home directory exists, because macOS deletes unaccessed regular files +there after three days. The socket stays under `/tmp/opencode-pty-/` to fit +socket path limits; temporary cleaners skip sockets. A stopping daemon removes +its own files and runtime directory, and a starting daemon removes abandoned +sibling runtime directories older than ten minutes. + ## Architecture ```text @@ -67,15 +78,16 @@ user input without blocking inside the callback. ## Playground ```sh -cargo run -- play +cargo run -- play --name play ``` Other service commands: ```sh -cargo run -- status -cargo run -- list -cargo run -- stop # destructive: terminates every terminal +cargo run -- status --name play +cargo run -- list --name play +cargo run -- watch 1 --name play +cargo run -- stop --name play # destructive: terminates every terminal ``` `play` starts and owns a new daemon; it cannot adopt an already running daemon. @@ -146,7 +158,7 @@ libclang and do not depend on a third-party Ghostty Rust crate. cargo fmt --check cargo clippy --all-targets --all-features -- -D warnings cargo test -printf 'demo\nlist\nquit\n' | cargo run -- play +printf 'demo\nlist\nquit\n' | cargo run -- play --name play ``` ### Direct Ghostty bindings diff --git a/SPEC.md b/SPEC.md index 84ec5c9..41f634b 100644 --- a/SPEC.md +++ b/SPEC.md @@ -100,8 +100,9 @@ terminals may die; live replacement of the daemon is not supported. Every daemon requires one authenticated owner connection within five seconds of startup. Owner loss stops the daemon unless a handoff was prepared on that connection. Handoff tickets expire 120 seconds after preparation and are -consumed by successful replacement ownership. A connected owner cannot be -displaced. The playground holds ownership until it exits; other CLI commands +consumed by successful replacement ownership. A valid ticket replaces even a +connected owner; the superseded connection can no longer prepare handoffs or +stop the daemon, and its disconnect does not affect the new owner. The playground holds ownership until it exits; other CLI commands only observe or operate an existing daemon and never start one. Protocol v7 uses four-byte big-endian framing with bounded UTF-8 JSON control @@ -116,6 +117,19 @@ service lock elects one process and protects stale socket cleanup. On Unix, the socket uses a fixed-length hash of the canonical runtime path under a private per-user `/tmp` directory to stay below platform path limits. +Every command requires `--name NAME`; the runtime directory is `DIR/NAME`, +where `--runtime-dir DIR` defaults to OpenCode's state directory, +`${XDG_STATE_HOME:-~/.local/state}/opencode/pty`. A name is one path component +of letters, digits, `.`, `_`, or `-`. Registration defaults to a temporary directory only when no home +directory can be found: macOS deletes regular files there that are unaccessed for three +days, even while the daemon runs. Temporary cleaners skip sockets, so the +socket stays in `/tmp`. On exit the daemon removes its registration and socket +only if they are still its own, then its empty runtime directory. At startup +it removes sibling runtime directories in `DIR` that are over ten +minutes old, contain only registration files, and have either a lock it can +acquire or, without a lock file, a valid registration whose PID no longer +exists. Directories without that evidence, including empty ones, are kept. + OpenCode chooses a fresh UUID runtime directory for each server, independent of the database. It starts the daemon only when the first terminal is created. Only an explicit restart handoff descriptor lets a replacement server reuse diff --git a/npm/package.json b/npm/package.json index b8a1d78..d7bd29e 100644 --- a/npm/package.json +++ b/npm/package.json @@ -1,6 +1,6 @@ { "name": "@opencode-ai/pty", - "version": "0.1.13", + "version": "0.2.0", "description": "Persistent PTY service for OpenCode", "type": "module", "license": "MIT", @@ -24,12 +24,12 @@ "index.js" ], "optionalDependencies": { - "@opencode-ai/pty-darwin-arm64": "0.1.13", - "@opencode-ai/pty-darwin-x64": "0.1.13", - "@opencode-ai/pty-linux-arm64-gnu": "0.1.13", - "@opencode-ai/pty-linux-arm64-musl": "0.1.13", - "@opencode-ai/pty-linux-x64-gnu": "0.1.13", - "@opencode-ai/pty-linux-x64-musl": "0.1.13" + "@opencode-ai/pty-darwin-arm64": "0.2.0", + "@opencode-ai/pty-darwin-x64": "0.2.0", + "@opencode-ai/pty-linux-arm64-gnu": "0.2.0", + "@opencode-ai/pty-linux-arm64-musl": "0.2.0", + "@opencode-ai/pty-linux-x64-gnu": "0.2.0", + "@opencode-ai/pty-linux-x64-musl": "0.2.0" }, "publishConfig": { "access": "public" diff --git a/src/client.rs b/src/client.rs index 1c5515d..c6a557a 100644 --- a/src/client.rs +++ b/src/client.rs @@ -1,12 +1,13 @@ use std::env; use std::io::{self, BufRead, Read, Write}; +use std::path::{Path, PathBuf}; use std::thread; use std::time::{Duration, Instant}; use anyhow::{Context, Result, anyhow, bail}; use base64::Engine; -use crate::daemon::{Registration, read_registration}; +use crate::daemon::{Registration, read_registration, runtime_dir}; use crate::protocol::{ AttachmentRole, Envelope, PROTOCOL_VERSION, Request, Response, SubscriptionEvent, read_frame, read_subscription_event, write_frame, @@ -57,14 +58,24 @@ impl TerminalSubscription { impl TerminalClient { #[cfg(unix)] - pub fn start() -> Result { + pub fn start(directory: &Path) -> Result { use std::os::unix::net::UnixStream; use std::os::unix::process::CommandExt; use std::process::{Command, Stdio}; + let invalid = || anyhow!("invalid runtime directory {}", directory.display()); + let name = directory + .file_name() + .and_then(|name| name.to_str()) + .ok_or_else(invalid)?; + let root = directory.parent().ok_or_else(invalid)?; let mut command = Command::new(env::current_exe()?); command .arg("daemon") + .arg("--name") + .arg(name) + .arg("--runtime-dir") + .arg(root) .stdin(Stdio::null()) .stdout(Stdio::null()) .stderr(Stdio::null()); @@ -79,7 +90,7 @@ impl TerminalClient { if let Some(status) = child.try_wait()? { bail!("opencode-pty daemon exited before ownership acquisition: {status}"); } - if let Ok(registration) = read_registration() + if let Ok(registration) = read_registration(directory) && registration.pid == child.id() { let mut stream = UnixStream::connect(®istration.socket)?; @@ -109,12 +120,12 @@ impl TerminalClient { } #[cfg(not(unix))] - pub fn start() -> Result { + pub fn start(_directory: &Path) -> Result { bail!("opencode-pty client transport is not implemented on this platform") } - pub fn discover() -> Result { - let registration = read_registration()?; + pub fn discover(directory: &Path) -> Result { + let registration = read_registration(directory)?; if registration.protocol != PROTOCOL_VERSION { bail!( "opencode-pty protocol mismatch: service={}, client={PROTOCOL_VERSION}", @@ -404,25 +415,52 @@ impl Drop for TerminalClient { pub fn run_cli() -> Result<()> { let mut args = env::args().skip(1); - match args.next().as_deref() { - Some("daemon") => { - if args.next().is_some() { - bail!("usage: opencode-pty daemon"); + let command = args.next(); + let mut name = None; + let mut root = None; + let mut positional = Vec::new(); + while let Some(arg) = args.next() { + match arg.as_str() { + "--name" => name = Some(args.next().ok_or_else(|| anyhow!("--name needs a value"))?), + "--runtime-dir" => { + root = Some(PathBuf::from( + args.next() + .ok_or_else(|| anyhow!("--runtime-dir needs a value"))?, + )) } - crate::daemon::run() + _ => positional.push(arg), + } + } + let directory = || -> Result { + let name = name + .as_deref() + .ok_or_else(|| anyhow!("--name is required; use `opencode-pty help`"))?; + runtime_dir(root.as_deref(), name) + }; + let no_arguments = |usage: &str| { + if positional.is_empty() { + Ok(()) + } else { + Err(anyhow!("usage: opencode-pty {usage}")) + } + }; + match command.as_deref() { + Some("daemon") => { + no_arguments("daemon --name NAME [--runtime-dir DIR]")?; + crate::daemon::run(&directory()?) } Some("fixture") => run_fixture(), - None | Some("play") => play(), - Some("status") => status(), - Some("stop") => stop(), - Some("list") => print_terminals(&TerminalClient::discover()?), + Some("play") => play(&directory()?), + Some("status") => status(&directory()?), + Some("stop") => stop(&directory()?), + Some("list") => print_terminals(&TerminalClient::discover(&directory()?)?), Some("watch") => { - let id = args - .next() - .ok_or_else(|| anyhow!("usage: opencode-pty watch TERMINAL_ID"))? + let id = positional + .first() + .ok_or_else(|| anyhow!("usage: opencode-pty watch TERMINAL_ID --name NAME"))? .parse() .context("expected terminal ID")?; - watch(id) + watch(&directory()?, id) } Some("version" | "--version" | "-V") => { println!( @@ -432,7 +470,7 @@ pub fn run_cli() -> Result<()> { ); Ok(()) } - Some("help" | "--help" | "-h") => { + None | Some("help" | "--help" | "-h") => { print_usage(); Ok(()) } @@ -440,8 +478,8 @@ pub fn run_cli() -> Result<()> { } } -fn status() -> Result<()> { - let client = TerminalClient::discover()?; +fn status(directory: &Path) -> Result<()> { + let client = TerminalClient::discover(directory)?; println!( "opencode-pty running: pid={} instance={} terminals={}", client.registration.pid, @@ -451,13 +489,13 @@ fn status() -> Result<()> { Ok(()) } -fn stop() -> Result<()> { - let client = TerminalClient::discover()?; +fn stop(directory: &Path) -> Result<()> { + let client = TerminalClient::discover(directory)?; let pid = client.registration.pid; client.shutdown()?; let deadline = Instant::now() + START_TIMEOUT; while Instant::now() < deadline { - if TerminalClient::discover().is_err() { + if TerminalClient::discover(directory).is_err() { println!("stopped opencode-pty pid={pid} (all terminals exited)"); return Ok(()); } @@ -467,8 +505,8 @@ fn stop() -> Result<()> { } #[cfg(unix)] -fn watch(id: TerminalId) -> Result<()> { - let client = TerminalClient::discover()?; +fn watch(directory: &Path, id: TerminalId) -> Result<()> { + let client = TerminalClient::discover(directory)?; let attachment_id = format!( "watch-{}-{:016x}", std::process::id(), @@ -503,12 +541,12 @@ fn watch(id: TerminalId) -> Result<()> { } #[cfg(not(unix))] -fn watch(_id: TerminalId) -> Result<()> { +fn watch(_directory: &Path, _id: TerminalId) -> Result<()> { bail!("streaming transport is not implemented on this platform") } -fn play() -> Result<()> { - let client = TerminalClient::start()?; +fn play(directory: &Path) -> Result<()> { + let client = TerminalClient::start(directory)?; let mut active = client.list()?.first().map(|terminal| terminal.id); println!("\n opencode-pty playground"); println!( @@ -682,7 +720,13 @@ fn set_stdin_raw() -> Result<()> { } fn print_usage() { - println!("usage: opencode-pty [play|status|list|watch ID|stop|daemon|--version]"); + println!( + "usage: opencode-pty [play|status|list|watch ID|stop|daemon] --name NAME [--runtime-dir DIR]" + ); + println!(" opencode-pty --version"); + println!( + " NAME selects DIR/NAME; DIR defaults to ${{XDG_STATE_HOME:-~/.local/state}}/opencode/pty" + ); } fn print_help() { diff --git a/src/daemon.rs b/src/daemon.rs index 2a45e15..4062ed2 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -1,5 +1,6 @@ -use std::path::PathBuf; +use std::path::{Path, PathBuf}; +use anyhow::{Context, Result, bail}; use serde::{Deserialize, Serialize}; pub const REGISTRATION_FILE: &str = "service.json"; @@ -14,14 +15,60 @@ pub struct Registration { pub token: String, } +/// Shared parent of runtime directories, inside OpenCode's state directory. +/// Registration files must not live in temporary directories, which macOS purges +/// while daemons run. +pub fn default_runtime_root() -> PathBuf { + resolve_runtime_root( + std::env::var_os("XDG_STATE_HOME").map(PathBuf::from), + std::env::home_dir(), + ) +} + +fn resolve_runtime_root(state: Option, home: Option) -> PathBuf { + state + .filter(|path| path.is_absolute()) + .or_else(|| home.map(|home| home.join(".local").join("state"))) + .unwrap_or_else(std::env::temp_dir) + .join("opencode") + .join("pty") +} + +/// Resolves `/`, where the name is a single plain path component. +pub fn runtime_dir(root: Option<&Path>, name: &str) -> Result { + let valid = !name.is_empty() + && name != "." + && name != ".." + && name + .bytes() + .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-')); + if !valid { + bail!("invalid runtime name {name:?}; use letters, digits, '.', '_', or '-'"); + } + Ok(root + .map(Path::to_path_buf) + .unwrap_or_else(default_runtime_root) + .join(name)) +} + +pub fn registration_path(directory: &Path) -> PathBuf { + directory.join(REGISTRATION_FILE) +} + +pub fn read_registration(directory: &Path) -> Result { + let data = std::fs::read(registration_path(directory)) + .context("opencode-pty registration is unavailable")?; + serde_json::from_slice(&data).context("invalid opencode-pty registration") +} + #[cfg(unix)] mod unix { use std::fs::{self, OpenOptions}; use std::net::Shutdown; use std::os::unix::ffi::OsStrExt; - use std::os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt}; + use std::os::unix::fs::{FileTypeExt, MetadataExt, OpenOptionsExt, PermissionsExt}; use std::os::unix::net::{UnixListener, UnixStream}; - use std::path::PathBuf; + use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::thread; @@ -32,36 +79,17 @@ mod unix { use fs2::FileExt; use sha2::{Digest, Sha256}; - use super::{LOCK_FILE, REGISTRATION_FILE, Registration}; + use super::{LOCK_FILE, REGISTRATION_FILE, Registration, read_registration, registration_path}; use crate::ownership::Ownership; use crate::protocol::{ Envelope, PROTOCOL_VERSION, Request, Response, read_frame, write_frame, write_output_frame, }; use crate::service::{CreateTerminal, StreamEvent, TerminalService}; - pub fn service_dir() -> PathBuf { - if let Some(path) = std::env::var_os("OPENCODE_PTY_RUNTIME_DIR") { - return PathBuf::from(path); - } - if let Some(path) = std::env::var_os("XDG_RUNTIME_DIR") { - return PathBuf::from(path).join("opencode-pty"); - } - let uid = nix::unistd::Uid::effective().as_raw(); - std::env::temp_dir().join(format!("opencode-pty-{uid}")) - } - - pub fn registration_path() -> PathBuf { - service_dir().join(REGISTRATION_FILE) - } - - pub fn read_registration() -> Result { - let data = - fs::read(registration_path()).context("opencode-pty registration is unavailable")?; - serde_json::from_slice(&data).context("invalid opencode-pty registration") - } + const STALE_RUNTIME_AGE: Duration = Duration::from_secs(10 * 60); - pub fn run() -> Result<()> { - let directory = service_dir(); + pub fn run(directory: &Path) -> Result<()> { + let directory = directory.to_path_buf(); fs::create_dir_all(&directory)?; fs::set_permissions(&directory, fs::Permissions::from_mode(0o700))?; let lock_path = directory.join(LOCK_FILE); @@ -81,16 +109,21 @@ mod unix { let listener = UnixListener::bind(&socket_path)?; fs::set_permissions(&socket_path, fs::Permissions::from_mode(0o600))?; listener.set_nonblocking(true)?; + let socket = SocketFile::bound(socket_path)?; let registration = Registration { instance_id: random_id(), pid: std::process::id(), protocol: PROTOCOL_VERSION, - socket: socket_path.clone(), + socket: socket.path.clone(), token: random_id(), }; let ownership = Arc::new(Mutex::new(Ownership::new(Instant::now()))); write_registration(&directory, ®istration)?; + let sweeper = { + let directory = directory.clone(); + thread::spawn(move || sweep_stale_runtimes(&directory)) + }; let service = Arc::new(TerminalService::default()); let shutdown = Arc::new(AtomicBool::new(false)); @@ -144,11 +177,14 @@ mod unix { } let (cleanup_tx, cleanup_rx) = crossbeam_channel::bounded::<()>(1); let cleanup_registration = registration.clone(); + let cleanup_socket = socket.clone(); + let cleanup_directory = directory.clone(); let watchdog = thread::spawn(move || { if cleanup_rx.recv_timeout(Duration::from_secs(5)).is_err() { eprintln!("opencode-pty cleanup timed out; forcing exit"); - let _ = remove_if_current(&cleanup_registration); - let _ = fs::remove_file(&cleanup_registration.socket); + let _ = remove_if_current(&cleanup_directory, &cleanup_registration); + cleanup_socket.remove_if_current(); + remove_runtime_directory(&cleanup_directory); std::process::exit(1); } }); @@ -156,29 +192,177 @@ mod unix { for (_, handler) in handlers { let _ = handler.join(); } + let _ = sweeper.join(); drop(service); - remove_if_current(®istration)?; - let _ = fs::remove_file(&socket_path); + let result = remove_if_current(&directory, ®istration); + socket.remove_if_current(); + // Remove the lock file before releasing it: a daemon that already opened it + // fails to lock it, and a later one creates a new lock file. + remove_runtime_directory(&directory); let _ = cleanup_tx.send(()); let _ = watchdog.join(); drop(lock); + result + } + + /// The bound socket's inode, so cleanup never unlinks a successor's socket. + #[derive(Clone)] + struct SocketFile { + path: PathBuf, + device: u64, + inode: u64, + } + + impl SocketFile { + fn bound(path: PathBuf) -> Result { + let metadata = fs::symlink_metadata(&path)?; + Ok(Self { + path, + device: metadata.dev(), + inode: metadata.ino(), + }) + } + + fn remove_if_current(&self) { + if fs::symlink_metadata(&self.path) + .is_ok_and(|metadata| metadata.dev() == self.device && metadata.ino() == self.inode) + { + let _ = fs::remove_file(&self.path); + } + } + } + + fn remove_runtime_directory(directory: &Path) { + let _ = fs::remove_file(directory.join(LOCK_FILE)); + let _ = fs::remove_dir(directory); + } + + /// Removes sibling runtime directories left by daemons that crashed. + fn sweep_stale_runtimes(directory: &Path) { + let Some(parent) = directory.parent() else { + return; + }; + let Ok(entries) = fs::read_dir(parent) else { + return; + }; + for entry in entries.flatten() { + let candidate = entry.path(); + if same_path(&candidate, directory) { + continue; + } + if let Err(error) = sweep_runtime(&candidate, SystemTime::now()) { + eprintln!( + "opencode-pty could not remove stale runtime {}: {error:#}", + candidate.display() + ); + } + } + } + + fn same_path(left: &Path, right: &Path) -> bool { + matches!( + (fs::canonicalize(left), fs::canonicalize(right)), + (Ok(left), Ok(right)) if left == right + ) + } + + fn sweep_runtime(directory: &Path, now: SystemTime) -> Result<()> { + let metadata = fs::symlink_metadata(directory)?; + if !metadata.is_dir() { + return Ok(()); + } + let age = now.duration_since(metadata.modified()?).unwrap_or_default(); + if age < STALE_RUNTIME_AGE { + return Ok(()); + } + let mut files = Vec::new(); + for entry in fs::read_dir(directory)? { + let entry = entry?; + let name = entry.file_name(); + let name = name.to_string_lossy(); + let known = name == LOCK_FILE + || name == REGISTRATION_FILE + || (name.starts_with("service.") && name.ends_with(".tmp")); + if !known || !entry.file_type()?.is_file() { + // Not a runtime directory this daemon created. + return Ok(()); + } + files.push(entry.path()); + } + + // Only a lock file or a parseable registration proves a daemon created this + // directory; anything else, including an empty directory, is left alone. + let _lock = match OpenOptions::new() + .read(true) + .write(true) + .open(directory.join(LOCK_FILE)) + { + Ok(lock) => { + if lock.try_lock_exclusive().is_err() { + return Ok(()); + } + Some(lock) + } + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + let Some(registration) = fs::read(registration_path(directory)) + .ok() + .and_then(|data| serde_json::from_slice::(&data).ok()) + else { + return Ok(()); + }; + if process_exists(registration.pid) { + return Ok(()); + } + None + } + Err(error) => return Err(error.into()), + }; + + let socket = socket_file(&socket_root(), &fs::canonicalize(directory)?); + if fs::symlink_metadata(&socket).is_ok_and(|metadata| metadata.file_type().is_socket()) { + fs::remove_file(&socket)?; + } + // Remove the lock last so a racing daemon cannot claim a half-removed directory. + files.sort_by_key(|file| file.ends_with(LOCK_FILE)); + for file in files { + fs::remove_file(file)?; + } + fs::remove_dir(directory)?; Ok(()) } - fn socket_path(directory: &std::path::Path) -> Result { + fn process_exists(pid: u32) -> bool { + let Ok(pid) = i32::try_from(pid) else { + return false; + }; + // SAFETY: signal zero only checks whether the PID exists. + let exists = unsafe { libc::kill(pid, 0) } == 0; + exists || std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM) + } + + fn socket_path(directory: &Path) -> Result { let directory = fs::canonicalize(directory).context("failed to resolve PTY runtime directory")?; - let digest = Sha256::digest(directory.as_os_str().as_bytes()); + let root = socket_root(); + ensure_private_directory(&root)?; + Ok(socket_file(&root, &directory)) + } + + // Sockets stay in /tmp so their paths fit sun_path; temp cleaners skip sockets. + fn socket_root() -> PathBuf { + PathBuf::from("/tmp").join(format!( + "opencode-pty-{}", + nix::unistd::Uid::effective().as_raw() + )) + } + + fn socket_file(root: &Path, canonical_directory: &Path) -> PathBuf { + let digest = Sha256::digest(canonical_directory.as_os_str().as_bytes()); let name = digest[..16] .iter() .map(|byte| format!("{byte:02x}")) .collect::(); - let root = PathBuf::from("/tmp").join(format!( - "opencode-pty-{}", - nix::unistd::Uid::effective().as_raw() - )); - ensure_private_directory(&root)?; - Ok(root.join(format!("{name}.sock"))) + root.join(format!("{name}.sock")) } fn ensure_private_directory(directory: &std::path::Path) -> Result<()> { @@ -588,14 +772,15 @@ mod unix { use std::io::Write; file.write_all(&data)?; file.sync_all()?; - fs::rename(&temporary, registration_path())?; + fs::rename(&temporary, registration_path(directory))?; Ok(()) } - fn remove_if_current(registration: &Registration) -> Result<()> { - if read_registration().is_ok_and(|current| current.instance_id == registration.instance_id) + fn remove_if_current(directory: &Path, registration: &Registration) -> Result<()> { + if read_registration(directory) + .is_ok_and(|current| current.instance_id == registration.instance_id) { - fs::remove_file(registration_path())?; + fs::remove_file(registration_path(directory))?; } Ok(()) } @@ -626,23 +811,158 @@ mod unix { fs::remove_dir_all(base).unwrap(); } + + #[test] + fn runtime_root_avoids_temporary_directories() { + let home = Some(PathBuf::from("/home/user")); + assert_eq!( + super::super::resolve_runtime_root(Some("/state".into()), home.clone()), + PathBuf::from("/state/opencode/pty") + ); + assert_eq!( + super::super::resolve_runtime_root(Some("relative".into()), home.clone()), + PathBuf::from("/home/user/.local/state/opencode/pty") + ); + assert_eq!( + super::super::resolve_runtime_root(None, home), + PathBuf::from("/home/user/.local/state/opencode/pty") + ); + } + + #[test] + fn runtime_names_are_single_components() { + let root = Path::new("/state"); + assert_eq!( + super::super::runtime_dir(Some(root), "ee7511b8-7db0.x_1").unwrap(), + root.join("ee7511b8-7db0.x_1") + ); + for name in ["", ".", "..", "a/b", "../x", "with space"] { + assert!( + super::super::runtime_dir(Some(root), name).is_err(), + "{name:?}" + ); + } + } + + struct Root(PathBuf); + + impl Root { + fn new() -> Self { + let root = + std::env::temp_dir().join(format!("opencode-pty-sweep-test-{}", random_id())); + fs::create_dir_all(&root).unwrap(); + Self(root) + } + + fn runtime(&self, files: &[&str], pid: Option) -> PathBuf { + let directory = self.0.join(random_id()); + fs::create_dir(&directory).unwrap(); + for file in files { + fs::write(directory.join(file), b"").unwrap(); + } + if let Some(pid) = pid { + let registration = Registration { + instance_id: random_id(), + pid, + protocol: PROTOCOL_VERSION, + socket: PathBuf::from("/nonexistent.sock"), + token: random_id(), + }; + fs::write( + directory.join(REGISTRATION_FILE), + serde_json::to_vec(®istration).unwrap(), + ) + .unwrap(); + } + directory + } + } + + impl Drop for Root { + fn drop(&mut self) { + let _ = fs::remove_dir_all(&self.0); + } + } + + fn later() -> SystemTime { + SystemTime::now() + STALE_RUNTIME_AGE + Duration::from_secs(1) + } + + fn dead_pid() -> u32 { + let mut child = std::process::Command::new("/usr/bin/true").spawn().unwrap(); + let pid = child.id(); + child.wait().unwrap(); + pid + } + + #[test] + fn sweep_removes_only_abandoned_old_runtimes() { + let root = Root::new(); + let abandoned = root.runtime(&[LOCK_FILE, REGISTRATION_FILE], None); + let young = root.runtime(&[LOCK_FILE], None); + let foreign = root.runtime(&[LOCK_FILE, "notes.txt"], None); + let locked = root.runtime(&[LOCK_FILE], None); + let held = fs::File::open(locked.join(LOCK_FILE)).unwrap(); + held.try_lock_exclusive().unwrap(); + + sweep_runtime(&abandoned, later()).unwrap(); + sweep_runtime(&young, SystemTime::now()).unwrap(); + sweep_runtime(&foreign, later()).unwrap(); + sweep_runtime(&locked, later()).unwrap(); + + assert!(!abandoned.exists()); + assert!(young.join(LOCK_FILE).exists()); + assert!(foreign.join("notes.txt").exists()); + assert!(locked.join(LOCK_FILE).exists()); + drop(held); + } + + #[test] + fn sweep_without_lock_file_checks_registered_pid() { + let root = Root::new(); + let live = root.runtime(&[], Some(std::process::id())); + let dead = root.runtime(&[], Some(dead_pid())); + + sweep_runtime(&live, later()).unwrap(); + sweep_runtime(&dead, later()).unwrap(); + + assert!(live.join(REGISTRATION_FILE).exists()); + assert!(!dead.exists()); + } + + #[test] + fn sweep_keeps_directories_without_daemon_evidence() { + let root = Root::new(); + let empty = root.runtime(&[], None); + let foreign = root.runtime(&[], None); + fs::write(foreign.join(REGISTRATION_FILE), br#"{"not":"ours"}"#).unwrap(); + let temporary = root.runtime(&["service.important.tmp"], None); + + for directory in [&empty, &foreign, &temporary] { + sweep_runtime(directory, later()).unwrap(); + assert!(directory.exists(), "{}", directory.display()); + } + } + + #[test] + fn sweep_removes_abandoned_socket() { + let root = Root::new(); + let abandoned = root.runtime(&[LOCK_FILE], None); + let socket = socket_path(&abandoned).unwrap(); + drop(UnixListener::bind(&socket).unwrap()); + + sweep_runtime(&abandoned, later()).unwrap(); + + assert!(!abandoned.exists()); + assert!(!socket.exists()); + } } } #[cfg(unix)] -pub use unix::{read_registration, registration_path, run, service_dir}; +pub use unix::run; #[cfg(not(unix))] -pub fn run() -> anyhow::Result<()> { +pub fn run(_directory: &Path) -> Result<()> { anyhow::bail!("persistent opencode-pty transport is not implemented on this platform") } - -#[cfg(not(unix))] -pub fn read_registration() -> anyhow::Result { - anyhow::bail!("persistent opencode-pty transport is not implemented on this platform") -} - -#[cfg(not(unix))] -pub fn registration_path() -> PathBuf { - PathBuf::from("opencode-pty-service.json") -} diff --git a/tests/ownership.rs b/tests/ownership.rs index 34e1204..d00b58f 100644 --- a/tests/ownership.rs +++ b/tests/ownership.rs @@ -1,7 +1,7 @@ #![cfg(unix)] use std::io::{Read, Write}; -use std::os::unix::net::UnixStream; +use std::os::unix::net::{UnixListener, UnixStream}; use std::path::PathBuf; use std::process::{Child, Command, Stdio}; use std::sync::Barrier; @@ -17,19 +17,21 @@ use opencode_pty::service::TerminalInfo; struct Daemon { child: Child, + root: PathBuf, directory: PathBuf, registration: Registration, } impl Daemon { fn start() -> Self { - let directory = std::env::temp_dir().join(format!( + let root = std::env::temp_dir().join(format!( "opencode-pty-ownership-{:032x}", rand::random::() )); + let directory = root.join("test"); let mut child = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) - .arg("daemon") - .env("OPENCODE_PTY_RUNTIME_DIR", &directory) + .args(["daemon", "--name", "test", "--runtime-dir"]) + .arg(&root) .stdin(Stdio::null()) .stdout(Stdio::null()) .spawn() @@ -53,6 +55,7 @@ impl Daemon { assert_eq!(registration.protocol, 7); Self { child, + root, directory, registration, } @@ -114,12 +117,19 @@ impl Daemon { } fn wait(&mut self) { + self.exit(); + assert!(!self.registration.socket.exists()); + } + + fn exit(&mut self) { let deadline = Instant::now() + Duration::from_secs(7); loop { if let Some(status) = self.child.try_wait().unwrap() { assert!(status.success(), "daemon exit: {status}"); - assert!(!self.directory.join("service.json").exists()); - assert!(!self.registration.socket.exists()); + assert!( + !self.directory.exists(), + "runtime directory was not removed" + ); return; } assert!(Instant::now() < deadline, "daemon did not stop"); @@ -148,7 +158,7 @@ impl Drop for Daemon { let _ = self.child.kill(); let _ = self.child.wait(); } - let _ = std::fs::remove_dir_all(&self.directory); + let _ = std::fs::remove_dir_all(&self.root); } } @@ -348,3 +358,21 @@ fn unclaimed_daemon_times_out() { let _partial = daemon.connect(); daemon.wait(); } + +#[test] +fn shutdown_keeps_a_replaced_socket() { + let mut daemon = Daemon::start(); + let (mut owner, response) = daemon.own(None); + assert!(matches!(response, Response::Owned)); + // A successor bound the same path; the exiting daemon must not unlink it. + std::fs::remove_file(&daemon.registration.socket).unwrap(); + let replacement = UnixListener::bind(&daemon.registration.socket).unwrap(); + assert!(matches!( + daemon.send(&mut owner, Request::Shutdown), + Response::Ok + )); + daemon.exit(); + assert!(daemon.registration.socket.exists()); + drop(replacement); + std::fs::remove_file(&daemon.registration.socket).unwrap(); +} diff --git a/tests/playground.rs b/tests/playground.rs index f407570..bcc404d 100644 --- a/tests/playground.rs +++ b/tests/playground.rs @@ -1,12 +1,13 @@ #![cfg(unix)] +use std::ffi::OsStr; use std::io::{BufRead, BufReader, Write}; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; use std::time::{Duration, Instant}; use std::time::{SystemTime, UNIX_EPOCH}; -fn runtime_dir(name: &str) -> PathBuf { +fn runtime_root(name: &str) -> PathBuf { std::env::temp_dir().join(format!( "opencode-pty-test-{name}-{}-{}", std::process::id(), @@ -17,6 +18,15 @@ fn runtime_dir(name: &str) -> PathBuf { )) } +fn target(root: &Path) -> [&OsStr; 4] { + [ + OsStr::new("--name"), + OsStr::new("play"), + OsStr::new("--runtime-dir"), + root.as_os_str(), + ] +} + fn output_with_timeout(mut child: std::process::Child) -> std::process::Output { let deadline = Instant::now() + Duration::from_secs(10); while Instant::now() < deadline { @@ -58,10 +68,11 @@ fn read_created_terminal(child: &mut std::process::Child) -> String { #[test] fn playground_proves_authoritative_query_response() { - let runtime = runtime_dir("query"); + let root = runtime_root("query"); + let runtime = root.join("play"); let mut child = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .arg("play") - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .stdin(Stdio::piped()) .stdout(Stdio::piped()) .spawn() @@ -74,7 +85,8 @@ fn playground_proves_authoritative_query_response() { .expect("commands written"); let output = output_with_timeout(child); assert!(!runtime.join("service.json").exists()); - std::fs::remove_dir_all(&runtime).unwrap(); + assert!(!runtime.exists(), "runtime directory was not removed"); + std::fs::remove_dir_all(&root).unwrap(); assert!( output.status.success(), "{}", @@ -86,10 +98,11 @@ fn playground_proves_authoritative_query_response() { #[test] fn playground_exit_stops_all_terminals_and_observers_do_not_start_daemons() { - let runtime = runtime_dir("ownership"); + let root = runtime_root("ownership"); + let runtime = root.join("play"); let mut owner = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .arg("play") - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .stdin(Stdio::piped()) .stdout(Stdio::piped()) .spawn() @@ -115,7 +128,7 @@ fn playground_exit_stops_all_terminals_and_observers_do_not_start_daemons() { let second = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .arg("play") - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .stdin(Stdio::null()) .stdout(Stdio::piped()) .stderr(Stdio::piped()) @@ -130,7 +143,7 @@ fn playground_exit_stops_all_terminals_and_observers_do_not_start_daemons() { assert!( Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .arg(command) - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .output() .unwrap() .status @@ -152,7 +165,7 @@ fn playground_exit_stops_all_terminals_and_observers_do_not_start_daemons() { assert!( !Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .args(args) - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .output() .unwrap() .status @@ -160,15 +173,17 @@ fn playground_exit_stops_all_terminals_and_observers_do_not_start_daemons() { ); assert!(!runtime.join("service.json").exists()); } - std::fs::remove_dir_all(&runtime).unwrap(); + assert!(!runtime.exists(), "runtime directory was not removed"); + std::fs::remove_dir_all(&root).unwrap(); } #[test] fn observer_stream_replays_and_follows_until_exit() { - let runtime = runtime_dir("stream"); + let root = runtime_root("stream"); + let runtime = root.join("play"); let mut owner = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .arg("play") - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .stdin(Stdio::piped()) .stdout(Stdio::piped()) .spawn() @@ -183,7 +198,7 @@ fn observer_stream_replays_and_follows_until_exit() { let watched = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) .args(["watch", &terminal_id]) - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(target(&root)) .stdout(Stdio::piped()) .stderr(Stdio::piped()) .spawn() @@ -192,7 +207,8 @@ fn observer_stream_replays_and_follows_until_exit() { owner.stdin.as_mut().unwrap().write_all(b"quit\n").unwrap(); assert!(output_with_timeout(owner).status.success()); assert!(!runtime.join("service.json").exists()); - std::fs::remove_dir_all(&runtime).unwrap(); + assert!(!runtime.exists(), "runtime directory was not removed"); + std::fs::remove_dir_all(&root).unwrap(); assert!( watched.status.success(), "watch failed with {:?}: {}", diff --git a/tests/rows.rs b/tests/rows.rs index 6b25101..1e22253 100644 --- a/tests/rows.rs +++ b/tests/rows.rs @@ -13,8 +13,9 @@ use opencode_pty::service::CreateTerminal; fn daemon_client_read_rows_roundtrip() { // Run the Rust client in a child test process so discovery's environment is // isolated without mutating this multithreaded test process's environment. - if std::env::var_os("OPENCODE_PTY_ROWS_TEST_CHILD").is_some() { - let client = TerminalClient::discover().unwrap(); + if let Some(directory) = std::env::var_os("OPENCODE_PTY_ROWS_TEST_CHILD") { + let directory = std::path::PathBuf::from(directory); + let client = TerminalClient::discover(&directory).unwrap(); let mut request = CreateTerminal::shell().unwrap(); request.program = "/bin/sh".to_string(); request.args = vec!["-c".to_string(), "printf 'one\ntwo\n\nfour\n'".to_string()]; @@ -45,7 +46,7 @@ fn daemon_client_read_rows_roundtrip() { .contains("positive") ); - let registration = opencode_pty::daemon::read_registration().unwrap(); + let registration = opencode_pty::daemon::read_registration(&directory).unwrap(); assert_eq!(registration.protocol, 7); // Exercise omitted rows independently of the Rust client's null encoding. let mut stream = std::os::unix::net::UnixStream::connect(registration.socket).unwrap(); @@ -70,11 +71,12 @@ fn daemon_client_read_rows_roundtrip() { return; } - let runtime = + let root = std::env::temp_dir().join(format!("opencode-pty-rows-{:032x}", rand::random::())); + let runtime = root.join("rows"); let mut daemon = Command::new(env!("CARGO_BIN_EXE_opencode-pty")) - .arg("daemon") - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) + .args(["daemon", "--name", "rows", "--runtime-dir"]) + .arg(&root) .stdin(Stdio::null()) .stdout(Stdio::null()) .spawn() @@ -107,12 +109,11 @@ fn daemon_client_read_rows_roundtrip() { "daemon_client_read_rows_roundtrip", "--nocapture", ]) - .env("OPENCODE_PTY_RUNTIME_DIR", &runtime) - .env("OPENCODE_PTY_ROWS_TEST_CHILD", "1") + .env("OPENCODE_PTY_ROWS_TEST_CHILD", &runtime) .output(); drop(owner); assert!(daemon.wait().unwrap().success()); - let _ = std::fs::remove_dir_all(runtime); + let _ = std::fs::remove_dir_all(root); let output = result.unwrap(); assert!( output.status.success(),