diff --git a/src/auth/store.test.ts b/src/auth/store.test.ts index e74b8fa7e..e6b344b43 100644 --- a/src/auth/store.test.ts +++ b/src/auth/store.test.ts @@ -1,10 +1,32 @@ import { describe, expect, test } from "bun:test"; -import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { + mkdir, + mkdtemp, + readFile, + rm, + utimes, + writeFile, +} from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { type } from "arktype"; -import { createAuthStore, type BaseTokens } from "./store.js"; +import { withMockedModule } from "../../tests/helpers/mock-module.js"; +import type { BaseTokens } from "./store.js"; + +let failingUnlinkPath: string | undefined; +await withMockedModule( + import.meta.resolve("node:fs/promises"), + (real: typeof import("node:fs/promises")) => ({ + ...real, + unlink: async (...args: Parameters) => { + if (args[0] === failingUnlinkPath) throw new Error("lock release failed"); + return real.unlink(...args); + }, + }), +); + +const { createAuthStore } = await import("./store.js"); type TestTokens = BaseTokens & { accountId?: string }; @@ -294,6 +316,64 @@ describe("createAuthStore", () => { } }); + test("surfaces a credential lock release failure after a successful write", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-release-error-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + }); + failingUnlinkPath = `${store.authPath(home)}.lock`; + + await expect( + store.saveProfile( + { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }, + home, + ), + ).rejects.toThrow("lock release failed"); + } finally { + failingUnlinkPath = undefined; + await rm(home, { recursive: true, force: true }); + } + }); + + test("does not mask a read-modify-write callback failure during release", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-error-precedence-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + }); + const profile = { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }; + await store.saveProfile(profile, home); + + const failingStore = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: (_value: unknown): _value is TestTokens => { + throw new Error("validator failed"); + }, + }); + failingUnlinkPath = `${store.authPath(home)}.lock`; + await expect(failingStore.saveProfile(profile, home)).rejects.toThrow( + "validator failed", + ); + } finally { + failingUnlinkPath = undefined; + await rm(home, { recursive: true, force: true }); + } + }); + test("releases the credential lock when a read-modify-write callback fails", async () => { const home = await mkdtemp(join(tmpdir(), "oauth-store-error-")); try { @@ -328,6 +408,200 @@ describe("createAuthStore", () => { } }); + test("takes over a dead-holder lock instead of timing out", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-takeover-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + }); + const lockPath = `${store.authPath(home)}.lock`; + await mkdir(join(home, TEST_SETTINGS_DIR), { recursive: true }); + // The maximum pid_t can never be a live holder: kill(pid, 0) answers + // ESRCH (or EINVAL), both of which read as dead. + await writeFile(lockPath, `${2_147_483_647}`, { mode: 0o600 }); + + const profile = { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }; + await expect(store.saveProfile(profile, home)).resolves.toBeUndefined(); + expect(await store.loadProfile("work", home)).toEqual(profile); + await expect(readFile(lockPath, "utf8")).rejects.toThrow("ENOENT"); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + test("takes over a dead-holder pid:counter claim instead of timing out", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-takeover-claim-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + }); + const lockPath = `${store.authPath(home)}.lock`; + await mkdir(join(home, TEST_SETTINGS_DIR), { recursive: true }); + // The NEW pid:counter claim format with a certainly-dead PID: kill(pid, 0) + // answers ESRCH (or EINVAL), both of which read as dead. The counter leg + // must not stop holderPid from reading the pid leg. + await writeFile(lockPath, `${2_147_483_647}:99`, { mode: 0o600 }); + + const profile = { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }; + await expect(store.saveProfile(profile, home)).resolves.toBeUndefined(); + expect(await store.loadProfile("work", home)).toEqual(profile); + await expect(readFile(lockPath, "utf8")).rejects.toThrow("ENOENT"); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + test("waits on a live-holder lock and times out without touching it", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-live-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + lockTimeoutMs: 100, + }); + const lockPath = `${store.authPath(home)}.lock`; + await mkdir(join(home, TEST_SETTINGS_DIR), { recursive: true }); + await writeFile(lockPath, `${process.pid}`, { mode: 0o600 }); + + await expect( + store.saveProfile( + { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }, + home, + ), + ).rejects.toThrow( + `Timed out waiting for OAuth credential lock ${lockPath}. ` + + "If no Corbits process is running, remove this lock file manually and retry.", + ); + expect(await readFile(lockPath, "utf8")).toBe(`${process.pid}`); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + test("waits on a live-holder pid:counter claim and times out without touching it", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-live-claim-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + lockTimeoutMs: 100, + }); + const lockPath = `${store.authPath(home)}.lock`; + await mkdir(join(home, TEST_SETTINGS_DIR), { recursive: true }); + // The NEW pid:counter claim format held by this live process. Never + // signal it; the waiter must read the pid leg as alive, time out, and + // leave the claim byte-identical. + const claim = `${process.pid}:42`; + await writeFile(lockPath, claim, { mode: 0o600 }); + + await expect( + store.saveProfile( + { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }, + home, + ), + ).rejects.toThrow( + `Timed out waiting for OAuth credential lock ${lockPath}. ` + + "If no Corbits process is running, remove this lock file manually and retry.", + ); + expect(await readFile(lockPath, "utf8")).toBe(claim); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + test("takes over a stale legacy lock but waits on a fresh one", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-legacy-")); + try { + const staleStore = createAuthStore({ + filename: "stale-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + }); + const staleLockPath = `${staleStore.authPath(home)}.lock`; + await mkdir(join(home, TEST_SETTINGS_DIR), { recursive: true }); + await writeFile(staleLockPath, "legacy-orphan", { mode: 0o600 }); + await utimes(staleLockPath, new Date(), new Date(Date.now() - 60_000)); + + const profile = { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }; + await expect( + staleStore.saveProfile(profile, home), + ).resolves.toBeUndefined(); + expect(await staleStore.loadProfile("work", home)).toEqual(profile); + + const freshStore = createAuthStore({ + filename: "fresh-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + lockTimeoutMs: 100, + }); + const freshLockPath = `${freshStore.authPath(home)}.lock`; + await writeFile(freshLockPath, "legacy-orphan", { mode: 0o600 }); + await expect(freshStore.saveProfile(profile, home)).rejects.toThrow( + "Timed out waiting for OAuth credential lock", + ); + expect(await readFile(freshLockPath, "utf8")).toBe("legacy-orphan"); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + + test("saves when a contended lock vanishes mid-wait", async () => { + const home = await mkdtemp(join(tmpdir(), "oauth-store-vanish-")); + try { + const store = createAuthStore({ + filename: "test-auth.json", + settingsDirName: TEST_SETTINGS_DIR, + isTokens: isTestTokens, + }); + const lockPath = `${store.authPath(home)}.lock`; + await mkdir(join(home, TEST_SETTINGS_DIR), { recursive: true }); + await writeFile(lockPath, `${2_147_483_647}`, { mode: 0o600 }); + + // Yank the stale lock out from under the waiter: whether the waiter + // observes the dead PID, an ENOENT read, or an ENOENT unlink, it must + // retry the exclusive create and land the save — never throw ENOENT. + const pending = store.saveProfile( + { + name: "work", + tokens: { access: "a", refresh: "r", expiresAt: 1 }, + createdAt: 1, + }, + home, + ); + await rm(lockPath, { force: true }); + await expect(pending).resolves.toBeUndefined(); + expect((await store.loadProfile("work", home))?.tokens.access).toBe("a"); + } finally { + await rm(home, { recursive: true, force: true }); + } + }); + test("fails closed with manual recovery guidance when an orphan lock exists", async () => { const home = await mkdtemp(join(tmpdir(), "oauth-store-orphan-")); try { diff --git a/src/auth/store.ts b/src/auth/store.ts index 1f587c99e..52ac0f9ee 100644 --- a/src/auth/store.ts +++ b/src/auth/store.ts @@ -3,6 +3,7 @@ import { open, readFile, rename, + stat, unlink, writeFile, } from "node:fs/promises"; @@ -60,10 +61,21 @@ interface AuthFile { const LOCK_RETRY_MS = 25; const LOCK_TIMEOUT_MS = 1_000; +// Age at which a lock file carrying no holder PID (a foreign writer, or a +// lock predating PID tagging) is presumed orphaned by a crashed holder and +// taken over. Stays far above LOCK_TIMEOUT_MS so a waiter never declares a +// live holder stale mid-wait; PID-tagged locks ignore this horizon. +const LOCK_STALE_MS = 5_000; + // Per-call unique temp (pid + counter). Matches mcp/auth-store — pid alone is not // unique per call if writeAuthFile ever overlaps in-process. let tmpWriteCounter = 0; +// Per-waiter unique lock claim (pid + counter). Two contenders never share a +// claim, so a steal re-read that still matches names the same file, and the +// post-create ownership check can tell our claim from a winner's. +let lockClaimCounter = 0; + // Same-process ops on one auth file queue here so a caller's lock deadline // starts when it actually runs, not when it was invoked — otherwise one lock // held past LOCK_TIMEOUT_MS fails the whole burst, not just the first waiter. @@ -97,6 +109,77 @@ function isProfile( return isTokens(parsed.tokens); } +function holderPid(content: string): number | null { + const pid = Number(content.split(":")[0]?.trim()); + return Number.isInteger(pid) && pid > 0 ? pid : null; +} + +// kill(pid, 0) liveness: success means the process exists (alive); EPERM +// means it exists but belongs to another user (alive); ESRCH/EINVAL mean no +// such process (dead). Any other failure reads as alive — never steal a live +// holder's lock on a confused signal check; the waiter times out with a +// recovery hint instead. +function isPidAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch (err) { + const code = (err as NodeJS.ErrnoException)?.code; + return code !== "ESRCH" && code !== "EINVAL"; + } +} + +// Null when the lock vanished under the waiter (a release raced the read); +// anything but ENOENT propagates — permission and disk errors must surface, +// not read as an empty lock. +async function readLockContent(lockPath: string): Promise { + try { + return await readFile(lockPath, "utf8"); + } catch (err) { + if (isErrnoCode(err, "ENOENT")) return null; + throw err; + } +} + +// True when the observed lock is safe to take over: a tagged holder whose +// PID is dead (crashed — reachable under any timeout), or a legacy untagged +// lock older than the stale horizon. A vanished lock reads as not stale; +// the acquire loop retries the exclusive create instead. +async function isLockStale( + lockPath: string, + content: string, +): Promise { + const pid = holderPid(content); + if (pid !== null) return !isPidAlive(pid); + try { + const info = await stat(lockPath); + return Date.now() - info.mtimeMs > LOCK_STALE_MS; + } catch (err) { + if (isErrnoCode(err, "ENOENT")) return false; + throw err; + } +} + +function lockTimeoutError(lockPath: string, cause: unknown): Error { + return new Error( + `Timed out waiting for OAuth credential lock ${lockPath}. ` + + "If no Corbits process is running, remove this lock file manually and retry.", + { cause }, + ); +} + +// Pace one contention round: throw once the deadline passed, else sleep a +// retry interval. Every wait path funnels here so steal contention never +// hot-spins. +async function paceLockWait( + deadline: number, + lockPath: string, + cause: unknown, +): Promise { + if (Date.now() >= deadline) throw lockTimeoutError(lockPath, cause); + await delay(LOCK_RETRY_MS); +} + export function createAuthStore( options: AuthStoreOptions, ): AuthStore { @@ -149,34 +232,107 @@ export function createAuthStore( const lockPath = `${path}.lock`; await mkdir(dirname(path), { recursive: true, mode: 0o700 }); const deadline = Date.now() + (options.lockTimeoutMs ?? LOCK_TIMEOUT_MS); + // holderPid reads the pid leg; the counter leg keeps every waiter's + // claim distinct. + const claim = `${process.pid}:${(lockClaimCounter += 1)}`; let lock; while (true) { try { - lock = await open(lockPath, "wx", 0o600); + const handle = await open(lockPath, "wx", 0o600); + try { + await handle.writeFile(claim, "utf8"); + } catch (writeError) { + try { + await handle.close(); + } catch { + // Ignore close errors on the cleanup path; the write error below + // is the one the caller must see. + } + try { + await unlink(lockPath); + } catch (unlinkError) { + if (!isErrnoCode(unlinkError, "ENOENT")) throw unlinkError; + } + throw writeError; + } + // A steal may have unlinked our fresh file and created its own + // between our create and write; the path then names a live + // winner. Never unlink here — close and re-contend so only the + // winner proceeds. + if ((await readLockContent(lockPath)) !== claim) { + try { + await handle.close(); + } catch { + // The path already names a live winner; the close outcome + // must not mask the paced retry below. + } + await paceLockWait(deadline, lockPath, undefined); + continue; + } + lock = handle; break; } catch (error) { if (!isErrnoCode(error, "EEXIST")) throw error; - if (Date.now() >= deadline) { - throw new Error( - `Timed out waiting for OAuth credential lock ${lockPath}. ` + - "If no Corbits process is running, remove this lock file manually and retry.", - { cause: error }, - ); + // A crashed holder never releases: take over a stale lock rather + // than brick the store. A lock that vanished under the read + // (a release raced us) is not stale — retry the exclusive create. + const content = await readLockContent(lockPath); + if (content !== null && (await isLockStale(lockPath, content))) { + // Re-check before unlinking so a concurrent takeover winner's + // fresh claim is never mistaken for the stale entry just + // observed — claims are unique per waiter, so a match still + // names the same file. A steal can still interleave between + // this re-read and the unlink; the post-create ownership + // check above then detects the loser and re-contends instead + // of running two holders. + if ((await readLockContent(lockPath)) === content) { + try { + await unlink(lockPath); + } catch (unlinkError) { + // A concurrent winner unlinked first; retry the create. + if (!isErrnoCode(unlinkError, "ENOENT")) throw unlinkError; + } + } + // Fall through to the deadline/sleep path: a steal that just + // lost to a concurrent winner paces like any other + // contention instead of hot-spinning. } - await delay(LOCK_RETRY_MS); + await paceLockWait(deadline, lockPath, error); } } + let callbackOutcome: + | { ok: true; value: TResult } + | { ok: false; error: unknown }; try { - return await callback(); - } finally { - try { - await lock.close(); - } finally { - await unlink(lockPath); + callbackOutcome = { ok: true, value: await callback() }; + } catch (error) { + callbackOutcome = { ok: false, error }; + } + + try { + await lock.close(); + } catch { + // The callback outcome owns precedence; a close failure must not mask it. + // The unlink below still runs. + } + + let releaseOutcome: { ok: true } | { ok: false; error: unknown } = { + ok: true, + }; + try { + await unlink(lockPath); + } catch (error) { + // A stale-takeover steal may legitimately remove the file first. + if (!isErrnoCode(error, "ENOENT")) { + releaseOutcome = { ok: false, error }; } } + + if (!callbackOutcome.ok) throw callbackOutcome.error; + if (!releaseOutcome.ok) throw releaseOutcome.error; + return callbackOutcome.value; } function enqueueAuthFileOp(