mirror of
https://github.com/dtzp555-max/ocp.git
synced 2026-07-21 21:15:09 +00:00
fix(server): bounded wait-queue + HTTP 429 for -p concurrency overflow (⑥)
spawnClaudeProcess used `if (activeRequests >= MAX_CONCURRENT) throw` → the client
got an opaque 500 AND the rejection was uncounted (a 15-concurrent stress run returned
7×500 while /health stats.errors stayed 0). The TUI path already had a bounded-queue
semaphore (TuiSemaphore); the -p path did not.
Now the -p path reuses TuiSemaphore as `claudeSemaphore = new TuiSemaphore(MAX_CONCURRENT,
{ maxQueue: CLAUDE_MAX_QUEUE })`. Requests beyond MAX_CONCURRENT WAIT (up to CLAUDE_MAX_QUEUE,
default 16) instead of being rejected; only when the queue is ALSO full does the request get
HTTP 429 + Retry-After (rate_limit_error, NOT 500), a distinct `concurrency_queue_full` log,
and a stats.queueRejections counter surfaced on /health. callClaude / callClaudeStreaming now
acquire a slot (acquireClaudeSlot) before spawning; the release fn is wired into the existing
idempotent cleanup() so the slot is freed on EVERY exit path (close/error/timeout/abort) — the
#37/#40 slot-leak guard. MAX_CONCURRENT semantics (max concurrent claude procs) are unchanged;
only overflow handling changed from throw-500 to queue-then-429. claudeSemaphore.limit is kept
in sync with runtime /settings maxConcurrent changes.
Live-verified (fake claude, MAX_CONCURRENT=1 MAX_QUEUE=1): 3 concurrent → exactly one 429
(Retry-After: 7, rate_limit_error body) + two 200s, on BOTH the non-streaming and streaming
paths; /health stats.queueRejections=1 while stats.errors=0; after drain activeRequests=0 /
inflight=0 / queued=0 and a follow-up request returns 200 (no slot leak). 2 new unit tests
(247 passed, 0 failed).
ALIGNMENT.md Rule 2 justification: this is an INFRA / PROCESS-MANAGEMENT change (a concurrency
queue + backpressure status code in front of the existing spawn). It does NOT mirror a cli.js
wire operation, adds no new endpoint or wire header, and introduces no API token — so there is
no cli.js function to cite. (Retry-After is a standard HTTP response header on OCP's own 429,
not a claude wire header.)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
+126
-16
@@ -275,6 +275,15 @@ const SYSTEM_PROMPT = process.env.CLAUDE_SYSTEM_PROMPT || "";
|
||||
const MCP_CONFIG = process.env.CLAUDE_MCP_CONFIG || "";
|
||||
let SESSION_TTL = parseInt(process.env.CLAUDE_SESSION_TTL || "3600000", 10);
|
||||
let MAX_CONCURRENT = parseInt(process.env.CLAUDE_MAX_CONCURRENT || "8", 10);
|
||||
// FIX ⑥ (concurrency): bound on requests WAITING for a -p concurrency slot. Beyond
|
||||
// MAX_CONCURRENT, requests queue (up to CLAUDE_MAX_QUEUE) instead of being rejected; when the
|
||||
// queue is ALSO full, the request gets HTTP 429 + Retry-After (not an opaque 500). See
|
||||
// claudeSemaphore / acquireClaudeSlot below.
|
||||
const CLAUDE_MAX_QUEUE = parseInt(process.env.CLAUDE_MAX_QUEUE || "16", 10);
|
||||
// Retry-After seconds advertised on a 429 backpressure response. A claude turn is typically a
|
||||
// few seconds to tens of seconds; a small constant nudge keeps well-behaved clients from
|
||||
// hammering while the queue drains.
|
||||
const CLAUDE_QUEUE_RETRY_AFTER = parseInt(process.env.CLAUDE_QUEUE_RETRY_AFTER || "5", 10);
|
||||
const BREAKER_THRESHOLD = parseInt(process.env.CLAUDE_BREAKER_THRESHOLD || "6", 10);
|
||||
const BREAKER_COOLDOWN = parseInt(process.env.CLAUDE_BREAKER_COOLDOWN || "120000", 10);
|
||||
const BREAKER_WINDOW = parseInt(process.env.CLAUDE_BREAKER_WINDOW || "300000", 10);
|
||||
@@ -403,6 +412,54 @@ function getSpawnHomeMode() {
|
||||
return _spawnHomeMode;
|
||||
}
|
||||
|
||||
// ── FIX ⑥ (concurrency): bounded wait-queue for the -p / stream-json path ──────────────
|
||||
// PROBLEM (proven): spawnClaudeProcess used `if (activeRequests >= MAX_CONCURRENT) throw` →
|
||||
// the client got an opaque 500 AND the rejection was NOT counted in stats (a 15-concurrent
|
||||
// stress run returned 7×500 while /health stats.errors stayed 0). The TUI path already has a
|
||||
// bounded-queue semaphore (TuiSemaphore); the -p path did not.
|
||||
//
|
||||
// FIX: requests beyond MAX_CONCURRENT WAIT on this semaphore (up to CLAUDE_MAX_QUEUE) instead of
|
||||
// being rejected. Only when the queue is ALSO full do we reject — with HTTP 429 + Retry-After
|
||||
// (deterministic backpressure), a distinct `concurrency_queue_full` log, and a stats.queueRejections
|
||||
// counter that shows up on /health. The slot is released on EVERY exit path via the existing
|
||||
// idempotent cleanup() (proc exit/close/error/timeout) — the #37/#40 slot-leak guard.
|
||||
const claudeSemaphore = new TuiSemaphore(MAX_CONCURRENT, { maxQueue: CLAUDE_MAX_QUEUE });
|
||||
|
||||
// Tagged error so callers can map this single overflow case to HTTP 429 (every OTHER throw stays
|
||||
// a 500). Carries retryAfter for the Retry-After header.
|
||||
class ConcurrencyOverflowError extends Error {
|
||||
constructor(message) { super(message); this.name = "ConcurrencyOverflowError"; this.httpStatus = 429; this.retryAfter = CLAUDE_QUEUE_RETRY_AFTER; }
|
||||
}
|
||||
|
||||
// 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. 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 claudeSemaphore.acquire();
|
||||
} catch (e) {
|
||||
stats.queued = claudeSemaphore.queued;
|
||||
stats.queueRejections++;
|
||||
logEvent("warn", "concurrency_queue_full", {
|
||||
limit: claudeSemaphore.limit, maxQueue: claudeSemaphore.maxQueue,
|
||||
inflight: claudeSemaphore.inflight, queued: claudeSemaphore.queued,
|
||||
});
|
||||
throw new ConcurrencyOverflowError(
|
||||
`backpressure: concurrency limit (${claudeSemaphore.limit}) reached and wait queue ` +
|
||||
`(${claudeSemaphore.maxQueue}) is full — retry shortly`);
|
||||
}
|
||||
stats.queued = claudeSemaphore.queued;
|
||||
let released = false;
|
||||
return function releaseClaudeSlot() {
|
||||
if (released) return; // idempotent — cleanup() may be reached via multiple proc events
|
||||
released = true;
|
||||
claudeSemaphore.release();
|
||||
stats.queued = claudeSemaphore.queued;
|
||||
};
|
||||
}
|
||||
|
||||
// SECURITY fail-loud: TUI-mode is incompatible with any configuration that allows
|
||||
// non-operator prompts to reach the interactive claude session. Three cases:
|
||||
// 1. AUTH_MODE=multi — guest/anonymous keys can submit prompts.
|
||||
@@ -613,6 +670,8 @@ const stats = {
|
||||
sessionHits: 0,
|
||||
sessionMisses: 0,
|
||||
oneOffRequests: 0,
|
||||
queued: 0, // current requests waiting for a -p concurrency slot (FIX ⑥)
|
||||
queueRejections: 0, // total requests rejected with HTTP 429 because the wait-queue was full (FIX ⑥)
|
||||
};
|
||||
const recentErrors = []; // last 20 errors
|
||||
|
||||
@@ -840,11 +899,14 @@ function getModelTier(cliModel) {
|
||||
// (messagesToPrompt), so multi-turn correctness is preserved without sessions.
|
||||
// The sessions Map is retained for stats/logging but no longer drives --resume.
|
||||
// Reference: OLP ADR 0009 Amendment 1 + commit 97e7d16.
|
||||
function spawnClaudeProcess(model, messages, conversationId, keyName) {
|
||||
if (stats.activeRequests >= MAX_CONCURRENT) {
|
||||
throw new Error(`concurrency limit reached (${stats.activeRequests}/${MAX_CONCURRENT})`);
|
||||
}
|
||||
|
||||
// FIX ⑥: concurrency is now bounded by the claudeSemaphore via acquireClaudeSlot(), which the
|
||||
// caller MUST await before calling this, passing the resulting release fn as `releaseSlot`. The
|
||||
// old `if (activeRequests >= MAX_CONCURRENT) throw` gate (→ opaque 500, uncounted) is GONE: at
|
||||
// most MAX_CONCURRENT callers hold a slot when they reach here, so this spawn is always within
|
||||
// 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 = () => {}) {
|
||||
const cliModel = MODEL_MAP[model] || model;
|
||||
|
||||
// Circuit breaker: disabled (see comment at top of breaker section)
|
||||
@@ -911,6 +973,10 @@ function spawnClaudeProcess(model, messages, conversationId, keyName) {
|
||||
cleaned = true;
|
||||
clearTimeout(overallTimer);
|
||||
stats.activeRequests--;
|
||||
// FIX ⑥: free the concurrency slot for a queued waiter. releaseSlot is itself idempotent,
|
||||
// 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 */ }
|
||||
}
|
||||
|
||||
// Guarantee slot release on ANY exit path (normal close, error, timeout kill,
|
||||
@@ -980,12 +1046,18 @@ function spawnClaudeProcess(model, messages, conversationId, keyName) {
|
||||
// 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.
|
||||
function callClaude(model, messages, conversationId, keyName) {
|
||||
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). 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);
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, keyName, releaseSlot);
|
||||
} catch (err) {
|
||||
releaseSlot();
|
||||
return reject(err);
|
||||
}
|
||||
|
||||
@@ -1152,14 +1224,28 @@ function startHeartbeat(res, intervalMs, sessionId) {
|
||||
// We parse line-by-line and forward content_block_delta text events as SSE.
|
||||
// The result event triggers the stop/[DONE] sequence.
|
||||
// Reference: OLP ADR 0009 Amendment 1 + commits 97e7d16, 65f945c.
|
||||
function callClaudeStreaming(model, messages, conversationId, res, authInfo = {}) {
|
||||
async function callClaudeStreaming(model, messages, conversationId, res, authInfo = {}) {
|
||||
const id = `chatcmpl-${randomUUID()}`;
|
||||
const created = Math.floor(Date.now() / 1000);
|
||||
|
||||
// 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.
|
||||
let releaseSlot;
|
||||
try {
|
||||
releaseSlot = await acquireClaudeSlot();
|
||||
} catch (err) {
|
||||
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" } });
|
||||
}
|
||||
|
||||
let ctx;
|
||||
try {
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, authInfo.keyName);
|
||||
ctx = spawnClaudeProcess(model, messages, conversationId, authInfo.keyName, releaseSlot);
|
||||
} catch (err) {
|
||||
releaseSlot();
|
||||
return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } });
|
||||
}
|
||||
|
||||
@@ -1352,12 +1438,23 @@ function sanitizeError(msg) {
|
||||
}
|
||||
|
||||
// ── Response helpers ────────────────────────────────────────────────────
|
||||
function jsonResponse(res, status, data) {
|
||||
function jsonResponse(res, status, data, extraHeaders = null) {
|
||||
if (res.headersSent || res.writableEnded || res.destroyed) return;
|
||||
res.writeHead(status, { "Content-Type": "application/json" });
|
||||
// extraHeaders is optional + additive (e.g. Retry-After on a 429); Content-Type always wins.
|
||||
res.writeHead(status, { ...(extraHeaders || {}), "Content-Type": "application/json" });
|
||||
res.end(JSON.stringify(data));
|
||||
}
|
||||
|
||||
// FIX ⑥: map an upstream error to the right HTTP response. A ConcurrencyOverflowError (the
|
||||
// wait-queue was full) becomes HTTP 429 + Retry-After + rate_limit_error; every other error
|
||||
// stays a 500 proxy_error (byte-for-byte the pre-fix behaviour for non-overflow errors).
|
||||
function respondUpstreamError(res, err) {
|
||||
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" } });
|
||||
}
|
||||
|
||||
function sendSSE(res, data, hb) {
|
||||
hb?.reset();
|
||||
res.write(`data: ${JSON.stringify(data)}\n\n`);
|
||||
@@ -1777,7 +1874,9 @@ function applySettingUpdate(key, value) {
|
||||
|
||||
switch (key) {
|
||||
case "timeout": TIMEOUT = value; break;
|
||||
case "maxConcurrent": MAX_CONCURRENT = 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;
|
||||
@@ -1986,7 +2085,7 @@ async function handleChatCompletions(req, res) {
|
||||
try { res.end(); } catch {}
|
||||
return;
|
||||
}
|
||||
return jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } });
|
||||
return respondUpstreamError(res, err);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2003,8 +2102,9 @@ async function handleChatCompletions(req, res) {
|
||||
try { res.end(); } catch {}
|
||||
return;
|
||||
}
|
||||
// Sanitize error: strip internal file paths before sending to client
|
||||
jsonResponse(res, 500, { error: { message: sanitizeError(err.message), type: "proxy_error" } });
|
||||
// Sanitize error: strip internal file paths before sending to client.
|
||||
// FIX ⑥: ConcurrencyOverflowError → 429 + Retry-After; all other errors → 500 (unchanged).
|
||||
respondUpstreamError(res, err);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2180,6 +2280,16 @@ const server = createServer(async (req, res) => {
|
||||
reason: shm.reason,
|
||||
};
|
||||
})(),
|
||||
// ── FIX ⑥ -p concurrency wait-queue surface — ADDITIVE ──
|
||||
// inflight/queued are live; queueRejections is cumulative (also in stats.queueRejections).
|
||||
// Lets the operator see backpressure instead of guessing from opaque 500s.
|
||||
concurrency: {
|
||||
maxConcurrent: MAX_CONCURRENT,
|
||||
maxQueue: claudeSemaphore.maxQueue,
|
||||
inflight: claudeSemaphore.inflight,
|
||||
queued: claudeSemaphore.queued,
|
||||
queueRejections: stats.queueRejections,
|
||||
},
|
||||
// ── TUI observability (audit C-5) — ADDITIVE block (ADR 0007 PR-B amendment) ──
|
||||
// /health is a grandfathered B.2 endpoint (ADR 0006). This block is NEW fields only;
|
||||
// every existing field above is byte-identical → behaviour-preserving for existing
|
||||
@@ -2468,7 +2578,7 @@ server.listen(PORT, BIND_ADDRESS, () => {
|
||||
console.log(`Architecture: on-demand spawning (no pool)`);
|
||||
console.log(`Models: ${MODELS.map((m) => m.id).join(", ")}`);
|
||||
console.log(`Claude binary: ${CLAUDE}`);
|
||||
console.log(`Timeout: ${TIMEOUT / 1000}s | Max concurrent: ${MAX_CONCURRENT}`);
|
||||
console.log(`Timeout: ${TIMEOUT / 1000}s | Max concurrent: ${MAX_CONCURRENT} | Queue: ${CLAUDE_MAX_QUEUE} (429 on overflow)`);
|
||||
console.log(`Circuit breaker: disabled`);
|
||||
console.log(`Tools: ${SKIP_PERMISSIONS ? "all (skip-permissions)" : ALLOWED_TOOLS.join(", ")}`);
|
||||
console.log(`Sessions: TTL=${SESSION_TTL / 1000}s`);
|
||||
|
||||
@@ -2046,6 +2046,45 @@ await asyncTest("wait queue is bounded — run() rejects with tui_queue_full whe
|
||||
assert.equal(sem.inflight, 0);
|
||||
});
|
||||
|
||||
console.log("\n-p concurrency wait-queue (FIX ⑥ — same TuiSemaphore reused for the -p path):");
|
||||
|
||||
// server.mjs reuses TuiSemaphore as `claudeSemaphore = new TuiSemaphore(MAX_CONCURRENT,
|
||||
// { maxQueue: CLAUDE_MAX_QUEUE })` and wraps acquire()/release() in acquireClaudeSlot(). These
|
||||
// tests assert the contract that the 429-mapping depends on: requests beyond the limit QUEUE
|
||||
// (not reject), only an overflow past the queue rejects (→ HTTP 429 in server.mjs), and a
|
||||
// released slot is reusable (the #37/#40 slot-leak guard — no leak on normal completion).
|
||||
await asyncTest("FIX ⑥: requests beyond MAX_CONCURRENT queue, not reject (limit=1, queue=1)", async () => {
|
||||
const sem = new TuiSemaphore(1, { maxQueue: 1 }); // mirrors CLAUDE_MAX_CONCURRENT=1, CLAUDE_MAX_QUEUE=1
|
||||
const g1 = deferred();
|
||||
const inflightP = sem.run(async () => { await g1.p; }); // request 1 — holds the only slot
|
||||
await new Promise((r) => setImmediate(r));
|
||||
assert.equal(sem.inflight, 1, "req1 inflight");
|
||||
const queuedP = sem.run(async () => {}); // request 2 — WAITS (queued), does NOT reject
|
||||
await new Promise((r) => setImmediate(r));
|
||||
assert.equal(sem.queued, 1, "req2 queued (waits), not rejected → would be served, not 429");
|
||||
// request 3 — queue full → reject (server.mjs maps this single case to 429 + Retry-After)
|
||||
await assert.rejects(sem.run(async () => {}), /tui_queue_full|queue/, "req3 overflows → reject (→429)");
|
||||
g1.resolve();
|
||||
await inflightP; await queuedP;
|
||||
assert.equal(sem.inflight, 0, "all slots released after drain (no leak)");
|
||||
assert.equal(sem.queued, 0, "queue fully drained");
|
||||
});
|
||||
|
||||
await asyncTest("FIX ⑥: slot released on normal completion is immediately reusable (no #37/#40 leak)", async () => {
|
||||
const sem = new TuiSemaphore(1, { maxQueue: 16 }); // mirrors default CLAUDE_MAX_QUEUE=16
|
||||
for (let i = 0; i < 5; i++) {
|
||||
await sem.run(async () => { /* a normal, completing turn */ });
|
||||
assert.equal(sem.inflight, 0, `slot released after turn ${i}`);
|
||||
}
|
||||
// Prove the limit still binds after many acquire/release cycles.
|
||||
const g = deferred();
|
||||
const held = sem.run(async () => { await g.p; });
|
||||
await new Promise((r) => setImmediate(r));
|
||||
assert.equal(sem.inflight, 1, "limit still enforced after reuse cycles");
|
||||
g.resolve(); await held;
|
||||
assert.equal(sem.inflight, 0);
|
||||
});
|
||||
|
||||
console.log("\nTUI drift observability (C-5):");
|
||||
|
||||
test("recordTuiEntrypoint: observed 'cli' is NOT a mismatch and sets lastEntrypoint", () => {
|
||||
|
||||
Reference in New Issue
Block a user