From 5acd9a92cd720e89eddefaa95838f1144e2a3083 Mon Sep 17 00:00:00 2001 From: xiaoyaner0201 Date: Tue, 22 Sep 2026 16:00:01 +0800 Subject: [PATCH 1/2] fix(resilience): release admission lease on client disconnect (#14456) releaseChatAdmissionWhenDone wired lease release into three consumer-driven paths only: pull-to-done, pull-throws, and cancel. A client that disconnects mid-stream may stop pulling and never cancel, so none of them run and the heavyweight slot is charged for the lifetime of the process. Once activeHeavy reaches OMNIROUTE_CHAT_MAX_HEAVY_IN_FLIGHT every heavy request is shed with 503 chat_admission_busy, with waiting=0 and no real load. Observe the request signal as a fallback release path and route all five call sites through it. Release is idempotent so a late abort after a clean finish cannot double-decrement. --- .../14456-admission-lease-abort-release.md | 1 + src/app/api/v1/chat/completions/route.ts | 5 +- src/app/api/v1/responses/route.ts | 5 +- src/shared/middleware/chatBodyAdmission.ts | 48 ++++++- src/shared/middleware/withChatAdmission.ts | 3 +- ...t-admission-abandoned-stream-lease.test.ts | 121 ++++++++++++++++++ 6 files changed, 172 insertions(+), 11 deletions(-) create mode 100644 changelog.d/fixes/14456-admission-lease-abort-release.md create mode 100644 tests/unit/chat-admission-abandoned-stream-lease.test.ts 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/chatBodyAdmission.ts b/src/shared/middleware/chatBodyAdmission.ts index ef38f4086ad..3bb6a809da5 100644 --- a/src/shared/middleware/chatBodyAdmission.ts +++ b/src/shared/middleware/chatBodyAdmission.ts @@ -1153,20 +1153,31 @@ export async function admitChatRequest( /** Release a lease if a handler rejects; otherwise bind it to the returned response lifecycle. */ export async function releaseChatAdmissionAfterHandler( responsePromise: Promise, - lease: ChatAdmissionLease | null + lease: ChatAdmissionLease | null, + options: ReleaseChatAdmissionOptions = {} ): Promise { try { - return releaseChatAdmissionWhenDone(await responsePromise, lease); + return releaseChatAdmissionWhenDone(await responsePromise, lease, options); } catch (error) { lease?.release(); throw error; } } +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 below 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 + lease: ChatAdmissionLease | null, + options: ReleaseChatAdmissionOptions = {} ): Response { if (!lease) return response; const isStreaming = response.headers.get("content-type")?.includes("text/event-stream"); @@ -1176,23 +1187,48 @@ export function releaseChatAdmissionWhenDone( } const reader = response.body.getReader(); + + // Release paths below are all driven by the consumer. If the client vanishes + // mid-stream the runtime may never pull again and never cancel, so the slot + // would be held until the process restarts. The request signal is the only + // event that still fires in that case. + 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) { - lease.release(); + releaseOnce(); controller.close(); } else { controller.enqueue(value); } } catch (error) { - lease.release(); + releaseOnce(); controller.error(error); } }, async cancel(reason) { - lease.release(); + releaseOnce(); await reader.cancel(reason).catch(() => undefined); }, }); 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); +}); From 7ed95a5952b98383e65b848ca82c4c25c8bc9c50 Mon Sep 17 00:00:00 2001 From: xiaoyaner0201 Date: Tue, 22 Sep 2026 16:20:20 +0800 Subject: [PATCH 2/2] refactor(admission): extract lease release into chatAdmissionRelease check:file-size freezes chatBodyAdmission.ts at 1206 lines and the fix pushed it to 1242. Move the release binding into its own module and re-export it from the original path, so every existing import site is unchanged. File is now 1159 lines; check:file-size passes. --- src/shared/middleware/chatAdmissionRelease.ts | 96 +++++++++++++++++++ src/shared/middleware/chatBodyAdmission.ts | 96 ++----------------- 2 files changed, 103 insertions(+), 89 deletions(-) create mode 100644 src/shared/middleware/chatAdmissionRelease.ts 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 3bb6a809da5..9efe0a83b36 100644 --- a/src/shared/middleware/chatBodyAdmission.ts +++ b/src/shared/middleware/chatBodyAdmission.ts @@ -1150,92 +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, - options: ReleaseChatAdmissionOptions = {} -): Promise { - try { - return releaseChatAdmissionWhenDone(await responsePromise, lease, options); - } catch (error) { - lease?.release(); - throw error; - } -} - -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 below 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(); - - // Release paths below are all driven by the consumer. If the client vanishes - // mid-stream the runtime may never pull again and never cancel, so the slot - // would be held until the process restarts. The request signal is the only - // event that still fires in that case. - 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, - }); -} +// Lease release binding lives in ./chatAdmissionRelease. Re-exported here so +// existing import sites keep working. +export { + releaseChatAdmissionAfterHandler, + releaseChatAdmissionWhenDone, + type ReleaseChatAdmissionOptions, +} from "./chatAdmissionRelease";