diff --git a/changelog.d/fixes/14456-admission-lease-abort-release.md b/changelog.d/fixes/14456-admission-lease-abort-release.md new file mode 100644 index 00000000000..4fcefccf74c --- /dev/null +++ b/changelog.d/fixes/14456-admission-lease-abort-release.md @@ -0,0 +1 @@ +- Release the heavyweight chat admission slot when a client disconnects without cancelling the SSE body. Previously `activeHeavy` only ever increased on abandoned streams and pinned at `OMNIROUTE_CHAT_MAX_HEAVY_IN_FLIGHT`, shedding every subsequent heavy request with 503 `chat_admission_busy` until the process was recreated. ([#14456](https://github.com/diegosouzapw/OmniRoute/issues/14456)) diff --git a/src/app/api/v1/chat/completions/route.ts b/src/app/api/v1/chat/completions/route.ts index 7576d6a44ed..7140bef5b3b 100644 --- a/src/app/api/v1/chat/completions/route.ts +++ b/src/app/api/v1/chat/completions/route.ts @@ -124,7 +124,7 @@ export async function POST(request) { const admission = admissionResult; request = admission.request; const finishAdmission = (response: Response) => - releaseChatAdmissionWhenDone(response, admission.lease); + releaseChatAdmissionWhenDone(response, admission.lease, { signal: request.signal }); try { // One-line marker for diagnosing 413 / Server-Action interceptions. @@ -270,7 +270,8 @@ export async function POST(request) { // eventual handler body; only that confirmed cleanup releases heavyweight capacity. const handlerResponse = releaseChatAdmissionAfterHandler( handleChat(request, null, parsedBody, reqId), - admission.lease + admission.lease, + { signal: request.signal } ); const streamedResponse = await withEarlyStreamKeepalive(handlerResponse, { signal: request.signal, diff --git a/src/app/api/v1/responses/route.ts b/src/app/api/v1/responses/route.ts index 65673efcd79..58b07920c25 100644 --- a/src/app/api/v1/responses/route.ts +++ b/src/app/api/v1/responses/route.ts @@ -108,7 +108,7 @@ async function postHandler(request: any) { const admission = admissionResult; request = admission.request; const finishAdmission = (response: Response) => - releaseChatAdmissionWhenDone(response, admission.lease); + releaseChatAdmissionWhenDone(response, admission.lease, { signal: request.signal }); try { let parsedBody; @@ -187,7 +187,8 @@ async function postHandler(request: any) { const correlationId = generateRequestId(); const handlerResponse = releaseChatAdmissionAfterHandler( handleChat(resolved, null, resolvedBody, correlationId), - admission.lease + admission.lease, + { signal: request.signal } ); return await withEarlyStreamKeepalive(handlerResponse, { signal: request.signal, diff --git a/src/shared/middleware/chatAdmissionRelease.ts b/src/shared/middleware/chatAdmissionRelease.ts new file mode 100644 index 00000000000..61dde6e3a8d --- /dev/null +++ b/src/shared/middleware/chatAdmissionRelease.ts @@ -0,0 +1,96 @@ +/** + * Lease release binding for streamed chat responses. + * + * Extracted from `chatBodyAdmission.ts` so the admission controller module does + * not keep growing. The release paths on a wrapped SSE body are all driven by + * the consumer (pull-to-done, pull-throws, cancel). A client that disconnects + * mid-stream may stop pulling without ever cancelling, so the request signal is + * observed as a fallback and the slot is returned exactly once. + */ +import type { ChatAdmissionLease } from "./chatBodyAdmission"; + +export interface ReleaseChatAdmissionOptions { + /** + * The inbound request signal. A disconnecting client may simply stop pulling + * the wrapped stream without ever cancelling it, in which case none of the + * consumer-driven release paths run. Aborting releases the slot. + */ + readonly signal?: AbortSignal; +} + +/** Hold a heavyweight lease through an SSE response without buffering the response body. */ +export function releaseChatAdmissionWhenDone( + response: Response, + lease: ChatAdmissionLease | null, + options: ReleaseChatAdmissionOptions = {} +): Response { + if (!lease) return response; + const isStreaming = response.headers.get("content-type")?.includes("text/event-stream"); + if (!isStreaming || !response.body) { + lease.release(); + return response; + } + + const reader = response.body.getReader(); + + const { signal } = options; + let detachAbortListener = (): void => undefined; + const releaseOnce = (): void => { + detachAbortListener(); + if (!lease.released) lease.release(); + }; + + if (signal) { + const onAbort = (): void => { + releaseOnce(); + void reader.cancel("client disconnected").catch(() => undefined); + }; + if (signal.aborted) { + onAbort(); + } else { + signal.addEventListener("abort", onAbort, { once: true }); + detachAbortListener = () => signal.removeEventListener("abort", onAbort); + } + } + + const body = new ReadableStream({ + async pull(controller) { + try { + const { done, value } = await reader.read(); + if (done) { + releaseOnce(); + controller.close(); + } else { + controller.enqueue(value); + } + } catch (error) { + releaseOnce(); + controller.error(error); + } + }, + async cancel(reason) { + releaseOnce(); + await reader.cancel(reason).catch(() => undefined); + }, + }); + + return new Response(body, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }); +} + +/** Release a lease if a handler rejects; otherwise bind it to the returned response lifecycle. */ +export async function releaseChatAdmissionAfterHandler( + responsePromise: Promise, + lease: ChatAdmissionLease | null, + options: ReleaseChatAdmissionOptions = {} +): Promise { + try { + return releaseChatAdmissionWhenDone(await responsePromise, lease, options); + } catch (error) { + lease?.release(); + throw error; + } +} diff --git a/src/shared/middleware/chatBodyAdmission.ts b/src/shared/middleware/chatBodyAdmission.ts index ef38f4086ad..9efe0a83b36 100644 --- a/src/shared/middleware/chatBodyAdmission.ts +++ b/src/shared/middleware/chatBodyAdmission.ts @@ -1150,56 +1150,10 @@ export async function admitChatRequest( return { admit: true, request: rebuildRequest(request, body), lease }; } -/** Release a lease if a handler rejects; otherwise bind it to the returned response lifecycle. */ -export async function releaseChatAdmissionAfterHandler( - responsePromise: Promise, - lease: ChatAdmissionLease | null -): Promise { - try { - return releaseChatAdmissionWhenDone(await responsePromise, lease); - } catch (error) { - lease?.release(); - throw error; - } -} - -/** Hold a heavyweight lease through an SSE response without buffering the response body. */ -export function releaseChatAdmissionWhenDone( - response: Response, - lease: ChatAdmissionLease | null -): Response { - if (!lease) return response; - const isStreaming = response.headers.get("content-type")?.includes("text/event-stream"); - if (!isStreaming || !response.body) { - lease.release(); - return response; - } - - const reader = response.body.getReader(); - const body = new ReadableStream({ - async pull(controller) { - try { - const { done, value } = await reader.read(); - if (done) { - lease.release(); - controller.close(); - } else { - controller.enqueue(value); - } - } catch (error) { - lease.release(); - controller.error(error); - } - }, - async cancel(reason) { - lease.release(); - await reader.cancel(reason).catch(() => undefined); - }, - }); - - return new Response(body, { - status: response.status, - statusText: response.statusText, - headers: response.headers, - }); -} +// Lease release binding lives in ./chatAdmissionRelease. Re-exported here so +// existing import sites keep working. +export { + releaseChatAdmissionAfterHandler, + releaseChatAdmissionWhenDone, + type ReleaseChatAdmissionOptions, +} from "./chatAdmissionRelease"; diff --git a/src/shared/middleware/withChatAdmission.ts b/src/shared/middleware/withChatAdmission.ts index 7a762b6f9ea..e3cf319eb48 100644 --- a/src/shared/middleware/withChatAdmission.ts +++ b/src/shared/middleware/withChatAdmission.ts @@ -38,7 +38,8 @@ export function withChatAdmission( try { return await releaseChatAdmissionAfterHandler( Promise.resolve(handler(admission.request, ...args)), - admission.lease + admission.lease, + { signal: request.signal } ); } catch (error) { admission.lease?.release(); diff --git a/tests/unit/chat-admission-abandoned-stream-lease.test.ts b/tests/unit/chat-admission-abandoned-stream-lease.test.ts new file mode 100644 index 00000000000..39b4439683f --- /dev/null +++ b/tests/unit/chat-admission-abandoned-stream-lease.test.ts @@ -0,0 +1,121 @@ +/** + * Regression: a heavyweight admission lease must not leak when the client goes + * away without draining or cancelling the SSE body. + * + * `releaseChatAdmissionWhenDone` wires release into three places on the wrapped + * stream: `pull` reaching `done`, `pull` throwing, and `cancel`. All three are + * driven by the *consumer*. When a client disconnects mid-stream, the runtime + * may stop pulling and never call `cancel`, so none of the three fire and the + * slot is held for the lifetime of the process. + * + * These tests assert the invariant (the slot is returned once the request is + * aborted) rather than any particular release mechanism. + */ +import test from "node:test"; +import assert from "node:assert/strict"; +import { + ChatAdmissionController, + releaseChatAdmissionWhenDone, +} from "../../src/shared/middleware/chatBodyAdmission.ts"; + +/** An SSE body that stays open, like a real upstream mid-generation. */ +function openSseResponse(signal?: AbortSignal): Response { + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("data: first\n\n")); + // Intentionally left open: upstream is still generating. + }, + }); + return new Response(body, { + status: 200, + headers: { "content-type": "text/event-stream" }, + ...(signal ? {} : {}), + }); +} + +test("abandoned SSE stream returns the heavyweight slot once the request aborts", async () => { + const controller = new ChatAdmissionController(1); + const lease = controller.tryAcquireHeavy(); + assert.ok(lease, "precondition: a slot is available"); + assert.equal(controller.activeHeavy, 1); + + const abort = new AbortController(); + const wrapped = releaseChatAdmissionWhenDone(openSseResponse(), lease, { + signal: abort.signal, + }); + + // The client reads one chunk, then vanishes: no further pull, no cancel. + const reader = wrapped.body!.getReader(); + await reader.read(); + + // The connection drops. Next.js aborts the request signal. + abort.abort(); + await new Promise((resolve) => setTimeout(resolve, 10)); + + assert.equal( + controller.activeHeavy, + 0, + "an aborted client must not hold a heavyweight slot forever" + ); + assert.equal(lease.released, true); +}); + +test("an already-aborted request does not take a slot hostage", async () => { + const controller = new ChatAdmissionController(1); + const lease = controller.tryAcquireHeavy(); + assert.ok(lease); + + const abort = new AbortController(); + abort.abort(); + + releaseChatAdmissionWhenDone(openSseResponse(), lease, { signal: abort.signal }); + await new Promise((resolve) => setTimeout(resolve, 10)); + + assert.equal(controller.activeHeavy, 0); +}); + +test("normal completion still releases exactly once (no double release)", async () => { + const controller = new ChatAdmissionController(1); + const lease = controller.tryAcquireHeavy(); + assert.ok(lease); + + const finite = new Response( + new ReadableStream({ + start(c) { + c.enqueue(new TextEncoder().encode("data: done\n\n")); + c.close(); + }, + }), + { status: 200, headers: { "content-type": "text/event-stream" } } + ); + + const abort = new AbortController(); + const wrapped = releaseChatAdmissionWhenDone(finite, lease, { signal: abort.signal }); + const reader = wrapped.body!.getReader(); + while (!(await reader.read()).done) { + /* drain */ + } + + assert.equal(controller.activeHeavy, 0); + + // A late abort after a clean finish must not double-decrement. + abort.abort(); + await new Promise((resolve) => setTimeout(resolve, 10)); + assert.equal(controller.activeHeavy, 0); +}); + +test("explicit cancel still releases the slot", async () => { + const controller = new ChatAdmissionController(1); + const lease = controller.tryAcquireHeavy(); + assert.ok(lease); + + const abort = new AbortController(); + const wrapped = releaseChatAdmissionWhenDone(openSseResponse(), lease, { + signal: abort.signal, + }); + const reader = wrapped.body!.getReader(); + await reader.read(); + await reader.cancel("client gone"); + + assert.equal(controller.activeHeavy, 0); +});