diff --git a/changelog.d/fixes/12577-huggingchat-buffer-cap.md b/changelog.d/fixes/12577-huggingchat-buffer-cap.md new file mode 100644 index 0000000000..a50c56acad --- /dev/null +++ b/changelog.d/fixes/12577-huggingchat-buffer-cap.md @@ -0,0 +1 @@ +- fix(sse): cap HuggingChat NDJSON body size and bound the read loop with the fetch timeout so a stalled or hostile upstream cannot buffer unbounded memory (#12577) diff --git a/open-sse/config/constants.ts b/open-sse/config/constants.ts index ebe568d2d3..ece185e129 100644 --- a/open-sse/config/constants.ts +++ b/open-sse/config/constants.ts @@ -55,6 +55,14 @@ export const SSE_HEARTBEAT_INTERVAL_MS = upstreamTimeouts.sseHeartbeatIntervalMs // Defaults to FETCH_TIMEOUT_MS. Override with FETCH_BODY_TIMEOUT_MS env var. export const FETCH_BODY_TIMEOUT_MS = upstreamTimeouts.fetchBodyTimeoutMs; +// Hard byte cap on the HuggingChat NDJSON body accumulated by +// open-sse/executors/huggingchat/jsonlStream.ts. Prevents a stalled/hostile upstream that +// never emits a terminal `finalAnswer` / `status: finished` marker from buffering +// indefinitely (#12577). Sized generously for legitimate long completions while staying +// well below a heap-exhausting size — mirrors the readCappedBuffer/readBodyCapped pattern +// already used by veoaifree-web.ts and context7-fetch.ts. +export const HUGGINGCHAT_MAX_BODY_BYTES = 4 * 1024 * 1024; + // Provider configurations // OAuth credentials read from env vars with hardcoded fallbacks for backward compatibility. // Use provider-credentials.json or env vars to override in production. diff --git a/open-sse/executors/huggingchat.ts b/open-sse/executors/huggingchat.ts index 30ec7da0ea..38c3c11849 100644 --- a/open-sse/executors/huggingchat.ts +++ b/open-sse/executors/huggingchat.ts @@ -538,7 +538,7 @@ export class HuggingChatExecutor extends BaseExecutor { resolvedModel, id, created, - signal, + combinedSignal, streamCancellationController.signal ); @@ -626,7 +626,7 @@ export class HuggingChatExecutor extends BaseExecutor { let fullText: string; try { - fullText = await readJsonlResponse(upstreamResponse.body, signal); + fullText = await readJsonlResponse(upstreamResponse.body, combinedSignal); } catch (err) { if (!(err instanceof HuggingChatStreamError)) throw err; const message = err instanceof Error ? err.message : String(err); diff --git a/open-sse/executors/huggingchat/jsonlStream.ts b/open-sse/executors/huggingchat/jsonlStream.ts index 3d4980aebb..830f75d1db 100644 --- a/open-sse/executors/huggingchat/jsonlStream.ts +++ b/open-sse/executors/huggingchat/jsonlStream.ts @@ -1,5 +1,10 @@ // Pure JSONL stream translation (HuggingChat NDJSON -> OpenAI SSE). Verbatim from huggingchat.ts. +import { HUGGINGCHAT_MAX_BODY_BYTES } from "../../config/constants.ts"; + +const MAX_BODY_EXCEEDED_MESSAGE = + "HuggingChat response exceeded the maximum supported size before completing"; + export class HuggingChatStreamError extends Error { constructor(message: string) { super(message); @@ -74,15 +79,23 @@ export async function* streamJsonlToOpenAi( id: string, created: number, signal?: AbortSignal | null, - cancellationSignal?: AbortSignal | null + cancellationSignal?: AbortSignal | null, + maxBytes: number = HUGGINGCHAT_MAX_BODY_BYTES ): AsyncGenerator { const reader = body.getReader(); const unbindReaderCancellation = bindReaderCancellation(reader, cancellationSignal); + // Also bind the plain `signal` so an already-in-flight `reader.read()` unblocks the + // instant it aborts, instead of only being noticed the next time the loop polls + // `signal?.aborted` (#12577 — a stalled upstream can otherwise leave the read + // suspended forever even once a caller-supplied timeout signal has fired). + const unbindSignalCancellation = bindReaderCancellation(reader, signal); const decoder = new TextDecoder(); let buffer = ""; let emittedRole = false; let fullText = ""; let finished = false; + let totalBytes = 0; + let exceededCap = false; try { while (true) { @@ -91,6 +104,13 @@ export async function* streamJsonlToOpenAi( const { value, done } = await reader.read(); if (done) break; + totalBytes += value.byteLength; + if (totalBytes > maxBytes) { + exceededCap = true; + cancelReader(reader); + break; + } + buffer += decoder.decode(value, { stream: true }); const lines = buffer.split("\n"); @@ -163,7 +183,7 @@ export async function* streamJsonlToOpenAi( if (finished) break; } - if (!finished && buffer.trim()) { + if (!finished && !exceededCap && buffer.trim()) { const parsed = parseJsonlLine(buffer.trim()); if (parsed.error) { throw new HuggingChatStreamError(parsed.error); @@ -190,9 +210,26 @@ export async function* streamJsonlToOpenAi( } } finally { unbindReaderCancellation(); + unbindSignalCancellation(); reader.releaseLock(); } + if (exceededCap) { + yield sseChunk({ + id, + object: "chat.completion.chunk", + created, + model, + error: { + message: MAX_BODY_EXCEEDED_MESSAGE, + type: "upstream_error", + code: "huggingchat_payload_too_large", + }, + }); + yield "data: [DONE]\n\n"; + return; + } + if (!signal?.aborted && !cancellationSignal?.aborted) { yield sseChunk({ id, @@ -209,12 +246,19 @@ export async function* streamJsonlToOpenAi( export async function readJsonlResponse( body: ReadableStream, - signal?: AbortSignal | null + signal?: AbortSignal | null, + maxBytes: number = HUGGINGCHAT_MAX_BODY_BYTES ): Promise { const reader = body.getReader(); + // Bind the signal so an already-in-flight `reader.read()` unblocks the instant it + // aborts, instead of only being noticed the next time the loop polls `signal?.aborted` + // (#12577 — a stalled upstream can otherwise leave the read suspended forever even + // once a caller-supplied timeout signal has fired). + const unbindSignalCancellation = bindReaderCancellation(reader, signal); const decoder = new TextDecoder(); let buffer = ""; let fullText = ""; + let totalBytes = 0; try { while (true) { @@ -223,6 +267,12 @@ export async function readJsonlResponse( const { value, done } = await reader.read(); if (done) break; + totalBytes += value.byteLength; + if (totalBytes > maxBytes) { + cancelReader(reader); + throw new HuggingChatStreamError(MAX_BODY_EXCEEDED_MESSAGE); + } + buffer += decoder.decode(value, { stream: true }); const lines = buffer.split("\n"); @@ -249,6 +299,7 @@ export async function readJsonlResponse( if (parsed.error) throw new HuggingChatStreamError(parsed.error); } } finally { + unbindSignalCancellation(); reader.releaseLock(); } diff --git a/tests/unit/huggingchat-jsonlstream-unbounded-buffer.test.ts b/tests/unit/huggingchat-jsonlstream-unbounded-buffer.test.ts new file mode 100644 index 0000000000..b8e618957e --- /dev/null +++ b/tests/unit/huggingchat-jsonlstream-unbounded-buffer.test.ts @@ -0,0 +1,119 @@ +// Regression test for issue #12577: HuggingChat NDJSON executor buffered the +// upstream body with no byte ceiling and no timeout, so a stalled/hostile +// upstream that never emits a terminal marker (`finalAnswer` / `status: +// finished`) drove unbounded memory growth per in-flight request. + +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + streamJsonlToOpenAi, + readJsonlResponse, + HuggingChatStreamError, +} from "../../open-sse/executors/huggingchat/jsonlStream.ts"; + +const REASONABLE_CAP_BYTES = 2 * 1024 * 1024; // 2 MB +const TEST_SAFETY_CEILING_BYTES = REASONABLE_CAP_BYTES * 4; // 8 MB + +function makeUnboundedStream(): { + body: ReadableStream; + getTotalSent: () => number; + getClosedBySafetyCeiling: () => boolean; +} { + const encoder = new TextEncoder(); + const tokenChunk = "a".repeat(32 * 1024); // 32 KB token payload per line + const line = JSON.stringify({ type: "stream", token: tokenChunk }) + "\n"; + const lineBytes = encoder.encode(line).byteLength; + + let totalSent = 0; + let closedBySafetyCeiling = false; + + const body = new ReadableStream({ + pull(controller) { + if (totalSent >= TEST_SAFETY_CEILING_BYTES) { + closedBySafetyCeiling = true; + controller.close(); + return; + } + controller.enqueue(encoder.encode(line)); + totalSent += lineBytes; + // Deliberately never emit a finalAnswer/status:finished terminal marker. + }, + }); + + return { + body, + getTotalSent: () => totalSent, + getClosedBySafetyCeiling: () => closedBySafetyCeiling, + }; +} + +test("streamJsonlToOpenAi aborts once accumulated upstream body exceeds a size cap, instead of buffering forever", async () => { + const { body, getTotalSent, getClosedBySafetyCeiling } = makeUnboundedStream(); + const encoder = new TextEncoder(); + + let sawUpstreamErrorChunk = false; + let bytesReceivedByConsumer = 0; + + for await (const chunk of streamJsonlToOpenAi( + body, + "gpt-huggingchat", + "id-1", + 0, + undefined, + undefined, + REASONABLE_CAP_BYTES + )) { + bytesReceivedByConsumer += encoder.encode(chunk).byteLength; + if (/upstream_error|too_large|payload.*exceed/i.test(chunk)) { + sawUpstreamErrorChunk = true; + break; + } + } + + assert.ok( + sawUpstreamErrorChunk, + `expected streamJsonlToOpenAi to abort with an upstream-error chunk once the ` + + `accumulated body exceeded ~${REASONABLE_CAP_BYTES} bytes, but it kept consuming ` + + `upstream data with no ceiling (sent ${getTotalSent()} bytes before the TEST's own ` + + `safety ceiling stepped in: closedBySafetyCeiling=${getClosedBySafetyCeiling()}, ` + + `bytesReceivedByConsumer=${bytesReceivedByConsumer}). This confirms issue #12577: ` + + `no byte cap is enforced on the read loop.` + ); + assert.ok( + getTotalSent() < TEST_SAFETY_CEILING_BYTES, + "expected the cap to trip well before the test's own 8MB safety ceiling" + ); +}); + +test("readJsonlResponse throws a HuggingChatStreamError once accumulated upstream body exceeds a size cap", async () => { + const { body, getClosedBySafetyCeiling } = makeUnboundedStream(); + + await assert.rejects( + () => readJsonlResponse(body, undefined, REASONABLE_CAP_BYTES), + (err: unknown) => err instanceof HuggingChatStreamError + ); + assert.equal( + getClosedBySafetyCeiling(), + false, + "expected the cap to trip well before the test's own 8MB safety ceiling" + ); +}); + +test("streamJsonlToOpenAi terminates the read loop once an idle-timeout signal fires", async () => { + const body = new ReadableStream({ + pull() { + // Never enqueue and never close: simulates a stalled upstream connection + // that sends nothing at all after headers, relying solely on the caller's + // timeout signal (mirroring huggingchat.ts's combinedSignal) to unblock. + }, + }); + + const idleTimeout = AbortSignal.timeout(50); + const chunks: string[] = []; + + for await (const chunk of streamJsonlToOpenAi(body, "gpt-huggingchat", "id-2", 0, idleTimeout)) { + chunks.push(chunk); + } + + assert.ok(idleTimeout.aborted, "expected the idle-timeout signal to have fired"); +});