Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions changelog.d/fixes/14456-admission-lease-abort-release.md
Original file line number Diff line number Diff line change
@@ -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))
5 changes: 3 additions & 2 deletions src/app/api/v1/chat/completions/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand Down
5 changes: 3 additions & 2 deletions src/app/api/v1/responses/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
96 changes: 96 additions & 0 deletions src/shared/middleware/chatAdmissionRelease.ts
Original file line number Diff line number Diff line change
@@ -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<Uint8Array>({
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<Response>,
lease: ChatAdmissionLease | null,
options: ReleaseChatAdmissionOptions = {}
): Promise<Response> {
try {
return releaseChatAdmissionWhenDone(await responsePromise, lease, options);
} catch (error) {
lease?.release();
throw error;
}
}
60 changes: 7 additions & 53 deletions src/shared/middleware/chatBodyAdmission.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Response>,
lease: ChatAdmissionLease | null
): Promise<Response> {
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<Uint8Array>({
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";
3 changes: 2 additions & 1 deletion src/shared/middleware/withChatAdmission.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
121 changes: 121 additions & 0 deletions tests/unit/chat-admission-abandoned-stream-lease.test.ts
Original file line number Diff line number Diff line change
@@ -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<Uint8Array>({
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<Uint8Array>({
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);
});
Loading