diff --git a/docs/shared-core.md b/docs/shared-core.md index 87d377a8..9f7b0255 100644 --- a/docs/shared-core.md +++ b/docs/shared-core.md @@ -154,7 +154,9 @@ processes against the same store: - Every connection sets `busy_timeout=5000` and uses one pooled handle. - Every write goes through `withScopeWriteLock`. It nests a cross-process lock for each scope (`src/services/turso/cross-process-write-lock.ts`). - That lock checks if the owner PID is alive and takes over stale locks. + That lock publishes a complete PID payload through a hard link before + entering the write section. It checks if the owner PID is alive and takes + over stale locks. See [TDR-039](tdr/039-publish-complete-cross-process-write-locks.md). - So shard allocation, vector-count sync, insert, and increment run as one critical section for each scope, across processes. - Races to create a shard end at the `UNIQUE(scope, scope_hash, shard_index)` diff --git a/docs/tdr/039-publish-complete-cross-process-write-locks.md b/docs/tdr/039-publish-complete-cross-process-write-locks.md new file mode 100644 index 00000000..8b56824d --- /dev/null +++ b/docs/tdr/039-publish-complete-cross-process-write-locks.md @@ -0,0 +1,83 @@ +# TDR-039: Publish complete cross-process write locks + +- **Date:** 2026-10-07 +- **Status:** Proposed +- **Deciders:** OMMS maintainers +- **Tags:** storage, concurrency, release, filesystem + +## Context + +Release 4.13.2 failed on macOS 26 in `tests/two-process-storage.test.ts:258`. +All 50 memory rows survived, but shard metadata totalled 48 vectors. The other +five platforms passed. Publishing was skipped. + +### Root Cause Analysis + +`writeFileSync(path, payload, { flag: "wx" })` creates an empty file before +writing its process ID (PID). Another process can read that empty file. +`readLiveLock` then treats it as corrupt, deletes it, and acquires the lock. +The first process finishes writing through its open handle and also enters. +Overlapping writes can overwrite the reconciled shard count. + +The regression test forces this interleaving through the real filesystem. It +observed two active writer callbacks on the old code, where one was expected. +The CI log does not record the exact filesystem interleaving. + +## Decision + +Write each complete payload to a unique temporary file in `.write-locks`. +Use `linkSync(candidate, path)` to publish it. The hard link creates the lock +without replacing an existing one. A contender sees a complete payload or no +lock. Keep the existing lock path, JSON fields, dead-owner recovery, and +15-second contention deadline. + +Remove the temporary file in `finally`. Remove the lock only after acquisition. +Report publication errors other than `EEXIST` to the caller. On Windows, retry +`EPERM`, `EACCES`, and `EBUSY` with the bounded waits from TDR-034: 1, 2, 5, 10, +20, 50, 100, and 200 milliseconds. Persistent errors still fail. + +## Consequences + +### Positive + +- A contender cannot reclaim a new lock while its payload is being written. +- Deterministic tests cover exclusivity, error cleanup, Windows retries, and dead-owner recovery. + +### Negative + +- The storage filesystem must support hard links, as APFS, ext4, and NTFS do. +- A process killed before cleanup can leave an unused temporary file. + +### Neutral + +- All hosts receive the fix through the shared storage code. Memory data stays unchanged. +- This change repairs publication. The existing stale-lock takeover logic stays unchanged. + +## Alternatives Considered + +| Option | Rejected because | +| --------------------------------------- | --------------------------------------------------------- | +| Retry the failed smoke test | Leaves overlapping writers possible. | +| Rename a prepared payload over the lock | Can replace another process's live lock. | +| Delay before removing an empty lock | Depends on how long the operating system pauses a writer. | + +## How to Recognise / Handle This Again + +1. Compare readable memory rows with the shard metadata counts. +2. Run `bash scripts/run-tests-isolated.sh tests/cross-process-write-lock.test.ts tests/two-process-storage.test.ts`. +3. Run `bun run ci:local` in the fix worktree before pushing. +4. Run the six-platform smoke on the fix branch before merging, with maintainer approval. + +## Revisit Triggers + +Reassess if counts still drift, stale-owner takeover permits overlapping writers, +or storage moves to a filesystem without hard links. + +## References + +- [Failed release run](https://github.com/cmdaltctr/omms/actions/runs/37684125173) +- [TDR-034: Windows hard-link retries](034-retry-windows-start-lock-hard-links.md) +- [Shared-store specification](../../openspec/specs/host-neutral-memory-core/spec.md) +- [Node filesystem API](https://nodejs.org/docs/latest-v24.x/api/fs.html#fslinksyncexistingpath-newpath) +- `src/services/turso/cross-process-write-lock.ts` +- `tests/cross-process-write-lock.test.ts` diff --git a/docs/tdr/README.md b/docs/tdr/README.md index eebf94ed..e10d59d2 100644 --- a/docs/tdr/README.md +++ b/docs/tdr/README.md @@ -60,6 +60,7 @@ TDRs capture **implementation-level technical decisions** such as platform-speci | [036](./036-restart-waits-for-its-own-copy.md) | Restart waits for its own copy | Proposed | 2026-10-06 | | [037](./037-skip-slow-tests-on-windows.md) | Skip slow tests on Windows | Proposed | 2026-10-07 | | [038](./038-resume-release-approval-after-npm.md) | Resume the release approval after npm shows the version | Proposed | 2026-10-07 | +| [039](./039-publish-complete-cross-process-write-locks.md) | Publish complete cross-process write locks | Proposed | 2026-10-07 | ## Status values diff --git a/src/services/turso/cross-process-write-lock.ts b/src/services/turso/cross-process-write-lock.ts index c86ab7c2..be07f49c 100644 --- a/src/services/turso/cross-process-write-lock.ts +++ b/src/services/turso/cross-process-write-lock.ts @@ -1,4 +1,5 @@ -import { existsSync, mkdirSync, readFileSync, unlinkSync, writeFileSync } from "node:fs"; +import { randomUUID } from "node:crypto"; +import { existsSync, linkSync, mkdirSync, readFileSync, unlinkSync, writeFileSync } from "node:fs"; import { join } from "node:path"; import { CONFIG } from "../../config.js"; @@ -58,6 +59,28 @@ function lockFilePath(scope: string, hash: string): string { return join(CONFIG.storagePath, WRITE_LOCK_DIR, `${scope}_${hash}.lock`); } +async function publishLock(candidate: string, path: string): Promise { + // Match the bounded Windows hard-link retries used by the web start lock (TDR-034). + const waits = [1, 2, 5, 10, 20, 50, 100, 200]; + for (let attempt = 0; ; attempt++) { + try { + linkSync(candidate, path); + return; + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if ( + process.platform !== "win32" || + !code || + !["EPERM", "EACCES", "EBUSY"].includes(code) || + waits[attempt] === undefined + ) { + throw error; + } + await new Promise((resolve) => setTimeout(resolve, waits[attempt])); + } + } +} + export async function withCrossProcessWriteLock( scope: "user" | "project", hash: string, @@ -71,30 +94,46 @@ export async function withCrossProcessWriteLock( const state: LockState = { pid: process.pid, timestamp: new Date().toISOString() }; const deadline = Date.now() + LOCK_TIMEOUT_MS; - for (;;) { - try { - writeFileSync(path, JSON.stringify(state), { flag: "wx" }); - break; - } catch { - const holder = readLiveLock(path); - if (!holder) continue; - if (Date.now() > deadline) { - throw new Error( - `Timed out acquiring the cross-process write lock for ${scope}/${hash}: ` + - `held by pid ${holder.pid} since ${holder.timestamp}` - ); + const candidate = `${path}.${process.pid}.${randomUUID()}.tmp`; + let acquired = false; + + try { + // Exclusive open exposes an empty file before its PID is written. Publish + // the complete payload with a hard link, which cannot replace a live lock. + writeFileSync(candidate, JSON.stringify(state), { flag: "wx" }); + for (;;) { + try { + await publishLock(candidate, path); + acquired = true; + break; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; + const holder = readLiveLock(path); + if (!holder) continue; + if (Date.now() > deadline) { + throw new Error( + `Timed out acquiring the cross-process write lock for ${scope}/${hash}: ` + + `held by pid ${holder.pid} since ${holder.timestamp}`, + { cause: error } + ); + } + await new Promise((resolve) => setTimeout(resolve, LOCK_POLL_MS + Math.random() * 10)); } - await new Promise((resolve) => setTimeout(resolve, LOCK_POLL_MS + Math.random() * 10)); } - } - try { return await fn(); } finally { + if (acquired) { + try { + unlinkSync(path); + } catch { + // Best effort: a stale file is cleaned up by the next acquirer. + } + } try { - unlinkSync(path); + unlinkSync(candidate); } catch { - // Best effort: a stale file is cleaned up by the next acquirer. + // A failed publication must not leave its temporary payload behind. } } } diff --git a/tests/cross-process-write-lock.test.ts b/tests/cross-process-write-lock.test.ts new file mode 100644 index 00000000..6cce6b97 --- /dev/null +++ b/tests/cross-process-write-lock.test.ts @@ -0,0 +1,171 @@ +import { afterEach, describe, expect, it, mock } from "bun:test"; +import * as fs from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +const realFs = { ...fs }; +let duringExclusiveCreate: (() => void) | undefined; +let linkErrors: NodeJS.ErrnoException[] = []; +const platformDescriptor = Object.getOwnPropertyDescriptor(process, "platform")!; + +function linkSync(from: fs.PathLike, to: fs.PathLike): void { + const error = linkErrors.shift(); + if (error) throw error; + realFs.linkSync(from, to); +} + +function writeFileSync(...args: Parameters): void { + const options = args[2]; + if (duringExclusiveCreate && typeof options === "object" && options?.flag === "wx") { + // Pause after open, before writing the PID, as another process can do. + const fd = realFs.openSync(args[0], "wx"); + const contender = duringExclusiveCreate; + duringExclusiveCreate = undefined; + try { + contender(); + realFs.writeFileSync(fd, args[1], options); + } finally { + realFs.closeSync(fd); + } + return; + } + realFs.writeFileSync(...args); +} + +mock.module("node:fs", () => ({ + ...realFs, + default: { ...realFs, writeFileSync, linkSync }, + writeFileSync, + linkSync, +})); + +const { CONFIG } = await import("../src/config.js"); +const { withCrossProcessWriteLock } = + await import("../src/services/turso/cross-process-write-lock.js"); +const previousStorage = CONFIG.storagePath; +const directories: string[] = []; +const hash = "a1b2c3d4e5f60718"; + +function storage(): string { + const dir = realFs.mkdtempSync(join(tmpdir(), "omms-write-lock-")); + directories.push(dir); + CONFIG.storagePath = dir; + return dir; +} + +function lockPath(dir: string): string { + return join(dir, ".write-locks", `project_${hash}.lock`); +} + +afterEach(() => { + duringExclusiveCreate = undefined; + linkErrors = []; + Object.defineProperty(process, "platform", platformDescriptor); + CONFIG.storagePath = previousStorage; + for (const dir of directories.splice(0)) realFs.rmSync(dir, { recursive: true, force: true }); +}); + +describe("cross-process write lock", () => { + it("keeps writers exclusive when a contender arrives before the PID is written", async () => { + const dir = storage(); + let active = 0; + let maximumActive = 0; + let release!: () => void; + const held = new Promise((resolve) => { + release = resolve; + }); + let contender: Promise | undefined; + + duringExclusiveCreate = () => { + contender = withCrossProcessWriteLock("project", hash, async () => { + active++; + maximumActive = Math.max(maximumActive, active); + await held; + active--; + }); + }; + const first = withCrossProcessWriteLock("project", hash, async () => { + active++; + maximumActive = Math.max(maximumActive, active); + active--; + }); + try { + await Promise.resolve(); + } finally { + release(); + await Promise.all([first, contender]); + } + + expect(contender).toBeDefined(); + expect(maximumActive).toBe(1); + expect(realFs.existsSync(lockPath(dir))).toBe(false); + expect(realFs.readdirSync(join(dir, ".write-locks"))).toEqual([]); + }); + + it("releases the lock after the writer throws", async () => { + const dir = storage(); + await expect( + withCrossProcessWriteLock("project", hash, async () => { + throw new Error("write failed"); + }) + ).rejects.toThrow("write failed"); + expect(realFs.existsSync(lockPath(dir))).toBe(false); + await expect(withCrossProcessWriteLock("project", hash, async () => "next")).resolves.toBe( + "next" + ); + expect(realFs.readdirSync(join(dir, ".write-locks"))).toEqual([]); + }); + + it("reports a publication error without entering the writer or leaving files", async () => { + const dir = storage(); + const error = Object.assign(new Error("disk error"), { code: "EIO" }); + linkErrors = [error]; + let entered = false; + await expect( + withCrossProcessWriteLock("project", hash, async () => { + entered = true; + }) + ).rejects.toThrow("disk error"); + expect(entered).toBe(false); + expect(realFs.readdirSync(join(dir, ".write-locks"))).toEqual([]); + }); + + it("retries temporary Windows refusals while publishing the complete lock", async () => { + const dir = storage(); + Object.defineProperty(process, "platform", { value: "win32" }); + linkErrors = ["EPERM", "EACCES", "EBUSY"].map((code) => + Object.assign(new Error("busy"), { code }) + ); + await expect(withCrossProcessWriteLock("project", hash, async () => "written")).resolves.toBe( + "written" + ); + expect(linkErrors).toHaveLength(0); + expect(realFs.readdirSync(join(dir, ".write-locks"))).toEqual([]); + }); + + it("bounds persistent Windows publication refusals and cleans up", async () => { + const dir = storage(); + Object.defineProperty(process, "platform", { value: "win32" }); + linkErrors = Array.from({ length: 9 }, () => + Object.assign(new Error("permission refused"), { code: "EPERM" }) + ); + await expect( + withCrossProcessWriteLock("project", hash, async () => "unreachable") + ).rejects.toThrow("permission refused"); + expect(linkErrors).toHaveLength(0); + expect(realFs.readdirSync(join(dir, ".write-locks"))).toEqual([]); + }); + + it("reclaims a lock whose owner process is dead", async () => { + const dir = storage(); + realFs.mkdirSync(join(dir, ".write-locks")); + realFs.writeFileSync( + lockPath(dir), + JSON.stringify({ pid: 2147483647, timestamp: new Date().toISOString() }) + ); + await expect(withCrossProcessWriteLock("project", hash, async () => "recovered")).resolves.toBe( + "recovered" + ); + expect(realFs.readdirSync(join(dir, ".write-locks"))).toEqual([]); + }); +});