diff --git a/MODULE.bazel.lock b/MODULE.bazel.lock index 79850381f..dd455e539 100644 --- a/MODULE.bazel.lock +++ b/MODULE.bazel.lock @@ -505,7 +505,7 @@ "FILE:@@//Cargo.toml 882e091720fde608c21c051c8efec1ddf8dda31f2274416699bebe2f2a5e66ec", "FILE:@@//crates/pi-ast/Cargo.toml 4467b5854e91088ad1669756f7a0aaec3c2208fc28057a29982ccc902ecaf356", "FILE:@@//crates/pi-iso/Cargo.toml 1cd7ed7e5e8cbe9b1683f38f1250d475d6bec29cd84c757af817eb5a79cffdb2", - "FILE:@@//crates/pi-natives/Cargo.toml fd137965d107e533777486cacd6142c5199cd6761de33dc4e43437149ba29d44", + "FILE:@@//crates/pi-natives/Cargo.toml 8e2e2be6de8b3b765683fc7a75617ec869bafb8cff45741bc4589d7d0dddf3c0", "FILE:@@//crates/pi-shell/Cargo.toml dff7badd2729fe8ba492500bab6119cb994f3d0bb21d559587b2c6a17168408a", "FILE:@@//crates/vendor/brush-builtins/Cargo.toml 9799575fe609133927291fbca4c4e190e0378c7efe20aa9f66bb0ae6a08cb6ef", "FILE:@@//crates/vendor/brush-core/Cargo.toml a15b65352d399a599153c739a69064c82c12170fd989bf96adeaee87bed98c2c", diff --git a/crates/pi-natives/Cargo.toml b/crates/pi-natives/Cargo.toml index 3f48a8546..a0b56e945 100644 --- a/crates/pi-natives/Cargo.toml +++ b/crates/pi-natives/Cargo.toml @@ -91,7 +91,7 @@ objc2-core-foundation = "=0.3.2" libc.workspace = true [target.'cfg(windows)'.dependencies] -windows-sys = { workspace = true, features = ["Wdk_Storage_FileSystem", "Win32_Foundation", "Win32_Graphics_Gdi", "Win32_Security", "Win32_UI_Input_KeyboardAndMouse", "Win32_UI_WindowsAndMessaging"] } +windows-sys = { workspace = true, features = ["Wdk_Storage_FileSystem", "Win32_Foundation", "Win32_Graphics_Gdi", "Win32_Security", "Win32_System_Threading", "Win32_UI_Input_KeyboardAndMouse", "Win32_UI_WindowsAndMessaging"] } clipboard-win.workspace = true uiautomation = "=0.25.0" winreg.workspace = true diff --git a/crates/pi-natives/src/file_lock/linux.rs b/crates/pi-natives/src/file_lock/linux.rs new file mode 100644 index 000000000..5ae538d5f --- /dev/null +++ b/crates/pi-natives/src/file_lock/linux.rs @@ -0,0 +1,29 @@ +use std::{ + io, + os::{ + linux::net::SocketAddrExt, + unix::net::{SocketAddr, UnixDatagram}, + }, +}; + +/// Linux lock held by an abstract Unix-domain socket binding. +pub struct PlatformFileLock { + socket: Option, +} + +pub fn try_acquire(path: &str) -> io::Result> { + let name = super::memory_lock_name(path); + let address = SocketAddr::from_abstract_name(name.as_bytes())?; + match UnixDatagram::bind_addr(&address) { + Ok(socket) => Ok(Some(PlatformFileLock { socket: Some(socket) })), + Err(error) if error.kind() == io::ErrorKind::AddrInUse => Ok(None), + Err(error) => Err(error), + } +} + +impl PlatformFileLock { + pub fn release(&mut self) -> io::Result<()> { + drop(self.socket.take()); + Ok(()) + } +} diff --git a/crates/pi-natives/src/file_lock/mod.rs b/crates/pi-natives/src/file_lock/mod.rs new file mode 100644 index 000000000..43e622762 --- /dev/null +++ b/crates/pi-natives/src/file_lock/mod.rs @@ -0,0 +1,86 @@ +//! Cross-process advisory locks backed by platform ownership primitives. +//! +//! Linux uses abstract Unix sockets and Windows uses named mutexes, so neither +//! platform leaves a filesystem artifact. Other Unix platforms use `flock(2)` +//! on a persistent sidecar because they lack a process-owned in-memory name +//! registry with automatic crash recovery. + +use napi_derive::napi; + +#[cfg(target_os = "linux")] +mod linux; +#[cfg(all(unix, not(target_os = "linux")))] +mod unix; +#[cfg(target_os = "windows")] +mod windows; + +#[cfg(target_os = "linux")] +use linux as platform; +#[cfg(all(unix, not(target_os = "linux")))] +use unix as platform; +#[cfg(target_os = "windows")] +use windows as platform; + +#[cfg(not(any(unix, target_os = "windows")))] +compile_error!("pi-natives file locks require Unix or Windows"); + +#[cfg(any(target_os = "linux", target_os = "windows"))] +fn memory_lock_name(path: &str) -> String { + const HIGH_SEED: u64 = 0x4f4d_502d_4c4f_434b; + const LOW_SEED: u64 = 0x5049_2d46_494c_454c; + let bytes = path.as_bytes(); + let high = xxhash_rust::xxh64::xxh64(bytes, HIGH_SEED); + let low = xxhash_rust::xxh64::xxh64(bytes, LOW_SEED); + format!("omp-file-lock-{high:016x}{low:016x}") +} + +/// Process-owned cross-platform advisory lock. +/// +/// `tryAcquire()` is non-blocking; its returned handle reports whether it won +/// through `acquired`. Ownership ends on `release()`, garbage collection, or +/// process exit; `release()` is idempotent. +#[napi(js_name = "FileLock")] +pub struct FileLock { + inner: Option, +} + +#[napi] +impl FileLock { + /// Try to acquire `path` without blocking. + #[napi(factory)] + pub fn try_acquire(path: String) -> napi::Result { + let inner = platform::try_acquire(&path).map_err(|error| { + napi::Error::from_reason(format!("Failed to acquire native file lock for {path}: {error}")) + })?; + Ok(Self { inner }) + } + + /// Whether this handle owns the requested lock. + #[napi(getter)] + pub fn acquired(&self) -> bool { + self.inner.is_some() + } + + /// Release this handle's ownership without affecting a successor. + #[napi] + pub fn release(&mut self) -> napi::Result<()> { + let Some(mut inner) = self.inner.take() else { + return Ok(()); + }; + if let Err(error) = inner.release() { + self.inner = Some(inner); + return Err(napi::Error::from_reason(format!( + "Failed to release native file lock: {error}" + ))); + } + Ok(()) + } +} + +impl Drop for FileLock { + fn drop(&mut self) { + if let Some(inner) = self.inner.as_mut() { + let _ = inner.release(); + } + } +} diff --git a/crates/pi-natives/src/file_lock/unix.rs b/crates/pi-natives/src/file_lock/unix.rs new file mode 100644 index 000000000..da449802d --- /dev/null +++ b/crates/pi-natives/src/file_lock/unix.rs @@ -0,0 +1,38 @@ +use std::{ + fs::{File, OpenOptions}, + io, + os::{fd::AsRawFd, unix::fs::OpenOptionsExt}, +}; + +/// Unix lock held by a persistent sidecar's open file description. +pub struct PlatformFileLock { + file: Option, +} + +pub fn try_acquire(path: &str) -> io::Result> { + let file = OpenOptions::new() + .read(true) + .write(true) + .create(true) + .mode(0o600) + .open(path)?; + // SAFETY: `file` owns a live descriptor for the duration of this call. The + // non-blocking operation only changes the advisory lock attached to that open + // file description. + let status = unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) }; + if status == 0 { + return Ok(Some(PlatformFileLock { file: Some(file) })); + } + let error = io::Error::last_os_error(); + if error.kind() == io::ErrorKind::WouldBlock { + return Ok(None); + } + Err(error) +} + +impl PlatformFileLock { + pub fn release(&mut self) -> io::Result<()> { + drop(self.file.take()); + Ok(()) + } +} diff --git a/crates/pi-natives/src/file_lock/windows.rs b/crates/pi-natives/src/file_lock/windows.rs new file mode 100644 index 000000000..b7d41d9ef --- /dev/null +++ b/crates/pi-natives/src/file_lock/windows.rs @@ -0,0 +1,59 @@ +use std::{ + io, + os::windows::io::{AsRawHandle, FromRawHandle, OwnedHandle}, + ptr, +}; + +use windows_sys::Win32::{ + Foundation::{ERROR_ALREADY_EXISTS, GetLastError, SetLastError}, + System::Threading::{CreateMutexW, ReleaseMutex}, +}; + +/// Windows lock held by a named kernel mutex. +pub struct PlatformFileLock { + handle: Option, +} + +pub fn try_acquire(path: &str) -> io::Result> { + let name = super::memory_lock_name(path); + let wide_name: Vec = format!(r"Global\{name}") + .encode_utf16() + .chain(std::iter::once(0)) + .collect(); + + // `bInitialOwner` only grants ownership when this call creates the mutex. + // Existing mutexes return `ERROR_ALREADY_EXISTS` without changing ownership, + // which also prevents Win32's same-thread recursive acquisition behavior. + // SAFETY: the attributes pointer is null, and `wide_name` is a live, + // NUL-terminated UTF-16 string for the duration of the call. + unsafe { SetLastError(0) }; + let raw_handle = unsafe { CreateMutexW(ptr::null(), 1, wide_name.as_ptr()) }; + if raw_handle.is_null() { + return Err(io::Error::last_os_error()); + } + // SAFETY: `CreateMutexW` returned a fresh owned handle. `OwnedHandle` closes + // it exactly once on every return path. + let handle = unsafe { OwnedHandle::from_raw_handle(raw_handle) }; + // SAFETY: this immediately observes the last-error value set by + // `CreateMutexW`; no intervening system call can overwrite it. + if unsafe { GetLastError() } == ERROR_ALREADY_EXISTS { + return Ok(None); + } + Ok(Some(PlatformFileLock { handle: Some(handle) })) +} + +impl PlatformFileLock { + pub fn release(&mut self) -> io::Result<()> { + let Some(handle) = self.handle.as_ref() else { + return Ok(()); + }; + // SAFETY: this handle was created with initial ownership on the calling + // N-API thread. On failure the handle stays live so ownership cannot be + // silently transferred. + if unsafe { ReleaseMutex(handle.as_raw_handle()) } == 0 { + return Err(io::Error::last_os_error()); + } + drop(self.handle.take()); + Ok(()) + } +} diff --git a/crates/pi-natives/src/lib.rs b/crates/pi-natives/src/lib.rs index 793f2ce87..2542d09a6 100644 --- a/crates/pi-natives/src/lib.rs +++ b/crates/pi-natives/src/lib.rs @@ -32,6 +32,7 @@ pub mod desktop; pub mod devicecheck; pub mod diff; pub mod fd; +pub mod file_lock; pub mod glob; pub mod glob_util; pub mod grep; diff --git a/packages/coding-agent/src/security/store.ts b/packages/coding-agent/src/security/store.ts index 7321e6825..d4c0dcc7e 100644 --- a/packages/coding-agent/src/security/store.ts +++ b/packages/coding-agent/src/security/store.ts @@ -34,7 +34,7 @@ const SECURITY_STORE_WRITE_CHAINS = new Map>(); async function withSecurityStoreWrite(key: string, operation: () => Promise): Promise { const lockTarget = path.join(key, "index.json"); const run = (SECURITY_STORE_WRITE_CHAINS.get(key) ?? Promise.resolve()).then(() => - withFileLock(lockTarget, operation, { staleMs: 60_000, retries: 200, retryDelayMs: 50 }), + withFileLock(lockTarget, operation, { retries: 200, retryDelayMs: 50 }), ); const guarded = run.catch(() => undefined); SECURITY_STORE_WRITE_CHAINS.set(key, guarded); diff --git a/packages/natives/CHANGELOG.md b/packages/natives/CHANGELOG.md index 514d826b5..5262a1fe5 100644 --- a/packages/natives/CHANGELOG.md +++ b/packages/natives/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Added + +- Added non-blocking, process-owned `FileLock` bindings using abstract Unix sockets on Linux, named mutexes on Windows, and persistent `flock(2)` sidecars on other Unix platforms. + ## [17.2.5] - 2026-08-03 ### Breaking Changes diff --git a/packages/natives/README.md b/packages/natives/README.md index b14d0e0df..55a5915ae 100644 --- a/packages/natives/README.md +++ b/packages/natives/README.md @@ -9,6 +9,7 @@ Native Rust functionality via N-API. - **SIXEL**: Terminal image encoding for SIXEL-capable terminals (decode, resize, encode in one pass) - **Audio**: Cross-platform low-latency microphone capture and gapless speaker playback - **WebRTC**: Native Opus media, SDP offer/answer negotiation, and data-channel events for live sessions +- **File locking**: Process-owned cross-process locks with in-memory kernel names on Linux/Windows and `flock(2)` sidecars on other Unix platforms General-purpose image processing (decode/resize/encode for files and buffers) lives in [`Bun.Image`](https://bun.com/docs/runtime/image) on the JS side; this diff --git a/packages/natives/native/index.d.ts b/packages/natives/native/index.d.ts index 67c5a6770..c35f1236d 100644 --- a/packages/natives/native/index.d.ts +++ b/packages/natives/native/index.d.ts @@ -61,6 +61,22 @@ export declare class DesktopSession { close(): Promise } +/** + * Process-owned cross-platform advisory lock. + * + * `tryAcquire()` is non-blocking; its returned handle reports whether it won + * through `acquired`. Ownership ends on `release()`, garbage collection, or + * process exit; `release()` is idempotent. + */ +export declare class FileLock { + /** Try to acquire `path` without blocking. */ + static tryAcquire(path: string): FileLock + /** Whether this handle owns the requested lock. */ + get acquired(): boolean + /** Release this handle's ownership without affecting a successor. */ + release(): void +} + /** WebRTC peer that accepts 16 kHz mono PCM and renders remote Opus audio. */ export declare class LiveWebRtcPeer { /** diff --git a/packages/natives/native/index.js b/packages/natives/native/index.js index b34a359b3..ec64b34f2 100644 --- a/packages/natives/native/index.js +++ b/packages/natives/native/index.js @@ -19,6 +19,7 @@ const nativeBindings = loadNative(); export const AudioCapture = nativeBindings.AudioCapture; export const AudioPlayback = nativeBindings.AudioPlayback; export const DesktopSession = nativeBindings.DesktopSession; +export const FileLock = nativeBindings.FileLock; export const LiveWebRtcPeer = nativeBindings.LiveWebRtcPeer; export const MacAppearanceObserver = nativeBindings.MacAppearanceObserver; export const MacOSPowerAssertion = nativeBindings.MacOSPowerAssertion; diff --git a/packages/natives/test/file-lock.test.ts b/packages/natives/test/file-lock.test.ts new file mode 100644 index 000000000..3618da300 --- /dev/null +++ b/packages/natives/test/file-lock.test.ts @@ -0,0 +1,38 @@ +import { expect, test } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as os from "node:os"; +import * as path from "node:path"; +import { FileLock } from "../native/index.js"; + +test("FileLock binds release to one native owner", async () => { + const root = await fs.mkdtemp(path.join(os.tmpdir(), "pi-native-lock-")); + const lockPath = path.join(root, "resource.lock"); + try { + const first = FileLock.tryAcquire(lockPath); + expect(first.acquired).toBe(true); + + const blocked = FileLock.tryAcquire(lockPath); + expect(blocked.acquired).toBe(false); + blocked.release(); + + first.release(); + const successor = FileLock.tryAcquire(lockPath); + expect(successor.acquired).toBe(true); + + first.release(); + const third = FileLock.tryAcquire(lockPath); + expect(third.acquired).toBe(false); + third.release(); + expect(successor.acquired).toBe(true); + + const usesInMemoryName = process.platform === "linux" || process.platform === "win32"; + expect(await Bun.file(lockPath).exists()).toBe(!usesInMemoryName); + + successor.release(); + const finalOwner = FileLock.tryAcquire(lockPath); + expect(finalOwner.acquired).toBe(true); + finalOwner.release(); + } finally { + await fs.rm(root, { recursive: true, force: true }); + } +}); diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index b64193cc3..05b4d70ef 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -51,24 +51,20 @@ import type { import { computeUsageWindowStats, fetchUsageSnapshots } from "./usage-windows"; const STATS_SYNC_LOCK_RETRY_MS = 25; -const STATS_SYNC_LOCK_ACQUIRE_STALE_MS = 10 * 1000; -const STATS_SYNC_LOCK_STALE_MS = 60 * 60 * 1000; +const STATS_SYNC_LOCK_WAIT_MS = 60 * 60 * 1000; /** * Serialize stats ingestion and archive reconciliation across processes. * The lock covers file discovery, parsing, and the final SQLite write so a * parse result for a session moved by GC can never commit after cleanup. - * The lock lives at `${dbPath}.sync.lock`; a dead owner is reclaimed - * immediately, a live-but-wedged one after an hour, and a lock abandoned - * mid-acquisition after ten seconds. + * The native lock is owned by an operating-system primitive, so an interrupted + * owner is released automatically and a live owner is never displaced. */ export async function withStatsSyncLock(dbPath: string, fn: () => Promise): Promise { await fs.promises.mkdir(path.dirname(dbPath), { recursive: true }); return await withFileLock(`${dbPath}.sync`, fn, { - staleMs: STATS_SYNC_LOCK_STALE_MS, - acquireStaleMs: STATS_SYNC_LOCK_ACQUIRE_STALE_MS, retryDelayMs: STATS_SYNC_LOCK_RETRY_MS, - retries: Math.ceil(STATS_SYNC_LOCK_STALE_MS / STATS_SYNC_LOCK_RETRY_MS), + retries: Math.ceil(STATS_SYNC_LOCK_WAIT_MS / STATS_SYNC_LOCK_RETRY_MS), }); } diff --git a/packages/stats/test/sync-serial.test.ts b/packages/stats/test/sync-serial.test.ts index 20a1a3f60..f759587e9 100644 --- a/packages/stats/test/sync-serial.test.ts +++ b/packages/stats/test/sync-serial.test.ts @@ -1,7 +1,7 @@ import { afterEach, describe, expect, it, vi } from "bun:test"; import * as fs from "node:fs/promises"; import * as path from "node:path"; -import { syncAllSessions, withStatsSyncLock } from "@oh-my-pi/omp-stats/aggregator"; +import { syncAllSessions } from "@oh-my-pi/omp-stats/aggregator"; import { getOverallStats } from "@oh-my-pi/omp-stats/db"; import { getSessionsDir } from "@oh-my-pi/pi-utils"; import { installStatsTestIsolation } from "./helpers/temp-agent"; @@ -93,71 +93,4 @@ describe("stats sync serial mode", () => { await expect(syncAllSessions({ workers: 2 })).rejects.toBe(workerProbe); expect(workerSpy).toHaveBeenCalled(); }); - - it("reclaims a dead owner's abandoned lock", async () => { - const dbPath = path.join(getSessionsDir(), "stats-lock.db"); - const lockPath = `${dbPath}.sync.lock`; - const deadPid = 424_242; - await fs.mkdir(lockPath, { recursive: true }); - await Bun.write( - path.join(lockPath, "info"), - JSON.stringify({ pid: deadPid, timestamp: Date.now() - 60_000, token: "dead-owner" }), - ); - vi.spyOn(process, "kill").mockImplementation(pid => { - if (pid === deadPid) throw Object.assign(new Error("dead owner"), { code: "ESRCH" }); - return true; - }); - - vi.spyOn(Bun, "sleep").mockResolvedValue(undefined); - const result = await withStatsSyncLock(dbPath, async () => "acquired"); - - expect(result).toBe("acquired"); - expect(await fs.stat(lockPath).catch(() => null)).toBeNull(); - }); - - it("reclaims an unstamped lock after the acquisition grace period", async () => { - const dbPath = path.join(getSessionsDir(), "stats-unstamped-lock.db"); - const lockPath = `${dbPath}.sync.lock`; - await fs.mkdir(lockPath, { recursive: true }); - const staleTime = new Date(Date.now() - 11_000); - await fs.utimes(lockPath, staleTime, staleTime); - vi.spyOn(Bun, "sleep").mockResolvedValue(undefined); - - const result = await withStatsSyncLock(dbPath, async () => "acquired"); - - expect(result).toBe("acquired"); - expect(await fs.stat(lockPath).catch(() => null)).toBeNull(); - }); - - it("waits for a live recently-stamped owner instead of reclaiming it", async () => { - const dbPath = path.join(getSessionsDir(), "stats-live-lock.db"); - const lockPath = `${dbPath}.sync.lock`; - const infoPath = path.join(lockPath, "info"); - await fs.mkdir(lockPath, { recursive: true }); - await Bun.write(infoPath, JSON.stringify({ pid: process.pid, timestamp: Date.now(), token: "live-owner" })); - const retryScheduled = Promise.withResolvers(); - const resumeRetry = Promise.withResolvers(); - let lockReleased = false; - vi.spyOn(Bun, "sleep").mockImplementation(async () => { - if (lockReleased) return; - retryScheduled.resolve(); - await resumeRetry.promise; - }); - - let acquired = false; - const pending = withStatsSyncLock(dbPath, async () => { - acquired = true; - }); - try { - await retryScheduled.promise; - expect(acquired).toBe(false); - expect(JSON.parse(await Bun.file(infoPath).text())).toMatchObject({ token: "live-owner" }); - } finally { - lockReleased = true; - await fs.rm(lockPath, { recursive: true, force: true }); - resumeRetry.resolve(); - await pending; - } - expect(acquired).toBe(true); - }); }); diff --git a/packages/utils/CHANGELOG.md b/packages/utils/CHANGELOG.md index 754af9e65..8cad11208 100644 --- a/packages/utils/CHANGELOG.md +++ b/packages/utils/CHANGELOG.md @@ -4,7 +4,7 @@ ### Added -- Added a shared `file-lock` utility module featuring configurable stale-lock acquisition grace periods (`acquireStaleMs`) to safely handle abandoned locks. +- Added a shared `file-lock` utility backed by process-owned native OS locks with automatic crash release and bounded asynchronous retry. ## [17.2.5] - 2026-08-03 diff --git a/packages/utils/src/file-lock.ts b/packages/utils/src/file-lock.ts index ba60d4454..261c0c99f 100644 --- a/packages/utils/src/file-lock.ts +++ b/packages/utils/src/file-lock.ts @@ -1,226 +1,66 @@ /** - * Cross-process advisory file lock shared by every package that serializes - * access to an on-disk resource (settings writes, MCP config, security store, - * stats sync). Locking is a `${filePath}.lock` directory — `mkdir` is atomic - * on every platform — holding an `info` file with the owner pid, acquisition - * timestamp, and a release token. - * - * Staleness: a lock is reclaimable when its owner process is gone, when a - * stamped lock is older than `staleMs` (wedged-owner recovery), or when an - * unstamped lock (owner crashed between `mkdir` and stamping) is older than - * `acquireStaleMs`. + * Cross-process advisory lock for packages that serialize access to an + * on-disk resource. The native handle is process-owned and automatically + * released on exit: Linux uses abstract Unix sockets, Windows uses named + * mutexes, and other Unix platforms use `flock(2)` on `${filePath}.lock`. */ -import { randomUUID } from "node:crypto"; -import * as fs from "node:fs/promises"; - -import { isEnoent } from "./fs-error"; -import * as logger from "./logger"; +import * as path from "node:path"; +import { FileLock as NativeFileLock } from "@oh-my-pi/pi-natives"; +/** Controls bounded waiting when an advisory file lock is contended. */ export interface FileLockOptions { - /** Age after which a stamped lock is reclaimable even if its owner is alive. */ - staleMs?: number; - /** Age after which an unstamped lock (crash mid-acquisition) is reclaimable; defaults to `staleMs`. */ - acquireStaleMs?: number; + /** Maximum acquisition attempts, including the initial attempt. */ retries?: number; + /** Delay between acquisition attempts. */ retryDelayMs?: number; } -const DEFAULT_OPTIONS: Required> = { - staleMs: 10_000, +const DEFAULT_OPTIONS: Required = { retries: 50, retryDelayMs: 100, }; -interface LockInfo { - pid: number; - timestamp: number; - token: string; -} - function getLockPath(filePath: string): string { - return `${filePath}.lock`; + return `${path.resolve(filePath)}.lock`; } -async function writeLockInfo(lockPath: string, token: string): Promise { - const info: LockInfo = { pid: process.pid, timestamp: Date.now(), token }; - await Bun.write(`${lockPath}/info`, JSON.stringify(info)); +function tryAcquireLock(lockPath: string): NativeFileLock | null { + const lock = NativeFileLock.tryAcquire(lockPath); + return lock.acquired ? lock : null; } -async function readLockInfo(lockPath: string): Promise { - try { - const content = await fs.readFile(`${lockPath}/info`, "utf-8"); - return JSON.parse(content) as LockInfo; - } catch { - return null; - } -} - -function isProcessAlive(pid: number): boolean { - try { - process.kill(pid, 0); - return true; - } catch (err) { - // EPERM means the pid exists but belongs to another user — still alive. - return (err as NodeJS.ErrnoException).code === "EPERM"; - } -} - -/** - * Identity of a lock artifact judged stale, pinning exactly what a reaper is - * allowed to remove: the stamped release token (or null while unstamped) and - * the lock directory's mtime. A re-acquired lock never matches, so a slow - * reaper cannot destroy a fresh owner's lock. - */ -interface StaleLockIdentity { - token: string | null; - mtimeMs: number; -} - -async function getStaleLockIdentity( - lockPath: string, - staleMs: number, - acquireStaleMs: number, -): Promise { - let mtimeMs: number; - try { - mtimeMs = (await fs.stat(lockPath)).mtimeMs; - } catch (err) { - if (isEnoent(err)) return null; - throw err; - } - const info = await readLockInfo(lockPath); - if (info) { - if (!isProcessAlive(info.pid)) return { token: info.token, mtimeMs }; - if (Date.now() - info.timestamp > staleMs) return { token: info.token, mtimeMs }; - return null; - } - - // No info file. Either the lock holder is between mkdir and writeLockInfo - // (fresh dir, do not reap) or the dir was already removed (also do not - // reap — there is nothing to clean up). - return Date.now() - mtimeMs > acquireStaleMs ? { token: null, mtimeMs } : null; -} - -/** - * Remove a stale lock without ever destroying a re-acquired one. The lock - * directory is atomically renamed into a unique graveyard path — only one - * concurrent reaper can win that claim — and deleted only if the claimed - * artifact still matches the identity that was judged stale; a mismatch - * (the lock was reaped and re-acquired in between) is renamed back. - */ -async function reapStaleLock(lockPath: string, expected: StaleLockIdentity): Promise { - const grave = `${lockPath}.reap-${randomUUID()}`; - try { - await fs.rename(lockPath, grave); - } catch (err) { - // Another reaper already claimed it (or the owner released) — done. - if (isEnoent(err)) return; - throw err; - } - let matches = false; - try { - const stat = await fs.stat(grave); - const info = await readLockInfo(grave); - matches = (info?.token ?? null) === expected.token && stat.mtimeMs === expected.mtimeMs; - } catch { - // Unreadable graveyard: treat as non-matching and restore below. - } - if (matches) { - await fs.rm(grave, { recursive: true, force: true }); - return; - } - // We claimed a lock that was re-acquired between judgment and rename — - // put it back. If a newer lock already took the path, drop the graveyard; - // the displaced owner's token-checked release degrades to a no-op. - try { - await fs.rename(grave, lockPath); - } catch { - logger.debug("file-lock: dropping displaced lock after failed restore", { lockPath, grave }); - await fs.rm(grave, { recursive: true, force: true }).catch(() => {}); - } -} - -async function tryAcquireLock(lockPath: string): Promise { - try { - await fs.mkdir(lockPath); - const token = randomUUID(); - await writeLockInfo(lockPath, token); - return token; - } catch (error) { - if ((error as NodeJS.ErrnoException).code === "EEXIST") { - return null; - } - throw error; - } -} - -async function releaseLock(lockPath: string, expectedToken: string): Promise { - try { - const info = await readLockInfo(lockPath); - if (!info || info.token !== expectedToken) { - // We are not the owner. The lock either expired and was reaped - // or another process has reclaimed it. Do nothing — releasing - // here would wipe the rightful owner's lock. - logger.debug("file-lock: skipping release for non-owned lock", { - lockPath, - expectedToken, - actualToken: info?.token, - }); - return; - } - await fs.rm(lockPath, { recursive: true }); - } catch { - // Ignore errors on release. - } -} - -async function acquireLock(filePath: string, options: FileLockOptions = {}): Promise<() => Promise> { +async function acquireLock(filePath: string, options: FileLockOptions = {}): Promise { const opts = { ...DEFAULT_OPTIONS, ...options }; - const acquireStaleMs = options.acquireStaleMs ?? opts.staleMs; const lockPath = getLockPath(filePath); for (let attempt = 0; attempt < opts.retries; attempt++) { - const token = await tryAcquireLock(lockPath); - if (token !== null) { - return () => releaseLock(lockPath, token); - } - - const stale = await getStaleLockIdentity(lockPath, opts.staleMs, acquireStaleMs); - if (stale) { - await reapStaleLock(lockPath, stale); - continue; - } - - await Bun.sleep(opts.retryDelayMs); + const lock = tryAcquireLock(lockPath); + if (lock) return lock; + if (attempt + 1 < opts.retries) await Bun.sleep(opts.retryDelayMs); } throw new Error(`Failed to acquire lock for ${filePath} after ${opts.retries} attempts`); } -/** Run `fn` while holding the exclusive `${filePath}.lock` advisory lock. */ +/** Run `fn` while holding an OS-backed exclusive lock for `filePath`. */ export async function withFileLock( filePath: string, fn: () => Promise, options: FileLockOptions = {}, ): Promise { - const release = await acquireLock(filePath, options); + const lock = await acquireLock(filePath, options); try { return await fn(); } finally { - await release(); + lock.release(); } } /** - * Test-only handles for the internal lock primitives. These are NOT part of - * the public API — they exist so the contract tests can validate token-keyed - * release semantics and the mkdir-race window without re-implementing them. + * Test-only acquisition handle for forcing ownership handoffs. This is not + * part of the supported package API. */ export const __internalsForTesting = { tryAcquireLock, - releaseLock, - readLockInfo, - getStaleLockIdentity, - reapStaleLock, getLockPath, }; diff --git a/packages/utils/test/file-lock.test.ts b/packages/utils/test/file-lock.test.ts index 47978f44a..a44d12963 100644 --- a/packages/utils/test/file-lock.test.ts +++ b/packages/utils/test/file-lock.test.ts @@ -3,10 +3,10 @@ import * as fs from "node:fs/promises"; import * as os from "node:os"; import * as path from "node:path"; import { __internalsForTesting, withFileLock } from "../src/file-lock"; +import { isEnoent } from "../src/fs-error"; import { removeWithRetries } from "../src/temp"; -const { tryAcquireLock, releaseLock, readLockInfo, getStaleLockIdentity, reapStaleLock, getLockPath } = - __internalsForTesting; +const { tryAcquireLock, getLockPath } = __internalsForTesting; const ROOTS: string[] = []; @@ -22,91 +22,74 @@ afterAll(async () => { } }); -describe("file-lock token ownership (F1)", () => { - test("releaseLock with the wrong token leaves the lock intact", async () => { +describe("native file-lock ownership", () => { + test("process death hands ownership to B while excluding C", async () => { const root = await mkRoot(); - const target = path.join(root, "data.json"); + const target = path.join(root, "abandoned.json"); + const readyPath = path.join(root, "holder-ready"); const lockPath = getLockPath(target); - - const token = await tryAcquireLock(lockPath); - expect(token).not.toBeNull(); - expect(typeof token).toBe("string"); - - // A contender that lost a race calling release with a guessed/empty token - // must NOT remove the rightful owner's lock. - await releaseLock(lockPath, "not-the-real-token"); - - const info = await readLockInfo(lockPath); - expect(info).not.toBeNull(); - expect(info?.token).toBe(token!); - - // The rightful owner can still release. - await releaseLock(lockPath, token!); - expect(await readLockInfo(lockPath)).toBeNull(); - }); - - test("getStaleLockIdentity does NOT declare a freshly-created empty dir stale", async () => { - const root = await mkRoot(); - const target = path.join(root, "race.json"); - const lockPath = getLockPath(target); - - // Simulate the precise window: mkdir succeeded for the winner but the - // info file has not been written yet. - await fs.mkdir(lockPath); - - const stale = await getStaleLockIdentity(lockPath, 10_000, 10_000); - expect(stale).toBeNull(); - - await removeWithRetries(lockPath); - }); - - test("a slow second reaper cannot destroy a reaped-and-reacquired lock", async () => { - const root = await mkRoot(); - const target = path.join(root, "contested.json"); - const lockPath = getLockPath(target); - - // A dead owner's stale lock, judged stale by two contenders. - await fs.mkdir(lockPath); - await Bun.write( - path.join(lockPath, "info"), - JSON.stringify({ pid: 999_999_999, timestamp: Date.now() - 60_000, token: "dead-token" }), + const holder = Bun.spawn( + [process.execPath, path.join(import.meta.dir, "fixtures/file-lock-holder.ts"), target, readyPath], + { + cwd: path.resolve(import.meta.dir, "../../.."), + env: { HOME: process.env.HOME ?? "", PATH: process.env.PATH ?? "" }, + stdin: "ignore", + stdout: "ignore", + stderr: "pipe", + }, ); - const judged = await getStaleLockIdentity(lockPath, 10_000, 10_000); - expect(judged).toMatchObject({ token: "dead-token" }); - // Reaper 1 wins: reaps and immediately re-acquires. - await reapStaleLock(lockPath, judged!); - const freshToken = await tryAcquireLock(lockPath); - expect(freshToken).not.toBeNull(); + try { + for (;;) { + try { + await fs.access(readyPath); + break; + } catch (error) { + if (!isEnoent(error)) throw error; + if (holder.exitCode !== null) { + throw new Error( + `lock holder exited before readiness (${holder.exitCode}): ${await new Response(holder.stderr).text()}`, + ); + } + } + } - // Reaper 2 acts on its stale pre-reap judgment: it must not remove the - // fresh owner's lock. - await reapStaleLock(lockPath, judged!); + holder.kill(); + expect(await holder.exited).not.toBe(0); - const info = await readLockInfo(lockPath); - expect(info?.token).toBe(freshToken!); + const ownerB = tryAcquireLock(lockPath); + if (!ownerB) throw new Error("B failed to acquire the abandoned lock"); + const ownerC = tryAcquireLock(lockPath); + expect(ownerC).toBeNull(); + expect(ownerB.acquired).toBe(true); + ownerB.release(); + } finally { + if (holder.exitCode === null) { + holder.kill(); + await holder.exited; + } + } + }, 10_000); - // The fresh owner can still release normally. - await releaseLock(lockPath, freshToken!); - expect(await readLockInfo(lockPath)).toBeNull(); - }); - - test("concurrent reapers of a vanished lock are a no-op", async () => { + test("a former owner's late release cannot unlock its successor", async () => { const root = await mkRoot(); - const target = path.join(root, "gone.json"); - const lockPath = getLockPath(target); + const lockPath = getLockPath(path.join(root, "handoff.json")); + const formerOwner = tryAcquireLock(lockPath); + if (!formerOwner) throw new Error("former owner failed to acquire"); + formerOwner.release(); - await fs.mkdir(lockPath); - await Bun.write( - path.join(lockPath, "info"), - JSON.stringify({ pid: 999_999_999, timestamp: Date.now() - 60_000, token: "dead-token" }), - ); - const judged = await getStaleLockIdentity(lockPath, 10_000, 10_000); - await reapStaleLock(lockPath, judged!); - // Second reap of the already-removed lock must not throw or recreate it. - await reapStaleLock(lockPath, judged!); - expect(await readLockInfo(lockPath)).toBeNull(); - expect(await fs.stat(lockPath).catch(() => null)).toBeNull(); + const successor = tryAcquireLock(lockPath); + if (!successor) throw new Error("successor failed to acquire"); + + // Force the old release path after the successor owns the same name. + formerOwner.release(); + expect(tryAcquireLock(lockPath)).toBeNull(); + expect(successor.acquired).toBe(true); + + successor.release(); + const finalOwner = tryAcquireLock(lockPath); + if (!finalOwner) throw new Error("final owner failed to acquire"); + finalOwner.release(); }); test("withFileLock serializes N concurrent writers without lost updates", async () => { @@ -123,9 +106,7 @@ describe("file-lock token ownership (F1)", () => { const text = await fs.readFile(target, "utf-8"); const data = JSON.parse(text) as { counter: number }; data.counter += 1; - // Widen the critical-section window so any concurrency leak - // surfaces as a lost update. - await Bun.sleep(2); + await Promise.resolve(); await fs.writeFile(target, JSON.stringify(data)); }, { retries: 500, retryDelayMs: 5 }, diff --git a/packages/utils/test/fixtures/file-lock-holder.ts b/packages/utils/test/fixtures/file-lock-holder.ts new file mode 100644 index 000000000..4282b18e5 --- /dev/null +++ b/packages/utils/test/fixtures/file-lock-holder.ts @@ -0,0 +1,14 @@ +import { withFileLock } from "../../src/file-lock"; + +const target = Bun.argv[2]; +const readyPath = Bun.argv[3]; +if (!target || !readyPath) throw new Error("file-lock-holder requires target and readiness paths"); + +await withFileLock( + target, + async () => { + await Bun.write(readyPath, "ready"); + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0); + }, + { retries: 1, retryDelayMs: 0 }, +);