feat: introduced native file lock bindings for cross-process advisory locking
- Added native `FileLock` bindings supporting cross-process advisory locking on Linux, Unix, and Windows. - Replaced directory-based file locking and custom stale-lock reclamation with OS-backed native locks. - Updated TypeScript declarations, native bindings, and package documentation for the new API. - Added comprehensive unit tests and fixtures validating single-owner constraints and process death handoff.
This commit is contained in:
Generated
+1
-1
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<UnixDatagram>,
|
||||
}
|
||||
|
||||
pub fn try_acquire(path: &str) -> io::Result<Option<PlatformFileLock>> {
|
||||
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(())
|
||||
}
|
||||
}
|
||||
@@ -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<platform::PlatformFileLock>,
|
||||
}
|
||||
|
||||
#[napi]
|
||||
impl FileLock {
|
||||
/// Try to acquire `path` without blocking.
|
||||
#[napi(factory)]
|
||||
pub fn try_acquire(path: String) -> napi::Result<Self> {
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<File>,
|
||||
}
|
||||
|
||||
pub fn try_acquire(path: &str) -> io::Result<Option<PlatformFileLock>> {
|
||||
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(())
|
||||
}
|
||||
}
|
||||
@@ -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<OwnedHandle>,
|
||||
}
|
||||
|
||||
pub fn try_acquire(path: &str) -> io::Result<Option<PlatformFileLock>> {
|
||||
let name = super::memory_lock_name(path);
|
||||
let wide_name: Vec<u16> = 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(())
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -34,7 +34,7 @@ const SECURITY_STORE_WRITE_CHAINS = new Map<string, Promise<unknown>>();
|
||||
async function withSecurityStoreWrite<T>(key: string, operation: () => Promise<T>): Promise<T> {
|
||||
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);
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Vendored
+16
@@ -61,6 +61,22 @@ export declare class DesktopSession {
|
||||
close(): Promise<undefined>
|
||||
}
|
||||
|
||||
/**
|
||||
* 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 {
|
||||
/**
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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 });
|
||||
}
|
||||
});
|
||||
@@ -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<T>(dbPath: string, fn: () => Promise<T>): Promise<T> {
|
||||
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),
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -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<void>();
|
||||
const resumeRetry = Promise.withResolvers<void>();
|
||||
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);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
+23
-183
@@ -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<Omit<FileLockOptions, "acquireStaleMs">> = {
|
||||
staleMs: 10_000,
|
||||
const DEFAULT_OPTIONS: Required<FileLockOptions> = {
|
||||
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<void> {
|
||||
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<LockInfo | null> {
|
||||
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<StaleLockIdentity | null> {
|
||||
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<void> {
|
||||
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<string | null> {
|
||||
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<void> {
|
||||
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<void>> {
|
||||
async function acquireLock(filePath: string, options: FileLockOptions = {}): Promise<NativeFileLock> {
|
||||
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<T>(
|
||||
filePath: string,
|
||||
fn: () => Promise<T>,
|
||||
options: FileLockOptions = {},
|
||||
): Promise<T> {
|
||||
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,
|
||||
};
|
||||
|
||||
@@ -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 },
|
||||
|
||||
@@ -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 },
|
||||
);
|
||||
Reference in New Issue
Block a user