mirror of
https://github.com/dtzp555-max/ocp.git
synced 2026-07-22 05:25:08 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1978a375e3 | ||
|
|
d97cef7b30 | ||
|
|
66dc9949ed | ||
|
|
758c2d703f | ||
|
|
6cf5e950a5 |
@@ -1,29 +1,5 @@
|
||||
# Changelog
|
||||
|
||||
## 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.
|
||||
|
||||
### TUI dead-code / footgun cleanup
|
||||
|
||||
- **A1 — removed inert entrypoint-env path** (`lib/tui/session.mjs`): deleted `resolveTuiEntrypointEnv()` and the redundant env-strip block in `runTuiTurn`. The `{env}` object passed to `spawnSync` (tmux itself) was the wrong target — tmux does NOT forward the spawning process's environment to the pane; the pane's `claude` gets its env exclusively from the `env` prefix string built inside `buildTuiCmd` (verified live 2026-06-01). The spawnSync env is now intentionally minimal (`HOME` only). Behavior is unchanged: `buildTuiCmd` already handled all claude-specific env vars via its prefix string.
|
||||
- **A2 — removed test-only transcript helpers** (`lib/tui/transcript.mjs`): deleted `encodeCwd()` and `transcriptPath()` exports and the tests that pinned them. Production resolves transcripts exclusively via `findTranscriptPath()` (glob by session-id), which is immune to the exact path-encoding rule. No non-test importers existed (grep confirms). A `// TODO` comment near `findTranscriptPath()` notes that a CI fixture-contract test would make claude-schema drift fail loudly.
|
||||
- **A3 — removed headless-unusable `--dangerously-skip-permissions` branch** (`lib/tui/session.mjs` + `README.md`): `OCP_TUI_FULL_TOOLS=1` now always takes the `--allowedTools` path. The removed branch pushed `--dangerously-skip-permissions` when `CLAUDE_SKIP_PERMISSIONS=true`; on claude v2.1.x this triggers an interactive bypass-acceptance screen that a headless tmux pane cannot answer → the turn hangs to the wallclock cap and bricks the pane. The working path is `--allowedTools` + scratch-home `settings.json` `additionalDirectories`. `CLAUDE_SKIP_PERMISSIONS` for the `-p` path is unchanged (still used in `server.mjs`).
|
||||
|
||||
### Docs
|
||||
|
||||
- **Client-tools boundary** (README `§ How It Works`): OCP is a text-prompt bridge only — it does not pass OpenAI `tools`/`functions` or Anthropic `tool_use` blocks to the client. Clients receive assistant TEXT only; client-local tool execution is not supported by design (bypassing `cli.js` = out of scope per `ALIGNMENT.md`).
|
||||
- **ToS honesty** (README `§ Deployment model & security`): pooling one Claude subscription across multiple distinct people may violate Anthropic's Consumer ToS and risk account suspension by the abuse classifier. The defensible framing is "one person, your own devices" — friends/team sharing is not. The prior language ("account terms are your call") was accurate but understated the risk.
|
||||
- **"Why OCP" posture** (README `§ Why OCP?`): new bullet making explicit that OCP drives the official `claude` CLI as-is — no OAuth token extraction, no binary patching, no protocol invention — so traffic looks like genuine Claude Code (`cc_entrypoint=cli`).
|
||||
- **Promotion plan** (`docs/PROMOTION.md`): "stable & visible" strategy covering goal (polish + low-key OSS visibility, NOT growth-hacking given the live ToS/billing risk), pre-requisites (stability first), honest ToS disclosure requirement, items explicitly skipped (multi-backend routing → OLP; gateway model-discovery; raw API passthrough → ALIGNMENT.md scope), TUI toggle as billing-split insurance, and low-key visibility actions. Framed as a recommendation for the maintainer to review, not a committed plan.
|
||||
|
||||
### Previously shipped (v3.20.x) — documented here for completeness
|
||||
|
||||
- **Default `-p` spawn-home isolation** (v3.20.0 / PR-A): per-request `claude` spawns run in a credential-free minimal scratch HOME (`$HOME/.ocp/spawn-home`, no `.credentials.json`/`settings.json`/plugins) with a neutral cwd and the env token, cutting per-request latency (measured ~10–28s → ~3–7s). Kill-switch: `OCP_SPAWN_REAL_HOME=1`. Active mode shown at startup and on `/health.spawn`.
|
||||
- **Bounded concurrency wait-queue** (v3.20.0 / PR-B): excess `-p` requests queue (up to `CLAUDE_MAX_QUEUE`, default 16) instead of being rejected; a full queue returns `HTTP 429` + `Retry-After` (not an opaque 500). New env vars: `CLAUDE_MAX_QUEUE`, `CLAUDE_QUEUE_RETRY_AFTER`. Surfaced on `/health.concurrency` + `/health.stats.queueRejections`.
|
||||
- **`ocp restart`** macOS `bootout`+`bootstrap` (v3.20.0 / PR-B): safe restart command that forces launchd to re-read the plist (unlike `kickstart -k` which reuses the cached env).
|
||||
- **`/ocp` plugin OpenClaw-2026.5.27 compat** (v3.20.0 / PR-C): gateway plugin updated for the current OpenClaw API version.
|
||||
|
||||
## v3.20.1 — 2026-06-13
|
||||
|
||||
TUI-mode auth hardening: fixes the recurring `Please run /login · API Error: 401` (the PI231 incident) and reaps leaked defunct `claude` sessions. ([#141](https://github.com/dtzp555-max/ocp/pull/141))
|
||||
|
||||
@@ -31,7 +31,6 @@ There are several Claude proxy projects. OCP picks a specific lane: **align tigh
|
||||
- **SSE heartbeat for long reasoning** ([v3.12.0](https://github.com/dtzp555-max/ocp/releases/tag/v3.12.0), opt-in). If you've ever watched your IDE die at the 60s idle mark during a long Claude tool-use pause — that's nginx/Cloudflare default behavior. OCP emits an SSE comment frame to keep the connection alive without polluting the response. ([PR #49](https://github.com/dtzp555-max/ocp/pull/49))
|
||||
- **`cli.js` alignment + CI guardrail.** LLM-assisted code drifts easily — it's tempting to invent plausible-looking endpoints that `cli.js` doesn't actually use. [`ALIGNMENT.md`](./ALIGNMENT.md) is binding: every endpoint OCP exposes must cite a `cli.js` line. The [`alignment.yml`](./.github/workflows/alignment.yml) CI workflow blocks PRs that introduce known-hallucinated tokens. The payoff is boring: your setup keeps working when `cli.js` ships its next minor.
|
||||
- **`models.json` single source of truth** (v3.11.0). Adding a model is one file edit; both `/v1/models` and the OpenClaw bootstrap derive from it. ([PR #30](https://github.com/dtzp555-max/ocp/pull/30))
|
||||
- **Drives the official CLI as-is, no binary patching.** OCP spawns the official `claude` CLI (or hosts it in an interactive tmux pane for TUI mode) — it does not extract OAuth tokens from memory, patch the binary, or invent protocol extensions. Traffic therefore looks like genuine Claude Code to Anthropic's classifiers (`cc_entrypoint=cli`). See `ALIGNMENT.md` for why this constraint is load-bearing.
|
||||
|
||||
### Comparison
|
||||
|
||||
@@ -424,7 +423,7 @@ ocp keys revoke son-ipad # Revoke a key
|
||||
- The per-key modes (`shared` / `multi`) give per-key **usage tracking, quotas, and cache separation** — useful for seeing who used what and capping budgets.
|
||||
- They do **not** give a **security isolation boundary**. The spawned `claude` runs with the **operator's filesystem access** and is *not* sandboxed per key. **Only share with people you fully trust, on a trusted network.**
|
||||
- For simple trusted family sharing, the easiest setup is a single shared **anonymous key** (see [Anonymous Access](#anonymous-access-optional)) — no per-person separation, same trust assumption.
|
||||
- **Account terms and ToS — read before sharing with others.** Claude Pro/Max are *per-user* accounts. Pooling a single subscription across **multiple distinct people** may violate Anthropic's Consumer Terms of Service and risk account suspension by the abuse classifier. The defensible framing is **"one person, your own devices"** — sharing with friends or a team is not. OCP does not change your account terms, and whether any particular sharing setup complies with the ToS is the account holder's responsibility. Review Anthropic's Usage Policy before extending access to other people.
|
||||
- **Account terms are your call.** Claude Pro/Max are *per-user* accounts, and Anthropic's Usage Policy governs who may use them. OCP is a localhost protocol adapter for your own tools and devices — it does not change your account terms, and whether any particular sharing setup complies with the Usage Policy is the account holder's responsibility. Review it before extending access to other people.
|
||||
|
||||
**Real per-user isolation (sandboxed, multi-tenant-safe) is planned for after 2026-06-15** — per-key ephemeral home + tool lockdown + an OS sandbox. Until then, treat a multi-user OCP as a *trusted-group convenience*, not a security boundary. (This is also why `CLAUDE_TUI_MODE` is single-user-only — see [Subscription-pool (TUI) mode](#subscription-pool-tui-mode).)
|
||||
|
||||
@@ -701,14 +700,6 @@ Your IDE → OCP (localhost:3456) → claude --output-format stream-json CLI →
|
||||
|
||||
OCP translates OpenAI-compatible `/v1/chat/completions` requests into `claude --output-format stream-json` CLI calls. Anthropic sees normal Claude Code usage — no API billing, no separate key needed.
|
||||
|
||||
### Client-tools boundary
|
||||
|
||||
OCP is a **text-prompt bridge** to the official `claude` CLI. It does **not** pass through OpenAI `tools`/`functions` payloads or Anthropic `tool_use` blocks to the client. Clients (Cline, Cursor, OpenClaw, etc.) pointed at OCP receive **assistant TEXT only** — they never get `tool_calls` to execute locally.
|
||||
|
||||
Any tool use happens server-side, under the `--allowedTools` set configured on the OCP host. In default mode (no `CLAUDE_NO_CONTEXT`), the `claude` CLI's own built-in tools are available to the model; in TUI mode, the operator controls the tool surface via `OCP_TUI_FULL_TOOLS`. Either way, the tools run under the operator's credentials on the server, and the client sees only the final text output.
|
||||
|
||||
**Client-local tool execution is not supported by design.** Supporting it would require bypassing the `claude` CLI to call the raw Anthropic API directly — that is a different product, and is out of scope per `ALIGNMENT.md` (every OCP endpoint must correspond to something `cli.js` actually does).
|
||||
|
||||
## Available Models
|
||||
|
||||
| Model ID | Notes |
|
||||
@@ -956,7 +947,7 @@ See [Subscription-pool (TUI) mode](#subscription-pool-tui-mode) and ADR 0007 PR-
|
||||
| `OCP_TUI_ENTRYPOINT` | `cli` | (TUI-mode) Billing-classifier labeling: `cli` (default) pins `cc_entrypoint=cli` deterministically; `auto` lets claude self-classify via TTY detection; `off` leaves the inherited env untouched. Honest only when the spawn is a genuine interactive PTY — see ADR 0007. |
|
||||
| `OCP_TUI_MAX_CONCURRENT` | `2` | (TUI-mode) Max concurrent interactive TUI turns. **Independent** of `CLAUDE_MAX_CONCURRENT` (which bounds the `-p`/stream-json path; TUI never uses it). A TUI turn is heavy (per-request cold-boot of tmux+claude + up to `CLAUDE_TUI_WALLCLOCK_MS` wallclock), so the default is low to keep small hosts (e.g. a Pi 4) alive under a burst. Excess turns **queue** (bounded); a full queue yields a 503. See ADR 0007 PR-B amendment. |
|
||||
| `OCP_SKIP_AUTH_TEST` | *(unset)* | When `=1`, skip the `claude -p` auth probe during `setup.mjs`. After 2026-06-15 this probe draws from the Agent SDK credit pool; set this to avoid burning a metered credit on re-installs or `ocp update` runs. Auth is validated at the first real request. |
|
||||
| `OCP_TUI_FULL_TOOLS` | *(unset)* | (TUI-mode, **single-user only**) When `=1`, grant the interactive session the **same tool surface as the `-p` path** — `--allowedTools` (+ optional `--mcp-config`, read from `CLAUDE_ALLOWED_TOOLS` / `CLAUDE_MCP_CONFIG`) — instead of the default MCP-walled, built-in-tools-only set. Lets a trusted single-operator TUI deployment run a **tool-using / MCP agent** (e.g. an OpenClaw assistant) on the subscription pool. Safe because TUI **refuses to boot under `AUTH_MODE=multi`** (hard exit) — no guest key can ever reach the TUI path, so this gate cannot expose tools to an untrusted caller. (Under `AUTH_MODE=shared` + `OCP_TUI_ALLOW_LAN=1`, anyone holding the single shared key reaches it — that is the existing TUI trust model, unchanged.) Note: `--dangerously-skip-permissions` / `CLAUDE_SKIP_PERMISSIONS` is **not** supported for TUI — claude v2.1.x shows an interactive bypass-acceptance screen in headless tmux that cannot be answered, bricking the pane. Use scratch-home `settings.json` `additionalDirectories` instead. See [Subscription-pool (TUI) mode](#subscription-pool-tui-mode) and ADR 0007. |
|
||||
| `OCP_TUI_FULL_TOOLS` | *(unset)* | (TUI-mode, **single-user only**) When `=1`, grant the interactive session the **same tool surface as the `-p` path** — `--allowedTools` (+ optional `--mcp-config` / `--dangerously-skip-permissions`, read from `CLAUDE_ALLOWED_TOOLS` / `CLAUDE_MCP_CONFIG` / `CLAUDE_SKIP_PERMISSIONS`) — instead of the default MCP-walled, built-in-tools-only set. Lets a trusted single-operator TUI deployment run a **tool-using / MCP agent** (e.g. an OpenClaw assistant) on the subscription pool. Safe because TUI **refuses to boot under `AUTH_MODE=multi`** (hard exit) — no guest key can ever reach the TUI path, so this gate cannot expose tools to an untrusted caller. (Under `AUTH_MODE=shared` + `OCP_TUI_ALLOW_LAN=1`, anyone holding the single shared key reaches it — that is the existing TUI trust model, unchanged.) See [Subscription-pool (TUI) mode](#subscription-pool-tui-mode) and ADR 0007. |
|
||||
|
||||
### Streaming heartbeat
|
||||
|
||||
|
||||
@@ -1,100 +0,0 @@
|
||||
# OCP Promotion Strategy — "Stable & Visible"
|
||||
|
||||
> **This document is a recommendation for the maintainer to review and adjust, not a committed plan.**
|
||||
> It reflects the project's current posture (post-v3.21.0) and should be revisited whenever
|
||||
> the Anthropic billing / ToS environment changes significantly.
|
||||
|
||||
---
|
||||
|
||||
## 1. Goal: Polish + Low-Key OSS Visibility
|
||||
|
||||
The goal is **stability and quiet discoverability**, not growth-hacking. OCP is a personal power tool
|
||||
that has been open-sourced because others can benefit from it. The right audience finds it via GitHub
|
||||
search, issue threads in related projects, and word of mouth — not viral posts.
|
||||
|
||||
**Explicitly avoid:**
|
||||
|
||||
- HN / Reddit front-page pushes, influencer outreach, or any campaign that would attract a large
|
||||
influx of users before the ToS/billing situation has settled. Anthropic is actively tightening
|
||||
billing and enforcement on subscription-sharing (the June-15 Agent-SDK billing split is
|
||||
*paused*, not cancelled — and consumer-ToS enforcement on multi-person sharing is a live risk).
|
||||
A high-traffic spotlight right now would draw scrutiny that a low-profile project avoids.
|
||||
- Promising features that require bypassing the `claude` CLI (raw API calls, OAuth extraction, etc.)
|
||||
— that would violate `ALIGNMENT.md` and the ToS simultaneously.
|
||||
|
||||
---
|
||||
|
||||
## 2. Pre-Requisite: Stability First
|
||||
|
||||
Do not promote until the house is in order:
|
||||
|
||||
- [x] The concurrency / latency perf fixes are shipped (v3.20.x–v3.21.0).
|
||||
- [x] Docs honesty is complete (client-tools boundary, ToS sharing disclosure, this doc).
|
||||
- [ ] The June-15 Agent-SDK billing split is either confirmed cancelled or OCP has a confirmed
|
||||
stable path (TUI toggle as insurance — see §5 below).
|
||||
|
||||
Promoting a project that has known rough edges in docs or stability only generates support burden
|
||||
and negative first impressions.
|
||||
|
||||
---
|
||||
|
||||
## 3. Honest ToS Disclosure on Sharing
|
||||
|
||||
Any promotion materials must carry the same disclosure as `README.md § "Deployment model & security"`:
|
||||
|
||||
> Pooling a single Claude subscription across **multiple distinct people** may violate Anthropic's
|
||||
> Consumer Terms of Service and risk account suspension. The defensible framing is "one person,
|
||||
> your own devices". Friends/team sharing is not.
|
||||
|
||||
This framing should appear in any README badge, linked blog post, or issue comment that mentions
|
||||
LAN sharing. It is not a disclaimer that discourages usage — it is honest positioning that protects
|
||||
both the project and its users.
|
||||
|
||||
---
|
||||
|
||||
## 4. What to Explicitly Skip
|
||||
|
||||
These items are **not gaps in OCP** — they are deliberate stance decisions:
|
||||
|
||||
- **Multi-backend routing** (routing to OpenAI, Gemini, Llama, etc.) — that is the sibling [OLP
|
||||
project](https://github.com/dtzp555-max/olp)'s role. OCP stays Claude-only by design.
|
||||
- **Gateway model-discovery** (auto-detecting which models a remote server offers) — not needed
|
||||
for OCP's single-provider, single-subscription model. `models.json` is the SPOT.
|
||||
- **Raw Anthropic API passthrough** (bypassing the `claude` CLI) — out of scope per `ALIGNMENT.md`.
|
||||
|
||||
Do not add these to OCP roadmaps or respond to feature requests for them with "planned" — the
|
||||
correct answer is "that's OLP territory" or "out of scope per ALIGNMENT.md".
|
||||
|
||||
---
|
||||
|
||||
## 5. TUI Toggle as Insurance
|
||||
|
||||
The `CLAUDE_TUI_MODE` opt-in is the primary mitigation if the June-15 billing split reactivates
|
||||
and makes the default `-p` path draw from the metered Agent SDK credit pool.
|
||||
|
||||
Keep the TUI toggle:
|
||||
- Functional and tested across the three deployment hosts.
|
||||
- Documented in the README, including the security constraints (single-user only).
|
||||
- Easily discoverable for users who get unexpectedly metered.
|
||||
|
||||
If the split reactivates, the recommended operator path is: set `CLAUDE_TUI_MODE=true` +
|
||||
`CLAUDE_CODE_OAUTH_TOKEN` → credential-isolated scratch home → subscription pool. That path is
|
||||
already shipped and documented.
|
||||
|
||||
---
|
||||
|
||||
## 6. Low-Key Visibility Actions (when §2 pre-requisites are met)
|
||||
|
||||
- Keep the GitHub README polished and honest — it is the primary landing page.
|
||||
- Respond promptly to issues and PRs — the project's reputation is built on reliability, not
|
||||
marketing.
|
||||
- Add OCP to the `awesome-claude` / `awesome-llm-tools` lists if they exist and allow self-PRs
|
||||
— low-effort, targeted, reaches the right audience.
|
||||
- When related projects (Cline, OpenCode, OpenClaw, Continue.dev) post about local Claude proxies,
|
||||
a short factual comment linking to OCP is appropriate — not spam.
|
||||
- Maintain the `CHANGELOG.md` with clear, honest summaries — users who are already running OCP
|
||||
are the best vector for word-of-mouth.
|
||||
|
||||
---
|
||||
|
||||
*Last updated: v3.21.0 cleanup cycle. Maintainer should re-read before any external promotion.*
|
||||
@@ -382,25 +382,11 @@ export function getCacheStats() {
|
||||
// Per ADR 0005 / spec D4: in-process scope only (single Node process per host).
|
||||
const inflightMap = new Map();
|
||||
|
||||
// `retryIf` (optional, audit finding M1): a predicate applied on the FOLLOWER path only.
|
||||
// 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) {
|
||||
export function singleflight(hash, fn) {
|
||||
const existing = inflightMap.get(hash);
|
||||
if (existing) {
|
||||
existing.requesters++;
|
||||
if (!retryIf) return existing.promise;
|
||||
return existing.promise.catch((err) => {
|
||||
if (!retryIf(err)) throw err;
|
||||
return singleflight(hash, fn, retryIf);
|
||||
});
|
||||
return existing.promise;
|
||||
}
|
||||
// Wrap fn() in Promise.resolve().then() so synchronous throws don't escape.
|
||||
const promise = Promise.resolve().then(fn).finally(() => {
|
||||
|
||||
@@ -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)];
|
||||
}
|
||||
+11
-61
@@ -20,15 +20,6 @@
|
||||
//
|
||||
// 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 {
|
||||
// limit: max concurrent slots. maxQueue: max waiters before run() rejects with backpressure.
|
||||
constructor(limit, { maxQueue } = {}) {
|
||||
@@ -43,30 +34,9 @@ export class TuiSemaphore {
|
||||
get inflight() { return this._inflight; }
|
||||
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
|
||||
// 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
|
||||
// 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"));
|
||||
}
|
||||
acquire() {
|
||||
if (this._inflight < this.limit) {
|
||||
this._inflight++;
|
||||
return Promise.resolve();
|
||||
@@ -76,44 +46,24 @@ export class TuiSemaphore {
|
||||
`tui_queue_full: TUI concurrency limit (${this.limit}) reached and wait queue ` +
|
||||
`(${this.maxQueue}) is full`));
|
||||
}
|
||||
return new Promise((resolve, reject) => {
|
||||
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);
|
||||
});
|
||||
return new Promise((resolve) => { this._waiters.push(resolve); });
|
||||
}
|
||||
|
||||
// Release a slot. Always frees the caller's own slot first, then re-grants it to the next
|
||||
// waiter ONLY if the (post-decrement) inflight count is still under the current limit (F1
|
||||
// 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 a slot. If a waiter is queued, hand the slot directly to it (inflight stays
|
||||
// constant across the handoff); otherwise decrement.
|
||||
release() {
|
||||
if (this._inflight > 0) this._inflight--;
|
||||
if (this._inflight < this.limit) {
|
||||
const next = this._waiters.shift();
|
||||
if (next) {
|
||||
this._inflight++;
|
||||
next();
|
||||
}
|
||||
const next = this._waiters.shift();
|
||||
if (next) {
|
||||
next(); // the woken waiter already "owns" the slot — inflight unchanged
|
||||
} else if (this._inflight > 0) {
|
||||
this._inflight--;
|
||||
}
|
||||
}
|
||||
|
||||
// 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.
|
||||
// `signal` (optional, F2) is forwarded to acquire() so a queued run() can be cancelled.
|
||||
async run(fn, signal) {
|
||||
await this.acquire(signal);
|
||||
async run(fn) {
|
||||
await this.acquire();
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
|
||||
+56
-86
@@ -15,48 +15,14 @@ import { tmpdir } from "node:os";
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { readTuiTranscript } from "./transcript.mjs";
|
||||
|
||||
// F7 fix (audit finding, LOW): the prefix used to be a bare, host-wide constant
|
||||
// ("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}-`;
|
||||
}
|
||||
|
||||
export const SESSION_PREFIX = "ocp-tui-"; // per-proxy namespace (coexistence rule)
|
||||
const TMUX = process.env.OCP_TUI_TMUX_BIN || "tmux";
|
||||
|
||||
const defaultTmux = (args, opts = {}) =>
|
||||
spawnSync(TMUX, args, { encoding: "utf8", ...opts });
|
||||
|
||||
// Kill ONLY our own stale sessions. Scoped to sessionPrefixForPort(port) so a co-hosted
|
||||
// OLP test instance's `olp-tui-*` sessions — AND a co-hosted second OCP instance's
|
||||
// `ocp-tui-<otherPort>-*` sessions — are never touched (F7 fix).
|
||||
// Kill ONLY our own stale sessions. Scoped to SESSION_PREFIX so a co-hosted
|
||||
// OLP test instance's `olp-tui-*` sessions are never touched.
|
||||
//
|
||||
// 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`
|
||||
@@ -70,40 +36,23 @@ const defaultTmux = (args, opts = {}) =>
|
||||
// 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.
|
||||
//
|
||||
// `port` (required) is this instance's own listen port (server.mjs's PORT / lib/constants.mjs
|
||||
// DEFAULT_PORT resolution) — the SPOT for "which sessions are ours."
|
||||
//
|
||||
// `includeLegacy` (default false): when true, sessions matching the exact OLD bare-prefix
|
||||
// shape (LEGACY_SESSION_NAME_RE) are ALSO treated as ours for kill-session purposes. This is
|
||||
// the boot-time legacy migration: an operator upgrading past this fix could otherwise be left
|
||||
// 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 } = {}) {
|
||||
// So after killing our own sessions, if the server has NO sessions left of ANY prefix
|
||||
// (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
|
||||
// server running (coexistence rule, ADR 0007) and let the next boot/periodic sweep retry
|
||||
// once the server is otherwise idle.
|
||||
export function reapStaleTuiSessions({ tmux = defaultTmux } = {}) {
|
||||
const r = tmux(["list-sessions", "-F", "#{session_name}"]);
|
||||
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 ownPrefix = sessionPrefixForPort(port);
|
||||
let killed = 0;
|
||||
let othersRemain = false;
|
||||
for (const name of names) {
|
||||
const isOwn = name.startsWith(ownPrefix);
|
||||
const isLegacyOwn = includeLegacy && LEGACY_SESSION_NAME_RE.test(name);
|
||||
if (isOwn || isLegacyOwn) {
|
||||
if (name.startsWith(SESSION_PREFIX)) {
|
||||
tmux(["kill-session", "-t", name]);
|
||||
killed++;
|
||||
} else {
|
||||
othersRemain = true; // a session we do NOT own (olp-tui-*, a sibling ocp-tui-<otherPort>-*,
|
||||
// or — outside includeLegacy — a legacy-shaped name) — never kill-server
|
||||
othersRemain = true; // a session we do NOT own (e.g. olp-tui-*) — never kill-server
|
||||
}
|
||||
}
|
||||
// Reap defunct `claude` zombies: safe ONLY when the server is now ours-only/empty.
|
||||
@@ -295,6 +244,23 @@ export function prepareTuiHome(realHome, tuiHome, cwd, { envTokenMode = false }
|
||||
ensureTuiCwdTrusted(tuiHome, cwd);
|
||||
}
|
||||
|
||||
// ── Billing-classifier labeling ─────────────────────────────────────────
|
||||
// Resolve CLAUDE_CODE_ENTRYPOINT on the spawn env per mode. ALWAYS deletes any
|
||||
// inherited value first (so a stray entrypoint from OCP's own parent env can never
|
||||
// leak into / mislabel the billing header). Then:
|
||||
// "cli" (default) → set "cli": deterministic subscription-pool classification.
|
||||
// HONEST ONLY because OCP's spawn is a genuine interactive PTY (tmux pane,
|
||||
// no -p, stdout not redirected). Never set "cli" on a non-interactive spawn.
|
||||
// "auto" → leave unset → claude self-classifies via its t$A (TTY → cli). Use to
|
||||
// observe/diagnose the real TTY-derived value.
|
||||
// "off" → leave the env exactly as inherited (diagnostics / honesty audit).
|
||||
export function resolveTuiEntrypointEnv(env, mode = "cli") {
|
||||
if (mode === "off") return env;
|
||||
delete env.CLAUDE_CODE_ENTRYPOINT;
|
||||
if (mode === "cli") env.CLAUDE_CODE_ENTRYPOINT = "cli";
|
||||
return env;
|
||||
}
|
||||
|
||||
// Build interactive claude argv: NO -p, NO --output-format (=> cc_entrypoint=cli).
|
||||
// MCP hard-disabled: --strict-mcp-config (no --mcp-config) is the only mechanism
|
||||
// that stops account-attached managed MCP from connecting (spec §5.2 / T6),
|
||||
@@ -355,24 +321,27 @@ export function buildTuiCmd(claudeBin, model, sessionId, ehome, entrypointMode)
|
||||
// DEFAULT (safe): hard-disable MCP (--strict-mcp-config + --disallowedTools mcp__*);
|
||||
// built-in tools stay on, acceptable for single-user A-path.
|
||||
// OCP_TUI_FULL_TOOLS=1: grant the SAME tool surface as the -p A-path
|
||||
// (--allowedTools [+ --mcp-config]), so a SINGLE-USER / trusted TUI deployment can
|
||||
// run a tool-using agent (e.g. an OpenClaw assistant that needs Bash/Read/Write/MCP)
|
||||
// on the subscription pool. ALWAYS uses --allowedTools (CLAUDE_SKIP_PERMISSIONS /
|
||||
// --dangerously-skip-permissions is intentionally removed: claude v2.1.x shows an
|
||||
// interactive bypass-acceptance screen in headless tmux that nothing can answer →
|
||||
// the turn hangs until the wallclock cap, bricks the pane; not recoverable without a
|
||||
// human at a keyboard). Use scratch-home settings.json additionalDirectories instead.
|
||||
// (--allowedTools [+ --mcp-config] [+ --dangerously-skip-permissions]), so a
|
||||
// SINGLE-USER / trusted TUI deployment can run a tool-using agent (e.g. an OpenClaw
|
||||
// assistant that needs Bash/Read/Write/MCP) on the subscription pool. This mirrors
|
||||
// buildCliArgs() in server.mjs. Safe to gate ON only because TUI is hard-incompatible
|
||||
// with AUTH_MODE=multi (server.mjs refuses to boot), so it can never widen a guest's
|
||||
// surface. Env mirrors server.mjs's CLAUDE_ALLOWED_TOOLS / _SKIP_PERMISSIONS / _MCP_CONFIG.
|
||||
let toolArgs;
|
||||
if (process.env.OCP_TUI_FULL_TOOLS === "1") {
|
||||
toolArgs = [];
|
||||
const allowed = (process.env.CLAUDE_ALLOWED_TOOLS ||
|
||||
"Bash,Read,Write,Edit,Glob,Grep,WebSearch,WebFetch,Agent")
|
||||
.split(",").map((s) => s.trim()).filter(Boolean);
|
||||
// shq EACH token: buildTuiCmd returns a SHELL STRING (run by tmux via sh -c), unlike
|
||||
// buildCliArgs which returns an argv array to spawn(). claude accepts scoped specifiers
|
||||
// like "Bash(npm run test:*)" / "Read(~/**)" whose ( ) * ~ would break/inject the shell
|
||||
// command if pasted bare. (operator-self-injection only — guests can't reach TUI.)
|
||||
if (allowed.length) toolArgs.push("--allowedTools", ...allowed.map(shq));
|
||||
if (process.env.CLAUDE_SKIP_PERMISSIONS === "true") {
|
||||
toolArgs.push("--dangerously-skip-permissions");
|
||||
} else {
|
||||
const allowed = (process.env.CLAUDE_ALLOWED_TOOLS ||
|
||||
"Bash,Read,Write,Edit,Glob,Grep,WebSearch,WebFetch,Agent")
|
||||
.split(",").map((s) => s.trim()).filter(Boolean);
|
||||
// shq EACH token: buildTuiCmd returns a SHELL STRING (run by tmux via sh -c), unlike
|
||||
// buildCliArgs which returns an argv array to spawn(). claude accepts scoped specifiers
|
||||
// like "Bash(npm run test:*)" / "Read(~/**)" whose ( ) * ~ would break/inject the shell
|
||||
// command if pasted bare. (operator-self-injection only — guests can't reach TUI.)
|
||||
if (allowed.length) toolArgs.push("--allowedTools", ...allowed.map(shq));
|
||||
}
|
||||
if (process.env.CLAUDE_MCP_CONFIG) toolArgs.push("--mcp-config", shq(process.env.CLAUDE_MCP_CONFIG));
|
||||
} else {
|
||||
toolArgs = ["--strict-mcp-config", "--disallowedTools", shq("mcp__*")];
|
||||
@@ -409,15 +378,12 @@ export async function runTuiTurn({
|
||||
home,
|
||||
realHome,
|
||||
cwd,
|
||||
port,
|
||||
wallclockMs = 120000,
|
||||
entrypointMode = "cli",
|
||||
tmux = defaultTmux,
|
||||
}) {
|
||||
const sessionId = randomUUID();
|
||||
// Port-scoped session name (F7 fix) — see sessionPrefixForPort / reapStaleTuiSessions
|
||||
// for why this instance's own listen port is the namespace discriminator.
|
||||
const tmuxName = sessionPrefixForPort(port) + sessionId.slice(0, 8);
|
||||
const tmuxName = SESSION_PREFIX + sessionId.slice(0, 8);
|
||||
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)
|
||||
|
||||
@@ -439,11 +405,15 @@ export async function runTuiTurn({
|
||||
const promptFile = `${tmpDir}/prompt.txt`;
|
||||
writeFileSync(promptFile, prompt, { mode: 0o600 });
|
||||
|
||||
// Minimal env for spawnSync (tmux itself). The pane's claude env comes exclusively
|
||||
// from the `env` prefix string built inside buildTuiCmd — tmux does NOT forward the
|
||||
// spawning process's env to the pane, so the {env} here is intentionally minimal.
|
||||
const env = { ...process.env };
|
||||
env.HOME = ehome; // tmux needs HOME; all claude-specific vars go via buildTuiCmd prefix
|
||||
// Build the env: disable marketplace auto-install, strip any Anthropic / CC
|
||||
// env vars that might interfere with interactive-mode classification.
|
||||
const env = { ...process.env, CLAUDE_CODE_DISABLE_OFFICIAL_MARKETPLACE_AUTOINSTALL: "1" };
|
||||
delete env.CLAUDECODE;
|
||||
delete env.ANTHROPIC_API_KEY;
|
||||
delete env.ANTHROPIC_BASE_URL;
|
||||
delete env.ANTHROPIC_AUTH_TOKEN;
|
||||
env.HOME = ehome; // claude reads credentials + writes the transcript under this HOME
|
||||
resolveTuiEntrypointEnv(env, entrypointMode);
|
||||
|
||||
try {
|
||||
// 1. Boot the interactive session inside tmux, rooted at the scratch cwd.
|
||||
|
||||
+15
-2
@@ -9,11 +9,24 @@ import { readFileSync, existsSync, readdirSync } from "node:fs";
|
||||
|
||||
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
|
||||
|
||||
// Project-dir encoding: claude replaces every "/" AND every "." with "-".
|
||||
// Verified live (claude v2.1.158): cwd /home/u/.ocp-tui/work is stored under
|
||||
// projects/-home-u--ocp-tui-work/ (the "." in ".ocp-tui" becomes "-", yielding
|
||||
// the double dash). The earlier "/"-only rule was wrong for dotted paths; the
|
||||
// fixture cwd /tmp/tui-test happened to have no dots so it never surfaced.
|
||||
// NOTE: prefer findTranscriptPath() (glob by session-id) for resolution — it is
|
||||
// immune to the exact encoding rule. This helper is kept for the known-path case.
|
||||
export function encodeCwd(cwd) {
|
||||
return cwd.replace(/[/.]/g, "-");
|
||||
}
|
||||
|
||||
export function transcriptPath(home, cwd, sessionId) {
|
||||
return `${home}/.claude/projects/${encodeCwd(cwd)}/${sessionId}.jsonl`;
|
||||
}
|
||||
|
||||
// Locate a session's transcript by its UUID across every projects subdir, without
|
||||
// reconstructing the encoded cwd. Robust to whatever encoding claude applies.
|
||||
// Returns the path, or null if not present yet (it appears once the turn starts).
|
||||
// TODO: add a CI fixture-contract test (a captured real transcript) so schema drift
|
||||
// in the claude JSONL format fails loudly rather than silently degrading.
|
||||
export function findTranscriptPath(home, sessionId) {
|
||||
if (!home || !sessionId) return null;
|
||||
const root = `${home}/.claude/projects`;
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "open-claude-proxy",
|
||||
"version": "3.21.0",
|
||||
"version": "3.20.1",
|
||||
"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",
|
||||
"bin": {
|
||||
|
||||
+86
-351
@@ -43,8 +43,7 @@ import { DEFAULT_PORT } from "./lib/constants.mjs";
|
||||
import { isLoopbackBind } from "./lib/net.mjs";
|
||||
import { runTuiTurn, reapStaleTuiSessions, resolveTuiHome } from "./lib/tui/session.mjs";
|
||||
import { detectTuiUpstreamError } from "./lib/tui/transcript.mjs";
|
||||
import { TuiSemaphore, SemaphoreAbortError, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs";
|
||||
import { createSerialMutex, createTtlCache, isTokenExpiring, orderLabelsLastGoodFirst } from "./lib/spawn-auth.mjs";
|
||||
import { TuiSemaphore, recordTuiEntrypoint, buildTuiHealthBlock } from "./lib/tui/semaphore.mjs";
|
||||
|
||||
const __dirname = dirname(fileURLToPath(import.meta.url));
|
||||
const _pkg = JSON.parse(readFileSync(join(__dirname, "package.json"), "utf8"));
|
||||
@@ -388,104 +387,29 @@ function prepareSpawnHome(dir = SPAWN_HOME_DIR) {
|
||||
} 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:false → legacy real-HOME spawn, no cwd override (no token, or kill-switch on).
|
||||
//
|
||||
// FIX F6 (2026-07-07): this decision is NO LONGER memoized permanently. The previous version
|
||||
// cached it forever at first call, which meant: (a) credentials appearing after startup never
|
||||
// enabled isolation; (b) `rm -rf ~/.ocp/spawn-home` at runtime made every isolated spawn ENOENT
|
||||
// until restart; (c) during a token-expiry stint /health reported isolated:true while spawns
|
||||
// actually ran real-HOME. Re-evaluating per spawn is cheap because F5's 30s keychain TTL cache
|
||||
// backs getOAuthCredentials(). This function is the CONFIG-level decision (isolated iff a token
|
||||
// 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).
|
||||
// NEVER logs/returns the token verbatim to any caller that logs; spawnClaudeProcess uses it only
|
||||
// to populate the spawn env. token is null when isolation is off.
|
||||
let _spawnHomeMode = null;
|
||||
function getSpawnHomeMode() {
|
||||
if (_spawnHomeMode) return _spawnHomeMode;
|
||||
if (SPAWN_REAL_HOME) {
|
||||
return { isolated: false, home: null, reason: "kill-switch (OCP_SPAWN_REAL_HOME=1)" };
|
||||
_spawnHomeMode = { isolated: false, home: null, token: null, reason: "kill-switch (OCP_SPAWN_REAL_HOME=1)" };
|
||||
return _spawnHomeMode;
|
||||
}
|
||||
let hasToken = false;
|
||||
try { hasToken = !!(getOAuthCredentials()?.accessToken); } catch { hasToken = false; }
|
||||
if (hasToken) return { isolated: true, home: SPAWN_HOME_DIR, reason: "oauth token resolved" };
|
||||
return { isolated: false, home: null, reason: "no oauth token resolvable" };
|
||||
}
|
||||
|
||||
// FIX F6: re-verify the scratch HOME exists before each isolated spawn and re-create it if it was
|
||||
// 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.
|
||||
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
|
||||
// / 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
|
||||
// fall back to real HOME, where the spawned claude refreshes the credential natively and self-heals
|
||||
// (the keychain token is then fresh again → next spawn is fast). The env-token path (Linux) carries
|
||||
// no expiresAt → never expiry-gated (those tokens are long-lived).
|
||||
function resolveSpawnToken() {
|
||||
try {
|
||||
const creds = getOAuthCredentials();
|
||||
if (!creds?.accessToken) return null;
|
||||
if (isTokenExpiring(creds)) return null; // 5-min buffer; applied to the CACHED creds every use
|
||||
return creds.accessToken;
|
||||
} 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();
|
||||
let token = null;
|
||||
try { token = getOAuthCredentials()?.accessToken || null; } catch { token = null; }
|
||||
if (token) {
|
||||
ensureSpawnHome(shm.home);
|
||||
return { isolated: true, home: shm.home, token, releaseFallback: null };
|
||||
prepareSpawnHome(SPAWN_HOME_DIR);
|
||||
_spawnHomeMode = { isolated: true, home: SPAWN_HOME_DIR, token, reason: "oauth token resolved" };
|
||||
} else {
|
||||
_spawnHomeMode = { isolated: false, home: null, token: null, reason: "no oauth token resolvable" };
|
||||
}
|
||||
// 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 };
|
||||
return _spawnHomeMode;
|
||||
}
|
||||
|
||||
// ── FIX ⑥ (concurrency): bounded wait-queue for the -p / stream-json path ──────────────
|
||||
@@ -507,61 +431,16 @@ class ConcurrencyOverflowError extends Error {
|
||||
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
|
||||
// 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
|
||||
// RequestDisconnectedError when `res` closes before a slot is granted (F2) — the caller must not
|
||||
// spawn claude in that case. `res` is optional (back-compat for any caller without a live response
|
||||
// object); omitting it just means a queued wait can't be cancelled early.
|
||||
//
|
||||
// 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
|
||||
// Rejects with ConcurrencyOverflowError when the wait-queue is full. Increments stats.queued while
|
||||
// waiting (decremented on acquire) and stats.queueRejections on overflow.
|
||||
async function acquireClaudeSlot() {
|
||||
stats.queued = claudeSemaphore.queued + 1; // reflect this waiter before we (maybe) block
|
||||
try {
|
||||
await slot;
|
||||
await claudeSemaphore.acquire();
|
||||
} catch (e) {
|
||||
detach();
|
||||
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++;
|
||||
logEvent("warn", "concurrency_queue_full", {
|
||||
limit: claudeSemaphore.limit, maxQueue: claudeSemaphore.maxQueue,
|
||||
@@ -571,7 +450,6 @@ async function acquireClaudeSlot(res) {
|
||||
`backpressure: concurrency limit (${claudeSemaphore.limit}) reached and wait queue ` +
|
||||
`(${claudeSemaphore.maxQueue}) is full — retry shortly`);
|
||||
}
|
||||
detach();
|
||||
stats.queued = claudeSemaphore.queued;
|
||||
let released = false;
|
||||
return function releaseClaudeSlot() {
|
||||
@@ -774,13 +652,7 @@ const TUI_REAP_INTERVAL_MS = 15 * 60 * 1000;
|
||||
const tuiReapInterval = TUI_MODE ? setInterval(() => {
|
||||
if (tuiSemaphore.inflight > 0 || tuiSemaphore.queued > 0) return; // a turn is live — defer
|
||||
try {
|
||||
// F7 fix: scope to THIS instance's own port; a sibling ocp-tui-<otherPort>-* session
|
||||
// (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 });
|
||||
const n = reapStaleTuiSessions();
|
||||
if (n) logEvent("info", "tui_reaped_stale_sessions", { count: n, trigger: "periodic" });
|
||||
} catch (e) { logEvent("error", "tui_periodic_reap_failed", { error: e.message }); }
|
||||
}, TUI_REAP_INTERVAL_MS) : null;
|
||||
@@ -1034,7 +906,7 @@ function getModelTier(cliModel) {
|
||||
// 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
|
||||
// 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;
|
||||
|
||||
// Circuit breaker: disabled (see comment at top of breaker section)
|
||||
@@ -1071,24 +943,22 @@ function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlo
|
||||
env.CLAUDE_CODE_DISABLE_AUTO_MEMORY = "1";
|
||||
}
|
||||
|
||||
// FIX ③ (latency) + F3 (concurrency): apply the pre-resolved per-spawn HOME/token decision.
|
||||
// The decision is resolved ASYNC in the caller (resolveSpawnDecision) so the real-HOME fallback
|
||||
// serialization can await its mutex; here we only apply the result. When isolated, run claude
|
||||
// under a credential-free minimal HOME with cwd = that same neutral dir, so it loads NONE of the
|
||||
// operator's global ~/.claude (plugins/skills/hooks) or the ~/ocp project CLAUDE.md/skills — the
|
||||
// measured 10–28s → 3–7s latency win. The env token is authoritative for `-p` (unlike
|
||||
// interactive claude). When no fresh token is resolvable, decision.isolated is false → real HOME
|
||||
// + inherited cwd (zero regression), and the spawned claude resolves+refreshes credentials
|
||||
// natively. The DISABLE_CLAUDE_MDS / AUTO_MEMORY flags are set unconditionally in isolated mode
|
||||
// (belt-and-braces; mirrors the TUI path).
|
||||
const decision = spawnDecision || { isolated: false, releaseFallback: null };
|
||||
// FIX ③ (latency): default-path spawn-home isolation. When a token is resolvable (and the
|
||||
// OCP_SPAWN_REAL_HOME kill-switch is off), run claude under a credential-free minimal HOME
|
||||
// with cwd = that same neutral dir, so it loads NONE of the operator's global ~/.claude
|
||||
// (plugins/skills/hooks) or the ~/ocp project CLAUDE.md/skills — the measured 10–28s → 3–7s
|
||||
// latency win. The env token is authoritative for `-p` (unlike interactive claude). When no
|
||||
// token is resolvable, falls back to real HOME + inherited cwd (zero regression). See
|
||||
// getSpawnHomeMode() / prepareSpawnHome() above. The DISABLE_CLAUDE_MDS / AUTO_MEMORY flags
|
||||
// are set unconditionally in isolated mode (belt-and-braces; mirrors the TUI path).
|
||||
const spawnHome = getSpawnHomeMode();
|
||||
const spawnOpts = { env, stdio: ["pipe", "pipe", "pipe"] };
|
||||
if (decision.isolated && decision.token) {
|
||||
env.HOME = decision.home;
|
||||
env.CLAUDE_CODE_OAUTH_TOKEN = decision.token; // env token is authoritative for -p
|
||||
if (spawnHome.isolated) {
|
||||
env.HOME = spawnHome.home;
|
||||
env.CLAUDE_CODE_OAUTH_TOKEN = spawnHome.token; // env token is authoritative for -p
|
||||
env.CLAUDE_CODE_DISABLE_CLAUDE_MDS = "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);
|
||||
@@ -1107,11 +977,6 @@ function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlo
|
||||
// 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).
|
||||
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,
|
||||
@@ -1181,37 +1046,18 @@ function spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlo
|
||||
// We accumulate full text across all content_block_delta events plus the
|
||||
// assistant-aggregate fallback, then resolve with the assembled string.
|
||||
// Reference: OLP ADR 0009 Amendment 1 + commit 97e7d16.
|
||||
// `res` (optional, F2) is the client's http.ServerResponse — passed through so a queued wait
|
||||
// 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) {
|
||||
async function callClaude(model, messages, conversationId, keyName) {
|
||||
// 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)
|
||||
// if the client goes away first). The release fn is passed into the spawn so the idempotent
|
||||
// cleanup() frees it on every exit path. If the spawn itself throws synchronously (before
|
||||
// cleanup is wired), release here so the slot never leaks.
|
||||
// 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;
|
||||
}
|
||||
// ConcurrencyOverflowError → 429 when the queue is full). The release fn is passed into the
|
||||
// spawn so the idempotent cleanup() frees it on every exit path. If the spawn itself throws
|
||||
// synchronously (before cleanup is wired), release here so the slot never leaks.
|
||||
const releaseSlot = await acquireClaudeSlot();
|
||||
return new Promise((resolve, reject) => {
|
||||
let ctx;
|
||||
try {
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot, spawnDecision);
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot);
|
||||
} catch (err) {
|
||||
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);
|
||||
}
|
||||
|
||||
@@ -1282,50 +1128,26 @@ async function callClaude(model, messages, conversationId, keyName, res) {
|
||||
// flag that could perturb cc_entrypoint classification.
|
||||
// Authority: claude CLI v2.1.158 interactive mode (cc_entrypoint=cli).
|
||||
// SECURITY: A-path single-user ONLY — home is NOT isolation (see ADR 0007).
|
||||
// `res` (optional, F2) is the client's http.ServerResponse — see closeSignalFor.
|
||||
async function callClaudeTui(model, messages, _conversationId, _keyName, res) {
|
||||
function callClaudeTui(model, messages, _conversationId, _keyName) {
|
||||
const cliModel = MODEL_MAP[model] || model;
|
||||
const prompt = messagesToPrompt(messages); // includes system as [System] inline
|
||||
recordModelRequest(cliModel, prompt.length);
|
||||
// C-4: gate the heavy interactive boot behind the TUI semaphore (queuing if all slots are
|
||||
// busy, up to maxQueue). F2: `signal` (tied to `res` "close") cancels a QUEUED wait the
|
||||
// instant the client disconnects, so a dead socket never triggers a cold-boot tmux+claude
|
||||
// spawn; detach() drops the "close" listener as soon as the wait settles rather than
|
||||
// holding it for the whole (up to 120s) turn.
|
||||
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,
|
||||
model: cliModel,
|
||||
claudeBin: CLAUDE,
|
||||
home: TUI_HOME,
|
||||
realHome: process.env.HOME,
|
||||
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,
|
||||
entrypointMode: TUI_ENTRYPOINT,
|
||||
});
|
||||
// C-4: gate the heavy interactive boot behind the TUI semaphore. run() acquires a slot
|
||||
// (queuing if all are busy, up to maxQueue), then releases 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.
|
||||
return tuiSemaphore.run(() => runTuiTurn({
|
||||
prompt,
|
||||
model: cliModel,
|
||||
claudeBin: CLAUDE,
|
||||
home: TUI_HOME,
|
||||
realHome: process.env.HOME,
|
||||
cwd: TUI_CWD,
|
||||
wallclockMs: TUI_WALLCLOCK_MS,
|
||||
entrypointMode: TUI_ENTRYPOINT,
|
||||
}).then(({ text, entrypoint, truncated }) => {
|
||||
// ── 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.
|
||||
|
||||
// C-2: the wall-clock cap hit with partial text and NO terminal marker — the turn
|
||||
@@ -1360,12 +1182,10 @@ async function callClaudeTui(model, messages, _conversationId, _keyName, res) {
|
||||
logEvent("warn", "tui_entrypoint_mismatch", { expected: "cli", got: entrypoint, model: cliModel });
|
||||
}
|
||||
return text;
|
||||
} catch (err) {
|
||||
}).catch((err) => {
|
||||
recordModelError(cliModel, false);
|
||||
throw err;
|
||||
} finally {
|
||||
tuiSemaphore.release();
|
||||
}
|
||||
}));
|
||||
}
|
||||
|
||||
// ── SSE heartbeat (opt-in idle watchdog) ────────────────────────────────
|
||||
@@ -1411,37 +1231,21 @@ async function callClaudeStreaming(model, messages, conversationId, res, authInf
|
||||
// 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
|
||||
// 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;
|
||||
try {
|
||||
releaseSlot = await acquireClaudeSlot(res);
|
||||
releaseSlot = await acquireClaudeSlot();
|
||||
} catch (err) {
|
||||
if (err instanceof RequestDisconnectedError) return; // client gone — nothing to write to
|
||||
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, 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;
|
||||
try {
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, authInfo.keyName, releaseSlot, spawnDecision);
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, authInfo.keyName, releaseSlot);
|
||||
} catch (err) {
|
||||
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" } });
|
||||
}
|
||||
|
||||
@@ -1713,51 +1517,6 @@ const OAUTH_REFRESH_MIN_BACKOFF = 60 * 1000;
|
||||
const OAUTH_REFRESH_MAX_BACKOFF = 3600 * 1000;
|
||||
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() {
|
||||
// 1. Env var fallback — highest precedence for explicit overrides.
|
||||
if (process.env.CLAUDE_CODE_OAUTH_TOKEN) {
|
||||
@@ -1771,8 +1530,17 @@ function getOAuthCredentials() {
|
||||
if (creds?.claudeAiOauth?.accessToken) return creds.claudeAiOauth;
|
||||
} catch { /* fall through to macOS keychain */ }
|
||||
|
||||
// 3. macOS keychain (both label formats) — F5: label-memoized + 30s TTL cached (see above).
|
||||
return readKeychainCreds();
|
||||
// 3. macOS keychain (both label formats)
|
||||
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) {
|
||||
@@ -2106,12 +1874,9 @@ function applySettingUpdate(key, value) {
|
||||
|
||||
switch (key) {
|
||||
case "timeout": TIMEOUT = value; break;
|
||||
// FIX ⑥ + F1: 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 —
|
||||
// in BOTH directions. setLimit() (not a bare `.limit =` assignment) is required: lowering
|
||||
// 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;
|
||||
// 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.
|
||||
case "maxConcurrent": MAX_CONCURRENT = value; claudeSemaphore.limit = Math.max(1, value); break;
|
||||
case "sessionTTL": SESSION_TTL = value; break;
|
||||
case "maxPromptChars": MAX_PROMPT_CHARS = value; break;
|
||||
case "cacheTTL": CACHE_TTL = value; break;
|
||||
@@ -2268,7 +2033,7 @@ async function handleChatCompletions(req, res) {
|
||||
const t0TuiStream = Date.now();
|
||||
const promptCharsTuiStream = messages.reduce((a, m) => a + contentToText(m.content).length, 0);
|
||||
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) {
|
||||
try { setCachedResponse(req._cacheHash, model, content); } catch (e) { logEvent("error", "cache_write_failed", { error: e.message }); }
|
||||
}
|
||||
@@ -2305,27 +2070,15 @@ async function handleChatCompletions(req, res) {
|
||||
// will re-read the freshly-populated cache entry here rather than spawning.
|
||||
const recheck = getCachedResponse(req._cacheHash, CACHE_TTL);
|
||||
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 }); }
|
||||
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()}`;
|
||||
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 }); }
|
||||
return;
|
||||
} 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 }); }
|
||||
console.error(`[proxy] error: ${err.message}`);
|
||||
if (res.headersSent || res.writableEnded || res.destroyed) {
|
||||
@@ -2338,14 +2091,11 @@ async function handleChatCompletions(req, res) {
|
||||
|
||||
// Fallback: cache disabled (CACHE_TTL=0) or no _cacheHash — original path untouched.
|
||||
try {
|
||||
const content = await upstreamCall(model, messages, conversationId, req._authKeyName, res);
|
||||
const content = await upstreamCall(model, messages, conversationId, req._authKeyName);
|
||||
const id = `chatcmpl-${randomUUID()}`;
|
||||
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 }); }
|
||||
} 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 }); }
|
||||
console.error(`[proxy] error: ${err.message}`);
|
||||
if (res.headersSent || res.writableEnded || res.destroyed) {
|
||||
@@ -2523,22 +2273,11 @@ const server = createServer(async (req, res) => {
|
||||
spawn: (() => {
|
||||
if (TUI_MODE) return { mode: "tui (default -p path unused)", isolated: false, home: null };
|
||||
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 {
|
||||
mode: effIsolated ? "isolated-scratch-home" : "real-home",
|
||||
isolated: effIsolated,
|
||||
home: effIsolated ? shm.home : null,
|
||||
reason: effIsolated
|
||||
? shm.reason
|
||||
: (shm.isolated
|
||||
? "oauth token within 5-min expiry window → real-HOME fallback (self-heals on next refresh)"
|
||||
: shm.reason),
|
||||
mode: shm.isolated ? "isolated-scratch-home" : "real-home",
|
||||
isolated: shm.isolated,
|
||||
home: shm.isolated ? shm.home : null,
|
||||
reason: shm.reason,
|
||||
};
|
||||
})(),
|
||||
// ── FIX ⑥ -p concurrency wait-queue surface — ADDITIVE ──
|
||||
@@ -2867,11 +2606,7 @@ server.listen(PORT, BIND_ADDRESS, () => {
|
||||
: "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}`);
|
||||
try {
|
||||
// F7 fix: scope to THIS instance's own port (see reapStaleTuiSessions). includeLegacy:
|
||||
// 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 });
|
||||
const n = reapStaleTuiSessions();
|
||||
if (n) logEvent("info", "tui_reaped_stale_sessions", { count: n });
|
||||
} catch {}
|
||||
}
|
||||
|
||||
+89
-503
@@ -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 { isLoopbackBind } from "./lib/net.mjs";
|
||||
import { createSerialMutex, createTtlCache, isTokenExpiring, orderLabelsLastGoodFirst } from "./lib/spawn-auth.mjs";
|
||||
import { createHash } from "node:crypto";
|
||||
import { strict as assert } from "node:assert";
|
||||
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");
|
||||
|
||||
// Initialize DB
|
||||
@@ -463,52 +451,6 @@ async function runSingleflightTests() {
|
||||
assert.equal(r1, 1);
|
||||
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();
|
||||
@@ -1398,12 +1340,23 @@ test("streamStringAsSSE empty content: role + stop + [DONE] only", () => {
|
||||
});
|
||||
|
||||
// ── Suite: TUI transcript reader ────────────────────────────────────────
|
||||
import { findTranscriptPath, parseTranscriptLines, isTerminalLine, extractLatestAssistantText, verifyEntrypoint, detectTuiUpstreamError } from "./lib/tui/transcript.mjs";
|
||||
import { encodeCwd, transcriptPath, findTranscriptPath, parseTranscriptLines, isTerminalLine, extractLatestAssistantText, verifyEntrypoint, detectTuiUpstreamError } from "./lib/tui/transcript.mjs";
|
||||
import { readFileSync as tuiReadFileSync, mkdtempSync as tuiMkdtemp0, mkdirSync as tuiMkdir0, writeFileSync as tuiWrite0 } from "node:fs";
|
||||
import { tmpdir as tuiTmp0 } from "node:os";
|
||||
|
||||
console.log("\nTUI transcript — path formula:");
|
||||
|
||||
test("encodeCwd replaces every slash AND every dot with dash", () => {
|
||||
// Verified live (claude v2.1.158): /home/u/.ocp-tui/work -> -home-u--ocp-tui-work
|
||||
assert.equal(encodeCwd("/home/u/.ocp-tui/work"), "-home-u--ocp-tui-work");
|
||||
assert.equal(encodeCwd("/tmp/tui-test"), "-tmp-tui-test"); // dot-free path still correct
|
||||
});
|
||||
test("transcriptPath composes HOME/.claude/projects/<enc>/<sid>.jsonl", () => {
|
||||
assert.equal(
|
||||
transcriptPath("/home/u", "/home/u/.ocp-tui/work", "abc-123"),
|
||||
"/home/u/.claude/projects/-home-u--ocp-tui-work/abc-123.jsonl"
|
||||
);
|
||||
});
|
||||
test("findTranscriptPath locates <sid>.jsonl across projects subdirs by UUID", () => {
|
||||
const home = tuiMkdtemp0(`${tuiTmp0()}/tui-home-`);
|
||||
const sid = "11111111-2222-3333-4444-555555555555";
|
||||
@@ -1710,23 +1663,12 @@ await asyncTest("readTuiTranscript throws when no text and cap elapses", async (
|
||||
});
|
||||
|
||||
// ── 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:");
|
||||
|
||||
// F7 fix: the session prefix is instance-scoped by listen port so a second OCP
|
||||
// instance on the same host (different port) is never mistaken for "ours".
|
||||
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");
|
||||
test("SESSION_PREFIX is ocp-tui-", () => {
|
||||
assert.equal(SESSION_PREFIX, "ocp-tui-");
|
||||
});
|
||||
|
||||
console.log("\nTUI command construction (proxy-purity / #4):");
|
||||
@@ -1798,7 +1740,7 @@ test("buildTuiCmd shq-escapes a token containing shell metacharacters (no inject
|
||||
test("buildTuiCmd OCP_TUI_FULL_TOOLS=1 grants -p-equivalent tool surface (single-user opt-in)", () => {
|
||||
const save = { ...process.env };
|
||||
const restore = () => {
|
||||
for (const k of ["OCP_TUI_FULL_TOOLS", "CLAUDE_MCP_CONFIG", "CLAUDE_ALLOWED_TOOLS"]) {
|
||||
for (const k of ["OCP_TUI_FULL_TOOLS", "CLAUDE_SKIP_PERMISSIONS", "CLAUDE_MCP_CONFIG", "CLAUDE_ALLOWED_TOOLS"]) {
|
||||
if (k in save) process.env[k] = save[k]; else delete process.env[k];
|
||||
}
|
||||
};
|
||||
@@ -1810,14 +1752,20 @@ test("buildTuiCmd OCP_TUI_FULL_TOOLS=1 grants -p-equivalent tool surface (single
|
||||
|
||||
// gate on: --allowedTools (default set incl Bash), MCP wall dropped
|
||||
process.env.OCP_TUI_FULL_TOOLS = "1";
|
||||
delete process.env.CLAUDE_SKIP_PERMISSIONS;
|
||||
delete process.env.CLAUDE_MCP_CONFIG;
|
||||
delete process.env.CLAUDE_ALLOWED_TOOLS;
|
||||
const full = buildTuiCmd("/usr/bin/claude", "m", "s", "/home/u", "cli");
|
||||
assert.ok(full.includes("--allowedTools") && full.includes("Bash"), "full-tools grants --allowedTools incl Bash");
|
||||
assert.ok(!full.includes("--strict-mcp-config") && !/--disallowedTools/.test(full), "full-tools drops the MCP wall");
|
||||
assert.ok(!full.includes("--dangerously-skip-permissions"), "skip-permissions branch is removed (bricks headless TUI)");
|
||||
|
||||
// skip-permissions supersedes --allowedTools
|
||||
process.env.CLAUDE_SKIP_PERMISSIONS = "true";
|
||||
const skip = buildTuiCmd("/usr/bin/claude", "m", "s", "/home/u", "cli");
|
||||
assert.ok(skip.includes("--dangerously-skip-permissions") && !skip.includes("--allowedTools"), "skip-permissions honored");
|
||||
|
||||
// mcp-config threaded through
|
||||
delete process.env.CLAUDE_SKIP_PERMISSIONS;
|
||||
process.env.CLAUDE_MCP_CONFIG = "/tmp/mcp.json";
|
||||
const mcp = buildTuiCmd("/usr/bin/claude", "m", "s", "/home/u", "cli");
|
||||
assert.ok(/--mcp-config '\/tmp\/mcp.json'/.test(mcp), "mcp-config passed through (shq'd)");
|
||||
@@ -1833,40 +1781,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 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 }; }
|
||||
return { status: 0, stdout: "" };
|
||||
};
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 });
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux });
|
||||
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");
|
||||
});
|
||||
|
||||
// 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)", () => {
|
||||
const fakeTmux = (_args) => ({ status: 1, stdout: "" });
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 });
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux });
|
||||
assert.equal(n, 0);
|
||||
});
|
||||
|
||||
@@ -1877,7 +1807,7 @@ test("reaper returns 0 for empty session list", () => {
|
||||
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 });
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux });
|
||||
assert.equal(n, 0);
|
||||
assert.equal(killed.length, 0);
|
||||
});
|
||||
@@ -1890,10 +1820,10 @@ test("reaper kill-servers when the server is ours-only (flush defunct claude zom
|
||||
const calls = [];
|
||||
const fakeTmux = (args) => {
|
||||
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: "" };
|
||||
};
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 });
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux });
|
||||
assert.equal(n, 2, "killed both of our sessions");
|
||||
assert.ok(calls.includes("kill-server"), "kill-server fired — reaps the defunct backlog");
|
||||
});
|
||||
@@ -1902,10 +1832,10 @@ test("reaper does NOT kill-server when a foreign (non-ocp) session remains (coex
|
||||
const calls = [];
|
||||
const fakeTmux = (args) => {
|
||||
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: "" };
|
||||
};
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux, port: 3456 });
|
||||
const n = reapStaleTuiSessions({ tmux: fakeTmux });
|
||||
assert.equal(n, 1, "killed only our own session");
|
||||
assert.ok(!calls.includes("kill-server"), "kill-server MUST NOT fire — would disrupt olp-tui-*");
|
||||
});
|
||||
@@ -1913,61 +1843,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)", () => {
|
||||
const calls = [];
|
||||
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)");
|
||||
});
|
||||
|
||||
// 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) ───────────────────────────────
|
||||
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";
|
||||
@@ -2050,8 +1929,58 @@ test("resolveTuiHome: explicit OCP_TUI_HOME wins regardless of env token (back-c
|
||||
assert.equal(resolveTuiHome({ realHome: "/home/u", configuredHome: "/custom/home", envTokenSet: false }), "/custom/home");
|
||||
});
|
||||
|
||||
// ── resolveTuiEntrypointEnv ───────────────────────────────────────────────
|
||||
import { resolveTuiEntrypointEnv } from "./lib/tui/session.mjs";
|
||||
|
||||
console.log("\nresolveTuiEntrypointEnv:");
|
||||
|
||||
test("mode 'cli' sets CLAUDE_CODE_ENTRYPOINT=cli", () => {
|
||||
const env = {};
|
||||
resolveTuiEntrypointEnv(env, "cli");
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, "cli");
|
||||
});
|
||||
|
||||
test("mode 'cli' overwrites an inherited CLAUDE_CODE_ENTRYPOINT value", () => {
|
||||
const env = { CLAUDE_CODE_ENTRYPOINT: "sdk-cli" };
|
||||
resolveTuiEntrypointEnv(env, "cli");
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, "cli");
|
||||
});
|
||||
|
||||
test("mode 'auto' deletes CLAUDE_CODE_ENTRYPOINT (leaves unset)", () => {
|
||||
const env = {};
|
||||
resolveTuiEntrypointEnv(env, "auto");
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, undefined);
|
||||
assert.ok(!Object.prototype.hasOwnProperty.call(env, "CLAUDE_CODE_ENTRYPOINT"));
|
||||
});
|
||||
|
||||
test("mode 'auto' deletes an inherited CLAUDE_CODE_ENTRYPOINT value", () => {
|
||||
const env = { CLAUDE_CODE_ENTRYPOINT: "sdk-cli" };
|
||||
resolveTuiEntrypointEnv(env, "auto");
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, undefined);
|
||||
assert.ok(!Object.prototype.hasOwnProperty.call(env, "CLAUDE_CODE_ENTRYPOINT"));
|
||||
});
|
||||
|
||||
test("mode 'off' leaves an inherited CLAUDE_CODE_ENTRYPOINT value untouched", () => {
|
||||
const env = { CLAUDE_CODE_ENTRYPOINT: "sdk-cli" };
|
||||
resolveTuiEntrypointEnv(env, "off");
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, "sdk-cli");
|
||||
});
|
||||
|
||||
test("mode 'off' with no inherited value leaves env unchanged", () => {
|
||||
const env = { OTHER: "x" };
|
||||
resolveTuiEntrypointEnv(env, "off");
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, undefined);
|
||||
assert.equal(env.OTHER, "x");
|
||||
});
|
||||
|
||||
test("default mode (no second arg) behaves like 'cli'", () => {
|
||||
const env = { CLAUDE_CODE_ENTRYPOINT: "sdk-cli" };
|
||||
resolveTuiEntrypointEnv(env);
|
||||
assert.equal(env.CLAUDE_CODE_ENTRYPOINT, "cli");
|
||||
});
|
||||
|
||||
// ── 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):");
|
||||
|
||||
@@ -2156,226 +2085,6 @@ await asyncTest("FIX ⑥: slot released on normal completion is immediately reus
|
||||
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):");
|
||||
|
||||
test("recordTuiEntrypoint: observed 'cli' is NOT a mismatch and sets lastEntrypoint", () => {
|
||||
@@ -2756,131 +2465,8 @@ test("isLoopbackBind: '100.64.0.1' → false (Tailscale IP)", () => {
|
||||
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 ──
|
||||
runAsyncTests().then(() => {
|
||||
closeDb();
|
||||
console.log(`\n=== Results: ${passed} passed, ${failed} failed ===\n`);
|
||||
process.exit(failed > 0 ? 1 : 0);
|
||||
}).catch((e) => {
|
||||
console.error("async test runner crashed:", e);
|
||||
closeDb();
|
||||
process.exit(1);
|
||||
});
|
||||
closeDb();
|
||||
|
||||
console.log(`\n=== Results: ${passed} passed, ${failed} failed ===\n`);
|
||||
process.exit(failed > 0 ? 1 : 0);
|
||||
|
||||
Reference in New Issue
Block a user