Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion docs/shared-core.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)`
Expand Down
83 changes: 83 additions & 0 deletions docs/tdr/039-publish-complete-cross-process-write-locks.md
Original file line number Diff line number Diff line change
@@ -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`
1 change: 1 addition & 0 deletions docs/tdr/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
75 changes: 57 additions & 18 deletions src/services/turso/cross-process-write-lock.ts
Original file line number Diff line number Diff line change
@@ -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";

Expand Down Expand Up @@ -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<void> {
// 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<T>(
scope: "user" | "project",
hash: string,
Expand All @@ -71,30 +94,46 @@ export async function withCrossProcessWriteLock<T>(
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.
}
}
}
171 changes: 171 additions & 0 deletions tests/cross-process-write-lock.test.ts
Original file line number Diff line number Diff line change
@@ -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<typeof fs.writeFileSync>): 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<void>((resolve) => {
release = resolve;
});
let contender: Promise<void> | 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([]);
});
});
Loading