diff --git a/Cargo.lock b/Cargo.lock index c2fda1a97..be5f4e63d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1961,6 +1961,7 @@ dependencies = [ "toml_edit", "uuid", "walkdir", + "windows-sys 0.59.0", "wiremock", "zip", ] diff --git a/crates/socket-patch-cli/CLI_CONTRACT.md b/crates/socket-patch-cli/CLI_CONTRACT.md index ebde249c6..d0c7544c7 100644 --- a/crates/socket-patch-cli/CLI_CONTRACT.md +++ b/crates/socket-patch-cli/CLI_CONTRACT.md @@ -32,7 +32,7 @@ Rows are in `--help` order (v5.0): the hosted/vendored workflow (`scan` → `vex **Removed in v4.0:** the `unlock` subcommand (a leftover lock from a crashed run never blocks acquisition — the OS releases a dead holder's advisory lock — so there is no stale-lock state to inspect or clear before a mutating command; `repair` briefly owned lock-file cleanup in v4.x, and since v5.0 every lock-taking command removes its own lock file on exit). -**Lock lifecycle (v5.0).** `<.socket>/apply.lock` never outlives the command that took it: acquisition creates `.socket/` when it is missing, the guard's drop unlinks the file WHILE the lock is still held (so a waiter can never lock an orphaned inode), releases it, and then removes `.socket/` itself if that left the directory empty — a run that had nothing to persist leaves no `.socket/` behind, and there is nothing to `.gitignore`. A leftover file from a crashed (SIGKILLed) run is reclaimed in place and removed by the next lock-taking command. The lock is taken by `apply`, `rollback`, `remove`, `repair`, `vendor`, agent-mode `get` and `scan --apply`/`--sync` (download → manifest write → nested apply is ONE lock window — the nested apply never re-acquires), and `scan`/`get` in vendored **and hosted** mode — hosted acquires it before its first wet write (the staged takeover reverts), never on `--dry-run` and never when the run would write nothing, so hosted previews and no-op runs create no `.socket/`. Dry runs of the other commands may still take the lock; it is residue-free either way. A live holder is `lock_held` (exit 1); a directory or special file squatting on `.socket/` or on the lock path is a lock I/O error — `lock_io` (exit 1, `failed to open lock file at : …`; a read-only project root surfaces the same code at the acquire, before any ledger or manifest write) — never `lock_held`. +**Lock lifecycle (v5.0).** `<.socket>/apply.lock` never outlives the command that took it: acquisition creates `.socket/` when it is missing, the guard's drop unlinks the file WHILE the lock is still held (so a waiter can never lock an orphaned inode), releases it, and then removes `.socket/` itself if that left the directory empty — a run that had nothing to persist leaves no `.socket/` behind, and there is nothing to `.gitignore`. An interrupted run (Ctrl-C, SIGTERM, SIGHUP; Ctrl-C, Ctrl-Break or console close on Windows) removes the file on its way out and still dies by that signal; only an uncatchable kill (SIGKILL, power loss) can leave `apply.lock`, and the next lock-taking command reclaims it in place and removes it. socket-patch never keeps a persistent lock file in the project. The lock is taken by `apply`, `rollback`, `remove`, `repair`, `vendor`, agent-mode `get` and `scan --apply`/`--sync` (download → manifest write → nested apply is ONE lock window — the nested apply never re-acquires), and `scan`/`get` in vendored **and hosted** mode — hosted acquires it before its first wet write (the staged takeover reverts), never on `--dry-run` and never when the run would write nothing, so hosted previews and no-op runs create no `.socket/`. Dry runs of the other commands may still take the lock; it is residue-free either way. A live holder is `lock_held` (exit 1); a directory or special file squatting on `.socket/` or on the lock path is a lock I/O error — `lock_io` (exit 1, `failed to open lock file at : …`; a read-only project root surfaces the same code at the acquire, before any ledger or manifest write) — never `lock_held`. **Bare-UUID fallback.** `socket-patch ` is rewritten to `socket-patch get `. The UUID shape checked is the standard 8-4-4-4-12 hex pattern (case-insensitive). See [`src/lib.rs::looks_like_uuid`](src/lib.rs). @@ -65,7 +65,7 @@ Every subcommand accepts the same set of "global" flags via a single shared `Glo | `--json` | `-j` | `SOCKET_JSON` | `false` | bool | Machine-readable output | | `--verbose` | `-v` | `SOCKET_VERBOSE` | `false` | bool | Extra detail | | `--silent` | `-s` | `SOCKET_SILENT` | `false` | bool | Errors only | -| `--dry-run` | — | `SOCKET_DRY_RUN` | `false` | bool | Preview, no mutations (a dry run may still take the transient `apply.lock`, removed again on exit — see "Lock lifecycle"; hosted and vendored previews never leave a `.socket/`) | +| `--dry-run` | — | `SOCKET_DRY_RUN` | `false` | bool | Preview, no mutations (a dry run may still take the transient `apply.lock`, removed again on exit or on interrupt — see "Lock lifecycle"; hosted and vendored previews never leave a `.socket/`) | | `--yes` | `-y` | `SOCKET_YES` | `false` | bool | Skip prompts (`scan` never prompts) | | `--lock-timeout` | — | `SOCKET_LOCK_TIMEOUT` | (none) | seconds (u64) | How long to wait for `<.socket>/apply.lock`. Unset and `0` both mean a single non-blocking try; a positive value retries with a 100 ms backoff. Only meaningful on the lock-taking subcommands — `apply`, `rollback`, `repair`, `remove`, `vendor`, and `scan`/`get` whenever they write (agent-mode download + apply, vendored, hosted) | | `--debug` | — | `SOCKET_DEBUG` | `false` | bool | Verbose debug logs to stderr | diff --git a/crates/socket-patch-cli/src/commands/lock_cli.rs b/crates/socket-patch-cli/src/commands/lock_cli.rs index 3e9d01b0a..d749f37c8 100644 --- a/crates/socket-patch-cli/src/commands/lock_cli.rs +++ b/crates/socket-patch-cli/src/commands/lock_cli.rs @@ -36,10 +36,15 @@ use crate::json_envelope::{Command, Envelope, EnvelopeError}; /// unlinks `apply.lock` while still holding the lock, releases it, and /// prunes an otherwise-empty `.socket/` — so no command leaves a lock /// file (or a bare `.socket/`) behind, and this wrapper never has to -/// touch the file. A leftover from a crashed run never contends: the -/// kernel released the dead holder's advisory lock along with its file -/// handle, so the acquire reclaims the file in place and removes it on -/// exit. `Held` therefore always means a *live* process. +/// touch the file. An interrupted run removes the file too: the +/// `interrupt` handlers (SIGINT/SIGTERM/SIGHUP, Windows console ctrl) +/// clean up the held lock before the signal ends the process. Only an +/// uncatchable kill (SIGKILL, power loss) can leave it, and such a +/// leftover never contends: the kernel released the dead holder's +/// advisory lock along with its file handle, so the acquire reclaims the +/// file in place and removes it on exit. `Held` therefore always means a +/// *live* process. socket-patch never keeps a persistent lock file in the +/// project; one that is ever needed lives outside it. pub(crate) fn acquire_or_emit( socket_dir: &Path, command: Command, diff --git a/crates/socket-patch-cli/src/interrupt.rs b/crates/socket-patch-cli/src/interrupt.rs new file mode 100644 index 000000000..f8330dfb4 --- /dev/null +++ b/crates/socket-patch-cli/src/interrupt.rs @@ -0,0 +1,96 @@ +//! Interrupt handling: an interrupted run must not leave `.socket/apply.lock` +//! behind. +//! +//! The default disposition of SIGINT, SIGTERM and SIGHUP (and of Ctrl-C, +//! Ctrl-Break and console close on Windows) ends the process without +//! running any destructor, so a lock-taking command killed that way never +//! reaches `LockGuard`'s drop. [`install`] puts a handler in front of the +//! default one that removes the held lock file +//! ([`cleanup_held_lock_on_interrupt`]) and then lets the signal end the +//! process exactly as before: same death-by-signal status (130 for Ctrl-C +//! in a shell), same exit code on Windows. +//! +//! Only an uncatchable kill (SIGKILL, power loss) can still leave the file, +//! and the next lock-taking command reclaims and removes it. + +use socket_patch_core::patch::apply_lock::cleanup_held_lock_on_interrupt; + +/// Install the interrupt handlers. Call once, first thing in `main`. +/// +/// Unix: a signal the process started with ignored (`nohup`, a launcher +/// that ignores Ctrl-C) keeps being ignored. The prompt's cursor guard +/// (`ui::prompt`) chains in front of this handler for SIGINT while a menu +/// is up and re-raises into it, so the cursor is restored first and the +/// lock removed second. +pub fn install() { + imp::install(); +} + +#[cfg(unix)] +mod imp { + use super::cleanup_held_lock_on_interrupt; + + const SIGNALS: [libc::c_int; 3] = [libc::SIGINT, libc::SIGTERM, libc::SIGHUP]; + + extern "C" fn on_signal(sig: libc::c_int) { + cleanup_held_lock_on_interrupt(); + // SAFETY: signal and raise are async-signal-safe. `sig` is blocked + // while this handler runs, so the re-raised signal is delivered + // once it returns, to the default disposition: the process dies by + // the same signal it would have without this handler. + unsafe { + libc::signal(sig, libc::SIG_DFL); + libc::raise(sig); + } + } + + pub(super) fn install() { + for sig in SIGNALS { + // SAFETY: plain sigaction calls with zeroed, then filled, + // structs; the handler only calls async-signal-safe code. + unsafe { + let mut old: libc::sigaction = std::mem::zeroed(); + if libc::sigaction(sig, std::ptr::null(), &mut old) != 0 + || old.sa_sigaction == libc::SIG_IGN + { + continue; + } + let mut new: libc::sigaction = std::mem::zeroed(); + new.sa_sigaction = on_signal as extern "C" fn(libc::c_int) as libc::sighandler_t; + new.sa_flags = libc::SA_RESTART; + libc::sigemptyset(&mut new.sa_mask); + libc::sigaction(sig, &new, std::ptr::null_mut()); + } + } + } +} + +#[cfg(windows)] +mod imp { + use super::cleanup_held_lock_on_interrupt; + use windows_sys::Win32::Foundation::{BOOL, FALSE, TRUE}; + use windows_sys::Win32::System::Console::{ + SetConsoleCtrlHandler, CTRL_BREAK_EVENT, CTRL_CLOSE_EVENT, CTRL_C_EVENT, + }; + + /// Runs on a thread the console creates. Returning FALSE hands the + /// event on to the default handler, which ends the process. + unsafe extern "system" fn on_ctrl(ctrl: u32) -> BOOL { + if matches!(ctrl, CTRL_C_EVENT | CTRL_BREAK_EVENT | CTRL_CLOSE_EVENT) { + cleanup_held_lock_on_interrupt(); + } + FALSE + } + + pub(super) fn install() { + // SAFETY: registers a handler with the documented signature. + unsafe { + SetConsoleCtrlHandler(Some(on_ctrl), TRUE); + } + } +} + +#[cfg(not(any(unix, windows)))] +mod imp { + pub(super) fn install() {} +} diff --git a/crates/socket-patch-cli/src/lib.rs b/crates/socket-patch-cli/src/lib.rs index 002b369c6..80a1d9ef7 100644 --- a/crates/socket-patch-cli/src/lib.rs +++ b/crates/socket-patch-cli/src/lib.rs @@ -8,6 +8,7 @@ pub mod args; pub mod commands; pub(crate) mod ecosystem_dispatch; +pub mod interrupt; /// The in-memory hosted engine, which lives in core /// ([`socket_patch_core::hosted::memory`]); re-exported under its old path /// for the `hosted-bundle` harness and the integration tests. diff --git a/crates/socket-patch-cli/src/main.rs b/crates/socket-patch-cli/src/main.rs index 13de81a3d..55cc522c6 100644 --- a/crates/socket-patch-cli/src/main.rs +++ b/crates/socket-patch-cli/src/main.rs @@ -27,6 +27,10 @@ async fn main() { // pipe. restore_default_sigpipe(); + // Ctrl-C / SIGTERM / SIGHUP remove a held `.socket/apply.lock` before + // the signal ends the process (see `interrupt`). + socket_patch_cli::interrupt::install(); + // Accept the JS socket-cli's SOCKET_CLI_* peer names (silently — // they are aliases, not deprecations) so `socket login` / socket-cli // env setups work for socket-patch unchanged. Canonical names win. diff --git a/crates/socket-patch-cli/src/ui/prompt.rs b/crates/socket-patch-cli/src/ui/prompt.rs index 70f8cbd0f..0eee32183 100644 --- a/crates/socket-patch-cli/src/ui/prompt.rs +++ b/crates/socket-patch-cli/src/ui/prompt.rs @@ -219,7 +219,9 @@ fn fit_menu(prompt: &str, options: &[String], width: usize) -> (String, Vec Duration { Duration::from_secs(0) } + +/// An interrupted holder removes `apply.lock` on its way out (#808). The +/// binary is parked by the debug-only `apply_lock.acquired~pause` +/// failpoint with the lock held, then signalled. It must die by that +/// signal (the exit status a shell reports is unchanged) and leave no lock +/// file, while the manifest, real state, survives. +#[cfg(unix)] +mod interrupted_holder { + use std::os::unix::process::{CommandExt, ExitStatusExt}; + use std::path::Path; + use std::process::{Child, Stdio}; + use std::time::{Duration, Instant}; + + use super::common::{binary, hermetic_command, jvm_env}; + use super::setup_socket_dir; + + /// Spawn `apply` in `root` and wait until it holds the lock. With + /// `ignore_sigint`, the child starts with SIGINT ignored, as under + /// `nohup` or a launcher that ignores Ctrl-C. + fn spawn_holding_lock(root: &Path, ignore_sigint: bool) -> Child { + spawn_parked_at(root, "apply_lock.acquired", ignore_sigint) + } + + /// Spawn `apply` in `root` and wait until it parks at `failpoint`. + fn spawn_parked_at(root: &Path, failpoint: &str, ignore_sigint: bool) -> Child { + let ready = root.join("failpoint-ready"); + let mut cmd = hermetic_command(&binary()); + jvm_env::isolate_cli(&mut cmd); + cmd.args(["apply", "--json"]) + .current_dir(root) + .env("SOCKET_PATCH_FAILPOINT", format!("{failpoint}~pause")) + .env("SOCKET_PATCH_FAILPOINT_READY", &ready) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + if ignore_sigint { + // SAFETY: only an async-signal-safe `signal` call between fork + // and exec. + unsafe { + cmd.pre_exec(|| { + libc::signal(libc::SIGINT, libc::SIG_IGN); + Ok(()) + }); + } + } + let mut child = cmd.spawn().expect("spawn socket-patch apply"); + let deadline = Instant::now() + Duration::from_secs(60); + while !ready.exists() { + if let Some(status) = child.try_wait().unwrap() { + panic!("apply exited ({status}) before reaching the lock failpoint"); + } + assert!( + Instant::now() < deadline, + "apply never reached the lock failpoint" + ); + std::thread::sleep(Duration::from_millis(20)); + } + child + } + + fn signal(child: &Child, sig: libc::c_int) { + // SAFETY: plain kill(2) on our own child's pid. + let rc = unsafe { libc::kill(child.id() as libc::pid_t, sig) }; + assert_eq!(rc, 0, "kill({sig}) failed"); + } + + fn assert_interrupt_cleans_up(sig: libc::c_int) { + assert_interrupt_at_cleans_up("apply_lock.acquired", sig); + } + + fn assert_interrupt_at_cleans_up(failpoint: &str, sig: libc::c_int) { + let dir = tempfile::tempdir().unwrap(); + let socket_dir = dir.path().join(".socket"); + setup_socket_dir(&socket_dir); + let lock = socket_dir.join("apply.lock"); + + let mut child = spawn_parked_at(dir.path(), failpoint, false); + assert!(lock.is_file(), "the parked apply holds apply.lock"); + + signal(&child, sig); + let status = child.wait().unwrap(); + assert_eq!( + status.signal(), + Some(sig), + "the process must still die by the signal; got {status}" + ); + assert!( + !lock.exists(), + "an interrupted run must not leave apply.lock behind (signal {sig})" + ); + assert!( + socket_dir.join("manifest.json").is_file(), + "real .socket/ state survives the interrupt" + ); + } + + #[test] + fn sigint_removes_the_lock_file() { + assert_interrupt_cleans_up(libc::SIGINT); + } + + #[test] + fn sigterm_removes_the_lock_file() { + assert_interrupt_cleans_up(libc::SIGTERM); + } + + #[test] + fn sighup_removes_the_lock_file() { + assert_interrupt_cleans_up(libc::SIGHUP); + } + + /// An interrupt that lands while the guard's drop is already under + /// way — before the drop's own unlink — still removes the file: the + /// handler ends the process without resuming the drop, so the guard + /// must stay in the interrupt table until the drop has unlinked it. + #[test] + fn interrupt_during_release_removes_the_lock_file() { + assert_interrupt_at_cleans_up("apply_lock.releasing", libc::SIGTERM); + } + + /// A process started with SIGINT ignored keeps ignoring it (no handler + /// is installed over SIG_IGN), and SIGTERM still cleans up. + #[test] + fn ignored_sigint_stays_ignored_and_sigterm_still_cleans_up() { + let dir = tempfile::tempdir().unwrap(); + let socket_dir = dir.path().join(".socket"); + setup_socket_dir(&socket_dir); + let lock = socket_dir.join("apply.lock"); + + let mut child = spawn_holding_lock(dir.path(), true); + signal(&child, libc::SIGINT); + std::thread::sleep(Duration::from_millis(300)); + assert!( + child.try_wait().unwrap().is_none(), + "an ignored SIGINT must not end the run" + ); + assert!(lock.is_file(), "and must not touch the held lock"); + + signal(&child, libc::SIGTERM); + let status = child.wait().unwrap(); + assert_eq!(status.signal(), Some(libc::SIGTERM), "got {status}"); + assert!(!lock.exists(), "SIGTERM removes apply.lock"); + assert!(socket_dir.join("manifest.json").is_file()); + } +} diff --git a/crates/socket-patch-core/Cargo.toml b/crates/socket-patch-core/Cargo.toml index 02c0cc31a..b331a0734 100644 --- a/crates/socket-patch-core/Cargo.toml +++ b/crates/socket-patch-core/Cargo.toml @@ -58,6 +58,10 @@ libc = { workspace = true } # do ourselves (mode-preserving, setuid-refusing — see update/swap.rs). [target.'cfg(windows)'.dependencies] self-replace = { workspace = true } +# apply_lock's interrupt cleanup unlinks the held lock file with POSIX +# semantics so the emptied `.socket/` can be removed before exit. Same +# windows-sys the CLI already builds, so no new crate. +windows-sys = { workspace = true, features = ["Win32_Foundation", "Win32_Storage_FileSystem"] } # sha2 0.10 compiles its aarch64 SHA-256 hardware path only under the `asm` # feature (x86/x86_64 already select SHA-NI at runtime without it); the diff --git a/crates/socket-patch-core/src/patch/apply_lock.rs b/crates/socket-patch-core/src/patch/apply_lock.rs index 48cefbb7c..221344b4c 100644 --- a/crates/socket-patch-core/src/patch/apply_lock.rs +++ b/crates/socket-patch-core/src/patch/apply_lock.rs @@ -34,14 +34,38 @@ //! for the next acquire to reclaim, rather than unlinked from under //! whoever locked it. //! +//! * An interrupted holder removes the file too. [`acquire`] registers +//! every held lock in a process-global table and the guard's drop +//! takes it out again; [`cleanup_held_lock_on_interrupt`] unlinks each +//! registered file (gated on the same inode identity as the drop) and +//! prunes an empty `.socket/`. The CLI calls it from its SIGINT, +//! SIGTERM and SIGHUP handlers (Ctrl-C, Ctrl-Break and console close +//! on Windows) before letting the signal end the process, because the +//! default disposition kills the process without running any drop. +//! Installing those handlers is the host's job: a host that embeds +//! this crate without the CLI (the Node addon) installs none, and an +//! interrupt there behaves like a crash. +//! The unlink precedes the process's death by a few syscalls (on +//! Windows, by the default ctrl handler's `ExitProcess`), during which +//! its other threads still run; a competing acquire that lands in that +//! window could briefly overlap them. +//! //! So no command leaves `apply.lock` behind (barring that non-cooperating //! replacement, which the next lock-taking command reclaims), and a //! project that had no //! `.socket/` before a run has none after it unless the run wrote real -//! state there. A leftover file from a crashed run needs no removal to -//! unblock anything — the kernel released the dead process's advisory -//! lock with its file handle — so the next acquire reclaims it in place -//! and removes it on exit. `Held` therefore always means a live process. +//! state there. Only an uncatchable kill (SIGKILL, power loss) can leave +//! the file. Such a leftover needs no removal to unblock anything — the +//! kernel released the dead process's advisory lock with its file +//! handle — so the next acquire reclaims it in place and removes it on +//! exit. `Held` therefore always means a live process. +//! +//! socket-patch never keeps a persistent lock file in the project. That +//! is why this module carries the deletion machinery (the identity +//! checks, the vanished/delete-pending retries, the unlink under the +//! lock). If a persistent lock is ever needed, it must live outside the +//! project (the user cache dir or a configurable path), never under +//! `.socket/`. //! //! # Windows //! @@ -156,6 +180,9 @@ pub struct LockGuard { handle: Option, path: PathBuf, socket_dir: PathBuf, + /// This guard's entry in the interrupt-cleanup table, if it got one + /// (see [`cleanup_held_lock_on_interrupt`]). + registration: Option, } impl LockGuard { @@ -178,6 +205,14 @@ impl Drop for LockGuard { run `socket-patch repair` after an unclean shutdown" ); } + // The guard stays in the interrupt table through R1, so an + // interrupt anywhere before the unlink below has happened (during + // the barrier, which can take seconds over a large vendored tree, + // or during the identity probe) still removes the file: the CLI's + // handler ends the process, so this drop never gets to finish. + // Debug-only hook: the interrupt tests park the process here and + // signal it. + crate::utils::failpoint::hit("apply_lock.releasing"); // R1: unlink while still holding the lock — but only the file we // hold. The unlink is gated on the path still naming the held // inode: after a non-cooperating `rm` + `touch`, the path names a @@ -198,12 +233,25 @@ impl Drop for LockGuard { if unlink { let _ = std::fs::remove_file(&self.path); } + // The file is dealt with: from here on an interrupt must only + // prune the directory, never unlink. Once the handle closes below, + // the inode number can be reused by the next holder's file and + // would pass the cleanup's identity check. (On Windows this also + // closes the table's duplicate handle, so it never outlives ours.) + if let Some(registration) = &self.registration { + interrupt::mark_released(registration); + } // R2: close the handle; the OS releases the advisory lock. self.handle = None; // R3: prune an otherwise-empty `.socket/`. Non-recursive, so it // fails harmlessly when anything else lives there — including a // concurrent acquirer's freshly created `apply.lock`. prune_empty_socket_dir(&self.socket_dir); + // Only now leave the interrupt table: an interrupt between R1 and + // here still prunes the `.socket/` this run left empty. + if let Some(registration) = self.registration.take() { + interrupt::unregister(registration); + } } } @@ -220,6 +268,303 @@ fn prune_empty_socket_dir(socket_dir: &Path) { } } +/// Remove every lock file this process holds, for a host's interrupt +/// handler (the CLI's SIGINT/SIGTERM/SIGHUP handler on Unix, its console +/// ctrl handler on Windows) to call just before the signal ends the +/// process. The default disposition kills the process without running +/// [`LockGuard`]'s drop, so without this an interrupted run would leave +/// `apply.lock` (and a `.socket/` the run created) behind. +/// +/// For each registered guard: if the lock path still names the inode the +/// guard holds, unlink it, then remove the `.socket/` directory if it is +/// now empty. The identity gate is the drop's: a replacement planted by a +/// non-cooperating `rm` + `touch` belongs to someone else and is left +/// alone. A guard whose drop already got past its own unlink is only +/// pruned, never unlinked. Nothing is unregistered, so a second call is a +/// harmless no-op. +/// +/// Unix: async-signal-safe. It reads only memory built at acquire time +/// and calls only `stat`, `unlink` and `rmdir`. Windows: the console ctrl +/// handler runs on an ordinary thread, so this takes a mutex and uses +/// std. A plain `DeleteFile` would leave the name delete-pending until +/// the process's lock handle closes at exit — too late to remove the +/// directory — so the unlink uses POSIX semantics (the name goes at once, +/// NTFS on Windows 10 1709 and later), falling back to `DeleteFile` +/// (and so to keeping an empty `.socket/`) where that is unsupported. +/// +/// Call it only when the process is about to end: the guards it cleaned +/// up still hold their locks, but their files are gone. +pub fn cleanup_held_lock_on_interrupt() { + interrupt::cleanup(None); +} + +/// The interrupt-cleanup table behind [`cleanup_held_lock_on_interrupt`]. +/// +/// Every command holds one guard at a time, but the in-process unit tests +/// hold many at once, so the table has a few slots. A guard that finds +/// them all taken is simply not registered: its drop still cleans up +/// normally, only an interrupt would leave its file behind. +#[cfg(unix)] +mod interrupt { + use std::ffi::{CStr, CString}; + use std::os::unix::ffi::OsStrExt; + use std::path::Path; + use std::ptr; + use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicUsize, Ordering::SeqCst}; + + use super::{LockGuard, SOCKET_DIR_NAME}; + + const SLOTS: usize = 8; + + /// Everything the handler needs, built before the slot is published so + /// the handler never allocates. + pub(super) struct Held { + lock: CString, + /// `None` unless the directory is literally named `.socket` (the + /// same gate as `prune_empty_socket_dir`). + socket_dir: Option, + dev: u64, + ino: u64, + /// Set by the guard's drop once it has dealt with the file itself: + /// the cleanup then only prunes the directory. + released: AtomicBool, + } + + static HELD: [AtomicPtr; SLOTS] = [const { AtomicPtr::new(ptr::null_mut()) }; SLOTS]; + + /// Number of cleanups currently reading the table. `unregister` frees + /// an entry only when none is: a cleanup that started before the + /// entry left the table may still be reading it. (With SeqCst, a + /// cleanup that starts after `unregister` read zero also loads the + /// slot after it was cleared, so it never sees the freed entry.) + static READERS: AtomicUsize = AtomicUsize::new(0); + + #[derive(Debug)] + pub(super) struct Registration(usize); + + fn c_path(path: &Path) -> Option { + CString::new(path.as_os_str().as_bytes()).ok() + } + + pub(super) fn register(guard: &LockGuard) -> Option { + let handle = guard.handle.as_ref()?; + let socket_dir = if guard + .socket_dir + .file_name() + .is_some_and(|name| name == SOCKET_DIR_NAME) + { + Some(c_path(&guard.socket_dir)?) + } else { + None + }; + let held = Box::into_raw(Box::new(Held { + lock: c_path(&guard.path)?, + socket_dir, + dev: handle.dev(), + ino: handle.ino(), + released: AtomicBool::new(false), + })); + for (i, slot) in HELD.iter().enumerate() { + if slot + .compare_exchange(ptr::null_mut(), held, SeqCst, SeqCst) + .is_ok() + { + return Some(Registration(i)); + } + } + // SAFETY: never published, so this is still the only pointer. + drop(unsafe { Box::from_raw(held) }); + None + } + + pub(super) fn mark_released(Registration(i): &Registration) { + let held = HELD[*i].load(SeqCst); + if !held.is_null() { + // SAFETY: only `unregister`, which takes the registration by + // value, frees the entry; we hold that registration. + unsafe { &*held }.released.store(true, SeqCst); + } + } + + pub(super) fn unregister(Registration(i): Registration) { + let held = HELD[i].swap(ptr::null_mut(), SeqCst); + if !held.is_null() && READERS.load(SeqCst) == 0 { + // SAFETY: the slot no longer publishes it and no cleanup is + // reading the table, so this is the last reference. When a + // cleanup is running, the process is about to die: leak it. + drop(unsafe { Box::from_raw(held) }); + } + } + + /// `only`: restrict the cleanup to one lock path (the unit tests, which + /// share the table with every other test in the binary). + pub(super) fn cleanup(only: Option<&CStr>) { + READERS.fetch_add(1, SeqCst); + for slot in &HELD { + let held = slot.load(SeqCst); + if held.is_null() { + continue; + } + // SAFETY: published entries are freed only after leaving the + // table while READERS is zero, and we hold READERS above zero. + let held = unsafe { &*held }; + if only.is_some_and(|only| only != held.lock.as_c_str()) { + continue; + } + // SAFETY: plain syscalls on NUL-terminated paths built at + // acquire time; `st` is a valid out-pointer. (The casts are + // no-ops on some targets only: `st_dev` is `i32` on macOS.) + #[allow(clippy::unnecessary_cast)] + unsafe { + let mut st: libc::stat = std::mem::zeroed(); + if !held.released.load(SeqCst) + && libc::stat(held.lock.as_ptr(), &mut st) == 0 + && st.st_dev as u64 == held.dev + && st.st_ino as u64 == held.ino + { + libc::unlink(held.lock.as_ptr()); + } + if let Some(dir) = &held.socket_dir { + // Non-recursive: fails harmlessly when not empty. + libc::rmdir(dir.as_ptr()); + } + } + } + READERS.fetch_sub(1, SeqCst); + } +} + +#[cfg(windows)] +mod interrupt { + use std::fs::OpenOptions; + use std::os::windows::fs::OpenOptionsExt; + use std::os::windows::io::AsRawHandle; + use std::path::{Path, PathBuf}; + use std::sync::Mutex; + + use same_file::Handle; + use windows_sys::Win32::Storage::FileSystem::{ + FileDispositionInfoEx, SetFileInformationByHandle, DELETE, FILE_DISPOSITION_FLAG_DELETE, + FILE_DISPOSITION_FLAG_POSIX_SEMANTICS, FILE_DISPOSITION_INFO_EX, FILE_READ_ATTRIBUTES, + FILE_SHARE_DELETE, FILE_SHARE_READ, FILE_SHARE_WRITE, + }; + + use super::{prune_empty_socket_dir, LockGuard}; + + struct Held { + id: u64, + lock: PathBuf, + socket_dir: PathBuf, + /// A duplicate of the guard's handle, kept only for the identity + /// comparison. The guard's drop clears it once it has dealt with + /// the file itself, before closing its own handle, so it never + /// keeps the lock or a delete-pending name alive. `None` then + /// means "prune only". + handle: Option, + } + + static HELD: Mutex> = Mutex::new(Vec::new()); + static NEXT_ID: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + + #[derive(Debug)] + pub(super) struct Registration(u64); + + fn table() -> std::sync::MutexGuard<'static, Vec> { + HELD.lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } + + pub(super) fn register(guard: &LockGuard) -> Option { + let dup = guard.handle.as_ref()?.as_file().try_clone().ok()?; + let handle = Handle::from_file(dup).ok()?; + let id = NEXT_ID.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + table().push(Held { + id, + lock: guard.path.clone(), + socket_dir: guard.socket_dir.clone(), + handle: Some(handle), + }); + Some(Registration(id)) + } + + pub(super) fn mark_released(Registration(id): &Registration) { + if let Some(held) = table().iter_mut().find(|held| held.id == *id) { + held.handle = None; + } + } + + pub(super) fn unregister(Registration(id): Registration) { + table().retain(|held| held.id != id); + } + + /// Unlink `path` if it still names `held`, with POSIX semantics so the + /// name is gone at once even though our lock handle stays open until + /// the process exits. Falls back to `DeleteFile` (delete-pending until + /// exit) where POSIX semantics are unsupported (FAT, older Windows). + fn unlink_if_held(path: &Path, held: &Handle) { + let Ok(file) = OpenOptions::new() + .access_mode(DELETE | FILE_READ_ATTRIBUTES) + .share_mode(FILE_SHARE_READ | FILE_SHARE_WRITE | FILE_SHARE_DELETE) + .open(path) + else { + return; + }; + let Ok(now) = Handle::from_file(file) else { + return; + }; + if now != *held { + return; + } + let info = FILE_DISPOSITION_INFO_EX { + Flags: FILE_DISPOSITION_FLAG_DELETE | FILE_DISPOSITION_FLAG_POSIX_SEMANTICS, + }; + // SAFETY: a live handle opened with DELETE access, and a correctly + // sized, initialized FILE_DISPOSITION_INFO_EX. + let ok = unsafe { + SetFileInformationByHandle( + now.as_file().as_raw_handle(), + FileDispositionInfoEx, + std::ptr::from_ref(&info).cast(), + std::mem::size_of::() as u32, + ) + }; + drop(now); + if ok == 0 { + let _ = std::fs::remove_file(path); + } + } + + pub(super) fn cleanup(only: Option<&Path>) { + for held in table().iter() { + if only.is_some_and(|only| only != held.lock) { + continue; + } + if let Some(handle) = &held.handle { + unlink_if_held(&held.lock, handle); + } + prune_empty_socket_dir(&held.socket_dir); + } + } +} + +#[cfg(not(any(unix, windows)))] +mod interrupt { + use super::LockGuard; + + #[derive(Debug)] + pub(super) struct Registration; + + pub(super) fn register(_: &LockGuard) -> Option { + None + } + + pub(super) fn mark_released(_: &Registration) {} + + pub(super) fn unregister(_: Registration) {} + + pub(super) fn cleanup(_: Option<&std::path::Path>) {} +} + /// Finish the group commit a crashed vendored run left half-written (see /// [`crate::utils::group_commit`]), now that no other command can be /// writing the files it covers. Every locked command runs it before reading @@ -319,9 +664,13 @@ pub fn acquire(socket_dir: &Path, timeout: Duration) -> Result { + Attempt::Acquired(mut guard) => { + guard.registration = interrupt::register(&guard); // On failure the guard drops here, releasing the lock. recover_group_commit(socket_dir)?; + // Debug-only hook: the interrupt tests park the process + // here, lock held, and signal it. + crate::utils::failpoint::hit("apply_lock.acquired"); return Ok(guard); } Attempt::Contended => { @@ -476,6 +825,7 @@ fn attempt(path: &Path, socket_dir: &Path) -> Attempt { handle: Some(held), path: path.to_path_buf(), socket_dir: socket_dir.to_path_buf(), + registration: None, }), // The path names a replacement: a releaser unlinked the inode we // locked between our open and our lock, and a newcomer created @@ -972,6 +1322,111 @@ mod tests { assert!(!socket.exists()); } + /// Run the interrupt cleanup for one lock path only: the table is + /// process-global, and the other tests in this binary hold locks of + /// their own while this runs. + #[cfg(unix)] + fn interrupt_cleanup_of(path: &Path) { + use std::os::unix::ffi::OsStrExt; + let path = std::ffi::CString::new(path.as_os_str().as_bytes()).unwrap(); + interrupt::cleanup(Some(&path)); + } + + #[cfg(windows)] + fn interrupt_cleanup_of(path: &Path) { + interrupt::cleanup(Some(path)); + } + + /// An interrupt removes the held lock file and the empty `.socket/` + /// the acquire created, and the guard's later drop (which a real + /// interrupt never reaches) still copes with the file being gone. + #[cfg(any(unix, windows))] + #[test] + fn interrupt_cleanup_removes_the_held_file_and_empty_socket_dir() { + let dir = tempfile::tempdir().unwrap(); + let socket = socket_dir(&dir); + let lock_path = socket.join("apply.lock"); + + let guard = acquire(&socket, Duration::ZERO).unwrap(); + assert!(guard.registration.is_some(), "acquire registers the lock"); + assert!(lock_path.is_file()); + + interrupt_cleanup_of(&lock_path); + assert!(!lock_path.exists(), "the interrupt removes apply.lock"); + assert!(!socket.exists(), "and the .socket/ it left empty"); + + drop(guard); + assert!(!lock_path.exists()); + assert!(!socket.exists()); + } + + /// The interrupt cleanup is gated on identity like the drop: a + /// replacement planted by `rm` + `touch` is someone else's file. + #[cfg(any(unix, windows))] + #[test] + fn interrupt_cleanup_leaves_a_replacement_file_alone() { + let dir = tempfile::tempdir().unwrap(); + let socket = socket_dir(&dir); + let lock_path = socket.join("apply.lock"); + + let guard = acquire(&socket, Duration::ZERO).unwrap(); + std::fs::remove_file(&lock_path).unwrap(); + std::fs::File::create(&lock_path).unwrap(); + + interrupt_cleanup_of(&lock_path); + assert!( + lock_path.is_file(), + "the cleanup must not unlink a file the guard does not hold" + ); + assert!(socket.is_dir()); + + drop(guard); + assert!(lock_path.is_file(), "nor may the drop"); + } + + /// Once the drop has dealt with the file (R1), an interrupt that lands + /// before the drop finishes only prunes: it must not unlink, because + /// after the handle closes the path may name the next holder's file. + #[cfg(any(unix, windows))] + #[test] + fn interrupt_cleanup_after_the_drops_unlink_only_prunes() { + let dir = tempfile::tempdir().unwrap(); + let socket = socket_dir(&dir); + let lock_path = socket.join("apply.lock"); + + let guard = acquire(&socket, Duration::ZERO).unwrap(); + interrupt::mark_released(guard.registration.as_ref().unwrap()); + interrupt_cleanup_of(&lock_path); + assert!( + lock_path.is_file(), + "a released guard's path is not unlinked" + ); + + drop(guard); + assert!(!lock_path.exists()); + assert!(!socket.exists()); + } + + /// The drop takes the guard out of the table: a later interrupt in + /// the same process leaves the next holder's file alone. + #[cfg(any(unix, windows))] + #[test] + fn dropped_guard_is_unregistered() { + let dir = tempfile::tempdir().unwrap(); + let socket = socket_dir(&dir); + let lock_path = socket.join("apply.lock"); + + drop(acquire(&socket, Duration::ZERO).unwrap()); + // Hold the next file through a plain handle, not a guard, so + // nothing re-registers this path. + std::fs::create_dir_all(&socket).unwrap(); + let other = std::fs::File::create(&lock_path).unwrap(); + other.try_lock_exclusive().unwrap(); + + interrupt_cleanup_of(&lock_path); + assert!(lock_path.is_file()); + } + /// Two threads hammering acquire/release on one `.socket/` — every /// release unlinking the file and pruning the directory, every /// acquire recreating both — must never observe two live guards and diff --git a/crates/socket-patch-core/src/utils/failpoint.rs b/crates/socket-patch-core/src/utils/failpoint.rs index 7b1522a87..ce29c491f 100644 --- a/crates/socket-patch-core/src/utils/failpoint.rs +++ b/crates/socket-patch-core/src/utils/failpoint.rs @@ -5,8 +5,16 @@ //! (status 86, no destructors, no further writes) at the `n`-th time //! (default: the first) it reaches [`hit`] with that name — the observable //! effect of a crash at that point: whatever was written before is on disk, -//! nothing after it is. Release builds compile [`hit`] to nothing, so no -//! environment variable can make a shipped binary stop half-way; the same +//! nothing after it is. +//! +//! `[@]~pause` parks the process there instead, forever, so a test +//! can signal it at a known point (the interrupt tests do this with the +//! apply lock held). Before parking it prints a `failpoint … paused` line +//! to stderr and, when `SOCKET_PATCH_FAILPOINT_READY=` is set, +//! creates that file as the test's ready marker. +//! +//! Release builds compile [`hit`] to nothing, so no environment variable +//! can make a shipped binary stop half-way; the same //! goes for [`switched_off`]. The `failpoints` feature compiles them into an //! optimized build as well; only the CLI's dev-dependencies enable it, so //! it reaches the optimized test binaries (`cargo test --release`) and @@ -32,10 +40,25 @@ pub fn hit(name: &str) { *count }; for point in spec.split(',') { - let (point, at) = match point.trim().split_once('@') { + let point = point.trim(); + let (point, pause) = match point.strip_suffix("~pause") { + Some(point) => (point, true), + None => (point, false), + }; + let (point, at) = match point.split_once('@') { Some((p, n)) => (p, n.parse::().unwrap_or(1)), - None => (point.trim(), 1), + None => (point, 1), }; + if point == name && at == count && pause { + drop(hits); + eprintln!("socket-patch: failpoint `{name}` #{count} paused"); + if let Some(ready) = std::env::var_os("SOCKET_PATCH_FAILPOINT_READY") { + let _ = std::fs::write(ready, b""); + } + loop { + std::thread::sleep(std::time::Duration::from_secs(3600)); + } + } if point == name && at == count { eprintln!("socket-patch: failpoint `{name}` #{count} hit; exiting"); std::process::exit(86);