diff --git a/packages/utils/src/stream.ts b/packages/utils/src/stream.ts index 22037127d..072a8df1f 100644 --- a/packages/utils/src/stream.ts +++ b/packages/utils/src/stream.ts @@ -1,3 +1,5 @@ +const trailingEvents = new WeakSet(); + import { abortableSource } from "./abortable"; const LF = 0x0a; @@ -210,19 +212,122 @@ function notifySseEventObserver(observer: SseEventObserver | undefined, event: S } } +function getClosingSuffix(data: string): string | null { + let trimmed = data.trim(); + if (!(trimmed.startsWith("{") || trimmed.startsWith("["))) { + return null; + } + + let suffix = ""; + if (trimmed.endsWith(",")) { + trimmed = trimmed.slice(0, -1).trim(); + } + + if (trimmed.endsWith(":")) { + suffix = "null"; + } else { + const match = trimmed.match(/(t|tr|tru|f|fa|fal|fals|n|nu|nul)$/i); + if (match) { + const partial = match[0].toLowerCase(); + const completions: Record = { + t: "rue", + tr: "ue", + tru: "e", + f: "alse", + fa: "lse", + fal: "se", + fals: "e", + n: "ull", + nu: "ll", + nul: "l", + }; + if (completions[partial]) { + suffix = completions[partial]; + } + } + } + + const stack: string[] = []; + let inString = false; + let escaped = false; + + for (let i = 0; i < trimmed.length; i++) { + const char = trimmed[i]; + if (escaped) { + escaped = false; + continue; + } + if (char === "\\") { + escaped = true; + continue; + } + if (char === '"') { + inString = !inString; + continue; + } + if (!inString) { + if (char === "{" || char === "[") { + stack.push(char); + } else if (char === "}") { + if (stack.pop() !== "{") return null; + } else if (char === "]") { + if (stack.pop() !== "[") return null; + } + } + } + + if (inString) { + return `"${stack + .reverse() + .map(c => (c === "{" ? "}" : "]")) + .join("")}`; + } + return ( + suffix + + stack + .reverse() + .map(c => (c === "{" ? "}" : "]")) + .join("") + ); +} + +function isJsonTruncated(data: string): boolean { + let trimmed = data.trim(); + if (trimmed.endsWith(",")) { + trimmed = trimmed.slice(0, -1).trim(); + } + const suffix = getClosingSuffix(data); + if (suffix === null || suffix === "") return false; + + try { + JSON.parse(trimmed + suffix); + return true; + } catch { + return false; + } +} + export async function* readSseJson( stream: ReadableStream, signal?: AbortSignal, onEvent?: SseEventObserver, ): AsyncGenerator { for await (const sse of readSseEvents(stream, signal)) { + const isTrailing = trailingEvents.has(sse); notifySseEventObserver(onEvent, sse); const data = sse.data; if (data === "" || data === "[DONE]") { if (data === "[DONE]") return; continue; } - yield JSON.parse(data) as T; + try { + yield JSON.parse(data) as T; + } catch (err) { + if (err instanceof SyntaxError && isTrailing && isJsonTruncated(data)) { + return; + } + throw err; + } } } @@ -353,12 +458,18 @@ export async function* readSseEvents( if (tail) { lineBuffer.clear(); const event = pushSseLine(tail, state); - if (event) yield event; + if (event) { + trailingEvents.add(event); + yield event; + } } } // Real services don't always close on a blank line — flush any pending event. const trailing = flushSseEvent(state); - if (trailing) yield trailing; + if (trailing) { + trailingEvents.add(trailing); + yield trailing; + } } catch (err) { if (signal?.aborted) return; throw err; diff --git a/packages/utils/test/stream.test.ts b/packages/utils/test/stream.test.ts index 618330030..503ab0edd 100644 --- a/packages/utils/test/stream.test.ts +++ b/packages/utils/test/stream.test.ts @@ -239,6 +239,167 @@ describe("readSseJson", () => { const output = await collectAsync(readSseJson(stream)); expect(output).toEqual([{ a: 1 }]); }); + + it("completes cleanly when the final data chunk is truncated JSON", async () => { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b":2')]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const output = await collectAsync(readSseJson(stream)); + expect(output).toEqual([{ a: 1 }]); + }); + + it("completes cleanly when the final data chunk is cut inside a JSON literal at EOF", async () => { + const testCases = ['data: {"finish_reason":nul', 'data: {"ok":tru', "data: [fal"]; + for (const dataChunk of testCases) { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode(dataChunk)]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + const output = await collectAsync(readSseJson(stream)); + expect(output).toEqual([{ a: 1 }]); + } + }); + + it("throws SyntaxError when a middle data chunk is malformed JSON", async () => { + const chunks = [encoder.encode('data: {"a":1\n\n'), encoder.encode('data: {"b":2}\n\n')]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event is complete but malformed JSON", async () => { + const chunks = [ + encoder.encode('data: {"a":1}\n\n'), + encoder.encode('data: {"b":2,}'), // balanced but malformed trailing comma + ]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event is plain text", async () => { + const chunks = [ + encoder.encode('data: {"a":1}\n\n'), + encoder.encode("data: Internal Server Error"), // plain text (not JSON) + ]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event is malformed with missing colons but balanced braces", async () => { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b" 2}')]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event is malformed with mismatched brackets/braces", async () => { + const chunks = [ + encoder.encode('data: {"a":1}\n\n'), + encoder.encode("data: [{]"), // mismatched closer + ]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event is malformed with syntax errors after values", async () => { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b": true garbage')]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event has invalid characters", async () => { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b": @')]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event has mismatched closers", async () => { + const chunks = [ + encoder.encode('data: {"a":1}\n\n'), + encoder.encode('data: {"b": ]'), // mismatched closer + ]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event has unbalanced braces but internal syntax errors", async () => { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode('data: {"b":1 "c":2')]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); + + it("throws SyntaxError when a final event has unbalanced braces but has complete invalid text inside", async () => { + const chunks = [encoder.encode('data: {"a":1}\n\n'), encoder.encode("data: {unterminated}")]; + const stream = new ReadableStream({ + start(controller) { + for (const chunk of chunks) controller.enqueue(chunk); + controller.close(); + }, + }); + + await expect(collectAsync(readSseJson(stream))).rejects.toThrow(SyntaxError); + }); }); function bytesStreamFromChunks(chunks: Uint8Array[]): ReadableStream {