Skip to content

Commit b16bac1

Browse files
committed
fix(codex): serialize shared-credential OAuth token refreshes
Concurrent headless runs sharing one Codex credential raced refresh-token rotation: overlapping grants revoked each other and the loser surfaced a bare not-found rejection. Refresh-and-persist now serializes on a lock file next to the home-level auth store (exclusive create + in-memory tail for same-process callers, stale takeover for crashed holders), so parallel agents observe the persisted result instead of racing. Refresh failures (missing profile, rejected refresh, lock timeout) throw CodexAuthError with a re-login hint and project onto the shared credential_failure shape, so the run surfaces a credential failure, never a bare not-found error.
1 parent 9b3c101 commit b16bac1

7 files changed

Lines changed: 495 additions & 2 deletions

File tree

Lines changed: 144 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,144 @@
1+
import { mkdtemp, rm, stat, utimes, writeFile } from "node:fs/promises";
2+
import { tmpdir } from "node:os";
3+
import { join } from "node:path";
4+
import { describe, expect, test } from "bun:test";
5+
import {
6+
CodexRefreshLockTimeoutError,
7+
withCodexRefreshLock,
8+
} from "./refresh-lock.js";
9+
10+
async function tempDir(): Promise<string> {
11+
return mkdtemp(join(tmpdir(), "cl8628-lock-"));
12+
}
13+
14+
describe("codex refresh lock", () => {
15+
test("concurrent holders serialize and the lock file is removed", async () => {
16+
const dir = await tempDir();
17+
try {
18+
const lock = join(dir, "refresh.lock");
19+
let active = 0;
20+
let maxActive = 0;
21+
const results = await Promise.all(
22+
[0, 1, 2, 3, 4].map((i) =>
23+
withCodexRefreshLock(lock, async () => {
24+
active += 1;
25+
maxActive = Math.max(maxActive, active);
26+
// The file exists while held so a second process contends on it.
27+
await stat(lock);
28+
await new Promise((resolve) => setTimeout(resolve, 20));
29+
active -= 1;
30+
return i;
31+
}),
32+
),
33+
);
34+
expect(results).toEqual([0, 1, 2, 3, 4]);
35+
expect(maxActive).toBe(1);
36+
await expect(stat(lock)).rejects.toThrow();
37+
} finally {
38+
await rm(dir, { recursive: true, force: true });
39+
}
40+
});
41+
42+
test("a foreign-held lock times out with a recovery hint", async () => {
43+
const dir = await tempDir();
44+
try {
45+
const lock = join(dir, "refresh.lock");
46+
// Simulate a lock held by another process: withCodexRefreshLock never
47+
// created it, so only the file path (not the in-memory chain) applies.
48+
await writeFile(lock, "");
49+
let failure: unknown;
50+
try {
51+
await withCodexRefreshLock(lock, async () => "never", {
52+
timeoutMs: 100,
53+
retryMs: 10,
54+
});
55+
} catch (err) {
56+
failure = err;
57+
}
58+
expect(failure).toBeInstanceOf(CodexRefreshLockTimeoutError);
59+
expect((failure as CodexRefreshLockTimeoutError).lockPath).toBe(lock);
60+
expect((failure as CodexRefreshLockTimeoutError).message).toContain(
61+
"remove this lock file manually and retry",
62+
);
63+
} finally {
64+
await rm(dir, { recursive: true, force: true });
65+
}
66+
});
67+
68+
test("a stale lock from a crashed holder is taken over", async () => {
69+
const dir = await tempDir();
70+
try {
71+
const lock = join(dir, "refresh.lock");
72+
await writeFile(lock, "");
73+
const ancient = new Date(Date.now() - 60_000);
74+
await utimes(lock, ancient, ancient);
75+
const result = await withCodexRefreshLock(
76+
lock,
77+
async () => "taken-over",
78+
{
79+
timeoutMs: 5_000,
80+
staleMs: 1_000,
81+
},
82+
);
83+
expect(result).toBe("taken-over");
84+
await expect(stat(lock)).rejects.toThrow();
85+
} finally {
86+
await rm(dir, { recursive: true, force: true });
87+
}
88+
});
89+
90+
test("a lock held by another process blocks acquisition until released", async () => {
91+
const dir = await tempDir();
92+
const holderPath = new URL(
93+
"../../../tests/fixtures/codex-refresh-lock/hold-lock.ts",
94+
import.meta.url,
95+
).pathname;
96+
const lock = join(dir, "refresh.lock");
97+
const proc = Bun.spawn(["bun", "run", holderPath, lock, "1500"], {
98+
stdout: "pipe",
99+
stderr: "pipe",
100+
});
101+
try {
102+
// Wait until the holder process actually holds the file lock.
103+
if (proc.stdout === null) throw new Error("holder has no stdout pipe");
104+
const reader = proc.stdout.getReader();
105+
const decoder = new TextDecoder();
106+
let output = "";
107+
const deadline = Date.now() + 10_000;
108+
try {
109+
while (!output.includes("held")) {
110+
if (Date.now() > deadline)
111+
throw new Error("lock holder never acquired the lock");
112+
const { value, done } = await reader.read();
113+
if (done) break;
114+
output += decoder.decode(value, { stream: true });
115+
}
116+
} finally {
117+
reader.releaseLock();
118+
}
119+
expect(output).toContain("held");
120+
121+
// While the other process holds it, acquisition times out instead of
122+
// overlapping the grant.
123+
let failure: unknown;
124+
try {
125+
await withCodexRefreshLock(lock, async () => "never", {
126+
timeoutMs: 300,
127+
retryMs: 10,
128+
});
129+
} catch (err) {
130+
failure = err;
131+
}
132+
expect(failure).toBeInstanceOf(CodexRefreshLockTimeoutError);
133+
134+
// Once the holder exits and releases, the lock is acquirable again.
135+
await proc.exited;
136+
expect(await withCodexRefreshLock(lock, async () => "acquired")).toBe(
137+
"acquired",
138+
);
139+
} finally {
140+
proc.kill();
141+
await rm(dir, { recursive: true, force: true });
142+
}
143+
});
144+
});

‎src/auth/codex/refresh-lock.ts‎

Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,114 @@
1+
import { mkdir, open, stat, unlink } from "node:fs/promises";
2+
import { dirname } from "node:path";
3+
4+
// Serializes Codex OAuth refresh-and-persist sections across processes that
5+
// share one credential store (concurrent headless runs). Same-process
6+
// callers chain on an in-memory tail; cross-process callers contend on an
7+
// exclusive lock file. At most one refresh grant is ever in flight, so a
8+
// second refresher observes the persisted result instead of racing the
9+
// authorization server's refresh-token rotation and revoking its sibling.
10+
const tails = new Map<string, Promise<void>>();
11+
12+
// Default stale horizon: crashed-holder recovery only. Kept far above any
13+
// legitimate hold and any contender timeout so a live holder is never
14+
// declared stale mid-wait (see the staleMs field docs).
15+
const DEFAULT_STALE_MS = 120_000;
16+
17+
export interface CodexRefreshLockOptions {
18+
/** Maximum wait for the lock before giving up. Defaults to 30_000. */
19+
timeoutMs?: number;
20+
/** Poll interval while contending. Defaults to 25. */
21+
retryMs?: number;
22+
/**
23+
* Lock-file age at which the holder is presumed crashed and the lock is
24+
* taken over. Defaults to 120_000. This must stay far above the longest
25+
* legitimate hold (the token endpoint timeout plus store I/O) — and, in
26+
* particular, far above any contender's own timeout — otherwise a waiter
27+
* declares the live holder stale mid-wait, steals the lock, and the two
28+
* overlapping refresh grants revoke each other under token rotation.
29+
*/
30+
staleMs?: number;
31+
}
32+
33+
// The lock could not be acquired in time: either a refresh is genuinely
34+
// stalled past the timeout, or a crashed holder's lock survived takeover.
35+
// Carries the lock path so the message can tell the operator how to recover.
36+
export class CodexRefreshLockTimeoutError extends Error {
37+
readonly lockPath: string;
38+
39+
constructor(lockPath: string, timeoutMs: number) {
40+
super(
41+
`Timed out after ${String(timeoutMs)}ms waiting for the Codex refresh lock ` +
42+
`at ${lockPath}. If no refresh is running, remove this lock file manually and retry.`,
43+
);
44+
this.name = "CodexRefreshLockTimeoutError";
45+
this.lockPath = lockPath;
46+
}
47+
}
48+
49+
async function acquireLockFile(
50+
lockPath: string,
51+
timeoutMs: number,
52+
retryMs: number,
53+
staleMs: number,
54+
): Promise<void> {
55+
await mkdir(dirname(lockPath), { recursive: true, mode: 0o700 });
56+
const start = Date.now();
57+
for (;;) {
58+
try {
59+
const handle = await open(lockPath, "wx", 0o600);
60+
await handle.close();
61+
return;
62+
} catch (err) {
63+
if ((err as NodeJS.ErrnoException)?.code !== "EEXIST") throw err;
64+
}
65+
try {
66+
// A crashed holder never releases: once its lock is older than the
67+
// stale horizon, unlink and take over rather than brick refreshes.
68+
const info = await stat(lockPath);
69+
if (Date.now() - info.mtimeMs > staleMs) {
70+
await unlink(lockPath).catch(() => undefined);
71+
continue;
72+
}
73+
} catch {
74+
// The lock vanished (or was never stat-able) between attempts: loop
75+
// around and contend for it again.
76+
}
77+
if (Date.now() - start >= timeoutMs)
78+
throw new CodexRefreshLockTimeoutError(lockPath, timeoutMs);
79+
await new Promise((resolve) => setTimeout(resolve, retryMs));
80+
}
81+
}
82+
83+
/**
84+
* Runs `work` while holding the refresh lock at `lockPath`. Concurrent
85+
* same-process callers queue in memory; concurrent processes queue on the
86+
* lock file. The lock file is always removed afterwards.
87+
*/
88+
export async function withCodexRefreshLock<T>(
89+
lockPath: string,
90+
work: () => Promise<T>,
91+
options: CodexRefreshLockOptions = {},
92+
): Promise<T> {
93+
const timeoutMs = options.timeoutMs ?? 30_000;
94+
const retryMs = options.retryMs ?? 25;
95+
const staleMs = options.staleMs ?? DEFAULT_STALE_MS;
96+
const previous = tails.get(lockPath) ?? Promise.resolve();
97+
let releaseTail!: () => void;
98+
const current = new Promise<void>((resolve) => {
99+
releaseTail = resolve;
100+
});
101+
tails.set(lockPath, current);
102+
try {
103+
await previous.catch(() => undefined);
104+
await acquireLockFile(lockPath, timeoutMs, retryMs, staleMs);
105+
try {
106+
return await work();
107+
} finally {
108+
await unlink(lockPath).catch(() => undefined);
109+
}
110+
} finally {
111+
if (tails.get(lockPath) === current) tails.delete(lockPath);
112+
releaseTail();
113+
}
114+
}
Lines changed: 117 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,117 @@
1+
import { mkdtemp, rm } from "node:fs/promises";
2+
import { tmpdir } from "node:os";
3+
import { join } from "node:path";
4+
import { describe, expect, test } from "bun:test";
5+
import { saveCodexProfile } from "../../config/oauth-stores.js";
6+
import {
7+
codexAuthFailureDiagnostic,
8+
CodexAuthError,
9+
getValidCodexToken,
10+
} from "./session.js";
11+
12+
async function tempHome(): Promise<string> {
13+
return mkdtemp(join(tmpdir(), "cl8628-failure-"));
14+
}
15+
16+
describe("codex auth failure surface", () => {
17+
test("a missing profile rejects with a re-login hint, never a bare not-found", async () => {
18+
const home = await tempHome();
19+
try {
20+
let failure: unknown;
21+
try {
22+
await getValidCodexToken("ghost", Date.now(), home);
23+
} catch (err) {
24+
failure = err;
25+
}
26+
expect(failure).toBeInstanceOf(CodexAuthError);
27+
const auth = failure as CodexAuthError;
28+
expect(auth.reason).toBe("missing");
29+
expect(auth.message).toContain('"ghost"');
30+
expect(auth.message).toContain("Log in again");
31+
expect(auth.message).not.toContain("No OAuth profile named");
32+
} finally {
33+
await rm(home, { recursive: true, force: true });
34+
}
35+
});
36+
37+
test("a rejected refresh rejects with a re-login hint", async () => {
38+
const home = await tempHome();
39+
const now = Date.now();
40+
await saveCodexProfile(
41+
{
42+
name: "shared",
43+
createdAt: now,
44+
tokens: {
45+
access: "access-1",
46+
refresh: "refresh-1",
47+
expiresAt: now - 300_000,
48+
accountId: "acct-1",
49+
},
50+
},
51+
home,
52+
);
53+
const originalFetch = globalThis.fetch;
54+
globalThis.fetch = (async () =>
55+
new Response(JSON.stringify({ error: "invalid_grant" }), {
56+
status: 400,
57+
})) as unknown as typeof fetch;
58+
try {
59+
let failure: unknown;
60+
try {
61+
await getValidCodexToken("shared", now, home);
62+
} catch (err) {
63+
failure = err;
64+
}
65+
expect(failure).toBeInstanceOf(CodexAuthError);
66+
expect((failure as CodexAuthError).reason).toBe("refresh-failed");
67+
expect((failure as CodexAuthError).message).toContain("Log in again");
68+
} finally {
69+
globalThis.fetch = originalFetch;
70+
await rm(home, { recursive: true, force: true });
71+
}
72+
});
73+
74+
test("fresh tokens resolve without touching the network", async () => {
75+
const home = await tempHome();
76+
const now = Date.now();
77+
await saveCodexProfile(
78+
{
79+
name: "shared",
80+
createdAt: now,
81+
tokens: {
82+
access: "access-1",
83+
refresh: "refresh-1",
84+
expiresAt: now + 3_600_000,
85+
accountId: "acct-1",
86+
},
87+
},
88+
home,
89+
);
90+
const originalFetch = globalThis.fetch;
91+
globalThis.fetch = (async () => {
92+
throw new Error("network must not be touched for fresh tokens");
93+
}) as unknown as typeof fetch;
94+
try {
95+
const access = await getValidCodexToken("shared", now, home);
96+
expect(access.access).toBe("access-1");
97+
expect(access.accountId).toBe("acct-1");
98+
} finally {
99+
globalThis.fetch = originalFetch;
100+
await rm(home, { recursive: true, force: true });
101+
}
102+
});
103+
104+
test("codexAuthFailureDiagnostic projects auth errors onto credential_failure", () => {
105+
const auth = new CodexAuthError(
106+
"personal",
107+
"refresh-failed",
108+
'Codex profile "personal" could not be refreshed (boom). Log in again.',
109+
);
110+
expect(codexAuthFailureDiagnostic(auth)).toEqual({
111+
category: "credential_failure",
112+
message: auth.message,
113+
});
114+
expect(codexAuthFailureDiagnostic(new Error("boom"))).toBeNull();
115+
expect(codexAuthFailureDiagnostic("boom")).toBeNull();
116+
});
117+
});

0 commit comments

Comments
 (0)