Compare commits

..
Author SHA1 Message Date
taodeng 90797864d5 fix(server): re-resolve -p spawn OAuth token per-spawn, expiry-aware (③ regression)
The FIX-③ spawn-home isolation memoized the OAuth token at startup. The macOS
keychain access token rotates (~hourly, refreshed by the operator's real claude),
so the startup snapshot went stale and every isolated -p spawn returned upstream
401 'Invalid authentication credentials' — a ~31h Mac-mini outage (PI231/oracle use
static long-lived env tokens, unaffected).

Fix: getSpawnHomeMode() now caches only the isolation DECISION; the token is
re-resolved FRESH per spawn via resolveSpawnToken(), which also returns null when a
known expiry has passed (5-min buffer) so the caller falls back to real HOME — where
the spawned claude refreshes the credential natively and self-heals. OCP still never
refreshes the token itself (a refresh-token grant would consume the single-use token
and log out the operator's real claude — issue #112). Env-token hosts carry no
expiresAt and are never expiry-gated. Infra/process change; no cli.js surface, no new
endpoint/header. 238 tests pass; live-verified sonnet 200 on a temp instance.
2026-06-26 20:33:59 +10:00
9 changed files with 148 additions and 1068 deletions
-10
View File
@@ -1,15 +1,5 @@
# Changelog # Changelog
## v3.21.1 — 2026-07-07
Patch release: three bug fixes from an independent concurrency/session-lifecycle audit, each its own PR with a fresh-context reviewer (Iron Rule 10). No new `cli.js` wire behavior, no new endpoint, header, or env var; the `/health` field set is unchanged (only value truthfulness improved).
### Fixed
- **TUI session-scope / boot-reap (#148)** — `lib/tui/session.mjs`'s tmux session prefix is now scoped per-instance by listen port (`ocp-tui-<port>-`) instead of a bare host-wide `ocp-tui-` constant, so a second OCP instance on the same host (e.g. a temporary verification instance) can no longer have its live TUI sessions reaped or `kill-server`'d by another instance's boot/periodic sweep. The one-time boot reap also claims exact-shape legacy `ocp-tui-<8hex>` sessions (pre-fix naming) once, to clean up zombies left behind across an in-place upgrade.
- **`-p` spawn-token mutex + keychain caching (#150)** — the real-HOME token fallback used when the keychain token is within its 5-minute expiry window is now serialized behind a mutex, so concurrent `-p` spawns no longer race the same single-use refresh token against each other (the credential-fork hazard). Added a 30s TTL cache + last-good-label memoization for the keychain read, cutting per-spawn event-loop blocking. The isolation decision (`/health` isolated/real-home reporting) is now re-evaluated per spawn instead of memoized forever, so `/health` no longer misreports a stale decision. New module `lib/spawn-auth.mjs` extracts the pure, unit-testable primitives (mutex, TTL cache, expiry gate, label ordering).
- **Concurrency queue / disconnect handling (#149)** — the shared semaphore now honors a runtime-lowered `maxConcurrent` immediately (previously a decrease was silently ignored until in-flight tasks finished on their own) and wakes queued waiters right away when the limit is raised. Queued `-p`/TUI requests are now linked to the client's HTTP connection via `AbortSignal`; a client that disconnects while queued is spliced out of the queue instead of still spawning `claude` once a slot frees. A singleflight follower whose leader disconnected now retries instead of inheriting a spurious 500, and a queued-then-disconnected request is no longer recorded as a usage failure or logged as an error (quiet disconnect handling).
## v3.21.0 — 2026-06-25 ## v3.21.0 — 2026-06-25
Cleanup + docs release: TUI dead-code removal, docs honesty, and release prep. No new `cli.js` wire behavior; the default path (`CLAUDE_TUI_MODE` unset) is byte-for-byte unchanged. Cleanup + docs release: TUI dead-code removal, docs honesty, and release prep. No new `cli.js` wire behavior; the default path (`CLAUDE_TUI_MODE` unset) is byte-for-byte unchanged.
-4
View File
@@ -894,10 +894,6 @@ node ~/ocp/scripts/sync-openclaw.mjs
This is read-only at startup; the warning never blocks the gateway from running. This is read-only at startup; the warning never blocks the gateway from running.
### A TUI session vanished right after upgrading OCP
If you ran a pre-3.21.1 OCP instance and a post-3.21.1 instance on the same host at the same time during an upgrade, the new instance's one-time boot reap can, once, kill an old-format (`ocp-tui-<8hex>`) live TUI session belonging to the still-running old instance — restart the affected session (`ocp restart` or re-run your TUI turn) and it will come back under the new instance's port-scoped naming.
### OpenClaw shows old models after `ocp update` (v3.10→v3.11 only) ### OpenClaw shows old models after `ocp update` (v3.10→v3.11 only)
One-time bootstrap quirk for the v3.10.0 → v3.11.0 jump only — the running shell had the old `cmd_update` cached. Run once manually: One-time bootstrap quirk for the v3.10.0 → v3.11.0 jump only — the running shell had the old `cmd_update` cached. Run once manually:
+2 -16
View File
@@ -382,25 +382,11 @@ export function getCacheStats() {
// Per ADR 0005 / spec D4: in-process scope only (single Node process per host). // Per ADR 0005 / spec D4: in-process scope only (single Node process per host).
const inflightMap = new Map(); const inflightMap = new Map();
// `retryIf` (optional, audit finding M1): a predicate applied on the FOLLOWER path only. export function singleflight(hash, fn) {
// When a follower joins an existing flight and the shared promise rejects with an error for
// which retryIf(err) is true (in practice: the LEADER's client disconnected while queued —
// an error that is personal to the leader, not a verdict about the upstream), the follower
// does NOT inherit that rejection. Instead it re-enters singleflight with its OWN fn: it
// either becomes the new leader (the map entry is already deleted — see the finally below,
// which runs before any follower's catch because it is attached upstream of the promise the
// followers await) or joins a flight another retrying follower just created. The leader's
// own rejection is never retried here — its error belongs to it (leader path returns the
// bare promise). Callers that pass no retryIf get the exact pre-M1 share-everything behavior.
export function singleflight(hash, fn, retryIf) {
const existing = inflightMap.get(hash); const existing = inflightMap.get(hash);
if (existing) { if (existing) {
existing.requesters++; existing.requesters++;
if (!retryIf) return existing.promise; return existing.promise;
return existing.promise.catch((err) => {
if (!retryIf(err)) throw err;
return singleflight(hash, fn, retryIf);
});
} }
// Wrap fn() in Promise.resolve().then() so synchronous throws don't escape. // Wrap fn() in Promise.resolve().then() so synchronous throws don't escape.
const promise = Promise.resolve().then(fn).finally(() => { const promise = Promise.resolve().then(fn).finally(() => {
-67
View File
@@ -1,67 +0,0 @@
// Pure, dependency-injected primitives for the `-p` spawn-token resolution + HOME-isolation
// layer. Extracted from server.mjs (findings F3 / F5 / F6, 2026-07-07) so the concurrency,
// caching and expiry logic is unit-testable WITHOUT booting the server or mocking execFileSync /
// child_process.spawn / fs. server.mjs owns all I/O (macOS keychain exec, process spawn, fs);
// this module owns only pure decision logic.
//
// ALIGNMENT NOTE: none of this touches the OAuth wire machinery (no endpoint / header / body).
// OCP still NEVER performs a refresh_token grant itself — these helpers only READ + GATE a token
// that some other process (the operator's real claude, or a spawned claude under the real HOME)
// refreshes. That property is load-bearing (issue #112) and preserved.
// Promise-chain mutex. `acquire()` resolves to a `release()` fn; the NEXT `acquire()` does not
// resolve until the current holder calls its `release()`. Serializes async critical sections
// without busy-waiting. release() is idempotent.
export function createSerialMutex() {
let tail = Promise.resolve();
return {
acquire() {
let release;
const gate = new Promise((r) => { release = r; });
const prev = tail;
tail = tail.then(() => gate);
// Hand the caller its release fn only after the previous holder has released.
return prev.then(() => {
let released = false;
return function releaseMutex() { if (!released) { released = true; release(); } };
});
},
};
}
// Short-TTL memo. `get(produce, now)` returns the cached value while `now - storedAt < ttlMs`,
// otherwise calls `produce()` and re-stores. A miss that produces null/undefined is STILL stored
// (so a genuinely-absent source is not re-probed on every call within the TTL window). `now` is
// injectable for testing.
export function createTtlCache({ ttlMs }) {
let value;
let at = -Infinity;
let has = false;
return {
get(produce, now = Date.now()) {
if (has && now - at < ttlMs) return value;
value = produce();
at = now;
has = true;
return value;
},
clear() { has = false; value = undefined; at = -Infinity; },
};
}
// Pure expiry gate. Returns true when `creds` carries a known expiry that is at/within `bufferMs`
// of `now`. Creds WITHOUT `expiresAt` (e.g. long-lived env tokens) are never treated as expiring.
// This gate is applied to the CACHED creds on EVERY use — which is precisely why a short-TTL
// keychain cache (createTtlCache) cannot reintroduce the #146 forever-stale-token regression: the
// cache bounds how often we re-READ the keychain, but the expiry decision is recomputed per use.
export function isTokenExpiring(creds, now = Date.now(), bufferMs = 300000) {
return !!(creds && creds.expiresAt && now + bufferMs >= creds.expiresAt);
}
// Order candidate keychain labels so the last-known-good label is tried first (avoids the
// wrong-label miss that doubles the `security` exec count on the hot path). Pure: performs no
// read. Returns a fresh array; input is not mutated.
export function orderLabelsLastGoodFirst(labels, lastGood) {
if (!lastGood || !labels.includes(lastGood)) return labels.slice();
return [lastGood, ...labels.filter((l) => l !== lastGood)];
}
+9 -59
View File
@@ -20,15 +20,6 @@
// //
// Pure + importable so test-features.mjs can assert the bound directly (no server boot). // Pure + importable so test-features.mjs can assert the bound directly (no server boot).
// Thrown by acquire() when the caller-supplied AbortSignal fires before a slot was granted
// (audit finding F2 — a client that disconnects while queued must never receive a slot; the
// queue entry is spliced out, not just flagged, so `queued` accounting stays exact). Distinct
// `name` lets callers (server.mjs acquireClaudeSlot) tell "client went away" apart from
// "queue is full" without string-matching the message.
export class SemaphoreAbortError extends Error {
constructor(message) { super(message); this.name = "SemaphoreAbortError"; }
}
export class TuiSemaphore { export class TuiSemaphore {
// limit: max concurrent slots. maxQueue: max waiters before run() rejects with backpressure. // limit: max concurrent slots. maxQueue: max waiters before run() rejects with backpressure.
constructor(limit, { maxQueue } = {}) { constructor(limit, { maxQueue } = {}) {
@@ -43,30 +34,9 @@ export class TuiSemaphore {
get inflight() { return this._inflight; } get inflight() { return this._inflight; }
get queued() { return this._waiters.length; } get queued() { return this._waiters.length; }
// Runtime-adjust the concurrency limit (audit finding F1 — a PATCH /settings maxConcurrent
// change must actually take effect, not just be ignored until every currently-inflight task
// happens to finish). Lowering the limit is handled lazily by release() (see below) — it
// simply stops re-granting until inflight drains under the new, lower limit. Raising the
// limit has immediate headroom, so we wake as many queued waiters as now fit.
setLimit(limit) {
this.limit = Math.max(1, parseInt(limit, 10) || 1);
while (this._inflight < this.limit && this._waiters.length > 0) {
const next = this._waiters.shift();
this._inflight++;
next();
}
}
// Acquire a slot. Resolves once a slot is free (immediately if under the limit, otherwise // Acquire a slot. Resolves once a slot is free (immediately if under the limit, otherwise
// when an in-flight task releases). Rejects synchronously-ish if the wait queue is full. // when an in-flight task releases). Rejects synchronously-ish if the wait queue is full.
// `signal` (optional AbortSignal, F2) lets the caller cancel a QUEUED wait — e.g. wired to acquire() {
// a client's socket "close" event so a request that disconnects before a slot is granted
// is removed from the queue instead of eventually being handed a slot for a dead socket.
// If `signal` is already aborted, reject immediately without ever touching the queue.
acquire(signal) {
if (signal?.aborted) {
return Promise.reject(new SemaphoreAbortError("acquire aborted before requesting a slot"));
}
if (this._inflight < this.limit) { if (this._inflight < this.limit) {
this._inflight++; this._inflight++;
return Promise.resolve(); return Promise.resolve();
@@ -76,44 +46,24 @@ export class TuiSemaphore {
`tui_queue_full: TUI concurrency limit (${this.limit}) reached and wait queue ` + `tui_queue_full: TUI concurrency limit (${this.limit}) reached and wait queue ` +
`(${this.maxQueue}) is full`)); `(${this.maxQueue}) is full`));
} }
return new Promise((resolve, reject) => { return new Promise((resolve) => { this._waiters.push(resolve); });
let waiter; // the FIFO entry — captured so onAbort can find + splice exactly this one
const onAbort = () => {
const idx = this._waiters.indexOf(waiter);
if (idx === -1) return; // already granted a slot (shifted out by release()/setLimit) — too late to cancel
this._waiters.splice(idx, 1); // remove, not just flag — keeps `queued` accounting exact
reject(new SemaphoreAbortError("acquire aborted while queued"));
};
waiter = () => {
signal?.removeEventListener("abort", onAbort);
resolve();
};
signal?.addEventListener("abort", onAbort, { once: true });
this._waiters.push(waiter);
});
} }
// Release a slot. Always frees the caller's own slot first, then re-grants it to the next // Release a slot. If a waiter is queued, hand the slot directly to it (inflight stays
// waiter ONLY if the (post-decrement) inflight count is still under the current limit (F1 // constant across the handoff); otherwise decrement.
// fix). This is what makes a runtime-lowered limit actually bite: if the limit was lowered
// while over-subscribed, releases stop re-granting and inflight drains toward the new limit
// instead of a freed slot being handed straight back out at the old, higher occupancy.
release() { release() {
if (this._inflight > 0) this._inflight--;
if (this._inflight < this.limit) {
const next = this._waiters.shift(); const next = this._waiters.shift();
if (next) { if (next) {
this._inflight++; next(); // the woken waiter already "owns" the slot — inflight unchanged
next(); } else if (this._inflight > 0) {
} this._inflight--;
} }
} }
// Run fn() under one slot. Releases in a finally so a throw (PR-A's honesty gates, // Run fn() under one slot. Releases in a finally so a throw (PR-A's honesty gates,
// wallclock truncation, paste-not-landed, tmux spawn failure) NEVER leaks a slot. // wallclock truncation, paste-not-landed, tmux spawn failure) NEVER leaks a slot.
// `signal` (optional, F2) is forwarded to acquire() so a queued run() can be cancelled. async run(fn) {
async run(fn, signal) { await this.acquire();
await this.acquire(signal);
try { try {
return await fn(); return await fn();
} finally { } finally {
+12 -66
View File
@@ -15,48 +15,14 @@ import { tmpdir } from "node:os";
import { randomUUID } from "node:crypto"; import { randomUUID } from "node:crypto";
import { readTuiTranscript } from "./transcript.mjs"; import { readTuiTranscript } from "./transcript.mjs";
// F7 fix (audit finding, LOW): the prefix used to be a bare, host-wide constant export const SESSION_PREFIX = "ocp-tui-"; // per-proxy namespace (coexistence rule)
// ("ocp-tui-"), so a SECOND OCP instance on the same host (e.g. a temporary
// verification instance stood up alongside production — a real pattern used during
// PR #144/#146 verification) would boot-reap and potentially kill-server the OTHER
// instance's LIVE sessions: the coexistence guard below only ever spared foreign
// PRODUCT prefixes (olp-tui-*), never a second ocp-tui-* instance on a different port.
//
// Fix: scope the prefix to the instance's own listen port. The port is the natural
// stable per-instance discriminator on one host (two OCP instances cannot share a
// port), so `ocp-tui-<port>-` uniquely namespaces this instance's sessions and makes
// a same-host sibling OCP instance look exactly like a foreign product (olp-tui-*) to
// the coexistence guard — its `ocp-tui-<otherPort>-*` sessions never match our own
// prefix and are therefore never reaped/kill-server'd by us.
//
// LEGACY_SESSION_PREFIX / LEGACY_SESSION_NAME_RE describe the OLD bare-prefix shape
// (pre-this-fix), retained ONLY for the boot-time legacy-zombie migration handled in
// reapStaleTuiSessions (see comment there). No code path in this version ever CREATES
// a legacy-shaped session name again — sessionPrefixForPort() is the only session-name
// prefix constructor used going forward.
export const LEGACY_SESSION_PREFIX = "ocp-tui-";
// Exact legacy shape: LEGACY_SESSION_PREFIX + sessionId.slice(0, 8), where sessionId is
// a randomUUID() — so the suffix is always exactly 8 lowercase hex characters with NO
// further separator. The new port-scoped shape always inserts a "-" between the port
// digits and the 8-hex suffix (see sessionPrefixForPort), so this regex can never match
// a new-shape name: a new-shape suffix is `<port digits>-<8 hex>` (contains a literal
// "-"), which `[0-9a-f]{8}$` anchored immediately after the prefix cannot satisfy.
export const LEGACY_SESSION_NAME_RE = /^ocp-tui-[0-9a-f]{8}$/;
// Build this instance's own session-name prefix, scoped by its listen port so a
// second OCP instance on the same host (different port) is never mistaken for "ours".
export function sessionPrefixForPort(port) {
return `ocp-tui-${port}-`;
}
const TMUX = process.env.OCP_TUI_TMUX_BIN || "tmux"; const TMUX = process.env.OCP_TUI_TMUX_BIN || "tmux";
const defaultTmux = (args, opts = {}) => const defaultTmux = (args, opts = {}) =>
spawnSync(TMUX, args, { encoding: "utf8", ...opts }); spawnSync(TMUX, args, { encoding: "utf8", ...opts });
// Kill ONLY our own stale sessions. Scoped to sessionPrefixForPort(port) so a co-hosted // Kill ONLY our own stale sessions. Scoped to SESSION_PREFIX so a co-hosted
// OLP test instance's `olp-tui-*` sessions — AND a co-hosted second OCP instance's // OLP test instance's `olp-tui-*` sessions are never touched.
// `ocp-tui-<otherPort>-*` sessions — are never touched (F7 fix).
// //
// Defunct-reaping (PI231 incident): the pane's `claude` process is a child of the // Defunct-reaping (PI231 incident): the pane's `claude` process is a child of the
// long-lived tmux SERVER daemon, NOT of the OCP node process — `tmux new-session -d` // long-lived tmux SERVER daemon, NOT of the OCP node process — `tmux new-session -d`
@@ -70,40 +36,23 @@ const defaultTmux = (args, opts = {}) =>
// merely re-signalling — is to stop the tmux server: when the server exits, the kernel // merely re-signalling — is to stop the tmux server: when the server exits, the kernel
// reparents its surviving children to init (PID 1), which reaps them immediately. // reparents its surviving children to init (PID 1), which reaps them immediately.
// //
// `port` (required) is this instance's own listen port (server.mjs's PORT / lib/constants.mjs // So after killing our own sessions, if the server has NO sessions left of ANY prefix
// DEFAULT_PORT resolution) — the SPOT for "which sessions are ours." // (i.e. nothing we could disrupt — no co-hosted `olp-tui-*` or other instance), we
// // `kill-server` to flush the defunct backlog. If ANY non-ocp session remains we leave the
// `includeLegacy` (default false): when true, sessions matching the exact OLD bare-prefix // server running (coexistence rule, ADR 0007) and let the next boot/periodic sweep retry
// shape (LEGACY_SESSION_NAME_RE) are ALSO treated as ours for kill-session purposes. This is // once the server is otherwise idle.
// the boot-time legacy migration: an operator upgrading past this fix could otherwise be left export function reapStaleTuiSessions({ tmux = defaultTmux } = {}) {
// with orphaned bare-prefix zombie sessions from the PREVIOUS (pre-fix) process generation of
// this SAME instance, since no live instance of the new version ever creates that shape again
// — a legacy-shaped session found at boot is therefore presumed to be this instance's own
// leftover, not a stranger's. Passed true ONLY from the one-time boot-reap call site in
// server.mjs; the periodic idle-reap sweep does NOT set it, so a lingering legacy session
// during steady-state is conservatively treated as foreign (correctly blocking kill-server)
// rather than assumed to be ours on every 15-minute tick. Residual (accepted, documented):
// if a genuinely-still-running PRE-FIX OCP instance is coexisting on the same host at the
// exact moment a new instance boots, its live legacy-shaped session could be reaped — the
// same class of residual risk the audit finding itself accepts ("no live instance of the new
// version creates them"); this PR does not regress that scenario, it only removes the far
// more common same-version collision (the actual F7 finding).
export function reapStaleTuiSessions({ tmux = defaultTmux, port, includeLegacy = false } = {}) {
const r = tmux(["list-sessions", "-F", "#{session_name}"]); const r = tmux(["list-sessions", "-F", "#{session_name}"]);
if (!r || r.status !== 0) return 0; // no tmux server / no sessions if (!r || r.status !== 0) return 0; // no tmux server / no sessions
const names = String(r.stdout || "").split("\n").map((s) => s.trim()).filter(Boolean); const names = String(r.stdout || "").split("\n").map((s) => s.trim()).filter(Boolean);
const ownPrefix = sessionPrefixForPort(port);
let killed = 0; let killed = 0;
let othersRemain = false; let othersRemain = false;
for (const name of names) { for (const name of names) {
const isOwn = name.startsWith(ownPrefix); if (name.startsWith(SESSION_PREFIX)) {
const isLegacyOwn = includeLegacy && LEGACY_SESSION_NAME_RE.test(name);
if (isOwn || isLegacyOwn) {
tmux(["kill-session", "-t", name]); tmux(["kill-session", "-t", name]);
killed++; killed++;
} else { } else {
othersRemain = true; // a session we do NOT own (olp-tui-*, a sibling ocp-tui-<otherPort>-*, othersRemain = true; // a session we do NOT own (e.g. olp-tui-*) — never kill-server
// or — outside includeLegacy — a legacy-shaped name) — never kill-server
} }
} }
// Reap defunct `claude` zombies: safe ONLY when the server is now ours-only/empty. // Reap defunct `claude` zombies: safe ONLY when the server is now ours-only/empty.
@@ -409,15 +358,12 @@ export async function runTuiTurn({
home, home,
realHome, realHome,
cwd, cwd,
port,
wallclockMs = 120000, wallclockMs = 120000,
entrypointMode = "cli", entrypointMode = "cli",
tmux = defaultTmux, tmux = defaultTmux,
}) { }) {
const sessionId = randomUUID(); const sessionId = randomUUID();
// Port-scoped session name (F7 fix) — see sessionPrefixForPort / reapStaleTuiSessions const tmuxName = SESSION_PREFIX + sessionId.slice(0, 8);
// for why this instance's own listen port is the namespace discriminator.
const tmuxName = sessionPrefixForPort(port) + sessionId.slice(0, 8);
const ehome = home || process.env.HOME; // HOME claude runs under (scratch or real) const ehome = home || process.env.HOME; // HOME claude runs under (scratch or real)
const rhome = realHome || process.env.HOME; // real home (OAuth + onboarded config source) const rhome = realHome || process.env.HOME; // real home (OAuth + onboarded config source)
+1 -1
View File
@@ -1,6 +1,6 @@
{ {
"name": "open-claude-proxy", "name": "open-claude-proxy",
"version": "3.21.1", "version": "3.21.0",
"description": "OCP (Open Claude Proxy) — use your Claude Pro/Max subscription as an OpenAI-compatible API for any IDE. Works with Cline, OpenCode, Aider, Continue.dev, OpenClaw, and more.", "description": "OCP (Open Claude Proxy) — use your Claude Pro/Max subscription as an OpenAI-compatible API for any IDE. Works with Cline, OpenCode, Aider, Continue.dev, OpenClaw, and more.",
"type": "module", "type": "module",
"bin": { "bin": {
+93 -333
View File
@@ -43,8 +43,7 @@ import { DEFAULT_PORT } from "./lib/constants.mjs";
import { isLoopbackBind } from "./lib/net.mjs"; import { isLoopbackBind } from "./lib/net.mjs";
import { runTuiTurn, reapStaleTuiSessions, resolveTuiHome } from "./lib/tui/session.mjs"; import { runTuiTurn, reapStaleTuiSessions, resolveTuiHome } from "./lib/tui/session.mjs";
import { detectTuiUpstreamError } from "./lib/tui/transcript.mjs"; import { detectTuiUpstreamError } from "./lib/tui/transcript.mjs";
import { TuiSemaphore, SemaphoreAbortError, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs"; import { TuiSemaphore, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs";
import { createSerialMutex, createTtlCache, isTokenExpiring, orderLabelsLastGoodFirst } from "./lib/spawn-auth.mjs";
const __dirname = dirname(fileURLToPath(import.meta.url)); const __dirname = dirname(fileURLToPath(import.meta.url));
const _pkg = JSON.parse(readFileSync(join(__dirname, "package.json"), "utf8")); const _pkg = JSON.parse(readFileSync(join(__dirname, "package.json"), "utf8"));
@@ -388,106 +387,50 @@ function prepareSpawnHome(dir = SPAWN_HOME_DIR) {
} catch { /* best effort — spawn will surface a hard error if the dir is truly unusable */ } } catch { /* best effort — spawn will surface a hard error if the dir is truly unusable */ }
} }
// Resolve the default-spawn HOME-isolation decision. Returns { isolated, home, reason }: // Resolve the default-spawn HOME-isolation decision ONCE, lazily + memoized (so it runs after
// getOAuthCredentials is defined regardless of source order, and the token probe happens at most
// once). Returns { isolated, home, token } where:
// - isolated:true → spawn under SPAWN_HOME_DIR with cwd=SPAWN_HOME_DIR + the env token. // - isolated:true → spawn under SPAWN_HOME_DIR with cwd=SPAWN_HOME_DIR + the env token.
// - isolated:false → legacy real-HOME spawn, no cwd override (no token, or kill-switch on). // - isolated:false → legacy real-HOME spawn, no cwd override (no token, or kill-switch on).
// // Caches only the isolation DECISION (isolated/home/reason), NOT the token — the token is
// FIX F6 (2026-07-07): this decision is NO LONGER memoized permanently. The previous version // re-resolved FRESH per spawn via resolveSpawnToken(). A memoized token goes stale when its
// cached it forever at first call, which meant: (a) credentials appearing after startup never // source rotates: the macOS keychain access token rotates (~hourly, refreshed by the operator's
// enabled isolation; (b) `rm -rf ~/.ocp/spawn-home` at runtime made every isolated spawn ENOENT // real claude), so a startup snapshot 401s once it expires (caused a ~31h Mac-mini 401 outage,
// until restart; (c) during a token-expiry stint /health reported isolated:true while spawns // 2026-06-26). OCP deliberately does NOT refresh the token itself — a refresh-token grant would
// actually ran real-HOME. Re-evaluating per spawn is cheap because F5's 30s keychain TTL cache // consume the single-use refresh token and log out the operator's real claude (issue #112).
// backs getOAuthCredentials(). This function is the CONFIG-level decision (isolated iff a token let _spawnHomeMode = null;
// resolves AND the kill-switch is off) and has NO fs side effects — the per-spawn EFFECTIVE
// decision additionally applies the expiry gate (resolveSpawnDecision), and scratch-HOME dir prep
// moved to ensureSpawnHome() at the isolated spawn site.
//
// The token itself is re-resolved FRESH per spawn via resolveSpawnToken(); a memoized token goes
// stale when its source rotates (the macOS keychain access token rotates ~hourly, refreshed by the
// operator's real claude), which 401'd every isolated spawn for ~31h on 2026-06-26 (#146). OCP
// deliberately does NOT refresh the token itself — a refresh-token grant would consume the
// single-use refresh token and log out the operator's real claude (issue #112).
function getSpawnHomeMode() { function getSpawnHomeMode() {
if (_spawnHomeMode) return _spawnHomeMode;
if (SPAWN_REAL_HOME) { if (SPAWN_REAL_HOME) {
return { isolated: false, home: null, reason: "kill-switch (OCP_SPAWN_REAL_HOME=1)" }; _spawnHomeMode = { isolated: false, home: null, reason: "kill-switch (OCP_SPAWN_REAL_HOME=1)" };
return _spawnHomeMode;
} }
let hasToken = false; let hasToken = false;
try { hasToken = !!(getOAuthCredentials()?.accessToken); } catch { hasToken = false; } try { hasToken = !!(getOAuthCredentials()?.accessToken); } catch { hasToken = false; }
if (hasToken) return { isolated: true, home: SPAWN_HOME_DIR, reason: "oauth token resolved" }; if (hasToken) {
return { isolated: false, home: null, reason: "no oauth token resolvable" }; prepareSpawnHome(SPAWN_HOME_DIR);
} _spawnHomeMode = { isolated: true, home: SPAWN_HOME_DIR, reason: "oauth token resolved" };
} else {
// FIX F6: re-verify the scratch HOME exists before each isolated spawn and re-create it if it was _spawnHomeMode = { isolated: false, home: null, reason: "no oauth token resolvable" };
// deleted at runtime (it used to be prepared once at startup, so a runtime deletion made every }
// isolated spawn fail ENOENT until restart). mkdirSync is recursive+idempotent → cheap to re-run. return _spawnHomeMode;
function ensureSpawnHome(dir = SPAWN_HOME_DIR) {
if (!existsSync(`${dir}/.claude`)) prepareSpawnHome(dir);
} }
// Resolve a FRESH OAuth access token for an isolated spawn. Read-only (keychain / credentials.json // Resolve a FRESH OAuth access token for an isolated spawn. Read-only (keychain / credentials.json
// / env) — NEVER refreshes/rotates (see getSpawnHomeMode note). Returns null if none resolvable OR // / env) — NEVER refreshes/rotates (see getSpawnHomeMode note). Returns null if none resolvable OR
// if a known expiry is within the 5-min buffer (isTokenExpiring): a null return makes the caller // if a known expiry has passed (5-min buffer): a null return makes the caller fall back to real
// fall back to real HOME, where the spawned claude refreshes the credential natively and self-heals // HOME, where the spawned claude refreshes the credential natively and self-heals (the keychain
// (the keychain token is then fresh again → next spawn is fast). The env-token path (Linux) carries // token is then fresh again → next spawn is fast). The env-token path (Linux) carries no expiresAt
// no expiresAt → never expiry-gated (those tokens are long-lived). // → never expiry-gated (those tokens are long-lived).
function resolveSpawnToken() { function resolveSpawnToken() {
try { try {
const creds = getOAuthCredentials(); const creds = getOAuthCredentials();
if (!creds?.accessToken) return null; if (!creds?.accessToken) return null;
if (isTokenExpiring(creds)) return null; // 5-min buffer; applied to the CACHED creds every use if (creds.expiresAt && Date.now() + 300000 >= creds.expiresAt) return null;
return creds.accessToken; return creds.accessToken;
} catch { return null; } } catch { return null; }
} }
// FIX F3 (2026-07-07): serializes ONLY the real-HOME fallback spawns. Isolated spawns (the common
// fast path) never touch this mutex.
const realHomeFallbackMutex = createSerialMutex();
// Resolve the EFFECTIVE per-spawn HOME/token decision. Returns
// { isolated, home, token, releaseFallback }
// `releaseFallback` is non-null ONLY for a real-HOME fallback holder — the caller MUST call it on
// spawn teardown (wired into cleanup()); it releases the serialization mutex. It is null (no-op)
// for isolated and stable real-HOME (kill-switch / no-token) spawns.
//
// This is async so the real-HOME fallback can `await` the mutex; the keychain reads inside stay
// synchronous (F5 keeps the call sites off async conversion).
async function resolveSpawnDecision() {
const shm = getSpawnHomeMode();
if (!shm.isolated) return { isolated: false, home: null, token: null, releaseFallback: null };
const token = resolveSpawnToken();
if (token) {
ensureSpawnHome(shm.home);
return { isolated: true, home: shm.home, token, releaseFallback: null };
}
// Token is present but within the 5-min expiry window → we would fall back to real HOME, where
// the spawned claude refreshes the credential natively. HAZARD PREVENTED: without serialization,
// every concurrent -p spawn inside this window runs claude under the real HOME simultaneously,
// and each spawned claude races a `refresh_token` grant against the SAME single-use refresh
// token — rotating it out from under the others AND the operator's own real claude (the
// credential-fork hazard; #112 / #146 class). Serialize: admit ONE real-HOME spawn at a time.
// When the next waiter is admitted (the prior holder torn down → its claude has had its lifetime
// to refresh the keychain), re-run resolveSpawnToken(): a now-fresh token means we proceed
// ISOLATED and release the mutex immediately, so the queue drains to the fast path instead of
// piling every request into the real HOME.
const release = await realHomeFallbackMutex.acquire();
try {
// Drop the 30s keychain TTL cache so the re-check reads FRESH keychain state — otherwise a
// waiter admitted right after the prior holder's claude refreshed the token could still see the
// stale (expiring) cached creds and needlessly fall back to real HOME again for up to ~30s.
invalidateKeychainReadCache();
const retry = resolveSpawnToken();
if (retry) {
release();
ensureSpawnHome(shm.home);
return { isolated: true, home: shm.home, token: retry, releaseFallback: null };
}
} catch (e) {
release();
throw e;
}
return { isolated: false, home: null, token: null, releaseFallback: release };
}
// ── FIX ⑥ (concurrency): bounded wait-queue for the -p / stream-json path ────────────── // ── FIX ⑥ (concurrency): bounded wait-queue for the -p / stream-json path ──────────────
// PROBLEM (proven): spawnClaudeProcess used `if (activeRequests >= MAX_CONCURRENT) throw` → // PROBLEM (proven): spawnClaudeProcess used `if (activeRequests >= MAX_CONCURRENT) throw` →
// the client got an opaque 500 AND the rejection was NOT counted in stats (a 15-concurrent // the client got an opaque 500 AND the rejection was NOT counted in stats (a 15-concurrent
@@ -507,61 +450,16 @@ class ConcurrencyOverflowError extends Error {
constructor(message) { super(message); this.name = "ConcurrencyOverflowError"; this.httpStatus = 429; this.retryAfter = CLAUDE_QUEUE_RETRY_AFTER; } constructor(message) { super(message); this.name = "ConcurrencyOverflowError"; this.httpStatus = 429; this.retryAfter = CLAUDE_QUEUE_RETRY_AFTER; }
} }
// Tagged error for audit finding F2: the client disconnected while queued (or was already gone
// before we even tried to queue it). Distinct from ConcurrencyOverflowError so callers never send
// a response on this path — there is no socket left to write to.
class RequestDisconnectedError extends Error {
constructor(message) { super(message); this.name = "RequestDisconnectedError"; }
}
// Build an AbortSignal that fires when `res` (an http.ServerResponse) closes — i.e. the client
// disconnected. Used to cancel a QUEUED concurrency-slot wait (F2) so a client that gives up
// before a slot is granted is spliced out of the wait queue instead of eventually spawning a
// claude process for a dead socket. If `res` has already closed by the time we get here (its
// underlying stream already torn down), the signal is returned pre-aborted so acquire() rejects
// immediately without ever touching the queue — the "close already fired before we attach" case.
// `detach()` MUST be called once the wait settles (granted or rejected) to avoid a listener leak.
function closeSignalFor(res) {
const controller = new AbortController();
if (!res || typeof res.on !== "function") return { signal: controller.signal, detach() {} };
if (res.destroyed) {
controller.abort();
return { signal: controller.signal, detach() {} };
}
const onClose = () => controller.abort();
res.on("close", onClose);
return { signal: controller.signal, detach() { res.removeListener("close", onClose); } };
}
// Acquire a -p concurrency slot, queuing if all are busy (up to CLAUDE_MAX_QUEUE). Resolves to a // Acquire a -p concurrency slot, queuing if all are busy (up to CLAUDE_MAX_QUEUE). Resolves to a
// release() fn that MUST be called exactly once on every exit path (wired into ctx.cleanup()). // release() fn that MUST be called exactly once on every exit path (wired into ctx.cleanup()).
// Rejects with ConcurrencyOverflowError when the wait-queue is full, or with // Rejects with ConcurrencyOverflowError when the wait-queue is full. Increments stats.queued while
// RequestDisconnectedError when `res` closes before a slot is granted (F2) — the caller must not // waiting (decremented on acquire) and stats.queueRejections on overflow.
// spawn claude in that case. `res` is optional (back-compat for any caller without a live response async function acquireClaudeSlot() {
// object); omitting it just means a queued wait can't be cancelled early. stats.queued = claudeSemaphore.queued + 1; // reflect this waiter before we (maybe) block
//
// F8 fix: stats.queued is set from claudeSemaphore.queued AFTER calling acquire() (not before) —
// acquire() synchronously updates _inflight/_waiters before its Promise ever resolves, so reading
// .queued right after the call already reflects reality. The old code set `queued + 1` BEFORE
// calling acquire() to account for "this waiter", which over-reported by 1 whenever the slot was
// granted immediately (the common case, not a queue at all).
async function acquireClaudeSlot(res) {
const { signal, detach } = closeSignalFor(res);
const slot = claudeSemaphore.acquire(signal);
stats.queued = claudeSemaphore.queued; // accurate: acquire() already updated the queue synchronously
try { try {
await slot; await claudeSemaphore.acquire();
} catch (e) { } catch (e) {
detach();
stats.queued = claudeSemaphore.queued; stats.queued = claudeSemaphore.queued;
if (e instanceof SemaphoreAbortError) {
// Client-driven cancellation, not backpressure — do NOT count it as a queueRejection or
// log it as concurrency_queue_full (that log/counter means "the queue itself is full").
logEvent("info", "concurrency_wait_cancelled", {
reason: "client_disconnected", inflight: claudeSemaphore.inflight, queued: claudeSemaphore.queued,
});
throw new RequestDisconnectedError("client disconnected while waiting for a concurrency slot");
}
stats.queueRejections++; stats.queueRejections++;
logEvent("warn", "concurrency_queue_full", { logEvent("warn", "concurrency_queue_full", {
limit: claudeSemaphore.limit, maxQueue: claudeSemaphore.maxQueue, limit: claudeSemaphore.limit, maxQueue: claudeSemaphore.maxQueue,
@@ -571,7 +469,6 @@ async function acquireClaudeSlot(res) {
`backpressure: concurrency limit (${claudeSemaphore.limit}) reached and wait queue ` + `backpressure: concurrency limit (${claudeSemaphore.limit}) reached and wait queue ` +
`(${claudeSemaphore.maxQueue}) is full — retry shortly`); `(${claudeSemaphore.maxQueue}) is full — retry shortly`);
} }
detach();
stats.queued = claudeSemaphore.queued; stats.queued = claudeSemaphore.queued;
let released = false; let released = false;
return function releaseClaudeSlot() { return function releaseClaudeSlot() {
@@ -774,13 +671,7 @@ const TUI_REAP_INTERVAL_MS = 15 * 60 * 1000;
const tuiReapInterval = TUI_MODE ? setInterval(() => { const tuiReapInterval = TUI_MODE ? setInterval(() => {
if (tuiSemaphore.inflight > 0 || tuiSemaphore.queued > 0) return; // a turn is live — defer if (tuiSemaphore.inflight > 0 || tuiSemaphore.queued > 0) return; // a turn is live — defer
try { try {
// F7 fix: scope to THIS instance's own port; a sibling ocp-tui-<otherPort>-* session const n = reapStaleTuiSessions();
// (a second OCP instance on the same host) is treated as foreign, same as olp-tui-*.
// includeLegacy is NOT set here — see reapStaleTuiSessions' comment: the periodic sweep
// conservatively treats any lingering bare-prefix legacy session as foreign so it can
// never trigger kill-server on a steady-state tick; only the one-time boot reap below
// claims legacy-shaped zombies.
const n = reapStaleTuiSessions({ port: PORT });
if (n) logEvent("info", "tui_reaped_stale_sessions", { count: n, trigger: "periodic" }); if (n) logEvent("info", "tui_reaped_stale_sessions", { count: n, trigger: "periodic" });
} catch (e) { logEvent("error", "tui_periodic_reap_failed", { error: e.message }); } } catch (e) { logEvent("error", "tui_periodic_reap_failed", { error: e.message }); }
}, TUI_REAP_INTERVAL_MS) : null; }, TUI_REAP_INTERVAL_MS) : null;
@@ -1034,7 +925,7 @@ function getModelTier(cliModel) {
// budget. releaseSlot is wired into the idempotent cleanup() so the slot is freed on EVERY exit // budget. releaseSlot is wired into the idempotent cleanup() so the slot is freed on EVERY exit
// path (close/error/timeout/abort). Back-compat: releaseSlot defaults to a no-op so any future // path (close/error/timeout/abort). Back-compat: releaseSlot defaults to a no-op so any future
// internal caller that does its own gating still works. // internal caller that does its own gating still works.
function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot = () => {}, spawnDecision = null) { function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot = () => {}) {
const cliModel = MODEL_MAP[model] || model; const cliModel = MODEL_MAP[model] || model;
// Circuit breaker: disabled (see comment at top of breaker section) // Circuit breaker: disabled (see comment at top of breaker section)
@@ -1071,24 +962,28 @@ function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlo
env.CLAUDE_CODE_DISABLE_AUTO_MEMORY = "1"; env.CLAUDE_CODE_DISABLE_AUTO_MEMORY = "1";
} }
// FIX ③ (latency) + F3 (concurrency): apply the pre-resolved per-spawn HOME/token decision. // FIX ③ (latency): default-path spawn-home isolation. When a token is resolvable (and the
// The decision is resolved ASYNC in the caller (resolveSpawnDecision) so the real-HOME fallback // OCP_SPAWN_REAL_HOME kill-switch is off), run claude under a credential-free minimal HOME
// serialization can await its mutex; here we only apply the result. When isolated, run claude // with cwd = that same neutral dir, so it loads NONE of the operator's global ~/.claude
// under a credential-free minimal HOME with cwd = that same neutral dir, so it loads NONE of the // (plugins/skills/hooks) or the ~/ocp project CLAUDE.md/skills — the measured 1028s → 37s
// operator's global ~/.claude (plugins/skills/hooks) or the ~/ocp project CLAUDE.md/skills — the // latency win. The env token is authoritative for `-p` (unlike interactive claude). When no
// measured 1028s → 37s latency win. The env token is authoritative for `-p` (unlike // token is resolvable, falls back to real HOME + inherited cwd (zero regression). See
// interactive claude). When no fresh token is resolvable, decision.isolated is false → real HOME // getSpawnHomeMode() / prepareSpawnHome() above. The DISABLE_CLAUDE_MDS / AUTO_MEMORY flags
// + inherited cwd (zero regression), and the spawned claude resolves+refreshes credentials // are set unconditionally in isolated mode (belt-and-braces; mirrors the TUI path).
// natively. The DISABLE_CLAUDE_MDS / AUTO_MEMORY flags are set unconditionally in isolated mode const spawnHome = getSpawnHomeMode();
// (belt-and-braces; mirrors the TUI path).
const decision = spawnDecision || { isolated: false, releaseFallback: null };
const spawnOpts = { env, stdio: ["pipe", "pipe", "pipe"] }; const spawnOpts = { env, stdio: ["pipe", "pipe", "pipe"] };
if (decision.isolated && decision.token) { if (spawnHome.isolated) {
env.HOME = decision.home; // Re-resolve the token FRESH per spawn (never a startup snapshot — keychain tokens rotate;
env.CLAUDE_CODE_OAUTH_TOKEN = decision.token; // env token is authoritative for -p // a stale snapshot 401s). If unresolvable right now, fall through to real HOME so the spawned
// claude resolves + refreshes credentials natively instead of 401ing on a stale/null token.
const freshToken = resolveSpawnToken();
if (freshToken) {
env.HOME = spawnHome.home;
env.CLAUDE_CODE_OAUTH_TOKEN = freshToken; // env token is authoritative for -p
env.CLAUDE_CODE_DISABLE_CLAUDE_MDS = "1"; env.CLAUDE_CODE_DISABLE_CLAUDE_MDS = "1";
env.CLAUDE_CODE_DISABLE_AUTO_MEMORY = "1"; env.CLAUDE_CODE_DISABLE_AUTO_MEMORY = "1";
spawnOpts.cwd = decision.home; // neutral cwd: no project CLAUDE.md/skills spawnOpts.cwd = spawnHome.home; // neutral cwd: no project CLAUDE.md/skills
}
} }
const proc = spawn(CLAUDE, cliArgs, spawnOpts); const proc = spawn(CLAUDE, cliArgs, spawnOpts);
@@ -1107,11 +1002,6 @@ function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlo
// and cleanup() is guarded by `cleaned`, so the slot is released exactly once on the first // and cleanup() is guarded by `cleaned`, so the slot is released exactly once on the first
// exit path reached (proc 'exit' fires before 'close'; 'error' covers spawn failure). // exit path reached (proc 'exit' fires before 'close'; 'error' covers spawn failure).
try { releaseSlot(); } catch { /* never let release throw out of cleanup */ } try { releaseSlot(); } catch { /* never let release throw out of cleanup */ }
// F3: release the real-HOME fallback serialization mutex (no-op for isolated/normal spawns).
// By now this spawn's claude has had its lifetime to refresh the keychain token, so the next
// queued fallback waiter re-checks resolveSpawnToken() and proceeds ISOLATED with the now-fresh
// token instead of piling into the real HOME. Idempotent; cleanup() is guarded by `cleaned`.
try { if (decision.releaseFallback) decision.releaseFallback(); } catch { /* never throw out of cleanup */ }
} }
// Guarantee slot release on ANY exit path (normal close, error, timeout kill, // Guarantee slot release on ANY exit path (normal close, error, timeout kill,
@@ -1181,37 +1071,18 @@ function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlo
// We accumulate full text across all content_block_delta events plus the // We accumulate full text across all content_block_delta events plus the
// assistant-aggregate fallback, then resolve with the assembled string. // assistant-aggregate fallback, then resolve with the assembled string.
// Reference: OLP ADR 0009 Amendment 1 + commit 97e7d16. // Reference: OLP ADR 0009 Amendment 1 + commit 97e7d16.
// `res` (optional, F2) is the client's http.ServerResponse — passed through so a queued wait async function callClaude(model, messages, conversationId, keyName) {
// can be cancelled the moment the client disconnects, instead of spawning claude for a dead
// socket once a slot finally frees up.
async function callClaude(model, messages, conversationId, keyName, res) {
// FIX ⑥: acquire a concurrency slot first (queues up to CLAUDE_MAX_QUEUE; rejects with a // FIX ⑥: acquire a concurrency slot first (queues up to CLAUDE_MAX_QUEUE; rejects with a
// ConcurrencyOverflowError → 429 when the queue is full, or a RequestDisconnectedError (F2) // ConcurrencyOverflowError → 429 when the queue is full). The release fn is passed into the
// if the client goes away first). The release fn is passed into the spawn so the idempotent // spawn so the idempotent cleanup() frees it on every exit path. If the spawn itself throws
// cleanup() frees it on every exit path. If the spawn itself throws synchronously (before // synchronously (before cleanup is wired), release here so the slot never leaks.
// cleanup is wired), release here so the slot never leaks. const releaseSlot = await acquireClaudeSlot();
// F2×F3 composition: the slot acquire comes FIRST and is the cancellable step — a client
// that disconnects while queued rejects here, BEFORE resolveSpawnDecision() runs, so a
// cancelled request can never acquire (or briefly hold) the real-HOME fallback mutex.
const releaseSlot = await acquireClaudeSlot(res);
// F3: resolve the per-spawn HOME/token decision (may serialize on the real-HOME fallback
// mutex). If it throws, release the just-acquired slot before propagating — cleanup() is
// not wired yet at this point.
let spawnDecision;
try {
spawnDecision = await resolveSpawnDecision();
} catch (err) {
releaseSlot();
throw err;
}
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
let ctx; let ctx;
try { try {
ctx = spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot, spawnDecision); ctx = spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot);
} catch (err) { } catch (err) {
releaseSlot(); releaseSlot();
// Spawn threw before cleanup() was wired → release the fallback mutex here so it never leaks.
try { spawnDecision.releaseFallback?.(); } catch { /* best effort */ }
return reject(err); return reject(err);
} }
@@ -1282,50 +1153,26 @@ async function callClaude(model, messages, conversationId, keyName, res) {
// flag that could perturb cc_entrypoint classification. // flag that could perturb cc_entrypoint classification.
// Authority: claude CLI v2.1.158 interactive mode (cc_entrypoint=cli). // Authority: claude CLI v2.1.158 interactive mode (cc_entrypoint=cli).
// SECURITY: A-path single-user ONLY — home is NOT isolation (see ADR 0007). // SECURITY: A-path single-user ONLY — home is NOT isolation (see ADR 0007).
// `res` (optional, F2) is the client's http.ServerResponse — see closeSignalFor. function callClaudeTui(model, messages, _conversationId, _keyName) {
async function callClaudeTui(model, messages, _conversationId, _keyName, res) {
const cliModel = MODEL_MAP[model] || model; const cliModel = MODEL_MAP[model] || model;
const prompt = messagesToPrompt(messages); // includes system as [System] inline const prompt = messagesToPrompt(messages); // includes system as [System] inline
recordModelRequest(cliModel, prompt.length); recordModelRequest(cliModel, prompt.length);
// C-4: gate the heavy interactive boot behind the TUI semaphore (queuing if all slots are // C-4: gate the heavy interactive boot behind the TUI semaphore. run() acquires a slot
// busy, up to maxQueue). F2: `signal` (tied to `res` "close") cancels a QUEUED wait the // (queuing if all are busy, up to maxQueue), then releases in a finally so any throw from
// instant the client disconnects, so a dead socket never triggers a cold-boot tmux+claude // runTuiTurn (tmux spawn failure, paste-not-landed) OR from the honesty gates below
// spawn; detach() drops the "close" listener as soon as the wait settles rather than // (truncation / error banner) can NEVER leak a slot. tuiSemaphore.inflight feeds /health.
// holding it for the whole (up to 120s) turn. return tuiSemaphore.run(() => runTuiTurn({
const { signal, detach } = closeSignalFor(res);
try {
await tuiSemaphore.acquire(signal);
} catch (err) {
detach();
if (err instanceof SemaphoreAbortError) {
// L1: client-driven cancellation, not an upstream failure — info, not error (mirrors
// acquireClaudeSlot's concurrency_wait_cancelled on the -p path).
logEvent("info", "concurrency_wait_cancelled", {
reason: "client_disconnected", path: "tui", inflight: tuiSemaphore.inflight, queued: tuiSemaphore.queued,
});
throw new RequestDisconnectedError("client disconnected while waiting for a TUI concurrency slot");
}
throw err;
}
detach();
// release() runs in a finally so any throw from runTuiTurn (tmux spawn failure,
// paste-not-landed) OR from the honesty gates below (truncation / error banner) can NEVER
// leak a slot. tuiSemaphore.inflight feeds /health.
try {
const { text, entrypoint, truncated } = await runTuiTurn({
prompt, prompt,
model: cliModel, model: cliModel,
claudeBin: CLAUDE, claudeBin: CLAUDE,
home: TUI_HOME, home: TUI_HOME,
realHome: process.env.HOME, realHome: process.env.HOME,
cwd: TUI_CWD, cwd: TUI_CWD,
port: PORT, // F7 fix: port-scopes the tmux session name so a sibling OCP instance on a
// different port never collides with this instance's reap/kill-server logic.
wallclockMs: TUI_WALLCLOCK_MS, wallclockMs: TUI_WALLCLOCK_MS,
entrypointMode: TUI_ENTRYPOINT, entrypointMode: TUI_ENTRYPOINT,
}); }).then(({ text, entrypoint, truncated }) => {
// ── Honesty gates (issue #133) ─ run BEFORE recordModelSuccess / cache write-back. // ── Honesty gates (issue #133) ─ run BEFORE recordModelSuccess / cache write-back.
// A throw here propagates to the catch below (recordModelError + reject), so the // A throw here propagates to the .catch below (recordModelError + reject), so the
// result never reaches the downstream setCachedResponse / singleflight / SUCCESS path. // result never reaches the downstream setCachedResponse / singleflight / SUCCESS path.
// C-2: the wall-clock cap hit with partial text and NO terminal marker — the turn // C-2: the wall-clock cap hit with partial text and NO terminal marker — the turn
@@ -1360,12 +1207,10 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) {
logEvent("warn", "tui_entrypoint_mismatch", { expected: "cli", got: entrypoint, model: cliModel }); logEvent("warn", "tui_entrypoint_mismatch", { expected: "cli", got: entrypoint, model: cliModel });
} }
return text; return text;
} catch (err) { }).catch((err) => {
recordModelError(cliModel, false); recordModelError(cliModel, false);
throw err; throw err;
} finally { }));
tuiSemaphore.release();
}
} }
// ── SSE heartbeat (opt-in idle watchdog) ──────────────────────────────── // ── SSE heartbeat (opt-in idle watchdog) ────────────────────────────────
@@ -1411,37 +1256,21 @@ async function callClaudeStreaming(model, messages, conversationId, res, authInf
// FIX ⑥: acquire a concurrency slot first (queues up to CLAUDE_MAX_QUEUE). On overflow, surface // FIX ⑥: acquire a concurrency slot first (queues up to CLAUDE_MAX_QUEUE). On overflow, surface
// HTTP 429 + Retry-After (NOT 500). Release is wired into cleanup() for every exit path; if the // HTTP 429 + Retry-After (NOT 500). Release is wired into cleanup() for every exit path; if the
// spawn throws synchronously before cleanup is wired, release here. // spawn throws synchronously before cleanup is wired, release here.
// F2: pass `res` so a queued wait is cancelled the instant this client disconnects — the client
// is already gone in that case, so there is no response to send back.
let releaseSlot; let releaseSlot;
try { try {
releaseSlot = await acquireClaudeSlot(res); releaseSlot = await acquireClaudeSlot();
} catch (err) { } catch (err) {
if (err instanceof RequestDisconnectedError) return; // client gone — nothing to write to
if (err instanceof ConcurrencyOverflowError) { if (err instanceof ConcurrencyOverflowError) {
return jsonResponse(res, 429, { error: { message: sanitizeError(err.message), type: "rate_limit_error" } }, { "Retry-After": String(err.retryAfter) }); return jsonResponse(res, 429, { error: { message: sanitizeError(err.message), type: "rate_limit_error" } }, { "Retry-After": String(err.retryAfter) });
} }
return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } }); return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } });
} }
// F3: resolve the per-spawn HOME/token decision (may serialize on the real-HOME fallback
// mutex). F2×F3 composition: this runs strictly AFTER the (cancellable) slot acquire, so a
// request cancelled while queued never touches the fallback mutex. If it throws, release
// the just-acquired slot before responding — cleanup() is not wired yet at this point.
let spawnDecision;
try {
spawnDecision = await resolveSpawnDecision();
} catch (err) {
releaseSlot();
return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } });
}
let ctx; let ctx;
try { try {
ctx = spawnClaudeProcess(model, messages, conversationId, authInfo.keyName, releaseSlot, spawnDecision); ctx = spawnClaudeProcess(model, messages, conversationId, authInfo.keyName, releaseSlot);
} catch (err) { } catch (err) {
releaseSlot(); releaseSlot();
// Spawn threw before cleanup() was wired → release the fallback mutex here so it never leaks.
try { spawnDecision.releaseFallback?.(); } catch { /* best effort */ }
return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } }); return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } });
} }
@@ -1713,51 +1542,6 @@ const OAUTH_REFRESH_MIN_BACKOFF = 60 * 1000;
const OAUTH_REFRESH_MAX_BACKOFF = 3600 * 1000; const OAUTH_REFRESH_MAX_BACKOFF = 3600 * 1000;
let oauthRefreshBackoff = { nextAttemptAt: 0, currentDelay: OAUTH_REFRESH_MIN_BACKOFF }; let oauthRefreshBackoff = { nextAttemptAt: 0, currentDelay: OAUTH_REFRESH_MIN_BACKOFF };
// FIX F5 (2026-07-07): the macOS keychain read (`security find-generic-password`, up to 5s × 2
// labels when the first label misses) ran on EVERY -p spawn's hot path, blocking the event loop
// (worst case 10s) and stalling all in-flight SSE streams. Two minimal, sync-preserving mitigations:
// (a) memoize the last-good keychain label and try it FIRST → one exec instead of two on the
// steady-state path (orderLabelsLastGoodFirst);
// (b) a short (30s) TTL cache of the keychain read result (createTtlCache).
// SAFETY vs the #146 regression: #146 was a token memoized FOREVER at startup that went stale and
// 401'd. This is a 30s TTL (not forever), AND resolveSpawnToken() re-applies the 5-min expiry gate
// (isTokenExpiring) to the CACHED creds on EVERY use — the creds object carries `expiresAt`, so a
// token expiring within the cache window is still rejected → real-HOME fallback. A short TTL bounds
// how often we re-READ the keychain; it does NOT bound how often we re-DECIDE expiry. This is why a
// short-TTL keychain cache + a per-use expiry check does not reintroduce the forever-stale bug.
const KEYCHAIN_LABELS = ["claude-code-credentials", "Claude Code-credentials"];
const KEYCHAIN_CACHE_TTL_MS = 30 * 1000;
const _keychainCache = createTtlCache({ ttlMs: KEYCHAIN_CACHE_TTL_MS });
let _lastGoodKeychainLabel = null;
// Read the macOS keychain credentials, label-memoized + short-TTL cached (F5). Sync (execFileSync);
// returns the `claudeAiOauth` creds object or null.
function readKeychainCreds() {
return _keychainCache.get(() => {
for (const label of orderLabelsLastGoodFirst(KEYCHAIN_LABELS, _lastGoodKeychainLabel)) {
try {
const raw = execFileSync("security", [
"find-generic-password", "-s", label, "-w"
], { encoding: "utf8", timeout: 5000 }).trim();
const creds = JSON.parse(raw);
if (creds?.claudeAiOauth?.accessToken) {
_lastGoodKeychainLabel = label; // remember the winner → try it first next time
return creds.claudeAiOauth;
}
} catch { /* try next label */ }
}
return null;
});
}
// F3 drain helper: drop the F5 keychain TTL cache so the NEXT getOAuthCredentials() re-reads the
// keychain from scratch. Called under the real-HOME fallback mutex just before the re-check, so a
// waiter admitted after the prior holder's claude refreshed the keychain sees the FRESH token
// immediately (and proceeds ISOLATED) instead of waiting out the ≤30s TTL on the stale creds.
function invalidateKeychainReadCache() {
_keychainCache.clear();
}
function getOAuthCredentials() { function getOAuthCredentials() {
// 1. Env var fallback — highest precedence for explicit overrides. // 1. Env var fallback — highest precedence for explicit overrides.
if (process.env.CLAUDE_CODE_OAUTH_TOKEN) { if (process.env.CLAUDE_CODE_OAUTH_TOKEN) {
@@ -1771,8 +1555,17 @@ function getOAuthCredentials() {
if (creds?.claudeAiOauth?.accessToken) return creds.claudeAiOauth; if (creds?.claudeAiOauth?.accessToken) return creds.claudeAiOauth;
} catch { /* fall through to macOS keychain */ } } catch { /* fall through to macOS keychain */ }
// 3. macOS keychain (both label formats) — F5: label-memoized + 30s TTL cached (see above). // 3. macOS keychain (both label formats)
return readKeychainCreds(); for (const label of ["claude-code-credentials", "Claude Code-credentials"]) {
try {
const raw = execFileSync("security", [
"find-generic-password", "-s", label, "-w"
], { encoding: "utf8", timeout: 5000 }).trim();
const creds = JSON.parse(raw);
if (creds?.claudeAiOauth?.accessToken) return creds.claudeAiOauth;
} catch { /* try next */ }
}
return null;
} }
async function refreshOAuthToken(refreshToken) { async function refreshOAuthToken(refreshToken) {
@@ -2106,12 +1899,9 @@ function applySettingUpdate(key, value) {
switch (key) { switch (key) {
case "timeout": TIMEOUT = value; break; case "timeout": TIMEOUT = value; break;
// FIX ⑥ + F1: keep the -p wait-queue semaphore's limit in sync with the runtime MAX_CONCURRENT // FIX ⑥: keep the -p wait-queue semaphore's limit in sync with the runtime MAX_CONCURRENT
// so a /settings change to maxConcurrent actually changes how many claude procs run at once // so a /settings change to maxConcurrent actually changes how many claude procs run at once.
// in BOTH directions. setLimit() (not a bare `.limit =` assignment) is required: lowering case "maxConcurrent": MAX_CONCURRENT = value; claudeSemaphore.limit = Math.max(1, value); break;
// needs release() to stop over-granting until inflight drains under the new cap, and raising
// needs queued waiters woken immediately to use the new headroom. See lib/tui/semaphore.mjs.
case "maxConcurrent": MAX_CONCURRENT = value; claudeSemaphore.setLimit(value); break;
case "sessionTTL": SESSION_TTL = value; break; case "sessionTTL": SESSION_TTL = value; break;
case "maxPromptChars": MAX_PROMPT_CHARS = value; break; case "maxPromptChars": MAX_PROMPT_CHARS = value; break;
case "cacheTTL": CACHE_TTL = value; break; case "cacheTTL": CACHE_TTL = value; break;
@@ -2268,7 +2058,7 @@ async function handleChatCompletions(req, res) {
const t0TuiStream = Date.now(); const t0TuiStream = Date.now();
const promptCharsTuiStream = messages.reduce((a, m) => a + contentToText(m.content).length, 0); const promptCharsTuiStream = messages.reduce((a, m) => a + contentToText(m.content).length, 0);
try { try {
const content = await callClaudeTui(model, messages, conversationId, req._authKeyName, res); const content = await callClaudeTui(model, messages, conversationId, req._authKeyName);
if (CACHE_TTL > 0 && req._cacheHash) { if (CACHE_TTL > 0 && req._cacheHash) {
try { setCachedResponse(req._cacheHash, model, content); } catch (e) { logEvent("error", "cache_write_failed", { error: e.message }); } try { setCachedResponse(req._cacheHash, model, content); } catch (e) { logEvent("error", "cache_write_failed", { error: e.message }); }
} }
@@ -2305,27 +2095,15 @@ async function handleChatCompletions(req, res) {
// will re-read the freshly-populated cache entry here rather than spawning. // will re-read the freshly-populated cache entry here rather than spawning.
const recheck = getCachedResponse(req._cacheHash, CACHE_TTL); const recheck = getCachedResponse(req._cacheHash, CACHE_TTL);
if (recheck) return recheck.response; if (recheck) return recheck.response;
const c = await upstreamCall(model, messages, conversationId, req._authKeyName, res); const c = await upstreamCall(model, messages, conversationId, req._authKeyName);
try { setCachedResponse(req._cacheHash, model, c); } catch (e) { logEvent("error", "cache_write_failed", { error: e.message }); } try { setCachedResponse(req._cacheHash, model, c); } catch (e) { logEvent("error", "cache_write_failed", { error: e.message }); }
return c; return c;
}, });
// M1: if the LEADER disconnected while queued (F2), its RequestDisconnectedError is
// personal to the leader — a live follower must not inherit it as a spurious 500.
// retryIf makes this follower re-enter singleflight with its OWN fn (own res, own
// disconnect signal), becoming the new leader or joining a retrying sibling's flight —
// but only while OUR client is still connected. If our client is also gone, the
// rejection propagates and the RDE early-return in the catch below ends it quietly.
(err) => err instanceof RequestDisconnectedError && !res.destroyed);
const id = `chatcmpl-${randomUUID()}`; const id = `chatcmpl-${randomUUID()}`;
completionResponse(res, id, model, content); completionResponse(res, id, model, content);
try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: content.length, elapsedMs: Date.now() - t0Usage, success: true }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); } try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: content.length, elapsedMs: Date.now() - t0Usage, success: true }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); }
return; return;
} catch (err) { } catch (err) {
// L1: a client disconnect while queued is NOT an upstream failure — mirror the
// streaming path (which returns without recording anything): no usage-failure row,
// no [proxy] error log, no error response (the socket is gone). The disconnect is
// already logged at info level (concurrency_wait_cancelled) by acquireClaudeSlot.
if (err instanceof RequestDisconnectedError) { try { res.end(); } catch {} return; }
try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: 0, elapsedMs: Date.now() - t0Usage, success: false }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); } try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: 0, elapsedMs: Date.now() - t0Usage, success: false }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); }
console.error(`[proxy] error: ${err.message}`); console.error(`[proxy] error: ${err.message}`);
if (res.headersSent || res.writableEnded || res.destroyed) { if (res.headersSent || res.writableEnded || res.destroyed) {
@@ -2338,14 +2116,11 @@ async function handleChatCompletions(req, res) {
// Fallback: cache disabled (CACHE_TTL=0) or no _cacheHash — original path untouched. // Fallback: cache disabled (CACHE_TTL=0) or no _cacheHash — original path untouched.
try { try {
const content = await upstreamCall(model, messages, conversationId, req._authKeyName, res); const content = await upstreamCall(model, messages, conversationId, req._authKeyName);
const id = `chatcmpl-${randomUUID()}`; const id = `chatcmpl-${randomUUID()}`;
completionResponse(res, id, model, content); completionResponse(res, id, model, content);
try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: content.length, elapsedMs: Date.now() - t0Usage, success: true }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); } try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: content.length, elapsedMs: Date.now() - t0Usage, success: true }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); }
} catch (err) { } catch (err) {
// L1: disconnect-while-queued — same quiet non-error outcome as the singleflight
// path above and the streaming path (see acquireClaudeSlot's info-level log).
if (err instanceof RequestDisconnectedError) { try { res.end(); } catch {} return; }
try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: 0, elapsedMs: Date.now() - t0Usage, success: false }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); } try { recordUsage({ keyId: req._authKeyId, keyName: req._authKeyName, model, promptChars, responseChars: 0, elapsedMs: Date.now() - t0Usage, success: false }); } catch (e) { logEvent("error", "usage_record_failed", { error: e.message }); }
console.error(`[proxy] error: ${err.message}`); console.error(`[proxy] error: ${err.message}`);
if (res.headersSent || res.writableEnded || res.destroyed) { if (res.headersSent || res.writableEnded || res.destroyed) {
@@ -2523,22 +2298,11 @@ const server = createServer(async (req, res) => {
spawn: (() => { spawn: (() => {
if (TUI_MODE) return { mode: "tui (default -p path unused)", isolated: false, home: null }; if (TUI_MODE) return { mode: "tui (default -p path unused)", isolated: false, home: null };
const shm = getSpawnHomeMode(); const shm = getSpawnHomeMode();
// FIX F6: report the EFFECTIVE current decision, not just token PRESENCE. During the
// 5-min pre-expiry window the token exists (shm.isolated=true) but resolveSpawnToken()
// returns null and spawns actually run real-HOME — so `isolated` MUST also reflect the
// expiry gate, or /health lies. The field SET is unchanged (grandfathered B.2 contract,
// ADR 0006 — HARD CONSTRAINT: no field add/remove/rename); only the VALUES are made
// truthful. resolveSpawnToken() is read-only + backed by F5's 30s keychain cache → cheap.
const effIsolated = shm.isolated && resolveSpawnToken() !== null;
return { return {
mode: effIsolated ? "isolated-scratch-home" : "real-home", mode: shm.isolated ? "isolated-scratch-home" : "real-home",
isolated: effIsolated, isolated: shm.isolated,
home: effIsolated ? shm.home : null, home: shm.isolated ? shm.home : null,
reason: effIsolated reason: shm.reason,
? shm.reason
: (shm.isolated
? "oauth token within 5-min expiry window → real-HOME fallback (self-heals on next refresh)"
: shm.reason),
}; };
})(), })(),
// ── FIX ⑥ -p concurrency wait-queue surface — ADDITIVE ── // ── FIX ⑥ -p concurrency wait-queue surface — ADDITIVE ──
@@ -2867,11 +2631,7 @@ server.listen(PORT, BIND_ADDRESS, () => {
: "credentials.json (no CLAUDE_CODE_OAUTH_TOKEN — see Troubleshooting #401)"; : "credentials.json (no CLAUDE_CODE_OAUTH_TOKEN — see Troubleshooting #401)";
console.log(` TUI-mode: ON home=${TUI_HOME} cwd=${TUI_CWD} auth=${tuiAuth} wallclock=${TUI_WALLCLOCK_MS}ms maxConcurrent=${TUI_MAX_CONCURRENT}`); console.log(` TUI-mode: ON home=${TUI_HOME} cwd=${TUI_CWD} auth=${tuiAuth} wallclock=${TUI_WALLCLOCK_MS}ms maxConcurrent=${TUI_MAX_CONCURRENT}`);
try { try {
// F7 fix: scope to THIS instance's own port (see reapStaleTuiSessions). includeLegacy: const n = reapStaleTuiSessions();
// true ONLY here — the one-time boot reap is the designated point to claim orphaned
// bare-prefix ("ocp-tui-<uuid8>") zombie sessions left by a PRE-fix process generation
// of this same instance (no live post-fix instance ever creates that shape again).
const n = reapStaleTuiSessions({ port: PORT, includeLegacy: true });
if (n) logEvent("info", "tui_reaped_stale_sessions", { count: n }); if (n) logEvent("info", "tui_reaped_stale_sessions", { count: n });
} catch {} } catch {}
} }
+19 -500
View File
@@ -5,7 +5,6 @@
*/ */
import { getDb, createKey, listKeys, validateKey, recordUsage, checkQuota, updateKeyQuota, getKeyQuota, findKey, cacheHash, getCachedResponse, setCachedResponse, clearCache, getCacheStats, closeDb, hasCacheControl, singleflight, getInflightStats } from "./keys.mjs"; import { getDb, createKey, listKeys, validateKey, recordUsage, checkQuota, updateKeyQuota, getKeyQuota, findKey, cacheHash, getCachedResponse, setCachedResponse, clearCache, getCacheStats, closeDb, hasCacheControl, singleflight, getInflightStats } from "./keys.mjs";
import { isLoopbackBind } from "./lib/net.mjs"; import { isLoopbackBind } from "./lib/net.mjs";
import { createSerialMutex, createTtlCache, isTokenExpiring, orderLabelsLastGoodFirst } from "./lib/spawn-auth.mjs";
import { createHash } from "node:crypto"; import { createHash } from "node:crypto";
import { strict as assert } from "node:assert"; import { strict as assert } from "node:assert";
import { unlinkSync } from "node:fs"; import { unlinkSync } from "node:fs";
@@ -34,17 +33,6 @@ function test(name, fn) {
} }
} }
async function testAsync(name, fn) {
try {
await fn();
passed++;
console.log(`${name}`);
} catch (e) {
failed++;
console.log(`${name}: ${e.message}`);
}
}
console.log("\n=== OCP Feature Tests (Quota + Cache) ===\n"); console.log("\n=== OCP Feature Tests (Quota + Cache) ===\n");
// Initialize DB // Initialize DB
@@ -463,52 +451,6 @@ async function runSingleflightTests() {
assert.equal(r1, 1); assert.equal(r1, 1);
assert.equal(r2, 2); assert.equal(r2, 2);
}); });
// 7. M1: leader disconnect while queued must not poison live followers. server.mjs passes
// retryIf = (err) => err instanceof RequestDisconnectedError && !res.destroyed — here we
// model that with a tagged error class. The leader (no retryIf on its own promise — the
// rejection is ITS OWN disconnect) sees the error; the live follower re-executes its OWN
// fn and gets a real result instead of a spurious inherited failure.
await asyncTest("M1: leader disconnects while queued → live follower re-executes and gets a real result", async () => {
class FakeDisconnectError extends Error {}
const leaderGate = Promise.withResolvers();
let leaderRuns = 0;
let followerRuns = 0;
const leaderFn = async () => { leaderRuns++; await leaderGate.promise; throw new FakeDisconnectError("leader client gone"); };
const followerFn = async () => { followerRuns++; return "real-execution"; };
const retryIf = (err) => err instanceof FakeDisconnectError;
const leaderP = singleflight("sf-m1-leader-dc", leaderFn); // becomes leader
const followerP = singleflight("sf-m1-leader-dc", followerFn, retryIf); // joins as follower
leaderGate.resolve(); // leader "disconnects" while holding the flight
await assert.rejects(leaderP, FakeDisconnectError, "the leader itself still sees its own disconnect");
assert.equal(await followerP, "real-execution", "follower got a REAL execution, not the leader's disconnect");
assert.equal(leaderRuns, 1, "leader fn ran once");
assert.equal(followerRuns, 1, "follower re-executed exactly once (as the new leader)");
assert.equal(getInflightStats().inflight, 0, "map fully cleaned up after the retry flight settles");
});
// 8. M1 guard: a follower whose retryIf returns false (server.mjs: its OWN client is also
// gone) inherits the rejection unchanged — no retry, no masked error. And a follower with
// NO retryIf keeps the exact pre-M1 share-everything behavior (test 2 pins the fan-out;
// this pins the predicate=false path specifically for the disconnect error).
await asyncTest("M1: follower with retryIf=false (own client also gone) inherits the leader's rejection, no retry", async () => {
class FakeDisconnectError extends Error {}
const gate = Promise.withResolvers();
let followerRuns = 0;
const leaderFn = async () => { await gate.promise; throw new FakeDisconnectError("leader client gone"); };
const followerFn = async () => { followerRuns++; return "should-never-run"; };
const leaderP = singleflight("sf-m1-both-dc", leaderFn);
const followerP = singleflight("sf-m1-both-dc", followerFn, () => false); // own client dead → no retry
gate.resolve();
await assert.rejects(leaderP, FakeDisconnectError);
await assert.rejects(followerP, FakeDisconnectError, "rejection propagates unchanged when retryIf says no");
assert.equal(followerRuns, 0, "follower fn never executed — no wasted spawn for a dead client");
assert.equal(getInflightStats().inflight, 0);
});
} }
await runSingleflightTests(); await runSingleflightTests();
@@ -1710,23 +1652,12 @@ await asyncTest("readTuiTranscript throws when no text and cap elapses", async (
}); });
// ── TUI session reaper ─────────────────────────────────────────────────── // ── TUI session reaper ───────────────────────────────────────────────────
import { reapStaleTuiSessions, sessionPrefixForPort, LEGACY_SESSION_PREFIX, LEGACY_SESSION_NAME_RE, buildTuiCmd } from "./lib/tui/session.mjs"; import { reapStaleTuiSessions, SESSION_PREFIX, buildTuiCmd } from "./lib/tui/session.mjs";
console.log("\nTUI session reaper:"); console.log("\nTUI session reaper:");
// F7 fix: the session prefix is instance-scoped by listen port so a second OCP test("SESSION_PREFIX is ocp-tui-", () => {
// instance on the same host (different port) is never mistaken for "ours". assert.equal(SESSION_PREFIX, "ocp-tui-");
test("sessionPrefixForPort embeds the port (F7 instance scoping)", () => {
assert.equal(sessionPrefixForPort(3456), "ocp-tui-3456-");
assert.equal(sessionPrefixForPort(4000), "ocp-tui-4000-");
assert.notEqual(sessionPrefixForPort(3456), sessionPrefixForPort(4000));
});
test("LEGACY_SESSION_NAME_RE matches only the exact old bare-prefix shape, never the new shape", () => {
assert.ok(LEGACY_SESSION_NAME_RE.test(`${LEGACY_SESSION_PREFIX}a1b2c3d4`), "legacy 8-hex shape matches");
assert.ok(!LEGACY_SESSION_NAME_RE.test("ocp-tui-3456-a1b2c3d4"), "new port-scoped shape must NOT match legacy regex");
assert.ok(!LEGACY_SESSION_NAME_RE.test("ocp-tui-a1b2c3"), "too-short suffix must not match");
assert.ok(!LEGACY_SESSION_NAME_RE.test("ocp-tui-a1b2c3d4extra"), "trailing extra chars must not match");
}); });
console.log("\nTUI command construction (proxy-purity / #4):"); console.log("\nTUI command construction (proxy-purity / #4):");
@@ -1833,40 +1764,22 @@ test("buildTuiCmd OCP_TUI_FULL_TOOLS=1 grants -p-equivalent tool surface (single
} }
}); });
test("reaper kills ONLY this instance's own port-scoped sessions, never olp-tui-", () => { test("reaper kills ONLY ocp-tui- sessions, never olp-tui-", () => {
const killed = []; const killed = [];
const fakeTmux = (args) => { const fakeTmux = (args) => {
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nolp-tui-bbbb\nmisc\nocp-tui-3456-cccc\n" }; if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-aaaa\nolp-tui-bbbb\nmisc\nocp-tui-cccc\n" };
if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; } if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; }
return { status: 0, stdout: "" }; return { status: 0, stdout: "" };
}; };
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 }); const n = reapStaleTuiSessions({ tmux: fakeTmux });
assert.equal(n, 2); assert.equal(n, 2);
assert.equal(killed.join(","), "ocp-tui-3456-aaaa,ocp-tui-3456-cccc"); assert.equal(killed.join(","), "ocp-tui-aaaa,ocp-tui-cccc");
assert.ok(!killed.includes("olp-tui-bbbb"), "olp-tui-bbbb must never be killed"); assert.ok(!killed.includes("olp-tui-bbbb"), "olp-tui-bbbb must never be killed");
}); });
// F7 fix: a second OCP instance on the same host (different port) must be treated exactly
// like a foreign product prefix — never reaped, never allowed to trigger kill-server.
test("reaper treats a sibling OCP instance on a DIFFERENT port as foreign (F7)", () => {
const killed = [];
const calls = [];
const fakeTmux = (args) => {
calls.push(args.join(" "));
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nocp-tui-9999-bbbb\n" };
if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; }
return { status: 0, stdout: "" };
};
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 });
assert.equal(n, 1, "killed only the own-port session");
assert.equal(killed.join(","), "ocp-tui-3456-aaaa");
assert.ok(!killed.includes("ocp-tui-9999-bbbb"), "sibling instance's session (port 9999) must NEVER be killed");
assert.ok(!calls.includes("kill-server"), "kill-server MUST NOT fire — sibling instance's session still live");
});
test("reaper returns 0 when tmux status !== 0 (no server)", () => { test("reaper returns 0 when tmux status !== 0 (no server)", () => {
const fakeTmux = (_args) => ({ status: 1, stdout: "" }); const fakeTmux = (_args) => ({ status: 1, stdout: "" });
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 }); const n = reapStaleTuiSessions({ tmux: fakeTmux });
assert.equal(n, 0); assert.equal(n, 0);
}); });
@@ -1877,7 +1790,7 @@ test("reaper returns 0 for empty session list", () => {
if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; } if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; }
return { status: 0, stdout: "" }; return { status: 0, stdout: "" };
}; };
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 }); const n = reapStaleTuiSessions({ tmux: fakeTmux });
assert.equal(n, 0); assert.equal(n, 0);
assert.equal(killed.length, 0); assert.equal(killed.length, 0);
}); });
@@ -1890,10 +1803,10 @@ test("reaper kill-servers when the server is ours-only (flush defunct claude zom
const calls = []; const calls = [];
const fakeTmux = (args) => { const fakeTmux = (args) => {
calls.push(args.join(" ")); calls.push(args.join(" "));
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nocp-tui-3456-bbbb\n" }; if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-aaaa\nocp-tui-bbbb\n" };
return { status: 0, stdout: "" }; return { status: 0, stdout: "" };
}; };
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 }); const n = reapStaleTuiSessions({ tmux: fakeTmux });
assert.equal(n, 2, "killed both of our sessions"); assert.equal(n, 2, "killed both of our sessions");
assert.ok(calls.includes("kill-server"), "kill-server fired — reaps the defunct backlog"); assert.ok(calls.includes("kill-server"), "kill-server fired — reaps the defunct backlog");
}); });
@@ -1902,10 +1815,10 @@ test("reaper does NOT kill-server when a foreign (non-ocp) session remains (coex
const calls = []; const calls = [];
const fakeTmux = (args) => { const fakeTmux = (args) => {
calls.push(args.join(" ")); calls.push(args.join(" "));
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nolp-tui-bbbb\n" }; if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-aaaa\nolp-tui-bbbb\n" };
return { status: 0, stdout: "" }; return { status: 0, stdout: "" };
}; };
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 }); const n = reapStaleTuiSessions({ tmux: fakeTmux });
assert.equal(n, 1, "killed only our own session"); assert.equal(n, 1, "killed only our own session");
assert.ok(!calls.includes("kill-server"), "kill-server MUST NOT fire — would disrupt olp-tui-*"); assert.ok(!calls.includes("kill-server"), "kill-server MUST NOT fire — would disrupt olp-tui-*");
}); });
@@ -1913,61 +1826,10 @@ test("reaper does NOT kill-server when a foreign (non-ocp) session remains (coex
test("reaper does NOT kill-server when there is no server (status !== 0)", () => { test("reaper does NOT kill-server when there is no server (status !== 0)", () => {
const calls = []; const calls = [];
const fakeTmux = (args) => { calls.push(args.join(" ")); return { status: 1, stdout: "" }; }; const fakeTmux = (args) => { calls.push(args.join(" ")); return { status: 1, stdout: "" }; };
reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 }); reapStaleTuiSessions({ tmux: fakeTmux });
assert.ok(!calls.includes("kill-server"), "no server → no kill-server (early return)"); assert.ok(!calls.includes("kill-server"), "no server → no kill-server (early return)");
}); });
// Legacy migration (F7): pre-fix versions created bare-prefix `ocp-tui-<uuid8>` sessions with
// no port segment. includeLegacy is the boot-only opt-in that claims these as our own leftover
// zombies; the periodic sweep never sets it, so a lingering legacy session cannot trigger
// kill-server on a routine 15-minute tick.
console.log("\nTUI legacy-prefix migration (boot-only reap, F7):");
test("reaper leaves legacy bare-prefix sessions untouched by default (includeLegacy unset)", () => {
const killed = [];
const calls = [];
const fakeTmux = (args) => {
calls.push(args.join(" "));
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nocp-tui-deadbeef\n" };
if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; }
return { status: 0, stdout: "" };
};
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 });
assert.equal(n, 1, "killed only the own-port session");
assert.ok(!killed.includes("ocp-tui-deadbeef"), "legacy session must NOT be reaped without includeLegacy");
assert.ok(!calls.includes("kill-server"), "legacy session blocks kill-server when not claimed");
});
test("reaper claims legacy bare-prefix sessions when includeLegacy=true (boot-time migration)", () => {
const killed = [];
const calls = [];
const fakeTmux = (args) => {
calls.push(args.join(" "));
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nocp-tui-deadbeef\n" };
if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; }
return { status: 0, stdout: "" };
};
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456, includeLegacy: true });
assert.equal(n, 2, "both own-port and legacy sessions reaped");
assert.ok(killed.includes("ocp-tui-deadbeef"), "legacy session claimed as our own leftover");
assert.ok(calls.includes("kill-server"), "kill-server fires once no foreign/unclaimed session remains");
});
test("reaper with includeLegacy=true still spares a sibling instance's port-scoped session", () => {
const killed = [];
const calls = [];
const fakeTmux = (args) => {
calls.push(args.join(" "));
if (args[0] === "list-sessions") return { status: 0, stdout: "ocp-tui-3456-aaaa\nocp-tui-deadbeef\nocp-tui-9999-zzzz\n" };
if (args[0] === "kill-session") { killed.push(args[args.indexOf("-t") + 1]); return { status: 0 }; }
return { status: 0, stdout: "" };
};
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456, includeLegacy: true });
assert.equal(n, 2, "own-port + legacy reaped, sibling instance untouched");
assert.ok(!killed.includes("ocp-tui-9999-zzzz"), "sibling instance session must never be claimed as legacy");
assert.ok(!calls.includes("kill-server"), "sibling instance's live session still blocks kill-server");
});
// ── TUI home preparation (scratch vs real) ─────────────────────────────── // ── TUI home preparation (scratch vs real) ───────────────────────────────
import { prepareTuiHome, ensureTuiCwdTrusted } from "./lib/tui/session.mjs"; import { prepareTuiHome, ensureTuiCwdTrusted } from "./lib/tui/session.mjs";
import { mkdtempSync as hMkdtemp, mkdirSync as hMkdir, writeFileSync as hWrite, readFileSync as hRead, existsSync as hExists, readlinkSync as hReadlink } from "node:fs"; import { mkdtempSync as hMkdtemp, mkdirSync as hMkdir, writeFileSync as hWrite, readFileSync as hRead, existsSync as hExists, readlinkSync as hReadlink } from "node:fs";
@@ -2051,7 +1913,7 @@ test("resolveTuiHome: explicit OCP_TUI_HOME wins regardless of env token (back-c
}); });
// ── TUI concurrency limiter + drift observability (PR-B: audit C-4 / C-5) ── // ── TUI concurrency limiter + drift observability (PR-B: audit C-4 / C-5) ──
import { TuiSemaphore, SemaphoreAbortError, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs"; import { TuiSemaphore, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs";
console.log("\nTUI concurrency limiter (C-4):"); console.log("\nTUI concurrency limiter (C-4):");
@@ -2156,226 +2018,6 @@ await asyncTest("FIX ⑥: slot released on normal completion is immediately reus
assert.equal(sem.inflight, 0); assert.equal(sem.inflight, 0);
}); });
// ── Audit F1 — runtime-lowered/raised limit must actually bite ──────────────
// server.mjs reuses this same TuiSemaphore as `claudeSemaphore`; a PATCH /settings
// maxConcurrent update now calls `claudeSemaphore.setLimit(value)` (see applySettingUpdate's
// "maxConcurrent" case). These tests pin the semaphore-level contract that fix depends on.
console.log("\nF1 — runtime concurrency-limit changes (setLimit / release honoring the current limit):");
await asyncTest("F1: lowering the limit mid-load — release() stops re-granting until inflight drains under the new limit", async () => {
const sem = new TuiSemaphore(3, { maxQueue: 16 });
const g = [deferred(), deferred(), deferred()];
const held = g.map((d) => sem.run(async () => { await d.p; }));
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 3, "3 tasks hold the 3 slots");
// A 4th arrives while at capacity — it queues.
const g4 = deferred();
const queued4 = sem.run(async () => { await g4.p; });
await new Promise((r) => setImmediate(r));
assert.equal(sem.queued, 1, "4th request queued");
// Operator lowers maxConcurrent from 3 to 1 while all 3 original slots are still inflight
// (mirrors a PATCH /settings maxConcurrent=1 hitting server.mjs mid-burst).
sem.setLimit(1);
assert.equal(sem.limit, 1);
// Releasing one of the 3 original holders must NOT hand the freed slot to the queued 4th
// request — before the F1 fix, release() handed slots off unconditionally, so inflight
// would have stayed pinned at the OLD higher occupancy forever.
g[0].resolve();
await held[0];
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 2, "inflight drains toward the new limit, not re-granted");
assert.equal(sem.queued, 1, "4th request is STILL queued — not over-admitted");
g[1].resolve();
await held[1];
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 1, "inflight now exactly at the new limit (1)");
assert.equal(sem.queued, 1, "still queued — inflight(1) is not < limit(1), so no grant yet");
// Releasing the LAST original holder finally drops inflight under the new limit — only
// now does the queued 4th request get granted.
g[2].resolve();
await held[2];
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 1, "queued 4th request now holds the single slot");
assert.equal(sem.queued, 0, "queue drained");
g4.resolve();
await queued4;
assert.equal(sem.inflight, 0);
});
await asyncTest("F1: raising the limit wakes queued waiters immediately, up to the new headroom", async () => {
const sem = new TuiSemaphore(1, { maxQueue: 16 });
const g1 = deferred();
const t1 = sem.run(async () => { await g1.p; }); // holds the only slot
await new Promise((r) => setImmediate(r));
const started = [];
const g2 = deferred(), g3 = deferred();
const t2 = sem.run(async () => { started.push(2); await g2.p; });
const t3 = sem.run(async () => { started.push(3); await g3.p; });
await new Promise((r) => setImmediate(r));
assert.equal(sem.queued, 2, "both queue behind the single holder");
assert.deepEqual(started, [], "neither queued task has started");
// Operator raises maxConcurrent from 1 to 3 (2 units of new headroom) — BOTH queued
// waiters must be woken immediately, without waiting for t1 to release.
sem.setLimit(3);
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 3, "t1 + both newly-woken waiters now hold slots");
assert.equal(sem.queued, 0, "queue drained by the limit raise");
assert.deepEqual(started.sort(), [2, 3], "both queued tasks started without waiting for t1's release");
g1.resolve(); g2.resolve(); g3.resolve();
await Promise.all([t1, t2, t3]);
assert.equal(sem.inflight, 0);
});
await asyncTest("F1: raising the limit wakes only as many waiters as the new headroom allows (FIFO)", async () => {
const sem = new TuiSemaphore(1, { maxQueue: 16 });
const g1 = deferred();
const t1 = sem.run(async () => { await g1.p; });
await new Promise((r) => setImmediate(r));
const started = [];
const g2 = deferred(), g3 = deferred();
const t2 = sem.run(async () => { started.push(2); await g2.p; });
const t3 = sem.run(async () => { started.push(3); await g3.p; });
await new Promise((r) => setImmediate(r));
assert.equal(sem.queued, 2);
sem.setLimit(2); // only 1 unit of new headroom (1 -> 2) — exactly one queued waiter wakes
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 2);
assert.equal(sem.queued, 1, "one waiter still queued — only one slot of headroom existed");
assert.deepEqual(started, [2], "FIFO: the earlier-queued waiter (t2) wakes, not t3");
// Freeing t1's slot afterward still honors the (now current) limit of 2 via release()'s
// normal path — the still-queued t3 gets in once a slot actually frees.
g1.resolve();
await t1;
await new Promise((r) => setImmediate(r));
assert.deepEqual(started, [2, 3], "t3 granted once a slot frees, honoring the raised limit");
assert.equal(sem.queued, 0);
g2.resolve(); g3.resolve();
await t2; await t3;
assert.equal(sem.inflight, 0);
});
// ── Audit F2 — queued waiters must be cancellable on client disconnect ──────
// server.mjs wires an AbortSignal derived from the client's res "close" event into
// claudeSemaphore.acquire()/tuiSemaphore.acquire() (see closeSignalFor + acquireClaudeSlot /
// callClaudeTui). These tests pin the semaphore-level cancellation contract that depends on.
console.log("\nF2 — queued-wait cancellation via AbortSignal (client disconnect while queued):");
await asyncTest("F2: aborting a QUEUED waiter rejects with SemaphoreAbortError and SPLICES it out (queued drops immediately, not just flagged)", async () => {
const sem = new TuiSemaphore(1, { maxQueue: 16 });
const g1 = deferred();
const t1 = sem.run(async () => { await g1.p; }); // holds the only slot
await new Promise((r) => setImmediate(r));
const controller = new AbortController();
const acquire2 = sem.acquire(controller.signal); // queues behind t1
await new Promise((r) => setImmediate(r));
assert.equal(sem.queued, 1, "second acquire queued");
controller.abort(); // simulates the client disconnecting while still queued
await assert.rejects(acquire2, SemaphoreAbortError, "cancelled waiter rejects with SemaphoreAbortError");
assert.equal(sem.queued, 0, "cancelled waiter is REMOVED — queue length drops immediately");
assert.equal(sem.inflight, 1, "t1's slot is untouched by the cancellation");
// Prove the cancelled waiter never later acquires a slot: free t1's slot and confirm
// nobody is waiting to receive it (the queue is genuinely empty, not just decremented).
g1.resolve();
await t1;
assert.equal(sem.inflight, 0, "slot freed with nobody queued — the cancelled waiter never got it");
});
await asyncTest("F2: an already-aborted signal rejects acquire() immediately, never touching the wait queue", async () => {
const sem = new TuiSemaphore(1, { maxQueue: 16 });
const g1 = deferred();
const t1 = sem.run(async () => { await g1.p; }); // holds the only slot
await new Promise((r) => setImmediate(r));
const controller = new AbortController();
controller.abort(); // client already gone before this request ever tries to acquire
await assert.rejects(sem.acquire(controller.signal), SemaphoreAbortError);
assert.equal(sem.queued, 0, "never entered the wait queue at all");
g1.resolve(); await t1;
});
await asyncTest("F2: cancelling one queued waiter preserves FIFO order for the others", async () => {
const sem = new TuiSemaphore(1, { maxQueue: 16 });
const g1 = deferred();
const t1 = sem.run(async () => { await g1.p; });
await new Promise((r) => setImmediate(r));
const started = [];
const cA = new AbortController();
const cB = new AbortController();
const accA = sem.acquire(cA.signal).then(() => started.push("A"));
const accB = sem.acquire(cB.signal).then(() => started.push("B"));
const g3 = deferred();
const t3 = sem.run(async () => { started.push("C"); await g3.p; });
await new Promise((r) => setImmediate(r));
assert.equal(sem.queued, 3, "A, B, C all queued behind t1");
cB.abort(); // B (the middle waiter) disconnects
await assert.rejects(accB, SemaphoreAbortError);
assert.equal(sem.queued, 2, "B removed; A and C remain, in original relative order");
g1.resolve();
await t1;
await new Promise((r) => setImmediate(r));
assert.deepEqual(started, ["A"], "A (queued first, still present) is granted next — FIFO preserved after B's removal");
assert.equal(sem.inflight, 1);
assert.equal(sem.queued, 1, "C still waiting");
sem.release(); // A was acquired directly (not via run()) — free its slot manually
await new Promise((r) => setImmediate(r));
assert.deepEqual(started, ["A", "C"], "C granted next");
g3.resolve();
await t3;
assert.equal(sem.inflight, 0);
});
await asyncTest("F2/L2: abort AFTER grant is a no-op — waiter keeps its slot, no rejection, slot released exactly once", async () => {
const sem = new TuiSemaphore(1, { maxQueue: 16 });
const g1 = deferred();
const t1 = sem.run(async () => { await g1.p; }); // holds the only slot
await new Promise((r) => setImmediate(r));
const controller = new AbortController();
let granted = false;
const acq = sem.acquire(controller.signal).then(() => { granted = true; });
await new Promise((r) => setImmediate(r));
assert.equal(sem.queued, 1, "waiter queued behind t1");
// t1 finishes → release() shifts the waiter out and grants it the slot (waiter() detaches
// the abort listener before resolving).
g1.resolve();
await t1;
await acq;
assert.equal(granted, true, "waiter was granted the slot");
assert.equal(sem.inflight, 1, "granted waiter holds the slot");
assert.equal(sem.queued, 0);
// The client disconnects AFTER the grant — the abort-after-grant race. onAbort must be a
// no-op (the waiter is no longer in _waiters; idx===-1 guard): no rejection materializes,
// the queue is untouched, and the slot is still owned by the (already-resolved) acquirer.
controller.abort();
await new Promise((r) => setImmediate(r));
assert.equal(sem.inflight, 1, "abort after grant did NOT revoke or double-free the slot");
assert.equal(sem.queued, 0, "abort after grant did not corrupt queue accounting");
// The slot is released exactly once via the normal path and is immediately reusable.
sem.release();
assert.equal(sem.inflight, 0, "slot released exactly once via the normal path");
await sem.run(async () => {}); // prove the semaphore is fully healthy afterward
assert.equal(sem.inflight, 0);
});
console.log("\nTUI drift observability (C-5):"); console.log("\nTUI drift observability (C-5):");
test("recordTuiEntrypoint: observed 'cli' is NOT a mismatch and sets lastEntrypoint", () => { test("recordTuiEntrypoint: observed 'cli' is NOT a mismatch and sets lastEntrypoint", () => {
@@ -2756,131 +2398,8 @@ test("isLoopbackBind: '100.64.0.1' → false (Tailscale IP)", () => {
assert.equal(isLoopbackBind("100.64.0.1"), false); assert.equal(isLoopbackBind("100.64.0.1"), false);
}); });
// ── Spawn-auth primitives (F3 / F5 / F6, lib/spawn-auth.mjs) ──
// Pure, dependency-injected primitives extracted from server.mjs so the spawn-token concurrency /
// caching / expiry logic is testable without booting the server or mocking execFileSync/spawn.
console.log("\nSpawn-auth (F3 mutex / F5 TTL cache + label memo / F6 expiry gate):");
// F5: expiry gate — the load-bearing invariant that lets a short-TTL keychain cache stay safe.
test("isTokenExpiring: creds within 5-min buffer → true", () => {
assert.equal(isTokenExpiring({ expiresAt: 1000 }, 1000 - 300000, 300000), true); // exactly at buffer edge
assert.equal(isTokenExpiring({ expiresAt: 1000 }, 900, 300000), true); // past the edge
});
test("isTokenExpiring: creds well beyond buffer → false", () => {
assert.equal(isTokenExpiring({ expiresAt: 10_000_000 }, 0, 300000), false);
});
test("isTokenExpiring: no expiresAt (long-lived env token) → never expiring", () => {
assert.equal(isTokenExpiring({ accessToken: "x" }, Date.now(), 300000), false);
assert.equal(isTokenExpiring(null, Date.now(), 300000), false);
});
// F5: last-good label ordering — one exec instead of two on the steady-state keychain path.
test("orderLabelsLastGoodFirst: last-good label is tried first", () => {
const labels = ["A", "B"];
assert.deepEqual(orderLabelsLastGoodFirst(labels, "B"), ["B", "A"]);
});
test("orderLabelsLastGoodFirst: null/unknown last-good → original order, fresh array", () => {
const labels = ["A", "B"];
assert.deepEqual(orderLabelsLastGoodFirst(labels, null), ["A", "B"]);
assert.deepEqual(orderLabelsLastGoodFirst(labels, "Z"), ["A", "B"]);
assert.notEqual(orderLabelsLastGoodFirst(labels, null), labels); // does not mutate/alias input
});
// F5: TTL cache — bounds how often we RE-READ the keychain (not how often we re-decide expiry).
test("createTtlCache: serves cached value within TTL, re-produces after TTL", () => {
const cache = createTtlCache({ ttlMs: 30000 });
let calls = 0;
const produce = () => { calls++; return `v${calls}`; };
assert.equal(cache.get(produce, 0), "v1");
assert.equal(cache.get(produce, 10000), "v1"); // within TTL → cached, producer NOT called
assert.equal(calls, 1);
assert.equal(cache.get(produce, 40000), "v2"); // past TTL → re-produced
assert.equal(calls, 2);
});
test("createTtlCache: caches a null miss (absent source not re-probed within TTL)", () => {
const cache = createTtlCache({ ttlMs: 30000 });
let calls = 0;
const produce = () => { calls++; return null; };
assert.equal(cache.get(produce, 0), null);
assert.equal(cache.get(produce, 5000), null);
assert.equal(calls, 1); // the null was cached, not re-probed
});
// F5 core safety property: a short-TTL cache CANNOT reintroduce the #146 forever-stale bug because
// the expiry gate is applied to the CACHED creds on every use. The cache keeps returning the same
// creds object, but isTokenExpiring flips to true the moment the clock crosses the expiry buffer.
test("TTL cache respects expiry gate: cached creds still rejected once clock passes expiry", () => {
const cache = createTtlCache({ ttlMs: 30000 });
const creds = { accessToken: "tok", expiresAt: 1_000_000 };
// t=980_000: cached AND not yet within the 5-min (300_000) buffer → usable.
const c1 = cache.get(() => creds, 980_000 - 300_000 - 1);
assert.equal(isTokenExpiring(c1, 980_000 - 300_000 - 1, 300000), false);
// t=800_000 later: SAME cached object returned (within TTL of the second read window), but now
// within the expiry buffer → gate rejects it → caller falls back to real HOME. No forever-stale.
const c2 = cache.get(() => creds, 990_000);
assert.equal(c2, c1, "cache returns the same creds object");
assert.equal(isTokenExpiring(c2, 990_000, 300000), true, "expiry gate still fires on cached creds");
});
// ── Async: F3 real-HOME fallback serialization mutex ──
async function runAsyncTests() {
await testAsync("createSerialMutex: second waiter blocks until first holder releases", async () => {
const mutex = createSerialMutex();
const order = [];
const rel1 = await mutex.acquire();
order.push("h1-enter");
let secondEntered = false;
const p2 = mutex.acquire().then((rel2) => { secondEntered = true; order.push("h2-enter"); return rel2; });
await new Promise((r) => setTimeout(r, 15));
assert.equal(secondEntered, false, "second waiter must NOT enter while first holds the mutex");
order.push("h1-release");
rel1();
const rel2 = await p2;
assert.equal(secondEntered, true, "second waiter enters only after release");
rel2();
assert.deepEqual(order, ["h1-enter", "h1-release", "h2-enter"]);
});
await testAsync("createSerialMutex: N acquires run strictly in FIFO order, never overlapping", async () => {
const mutex = createSerialMutex();
const events = [];
let active = 0;
async function critical(id) {
const rel = await mutex.acquire();
active++;
assert.equal(active, 1, `only one holder at a time (id=${id})`);
events.push(`start${id}`);
await new Promise((r) => setTimeout(r, 5));
events.push(`end${id}`);
active--;
rel();
}
await Promise.all([critical(1), critical(2), critical(3)]);
assert.deepEqual(events, ["start1", "end1", "start2", "end2", "start3", "end3"]);
});
await testAsync("createSerialMutex: release() is idempotent (double-release does not double-admit)", async () => {
const mutex = createSerialMutex();
const rel1 = await mutex.acquire();
rel1();
rel1(); // second call must be a no-op
const rel2 = await mutex.acquire(); // should acquire cleanly, exactly once
let thirdEntered = false;
const p3 = mutex.acquire().then((r) => { thirdEntered = true; return r; });
await new Promise((r) => setTimeout(r, 15));
assert.equal(thirdEntered, false, "double-release must not have leaked an extra admit slot");
rel2();
(await p3)();
});
}
// ── Cleanup ── // ── Cleanup ──
runAsyncTests().then(() => { closeDb();
closeDb();
console.log(`\n=== Results: ${passed} passed, ${failed} failed ===\n`); console.log(`\n=== Results: ${passed} passed, ${failed} failed ===\n`);
process.exit(failed > 0 ? 1 : 0); process.exit(failed > 0 ? 1 : 0);
}).catch((e) => {
console.error("async test runner crashed:", e);
closeDb();
process.exit(1);
});