From 50c93d266aff099b91bb63aec5913e2f8c0139f5 Mon Sep 17 00:00:00 2001 From: Brandon Bennett Date: Sun, 9 Aug 2026 18:33:32 -0700 Subject: [PATCH] fix(admission): complete REJECT_MAP, literal lane env read, split oversized test file MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Three CI-gate fixes surfaced by the post-merge check run (head 3de77166e): 1. open-sse-typecheck (TS2741): REJECT_MAP was missing the ADMISSION_LANE_EVICTED entry that controller.ts:662 emits on lane eviction. Add the 503 mapping so the Record is total. 2. Docs Gates fabricated-claim: OMNIROUTE_CHAT_VIRTUAL_LANES was read dynamically via ENV_KEYS.virtualLanes (env[key]), invisible to the literal env.X scanner. Read it literally — behavior-identical, doc claim now verifiable. 3. check:file-size: chat-body-admission.test.ts (1307 lines) exceeded the 1000-line new-file cap. Split the queue-wait/abort/heap-valve section into chat-body-admission-queue.test.ts (818 + 513 lines, both under cap). Suite: 125/125 across 8 files. All three checkers pass locally. --- open-sse/services/admission/runtime.ts | 8 +- tests/unit/chat-body-admission-queue.test.ts | 513 +++++++++++++++++++ tests/unit/chat-body-admission.test.ts | 489 ------------------ 3 files changed, 520 insertions(+), 490 deletions(-) create mode 100644 tests/unit/chat-body-admission-queue.test.ts diff --git a/open-sse/services/admission/runtime.ts b/open-sse/services/admission/runtime.ts index b5183bacda..eefa9115ae 100644 --- a/open-sse/services/admission/runtime.ts +++ b/open-sse/services/admission/runtime.ts @@ -121,7 +121,7 @@ export function resolveAdaptiveAdmissionConfigFromEnv( validateConfig(cfg); // Per-connection virtual admission lanes (#9654) — opt-in via OMNIROUTE_CHAT_VIRTUAL_LANES. - const vlRaw = env[ENV_KEYS.virtualLanes]; + const vlRaw = env.OMNIROUTE_CHAT_VIRTUAL_LANES; cfg.virtualLanes = vlRaw === "1" || vlRaw === "true"; return cfg; @@ -243,6 +243,12 @@ const REJECT_MAP: Record = { message: "Service temporarily unavailable", retryAfter: "1", }, + ADMISSION_LANE_EVICTED: { + status: 503, + code: "admission_lane_evicted", + message: "Connection lane evicted", + retryAfter: "1", + }, }; function isAdmissionRejectError( diff --git a/tests/unit/chat-body-admission-queue.test.ts b/tests/unit/chat-body-admission-queue.test.ts new file mode 100644 index 0000000000..caaa36c8e0 --- /dev/null +++ b/tests/unit/chat-body-admission-queue.test.ts @@ -0,0 +1,513 @@ +// #9654: queue-wait, AbortSignal cancellation, and the queued-bytes heap valve. +// Split from chat-body-admission.test.ts to stay under the 1000-line new-file cap. +import test from "node:test"; +import assert from "node:assert/strict"; + +const admissionModule = await import("../../src/shared/middleware/chatBodyAdmission.ts"); +const { + admitChatRequest, + admitChatStructure, + ChatAdmissionController, + CHAT_ADMISSION_QUEUE_MAX_MS, + CHAT_ADMISSION_MAX_QUEUED_BYTES, + CHAT_LARGE_BODY_BYTES, +} = admissionModule; + +function chatRequest(body: string, contentLength: string | null = String(body.length)): Request { + const headers: Record = { "content-type": "application/json" }; + if (contentLength !== null) headers["content-length"] = contentLength; + return new Request("http://x/v1/chat/completions", { + method: "POST", + headers, + body, + }); +} + +test("a heavy structural request waits for capacity instead of failing immediately", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const pending = admitChatStructure( + { + messages: [ + { role: "user", content: "one" }, + { role: "user", content: "two" }, + ], + }, + null, + { + controller, + maxMessages: 10, + heavyMessages: 2, + heavyTools: 10, + heavyTokens: 10_000, + queueMs: 500, + } + ); + + // Capacity is still busy: the request must not have resolved (admit/reject) yet. + let settled = false; + void pending.then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(settled, false, "must wait while capacity is busy"); + + held.release(); + const result = await pending; + assert.equal(result.admit, true); + if (result.admit) { + assert.equal(controller.activeHeavy, 1, "waiting request acquires the freed lease"); + result.lease?.release(); + } + assert.equal(controller.activeHeavy, 0); +}); + +test("waiting for admission times out into a retryable 503", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const started = Date.now(); + const result = await admitChatStructure( + { + messages: [ + { role: "user", content: "one" }, + { role: "user", content: "two" }, + ], + }, + null, + { + controller, + maxMessages: 10, + heavyMessages: 2, + heavyTools: 10, + heavyTokens: 10_000, + queueMs: 50, + } + ); + + assert.equal(result.admit, false); + if (!result.admit) { + assert.equal(result.response.status, 503); + assert.equal(result.response.headers.get("retry-after"), "1"); + assert.equal((await result.response.json()).error.code, "chat_admission_busy"); + } + assert.ok(Date.now() - started >= 40, "must wait for the queue deadline before rejecting"); + assert.equal(controller.activeHeavy, 1, "the holder keeps its lease"); + held.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("byte-heavy admission waits for capacity when queueMs is set", async () => { + const controller = new ChatAdmissionController(1); + const body = JSON.stringify({ messages: [{ role: "user", content: "x".repeat(40) }] }); + const options = { controller, largeBodyBytes: 32, hardMaxBytes: 1024, queueMs: 500 }; + + const first = await admitChatRequest(chatRequest(body), options); + assert.equal(first.admit, true); + if (!first.admit) return; + + const second = admitChatRequest(chatRequest(body), options); + let secondSettled = false; + void second.then(() => { + secondSettled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal( + secondSettled, + false, + "second heavy request must queue while the first holds capacity" + ); + + first.lease?.release(); + const secondResult = await second; + assert.equal(secondResult.admit, true, "second request acquires capacity after release"); + if (secondResult.admit) secondResult.lease?.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("expired admission queue keeps the legacy immediate 503 behaviour", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const result = await admitChatStructure( + { + messages: [ + { role: "user", content: "one" }, + { role: "user", content: "two" }, + ], + }, + null, + { + controller, + maxMessages: 10, + heavyMessages: 2, + heavyTools: 10, + heavyTokens: 10_000, + queueMs: 0, + } + ); + + assert.equal(result.admit, false); + if (!result.admit) assert.equal(result.response.status, 503); + held.release(); +}); + +test("admission waiters are served FIFO as capacity frees", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const body = { + messages: [ + { role: "user", content: "one" }, + { role: "user", content: "two" }, + ], + }; + const options = { + controller, + maxMessages: 10, + heavyMessages: 2, + heavyTools: 10, + heavyTokens: 10_000, + queueMs: 500, + }; + const first = admitChatStructure(body, null, options); + const second = admitChatStructure(body, null, options); + + held.release(); + const firstResult = await first; + assert.equal(firstResult.admit, true); + if (firstResult.admit) firstResult.lease?.release(); + const secondResult = await second; + assert.equal(secondResult.admit, true); + if (secondResult.admit) secondResult.lease?.release(); + assert.equal(controller.activeHeavy, 0); +}); + +// ── AbortSignal support in acquireHeavyWithin (#9654 / U2) ──────────────── +// A disconnected client must not keep parking in the admission queue for the +// full queueMs. On abort the waiter is removed from the FIFO immediately and +// the acquire resolves `null` early (the caller's 503 is dropped on the dead +// connection); no capacity is consumed and the freed slot does not wake it. + +test("aborting the admission wait settles early, grants no lease, and removes the waiter", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const abortController = new AbortController(); + const pending = controller.acquireHeavyWithin(2_000, abortController.signal); + + // Parked while capacity is busy. + let settled = false; + void pending.then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(settled, false, "must be parked while capacity is busy"); + + abortController.abort(); + + // Must settle well before the 2s deadline. + let settledAfterAbort = false; + void pending.then(() => { + settledAfterAbort = true; + }); + await new Promise((resolve) => setTimeout(resolve, 50)); + assert.equal(settledAfterAbort, true, "abort must settle the wait promptly, not park for queueMs"); + + const lease = await pending; + assert.equal(lease, null, "abort must not grant a lease"); + assert.equal(controller.activeHeavy, 1, "the holder keeps its lease; the aborted wait consumed nothing"); + + // Releasing must NOT wake the removed waiter: capacity stays free. + held.release(); + await new Promise((resolve) => setTimeout(resolve, 0)); + assert.equal( + controller.activeHeavy, + 0, + "releasing after abort must not wake the removed waiter" + ); +}); + +test("aborting the head waiter preserves FIFO order for remaining waiters", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const firstAbort = new AbortController(); + const first = controller.acquireHeavyWithin(2_000, firstAbort.signal); + const second = controller.acquireHeavyWithin(2_000); + + // Both are parked, head-first. + await new Promise((resolve) => setTimeout(resolve, 30)); + + // Abort the HEAD waiter: it must leave the queue without disturbing the rest. + firstAbort.abort(); + assert.equal(await first, null, "head waiter returns null on abort"); + + // The remaining waiter is now first in line and must get the freed capacity. + held.release(); + const secondLease = await second; + assert.ok(secondLease, "remaining waiter must acquire the freed capacity"); + secondLease?.release(); + assert.equal(controller.activeHeavy, 0); +}); + +// ── Heap-pressure safety valve (#9654 / U3) ─────────────────────────────── +// The queue-wait parks fully-buffered bodies; the queued-bytes cap bounds the +// total buffered memory parked per lane so the wait cannot recreate the #4380 +// heap amplification. Over-budget waits are rejected immediately (503). + +test("queued-bytes cap rejects an over-budget wait without parking", async () => { + const controller = new ChatAdmissionController(1, 200); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + // First waiter parks within budget. + const first = controller.acquireHeavyWithin(2_000, undefined, 150); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(controller.queuedBytes, 150); + + // Second waiter would push the total over the 200-byte budget → must NOT park. + const started = Date.now(); + const second = await controller.acquireHeavyWithin(2_000, undefined, 100); + assert.equal(second, null, "over-budget wait must be rejected"); + assert.ok(Date.now() - started < 500, "rejection must be immediate, not park for queueMs"); + assert.equal(controller.queuedBytes, 150, "rejected waiter must not be charged"); + assert.equal(controller.activeHeavy, 1, "holder keeps its lease"); + + // Free the slot: the parked waiter acquires and its bytes leave the queue. + held.release(); + const firstLease = await first; + assert.ok(firstLease, "in-budget waiter acquires the freed slot"); + assert.equal(controller.queuedBytes, 0, "acquired waiter's bytes must leave the queue"); + firstLease?.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("aborting a parked wait releases its queued bytes", async () => { + const controller = new ChatAdmissionController(1, 1_000); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const abortController = new AbortController(); + const pending = controller.acquireHeavyWithin(2_000, abortController.signal, 400); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(controller.queuedBytes, 400); + + abortController.abort(); + assert.equal(await pending, null); + assert.equal(controller.queuedBytes, 0, "abort must release the charged bytes"); + + held.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("a timed-out wait releases its queued bytes", async () => { + const controller = new ChatAdmissionController(1, 1_000); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const pending = controller.acquireHeavyWithin(50, undefined, 400); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(controller.queuedBytes, 400); + + assert.equal(await pending, null); + assert.equal(controller.queuedBytes, 0, "timeout must release the charged bytes"); + held.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("byte-heavy admission enforces the queued-bytes cap end-to-end", async () => { + const controller = new ChatAdmissionController(1, 100); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const body = JSON.stringify({ messages: [{ role: "user", content: "x".repeat(40) }] }); + const options = { controller, largeBodyBytes: 32, hardMaxBytes: 1024, queueMs: 2_000 }; + + // First request parks: declared length (~70B) fits the budget. + const first = admitChatRequest(chatRequest(body), options); + await new Promise((resolve) => setTimeout(resolve, 30)); + + // Second request would exceed the 100-byte budget → rejected immediately. + const started = Date.now(); + const second = await admitChatRequest(chatRequest(body), options); + assert.equal(second.admit, false, "over-budget byte-heavy wait must not admit"); + if (!second.admit) assert.equal(second.response.status, 503); + assert.ok(Date.now() - started < 500, "over-budget wait must reject immediately"); + + held.release(); + const firstResult = await first; + assert.equal(firstResult.admit, true); + if (firstResult.admit) firstResult.lease?.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("structural admission enforces the queued-bytes cap end-to-end", async () => { + const controller = new ChatAdmissionController(1, CHAT_LARGE_BODY_BYTES); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const structural = { + messages: [ + { role: "user", content: "one" }, + { role: "user", content: "two" }, + ], + }; + const options = { + controller, + maxMessages: 10, + heavyMessages: 2, + heavyTools: 10, + heavyTokens: 10_000, + queueMs: 2_000, + }; + + // First structural wait parks, charging the conservative 256KB weight. + const first = admitChatStructure(structural, null, options); + await new Promise((resolve) => setTimeout(resolve, 30)); + + // Second would double the charge → rejected immediately. + const started = Date.now(); + const second = await admitChatStructure(structural, null, options); + assert.equal(second.admit, false, "over-budget structural wait must not admit"); + if (!second.admit) assert.equal(second.response.status, 503); + assert.ok(Date.now() - started < 500, "over-budget structural wait must reject immediately"); + + held.release(); + const firstResult = await first; + assert.equal(firstResult.admit, true); + if (firstResult.admit) firstResult.lease?.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("queue-wait defaults are bounded (2s wait, 4MB queued-bytes budget)", () => { + if (process.env.OMNIROUTE_CHAT_ADMISSION_QUEUE_MS === undefined) { + assert.equal(CHAT_ADMISSION_QUEUE_MAX_MS, 2_000); + } + if (process.env.OMNIROUTE_CHAT_ADMISSION_MAX_QUEUED_BYTES === undefined) { + assert.equal(CHAT_ADMISSION_MAX_QUEUED_BYTES, 4 * 1024 * 1024); + } +}); + +test("a pre-aborted signal never parks in the admission queue", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const abortController = new AbortController(); + abortController.abort("client already disconnected"); + + const pending = controller.acquireHeavyWithin(2_000, abortController.signal); + let settled = false; + void pending.then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 50)); + assert.equal(settled, true, "a pre-aborted signal must settle immediately, not park"); + + const lease = await pending; + assert.equal(lease, null, "no lease is granted after abort"); + assert.equal(controller.activeHeavy, 1, "holder keeps capacity; aborted wait consumed nothing"); + held.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("aborting the request signal cancels a queued byte-heavy wait", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const abortController = new AbortController(); + const body = JSON.stringify({ messages: [{ role: "user", content: "x".repeat(40) }] }); + const request = new Request("http://x/v1/chat/completions", { + method: "POST", + headers: { "content-type": "application/json" }, + body, + signal: abortController.signal, + }); + const pending = admitChatRequest(request, { + controller, + largeBodyBytes: 32, + hardMaxBytes: 1024, + queueMs: 2_000, + }); + + let settled = false; + void pending.then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(settled, false, "must queue while capacity is busy"); + + abortController.abort(); + const started = Date.now(); + const result = await pending; + assert.ok( + Date.now() - started < 500, + "abort must cancel the queue-wait early, not park the full queueMs" + ); + assert.equal(result.admit, false, "abort must not admit"); + if (!result.admit) { + assert.equal(result.response.status, 503); + assert.equal((await result.response.json()).error.code, "chat_admission_busy"); + } + assert.equal(controller.activeHeavy, 1, "holder keeps capacity; aborted wait consumed nothing"); + held.release(); + assert.equal(controller.activeHeavy, 0); +}); + +test("aborting the signal cancels a structural queue-wait", async () => { + const controller = new ChatAdmissionController(1); + const held = controller.tryAcquireHeavy(); + assert.ok(held); + + const abortController = new AbortController(); + const pending = admitChatStructure( + { + messages: [ + { role: "user", content: "one" }, + { role: "user", content: "two" }, + ], + }, + null, + { + controller, + maxMessages: 10, + heavyMessages: 2, + heavyTools: 10, + heavyTokens: 10_000, + queueMs: 2_000, + signal: abortController.signal, + } + ); + + let settled = false; + void pending.then(() => { + settled = true; + }); + await new Promise((resolve) => setTimeout(resolve, 30)); + assert.equal(settled, false, "must queue while capacity is busy"); + + abortController.abort(); + const started = Date.now(); + const result = await pending; + assert.ok( + Date.now() - started < 500, + "abort must cancel the queue-wait early, not park the full queueMs" + ); + assert.equal(result.admit, false, "abort must not admit"); + if (!result.admit) { + assert.equal(result.response.status, 503); + assert.equal((await result.response.json()).error.code, "chat_admission_busy"); + } + assert.equal(controller.activeHeavy, 1, "holder keeps its lease"); + held.release(); + assert.equal(controller.activeHeavy, 0); +}); diff --git a/tests/unit/chat-body-admission.test.ts b/tests/unit/chat-body-admission.test.ts index bdd9c353c8..8544bc50c0 100644 --- a/tests/unit/chat-body-admission.test.ts +++ b/tests/unit/chat-body-admission.test.ts @@ -816,492 +816,3 @@ test("sk_omniroute sentinel is rejected once an env key is configured (REQUIRE_A restore(); } }); - -test("a heavy structural request waits for capacity instead of failing immediately", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const pending = admitChatStructure( - { - messages: [ - { role: "user", content: "one" }, - { role: "user", content: "two" }, - ], - }, - null, - { - controller, - maxMessages: 10, - heavyMessages: 2, - heavyTools: 10, - heavyTokens: 10_000, - queueMs: 500, - } - ); - - // Capacity is still busy: the request must not have resolved (admit/reject) yet. - let settled = false; - void pending.then(() => { - settled = true; - }); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(settled, false, "must wait while capacity is busy"); - - held.release(); - const result = await pending; - assert.equal(result.admit, true); - if (result.admit) { - assert.equal(controller.activeHeavy, 1, "waiting request acquires the freed lease"); - result.lease?.release(); - } - assert.equal(controller.activeHeavy, 0); -}); - -test("waiting for admission times out into a retryable 503", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const started = Date.now(); - const result = await admitChatStructure( - { - messages: [ - { role: "user", content: "one" }, - { role: "user", content: "two" }, - ], - }, - null, - { - controller, - maxMessages: 10, - heavyMessages: 2, - heavyTools: 10, - heavyTokens: 10_000, - queueMs: 50, - } - ); - - assert.equal(result.admit, false); - if (!result.admit) { - assert.equal(result.response.status, 503); - assert.equal(result.response.headers.get("retry-after"), "1"); - assert.equal((await result.response.json()).error.code, "chat_admission_busy"); - } - assert.ok(Date.now() - started >= 40, "must wait for the queue deadline before rejecting"); - assert.equal(controller.activeHeavy, 1, "the holder keeps its lease"); - held.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("byte-heavy admission waits for capacity when queueMs is set", async () => { - const controller = new ChatAdmissionController(1); - const body = JSON.stringify({ messages: [{ role: "user", content: "x".repeat(40) }] }); - const options = { controller, largeBodyBytes: 32, hardMaxBytes: 1024, queueMs: 500 }; - - const first = await admitChatRequest(chatRequest(body), options); - assert.equal(first.admit, true); - if (!first.admit) return; - - const second = admitChatRequest(chatRequest(body), options); - let secondSettled = false; - void second.then(() => { - secondSettled = true; - }); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal( - secondSettled, - false, - "second heavy request must queue while the first holds capacity" - ); - - first.lease?.release(); - const secondResult = await second; - assert.equal(secondResult.admit, true, "second request acquires capacity after release"); - if (secondResult.admit) secondResult.lease?.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("expired admission queue keeps the legacy immediate 503 behaviour", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const result = await admitChatStructure( - { - messages: [ - { role: "user", content: "one" }, - { role: "user", content: "two" }, - ], - }, - null, - { - controller, - maxMessages: 10, - heavyMessages: 2, - heavyTools: 10, - heavyTokens: 10_000, - queueMs: 0, - } - ); - - assert.equal(result.admit, false); - if (!result.admit) assert.equal(result.response.status, 503); - held.release(); -}); - -test("admission waiters are served FIFO as capacity frees", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const body = { - messages: [ - { role: "user", content: "one" }, - { role: "user", content: "two" }, - ], - }; - const options = { - controller, - maxMessages: 10, - heavyMessages: 2, - heavyTools: 10, - heavyTokens: 10_000, - queueMs: 500, - }; - const first = admitChatStructure(body, null, options); - const second = admitChatStructure(body, null, options); - - held.release(); - const firstResult = await first; - assert.equal(firstResult.admit, true); - if (firstResult.admit) firstResult.lease?.release(); - const secondResult = await second; - assert.equal(secondResult.admit, true); - if (secondResult.admit) secondResult.lease?.release(); - assert.equal(controller.activeHeavy, 0); -}); - -// ── AbortSignal support in acquireHeavyWithin (#9654 / U2) ──────────────── -// A disconnected client must not keep parking in the admission queue for the -// full queueMs. On abort the waiter is removed from the FIFO immediately and -// the acquire resolves `null` early (the caller's 503 is dropped on the dead -// connection); no capacity is consumed and the freed slot does not wake it. - -test("aborting the admission wait settles early, grants no lease, and removes the waiter", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const abortController = new AbortController(); - const pending = controller.acquireHeavyWithin(2_000, abortController.signal); - - // Parked while capacity is busy. - let settled = false; - void pending.then(() => { - settled = true; - }); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(settled, false, "must be parked while capacity is busy"); - - abortController.abort(); - - // Must settle well before the 2s deadline. - let settledAfterAbort = false; - void pending.then(() => { - settledAfterAbort = true; - }); - await new Promise((resolve) => setTimeout(resolve, 50)); - assert.equal(settledAfterAbort, true, "abort must settle the wait promptly, not park for queueMs"); - - const lease = await pending; - assert.equal(lease, null, "abort must not grant a lease"); - assert.equal(controller.activeHeavy, 1, "the holder keeps its lease; the aborted wait consumed nothing"); - - // Releasing must NOT wake the removed waiter: capacity stays free. - held.release(); - await new Promise((resolve) => setTimeout(resolve, 0)); - assert.equal( - controller.activeHeavy, - 0, - "releasing after abort must not wake the removed waiter" - ); -}); - -test("aborting the head waiter preserves FIFO order for remaining waiters", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const firstAbort = new AbortController(); - const first = controller.acquireHeavyWithin(2_000, firstAbort.signal); - const second = controller.acquireHeavyWithin(2_000); - - // Both are parked, head-first. - await new Promise((resolve) => setTimeout(resolve, 30)); - - // Abort the HEAD waiter: it must leave the queue without disturbing the rest. - firstAbort.abort(); - assert.equal(await first, null, "head waiter returns null on abort"); - - // The remaining waiter is now first in line and must get the freed capacity. - held.release(); - const secondLease = await second; - assert.ok(secondLease, "remaining waiter must acquire the freed capacity"); - secondLease?.release(); - assert.equal(controller.activeHeavy, 0); -}); - -// ── Heap-pressure safety valve (#9654 / U3) ─────────────────────────────── -// The queue-wait parks fully-buffered bodies; the queued-bytes cap bounds the -// total buffered memory parked per lane so the wait cannot recreate the #4380 -// heap amplification. Over-budget waits are rejected immediately (503). - -test("queued-bytes cap rejects an over-budget wait without parking", async () => { - const controller = new ChatAdmissionController(1, 200); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - // First waiter parks within budget. - const first = controller.acquireHeavyWithin(2_000, undefined, 150); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(controller.queuedBytes, 150); - - // Second waiter would push the total over the 200-byte budget → must NOT park. - const started = Date.now(); - const second = await controller.acquireHeavyWithin(2_000, undefined, 100); - assert.equal(second, null, "over-budget wait must be rejected"); - assert.ok(Date.now() - started < 500, "rejection must be immediate, not park for queueMs"); - assert.equal(controller.queuedBytes, 150, "rejected waiter must not be charged"); - assert.equal(controller.activeHeavy, 1, "holder keeps its lease"); - - // Free the slot: the parked waiter acquires and its bytes leave the queue. - held.release(); - const firstLease = await first; - assert.ok(firstLease, "in-budget waiter acquires the freed slot"); - assert.equal(controller.queuedBytes, 0, "acquired waiter's bytes must leave the queue"); - firstLease?.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("aborting a parked wait releases its queued bytes", async () => { - const controller = new ChatAdmissionController(1, 1_000); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const abortController = new AbortController(); - const pending = controller.acquireHeavyWithin(2_000, abortController.signal, 400); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(controller.queuedBytes, 400); - - abortController.abort(); - assert.equal(await pending, null); - assert.equal(controller.queuedBytes, 0, "abort must release the charged bytes"); - - held.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("a timed-out wait releases its queued bytes", async () => { - const controller = new ChatAdmissionController(1, 1_000); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const pending = controller.acquireHeavyWithin(50, undefined, 400); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(controller.queuedBytes, 400); - - assert.equal(await pending, null); - assert.equal(controller.queuedBytes, 0, "timeout must release the charged bytes"); - held.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("byte-heavy admission enforces the queued-bytes cap end-to-end", async () => { - const controller = new ChatAdmissionController(1, 100); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const body = JSON.stringify({ messages: [{ role: "user", content: "x".repeat(40) }] }); - const options = { controller, largeBodyBytes: 32, hardMaxBytes: 1024, queueMs: 2_000 }; - - // First request parks: declared length (~70B) fits the budget. - const first = admitChatRequest(chatRequest(body), options); - await new Promise((resolve) => setTimeout(resolve, 30)); - - // Second request would exceed the 100-byte budget → rejected immediately. - const started = Date.now(); - const second = await admitChatRequest(chatRequest(body), options); - assert.equal(second.admit, false, "over-budget byte-heavy wait must not admit"); - if (!second.admit) assert.equal(second.response.status, 503); - assert.ok(Date.now() - started < 500, "over-budget wait must reject immediately"); - - held.release(); - const firstResult = await first; - assert.equal(firstResult.admit, true); - if (firstResult.admit) firstResult.lease?.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("structural admission enforces the queued-bytes cap end-to-end", async () => { - const controller = new ChatAdmissionController(1, CHAT_LARGE_BODY_BYTES); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const structural = { - messages: [ - { role: "user", content: "one" }, - { role: "user", content: "two" }, - ], - }; - const options = { - controller, - maxMessages: 10, - heavyMessages: 2, - heavyTools: 10, - heavyTokens: 10_000, - queueMs: 2_000, - }; - - // First structural wait parks, charging the conservative 256KB weight. - const first = admitChatStructure(structural, null, options); - await new Promise((resolve) => setTimeout(resolve, 30)); - - // Second would double the charge → rejected immediately. - const started = Date.now(); - const second = await admitChatStructure(structural, null, options); - assert.equal(second.admit, false, "over-budget structural wait must not admit"); - if (!second.admit) assert.equal(second.response.status, 503); - assert.ok(Date.now() - started < 500, "over-budget structural wait must reject immediately"); - - held.release(); - const firstResult = await first; - assert.equal(firstResult.admit, true); - if (firstResult.admit) firstResult.lease?.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("queue-wait defaults are bounded (2s wait, 4MB queued-bytes budget)", () => { - if (process.env.OMNIROUTE_CHAT_ADMISSION_QUEUE_MS === undefined) { - assert.equal(CHAT_ADMISSION_QUEUE_MAX_MS, 2_000); - } - if (process.env.OMNIROUTE_CHAT_ADMISSION_MAX_QUEUED_BYTES === undefined) { - assert.equal(CHAT_ADMISSION_MAX_QUEUED_BYTES, 4 * 1024 * 1024); - } -}); - -test("a pre-aborted signal never parks in the admission queue", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const abortController = new AbortController(); - abortController.abort("client already disconnected"); - - const pending = controller.acquireHeavyWithin(2_000, abortController.signal); - let settled = false; - void pending.then(() => { - settled = true; - }); - await new Promise((resolve) => setTimeout(resolve, 50)); - assert.equal(settled, true, "a pre-aborted signal must settle immediately, not park"); - - const lease = await pending; - assert.equal(lease, null, "no lease is granted after abort"); - assert.equal(controller.activeHeavy, 1, "holder keeps capacity; aborted wait consumed nothing"); - held.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("aborting the request signal cancels a queued byte-heavy wait", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const abortController = new AbortController(); - const body = JSON.stringify({ messages: [{ role: "user", content: "x".repeat(40) }] }); - const request = new Request("http://x/v1/chat/completions", { - method: "POST", - headers: { "content-type": "application/json" }, - body, - signal: abortController.signal, - }); - const pending = admitChatRequest(request, { - controller, - largeBodyBytes: 32, - hardMaxBytes: 1024, - queueMs: 2_000, - }); - - let settled = false; - void pending.then(() => { - settled = true; - }); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(settled, false, "must queue while capacity is busy"); - - abortController.abort(); - const started = Date.now(); - const result = await pending; - assert.ok( - Date.now() - started < 500, - "abort must cancel the queue-wait early, not park the full queueMs" - ); - assert.equal(result.admit, false, "abort must not admit"); - if (!result.admit) { - assert.equal(result.response.status, 503); - assert.equal((await result.response.json()).error.code, "chat_admission_busy"); - } - assert.equal(controller.activeHeavy, 1, "holder keeps capacity; aborted wait consumed nothing"); - held.release(); - assert.equal(controller.activeHeavy, 0); -}); - -test("aborting the signal cancels a structural queue-wait", async () => { - const controller = new ChatAdmissionController(1); - const held = controller.tryAcquireHeavy(); - assert.ok(held); - - const abortController = new AbortController(); - const pending = admitChatStructure( - { - messages: [ - { role: "user", content: "one" }, - { role: "user", content: "two" }, - ], - }, - null, - { - controller, - maxMessages: 10, - heavyMessages: 2, - heavyTools: 10, - heavyTokens: 10_000, - queueMs: 2_000, - signal: abortController.signal, - } - ); - - let settled = false; - void pending.then(() => { - settled = true; - }); - await new Promise((resolve) => setTimeout(resolve, 30)); - assert.equal(settled, false, "must queue while capacity is busy"); - - abortController.abort(); - const started = Date.now(); - const result = await pending; - assert.ok( - Date.now() - started < 500, - "abort must cancel the queue-wait early, not park the full queueMs" - ); - assert.equal(result.admit, false, "abort must not admit"); - if (!result.admit) { - assert.equal(result.response.status, 503); - assert.equal((await result.response.json()).error.code, "chat_admission_busy"); - } - assert.equal(controller.activeHeavy, 1, "holder keeps its lease"); - held.release(); - assert.equal(controller.activeHeavy, 0); -});