From 9fb082bf1434b8df1797f416e7e54f34aaa23fda Mon Sep 17 00:00:00 2001 From: can1357 Date: Mon, 3 Aug 2026 15:21:10 +0200 Subject: [PATCH] refactor(utils): consolidated file locking into pi-utils file-lock - Moved the coding-agent lock-directory primitive to @oh-my-pi/pi-utils/file-lock and migrated settings, MCP config-writer, and security store imports. - Replaced the stats aggregator's parallel ~200-line token/breaker lock protocol with the shared primitive: dead owners reclaimed immediately, live-but-wedged owners after STATS_SYNC_LOCK_STALE_MS, unstamped acquisitions after the new acquireStaleMs grace (10s). - Shared primitive now treats EPERM kill probes as live owners. - Rewrote the stats lock-reclamation regressions against the shared protocol and moved the file-lock contract test into pi-utils. --- packages/coding-agent/src/config/file-lock.ts | 164 ------------- packages/coding-agent/src/config/settings.ts | 2 +- .../coding-agent/src/mcp/config-writer.ts | 2 +- packages/coding-agent/src/security/store.ts | 2 +- .../test/settings-manager.test.ts | 2 +- packages/stats/src/aggregator.ts | 201 +--------------- packages/stats/test/sync-serial.test.ts | 55 ++--- packages/utils/src/file-lock.ts | 226 ++++++++++++++++++ packages/utils/src/index.ts | 1 + .../test/file-lock.test.ts | 62 ++++- 10 files changed, 319 insertions(+), 398 deletions(-) delete mode 100644 packages/coding-agent/src/config/file-lock.ts create mode 100644 packages/utils/src/file-lock.ts rename packages/{coding-agent => utils}/test/file-lock.test.ts (51%) diff --git a/packages/coding-agent/src/config/file-lock.ts b/packages/coding-agent/src/config/file-lock.ts deleted file mode 100644 index b413407f7..000000000 --- a/packages/coding-agent/src/config/file-lock.ts +++ /dev/null @@ -1,164 +0,0 @@ -import { randomUUID } from "node:crypto"; -import * as fs from "node:fs/promises"; -import { isEnoent, logger } from "@oh-my-pi/pi-utils"; - -export interface FileLockOptions { - staleMs?: number; - retries?: number; - retryDelayMs?: number; -} - -const DEFAULT_OPTIONS: Required = { - staleMs: 10_000, - retries: 50, - retryDelayMs: 100, -}; - -interface LockInfo { - pid: number; - timestamp: number; - token: string; -} - -function getLockPath(filePath: string): string { - return `${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)); -} - -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 { - return false; - } -} - -async function isLockStale(lockPath: string, staleMs: number): Promise { - const info = await readLockInfo(lockPath); - if (info) { - if (!isProcessAlive(info.pid)) return true; - if (Date.now() - info.timestamp > staleMs) return true; - return false; - } - - // 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, and an unguarded fs.rm here would - // race with another contender's successful mkdir and wipe their dir). - try { - const stat = await fs.stat(lockPath); - return Date.now() - stat.mtimeMs > staleMs; - } catch (err) { - if (isEnoent(err)) return false; - throw err; - } -} - -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 { - if (expectedToken !== undefined) { - 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 lockExists(lockPath: string): Promise { - try { - await fs.stat(lockPath); - return true; - } catch (err) { - if (isEnoent(err)) return false; - throw err; - } -} - -async function acquireLock(filePath: string, options: FileLockOptions = {}): Promise<() => Promise> { - const opts = { ...DEFAULT_OPTIONS, ...options }; - const lockPath = getLockPath(filePath); - - for (let attempt = 0; attempt < opts.retries; attempt++) { - const token = await tryAcquireLock(lockPath); - if (token !== null) { - return () => releaseLock(lockPath, token); - } - - if ((await lockExists(lockPath)) && (await isLockStale(lockPath, opts.staleMs))) { - // Reaping a stale lock — no token because we didn't acquire it. The - // rightful owner is presumed dead; rm without ownership check. - await releaseLock(lockPath); - continue; - } - - await Bun.sleep(opts.retryDelayMs); - } - - throw new Error(`Failed to acquire lock for ${filePath} after ${opts.retries} attempts`); -} - -export async function withFileLock( - filePath: string, - fn: () => Promise, - options: FileLockOptions = {}, -): Promise { - const release = await acquireLock(filePath, options); - try { - return await fn(); - } finally { - await 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. - */ -export const __internalsForTesting = { - tryAcquireLock, - releaseLock, - readLockInfo, - isLockStale, - getLockPath, -}; diff --git a/packages/coding-agent/src/config/settings.ts b/packages/coding-agent/src/config/settings.ts index fe10e861b..09d4425b5 100644 --- a/packages/coding-agent/src/config/settings.ts +++ b/packages/coding-agent/src/config/settings.ts @@ -41,7 +41,7 @@ import { AUTO_IMAGE_PROVIDER_ORDER, isImageProviderId } from "../tools/image-pro import { type EditMode, normalizeEditMode } from "../utils/edit-mode"; import { INSPECT_IMAGE_MODES } from "../utils/inspect-image-mode"; import { isSearchProviderId, SEARCH_PROVIDER_ORDER } from "../web/search/types"; -import { withFileLock } from "./file-lock"; +import { withFileLock } from "@oh-my-pi/pi-utils/file-lock"; import { type BashInterceptorRule, type GroupPrefix, diff --git a/packages/coding-agent/src/mcp/config-writer.ts b/packages/coding-agent/src/mcp/config-writer.ts index 4d1855fc3..a7a95015c 100644 --- a/packages/coding-agent/src/mcp/config-writer.ts +++ b/packages/coding-agent/src/mcp/config-writer.ts @@ -8,7 +8,7 @@ import * as fs from "node:fs"; import * as path from "node:path"; import { isEnoent } from "@oh-my-pi/pi-utils"; import { invalidate as invalidateFsCache } from "../capability/fs"; -import { withFileLock } from "../config/file-lock"; +import { withFileLock } from "@oh-my-pi/pi-utils/file-lock"; import { validateServerConfig } from "./config"; import { MCP_CONFIG_SCHEMA_URL, type MCPConfigFile, type MCPServerConfig } from "./types"; diff --git a/packages/coding-agent/src/security/store.ts b/packages/coding-agent/src/security/store.ts index f253d452a..7321e6825 100644 --- a/packages/coding-agent/src/security/store.ts +++ b/packages/coding-agent/src/security/store.ts @@ -1,7 +1,7 @@ import * as fs from "node:fs/promises"; import * as path from "node:path"; import { getSecurityProjectDir, isEnoent } from "@oh-my-pi/pi-utils"; -import { withFileLock } from "../config/file-lock"; +import { withFileLock } from "@oh-my-pi/pi-utils/file-lock"; import * as git from "../utils/git"; import { compareSecurityLineage } from "./comparison"; import type { diff --git a/packages/coding-agent/test/settings-manager.test.ts b/packages/coding-agent/test/settings-manager.test.ts index 2f9bfd6db..884e8a9ce 100644 --- a/packages/coding-agent/test/settings-manager.test.ts +++ b/packages/coding-agent/test/settings-manager.test.ts @@ -20,7 +20,7 @@ import { AUTO_IMAGE_PROVIDER_ORDER } from "@oh-my-pi/pi-coding-agent/tools/image import { SEARCH_PROVIDER_ORDER } from "@oh-my-pi/pi-coding-agent/web/search/types"; import { getProjectAgentDir, TempDir } from "@oh-my-pi/pi-utils"; import { YAML } from "bun"; -import * as fileLock from "../src/config/file-lock"; +import * as fileLock from "@oh-my-pi/pi-utils/file-lock"; import { beginSettingsTest, restoreSettingsTestState, type SettingsTestState } from "./helpers/settings-test-state"; function context(): Context { diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index e831fe396..b64193cc3 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -1,6 +1,7 @@ import * as fs from "node:fs"; import * as path from "node:path"; import { getStatsDbPath, workerHostEntry } from "@oh-my-pi/pi-utils"; +import { withFileLock } from "@oh-my-pi/pi-utils/file-lock"; import { getRecentErrors as dbGetRecentErrors, getRecentRequests as dbGetRecentRequests, @@ -53,204 +54,22 @@ 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; -interface StatsLockSnapshot { - dev: number; - ino: number; - size: number; - mtimeMs: number; - text: string; -} - -function errorCode(error: unknown): string | undefined { - if (!error || typeof error !== "object" || !("code" in error)) return undefined; - return typeof error.code === "string" ? error.code : undefined; -} - -function sameStatsLock(left: StatsLockSnapshot, right: StatsLockSnapshot): boolean { - return ( - left.dev === right.dev && - left.ino === right.ino && - left.size === right.size && - left.mtimeMs === right.mtimeMs && - left.text === right.text - ); -} - -async function readStatsLockSnapshot(lockPath: string): Promise { - try { - const before = await fs.promises.stat(lockPath); - const text = await fs.promises.readFile(lockPath, "utf8"); - const after = await fs.promises.stat(lockPath); - const first = { - dev: before.dev, - ino: before.ino, - size: before.size, - mtimeMs: before.mtimeMs, - text, - }; - const second = { - dev: after.dev, - ino: after.ino, - size: after.size, - mtimeMs: after.mtimeMs, - text, - }; - return sameStatsLock(first, second) ? second : null; - } catch (error) { - if (errorCode(error) === "ENOENT") return null; - throw error; - } -} - -function statsLockOwnerIsRunning(snapshot: StatsLockSnapshot): boolean { - const pid = Number.parseInt(snapshot.text.split(/\r?\n/, 1)[0] ?? "", 10); - if (!Number.isSafeInteger(pid) || pid <= 0) { - return Date.now() - snapshot.mtimeMs < STATS_SYNC_LOCK_ACQUIRE_STALE_MS; - } - try { - process.kill(pid, 0); - return true; - } catch (error) { - return errorCode(error) !== "ESRCH"; - } -} - -function createStatsLockToken(): string { - return `${process.pid}\n${Date.now()}\n${Math.random()}\n`; -} - -async function removeAbandonedStatsLockFile(lockPath: string): Promise { - const snapshot = await readStatsLockSnapshot(lockPath); - if (!snapshot || statsLockOwnerIsRunning(snapshot)) return false; - - const current = await readStatsLockSnapshot(lockPath); - if (!current || !sameStatsLock(snapshot, current) || statsLockOwnerIsRunning(current)) return false; - try { - await fs.promises.unlink(lockPath); - return true; - } catch (error) { - if (errorCode(error) === "ENOENT") return false; - throw error; - } -} - -async function removeOwnedStatsLockFile(lockPath: string, owned: StatsLockSnapshot): Promise { - const current = await readStatsLockSnapshot(lockPath); - if (!current || !sameStatsLock(owned, current)) return; - try { - await fs.promises.unlink(lockPath); - } catch (error) { - if (errorCode(error) !== "ENOENT") throw error; - } -} - -async function removeAbandonedStatsLock(lockPath: string): Promise { - const snapshot = await readStatsLockSnapshot(lockPath); - if (!snapshot || statsLockOwnerIsRunning(snapshot)) return false; - - const breakerPath = `${lockPath}.break`; - let breaker: fs.promises.FileHandle; - try { - breaker = await fs.promises.open(breakerPath, "wx"); - } catch (error) { - if (errorCode(error) === "EEXIST") { - await removeAbandonedStatsLockFile(breakerPath); - return false; - } - throw error; - } - - const breakerToken = createStatsLockToken(); - let breakerSnapshot: StatsLockSnapshot | null = null; - let breakerTokenWritten = false; - let removed = false; - let cleanupError: unknown; - try { - await breaker.writeFile(breakerToken); - breakerTokenWritten = true; - breakerSnapshot = await readStatsLockSnapshot(breakerPath); - if (!breakerSnapshot || breakerSnapshot.text !== breakerToken) { - throw new Error(`Stats lock breaker changed while acquiring ${breakerPath}`); - } - - const current = await readStatsLockSnapshot(lockPath); - if (current && sameStatsLock(snapshot, current) && !statsLockOwnerIsRunning(current)) { - try { - await fs.promises.unlink(lockPath); - removed = true; - } catch (error) { - if (errorCode(error) !== "ENOENT") throw error; - } - } - } finally { - try { - await breaker.close(); - } catch (error) { - cleanupError = error; - } - try { - if (!breakerSnapshot && breakerTokenWritten) { - const current = await readStatsLockSnapshot(breakerPath); - if (current?.text === breakerToken) breakerSnapshot = current; - } - if (breakerSnapshot) await removeOwnedStatsLockFile(breakerPath, breakerSnapshot); - } catch (error) { - cleanupError ??= error; - } - } - if (cleanupError) throw cleanupError; - return removed; -} - /** * 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. */ export async function withStatsSyncLock(dbPath: string, fn: () => Promise): Promise { - const lockPath = `${dbPath}.sync.lock`; await fs.promises.mkdir(path.dirname(dbPath), { recursive: true }); - let handle: fs.promises.FileHandle | undefined; - while (!handle) { - try { - handle = await fs.promises.open(lockPath, "wx"); - } catch (error) { - if (errorCode(error) !== "EEXIST") throw error; - await removeAbandonedStatsLock(lockPath); - await Bun.sleep(STATS_SYNC_LOCK_RETRY_MS); - } - } - - const token = createStatsLockToken(); - let result: T | undefined; - let operationError: unknown; - let operationFailed = false; - let tokenWritten = false; - try { - await handle.writeFile(token); - tokenWritten = true; - result = await fn(); - } catch (error) { - operationFailed = true; - operationError = error; - } - - let cleanupError: unknown; - try { - await handle.close(); - } catch (error) { - cleanupError = error; - } - try { - const current = await fs.promises.readFile(lockPath, "utf8"); - if (!tokenWritten || current === token) await fs.promises.unlink(lockPath); - } catch (error) { - if (errorCode(error) !== "ENOENT") cleanupError ??= error; - } - - if (operationFailed) throw operationError; - if (cleanupError) throw cleanupError; - return result as T; + 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), + }); } /** diff --git a/packages/stats/test/sync-serial.test.ts b/packages/stats/test/sync-serial.test.ts index 063f82e6f..7b0b7e87c 100644 --- a/packages/stats/test/sync-serial.test.ts +++ b/packages/stats/test/sync-serial.test.ts @@ -94,17 +94,15 @@ describe("stats sync serial mode", () => { expect(workerSpy).toHaveBeenCalled(); }); - it("reclaims a dead owner's abandoned stale breaker", async () => { + it("reclaims a dead owner's abandoned lock", async () => { const dbPath = path.join(getSessionsDir(), "stats-lock.db"); const lockPath = `${dbPath}.sync.lock`; - const breakerPath = `${lockPath}.break`; const deadPid = 424_242; - await fs.mkdir(path.dirname(dbPath), { recursive: true }); - await Bun.write(lockPath, `${deadPid}\n1\nprimary\n`); - await Bun.write(breakerPath, `${deadPid}\n1\nbreaker\n`); - const staleTime = new Date(Date.now() - 2 * 60 * 60 * 1000); - await fs.utimes(lockPath, staleTime, staleTime); - await fs.utimes(breakerPath, staleTime, staleTime); + 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; @@ -114,15 +112,13 @@ describe("stats sync serial mode", () => { const result = await withStatsSyncLock(dbPath, async () => "acquired"); expect(result).toBe("acquired"); - expect(await Bun.file(lockPath).exists()).toBe(false); - expect(await Bun.file(breakerPath).exists()).toBe(false); + 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(path.dirname(dbPath), { recursive: true }); - await Bun.write(lockPath, ""); + 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); @@ -130,30 +126,23 @@ describe("stats sync serial mode", () => { const result = await withStatsSyncLock(dbPath, async () => "acquired"); expect(result).toBe("acquired"); - expect(await Bun.file(lockPath).exists()).toBe(false); + expect(await fs.stat(lockPath).catch(() => null)).toBeNull(); }); - it("never reclaims a live owner-stamped breaker even when its mtime is stale", async () => { - const dbPath = path.join(getSessionsDir(), "stats-live-breaker.db"); + 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 breakerPath = `${lockPath}.break`; - const deadPid = 424_243; - const liveBreakerToken = `${process.pid}\n1\nlive-breaker\n`; - await fs.mkdir(path.dirname(dbPath), { recursive: true }); - await Bun.write(lockPath, `${deadPid}\n1\nprimary\n`); - await Bun.write(breakerPath, liveBreakerToken); - const staleTime = new Date(Date.now() - 2 * 60 * 60 * 1000); - await fs.utimes(lockPath, staleTime, staleTime); - await fs.utimes(breakerPath, staleTime, staleTime); - vi.spyOn(process, "kill").mockImplementation(pid => { - if (pid === deadPid) throw Object.assign(new Error("dead owner"), { code: "ESRCH" }); - return true; - }); + 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 breakerReleased = false; + let lockReleased = false; vi.spyOn(Bun, "sleep").mockImplementation(async () => { - if (breakerReleased) return; + if (lockReleased) return; retryScheduled.resolve(); await resumeRetry.promise; }); @@ -165,10 +154,10 @@ describe("stats sync serial mode", () => { try { await retryScheduled.promise; expect(acquired).toBe(false); - expect(await Bun.file(breakerPath).text()).toBe(liveBreakerToken); + expect(JSON.parse(await Bun.file(infoPath).text())).toMatchObject({ token: "live-owner" }); } finally { - breakerReleased = true; - await fs.rm(breakerPath, { force: true }); + lockReleased = true; + await fs.rm(lockPath, { recursive: true, force: true }); resumeRetry.resolve(); await pending; } diff --git a/packages/utils/src/file-lock.ts b/packages/utils/src/file-lock.ts new file mode 100644 index 000000000..ba60d4454 --- /dev/null +++ b/packages/utils/src/file-lock.ts @@ -0,0 +1,226 @@ +/** + * 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`. + */ +import { randomUUID } from "node:crypto"; +import * as fs from "node:fs/promises"; + +import { isEnoent } from "./fs-error"; +import * as logger from "./logger"; + +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; + retries?: number; + retryDelayMs?: number; +} + +const DEFAULT_OPTIONS: Required> = { + staleMs: 10_000, + retries: 50, + retryDelayMs: 100, +}; + +interface LockInfo { + pid: number; + timestamp: number; + token: string; +} + +function getLockPath(filePath: string): string { + return `${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)); +} + +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> { + 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); + } + + throw new Error(`Failed to acquire lock for ${filePath} after ${opts.retries} attempts`); +} + +/** Run `fn` while holding the exclusive `${filePath}.lock` advisory lock. */ +export async function withFileLock( + filePath: string, + fn: () => Promise, + options: FileLockOptions = {}, +): Promise { + const release = await acquireLock(filePath, options); + try { + return await fn(); + } finally { + await 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. + */ +export const __internalsForTesting = { + tryAcquireLock, + releaseLock, + readLockInfo, + getStaleLockIdentity, + reapStaleLock, + getLockPath, +}; diff --git a/packages/utils/src/index.ts b/packages/utils/src/index.ts index 0c6726687..8a2ba3b8c 100644 --- a/packages/utils/src/index.ts +++ b/packages/utils/src/index.ts @@ -5,6 +5,7 @@ export * from "./color"; export * from "./dirs"; export * from "./env"; export * from "./fetch-retry"; +export * from "./file-lock"; export * from "./format"; export * from "./frontmatter"; export * from "./fs-error"; diff --git a/packages/coding-agent/test/file-lock.test.ts b/packages/utils/test/file-lock.test.ts similarity index 51% rename from packages/coding-agent/test/file-lock.test.ts rename to packages/utils/test/file-lock.test.ts index 5b37e30a8..47978f44a 100644 --- a/packages/coding-agent/test/file-lock.test.ts +++ b/packages/utils/test/file-lock.test.ts @@ -2,10 +2,11 @@ import { afterAll, describe, 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 { __internalsForTesting, withFileLock } from "@oh-my-pi/pi-coding-agent/config/file-lock"; -import { removeWithRetries } from "@oh-my-pi/pi-utils"; +import { __internalsForTesting, withFileLock } from "../src/file-lock"; +import { removeWithRetries } from "../src/temp"; -const { tryAcquireLock, releaseLock, readLockInfo, isLockStale, getLockPath } = __internalsForTesting; +const { tryAcquireLock, releaseLock, readLockInfo, getStaleLockIdentity, reapStaleLock, getLockPath } = + __internalsForTesting; const ROOTS: string[] = []; @@ -44,7 +45,7 @@ describe("file-lock token ownership (F1)", () => { expect(await readLockInfo(lockPath)).toBeNull(); }); - test("isLockStale does NOT declare a freshly-created empty dir stale", async () => { + 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); @@ -53,12 +54,61 @@ describe("file-lock token ownership (F1)", () => { // info file has not been written yet. await fs.mkdir(lockPath); - const stale = await isLockStale(lockPath, 10_000); - expect(stale).toBe(false); + 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 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(); + + // Reaper 2 acts on its stale pre-reap judgment: it must not remove the + // fresh owner's lock. + await reapStaleLock(lockPath, judged!); + + const info = await readLockInfo(lockPath); + expect(info?.token).toBe(freshToken!); + + // 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 () => { + const root = await mkRoot(); + const target = path.join(root, "gone.json"); + const lockPath = getLockPath(target); + + 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(); + }); + test("withFileLock serializes N concurrent writers without lost updates", async () => { const root = await mkRoot(); const target = path.join(root, "counter.json");