From 22fa8f6623e8106119824fb36eb90a36eb19f4f2 Mon Sep 17 00:00:00 2001 From: roboomp Date: Wed, 15 Jul 2026 02:55:02 +0000 Subject: [PATCH] perf(stream): batched response decoding - Batched complete SSE lines into one UTF-8 decode per source chunk. - Parsed Ollama NDJSON bytes directly without TextDecoder buffering. Fixes #5542 --- packages/ai/CHANGELOG.md | 4 ++ packages/ai/src/providers/ollama.ts | 34 +------------ packages/utils/CHANGELOG.md | 4 ++ packages/utils/src/stream.ts | 74 +++++++++++++++++++---------- packages/utils/test/stream.test.ts | 15 +++++- 5 files changed, 72 insertions(+), 59 deletions(-) diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 3464f09cd..5bd2f20d5 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Parsed Ollama NDJSON response bytes directly instead of decoding and buffering every network chunk as text. ([#5542](https://github.com/can1357/oh-my-pi/issues/5542)) + ## [16.5.2] - 2026-07-14 ### Added diff --git a/packages/ai/src/providers/ollama.ts b/packages/ai/src/providers/ollama.ts index bfcdb65e0..851b347d3 100644 --- a/packages/ai/src/providers/ollama.ts +++ b/packages/ai/src/providers/ollama.ts @@ -1,4 +1,4 @@ -import { fetchWithRetry, parseStreamingJson } from "@oh-my-pi/pi-utils"; +import { fetchWithRetry, parseStreamingJson, readJsonl } from "@oh-my-pi/pi-utils"; import * as AIError from "../error"; import { getEnvApiKey } from "../stream"; import type { @@ -347,36 +347,6 @@ async function captureHttpErrorResponse(response: Response): Promise): AsyncGenerator { - const reader = stream.getReader(); - const decoder = new TextDecoder(); - let buffer = ""; - while (true) { - const { done, value } = await reader.read(); - if (done) { - break; - } - buffer += decoder.decode(value, { stream: true }); - while (true) { - const newlineIndex = buffer.indexOf("\n"); - if (newlineIndex < 0) { - break; - } - const line = buffer.slice(0, newlineIndex).trim(); - buffer = buffer.slice(newlineIndex + 1); - if (!line) { - continue; - } - yield JSON.parse(line) as OllamaChatChunk; - } - } - buffer += decoder.decode(); - const tail = buffer.trim(); - if (tail) { - yield JSON.parse(tail) as OllamaChatChunk; - } -} - function createEmptyOutput(model: Model<"ollama-chat">): AssistantMessage { return { role: "assistant", @@ -622,7 +592,7 @@ const streamOllamaOnce = ( }); } stream.push({ type: "start", partial: output }); - for await (const chunk of iterateNdjson(response.body)) { + for await (const chunk of readJsonl(response.body)) { if (chunk.message?.thinking) { suppressHealedThinking = true; endActiveTextBlock(); diff --git a/packages/utils/CHANGELOG.md b/packages/utils/CHANGELOG.md index 100d252d2..1e3606a36 100644 --- a/packages/utils/CHANGELOG.md +++ b/packages/utils/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Batched complete SSE lines into one UTF-8 decode per source chunk, avoiding per-line decoder overhead during active response streaming. ([#5542](https://github.com/can1357/oh-my-pi/issues/5542)) + ## [16.5.2] - 2026-07-14 ### Fixed diff --git a/packages/utils/src/stream.ts b/packages/utils/src/stream.ts index 0cbac4a8d..9ae9a7d0c 100644 --- a/packages/utils/src/stream.ts +++ b/packages/utils/src/stream.ts @@ -150,6 +150,29 @@ class ConcatSink { } } } + + appendAndFlushText(chunk: Uint8Array, decoder: TextDecoder): string | undefined { + const lastNewline = chunk.lastIndexOf(LF); + if (lastNewline === -1) { + this.append(chunk); + return undefined; + } + + const completeEnd = lastNewline + 1; + let text: string; + if (this.isEmpty) { + const complete = completeEnd === chunk.length ? chunk : chunk.subarray(0, completeEnd); + text = decoder.decode(complete); + } else { + this.append(completeEnd === chunk.length ? chunk : chunk.subarray(0, completeEnd)); + text = decoder.decode(this.flush()); + this.clear(); + } + if (completeEnd < chunk.length) { + this.append(chunk.subarray(completeEnd)); + } + return text; + } *pullJSONL(chunk: Uint8Array, beg: number, end: number) { if (this.isEmpty) { const { values, error, read, done } = parseJsonlChunkCompat(chunk, beg, end); @@ -275,14 +298,9 @@ interface SseEventState { raw: string[]; } -// Single decoder reused for all line decodes. Safe because lines are split on -// LF (0x0a) which is always a single-byte ASCII char in UTF-8 and never appears -// inside a multi-byte sequence — so each line is itself a complete UTF-8 run. -const SSE_LINE_DECODER = new TextDecoder("utf-8"); - -function decodeSseLineBytes(line: Uint8Array, end: number): string { - return end === line.length ? SSE_LINE_DECODER.decode(line) : SSE_LINE_DECODER.decode(line.subarray(0, end)); -} +// Complete lines are decoded in one batch per source chunk. Each batch ends on +// LF, which cannot split a multi-byte UTF-8 sequence. +const SSE_DECODER = new TextDecoder("utf-8"); function flushSseEvent(state: SseEventState): ServerSentEvent | null { if (state.event === null && state.data === null) { @@ -300,25 +318,25 @@ function flushSseEvent(state: SseEventState): ServerSentEvent | null { return event; } -function pushSseLine(line: Uint8Array, state: SseEventState): ServerSentEvent | null { - // `appendAndFlushLines` splits on LF only; strip a trailing CR so CRLF sources +function pushSseLine(line: string, state: SseEventState): ServerSentEvent | null { + // Complete-line batches split on LF only; strip a trailing CR so CRLF sources // don't leak `\r` into field values. - let end = line.length; - if (end > 0 && line[end - 1] === 0x0d /* '\r' */) end--; - if (end === 0) return flushSseEvent(state); + if (line.charCodeAt(line.length - 1) === 0x0d /* '\r' */) { + line = line.slice(0, -1); + } + if (line.length === 0) return flushSseEvent(state); // Comment line: keep in `raw` for diagnostic context, skip parsing. - if (line[0] === 0x3a /* ':' */) { - state.raw.push(decodeSseLineBytes(line, end)); + if (line.charCodeAt(0) === 0x3a /* ':' */) { + state.raw.push(line); return null; } - const text = decodeSseLineBytes(line, end); - state.raw.push(text); + state.raw.push(line); - const colon = text.indexOf(":"); - const fieldName = colon === -1 ? text : text.slice(0, colon); - let value = colon === -1 ? "" : text.slice(colon + 1); + const colon = line.indexOf(":"); + const fieldName = colon === -1 ? line : line.slice(0, colon); + let value = colon === -1 ? "" : line.slice(colon + 1); if (value.charCodeAt(0) === 0x20 /* ' ' */) value = value.slice(1); if (fieldName === "event") { @@ -344,9 +362,8 @@ function pushSseLine(line: Uint8Array, state: SseEventState): ServerSentEvent | * Use `readSseJson` instead when every event is a single `data:` JSON object * and you don't need access to the `event:` field. * - * Internally backed by a Buffer-based line reader (`ConcatSink`) so chunk - * concatenation is O(n) and never triggers per-line string slicing of the - * accumulated buffer. + * Internally backed by a Buffer-based reader (`ConcatSink`) that batches all + * complete lines in each source chunk into one UTF-8 decode. * * @example * ```ts @@ -365,9 +382,14 @@ export async function* readSseEvents( const source = abortableSource(stream, signal); try { for await (const chunk of source) { - for (const line of lineBuffer.appendAndFlushLines(chunk)) { - const event = pushSseLine(line, state); + const text = lineBuffer.appendAndFlushText(chunk, SSE_DECODER); + if (text === undefined) continue; + let start = 0; + while (start < text.length) { + const newline = text.indexOf("\n", start); + const event = pushSseLine(text.slice(start, newline), state); if (event) yield event; + start = newline + 1; } } // Treat any trailing partial line (no terminating LF) as a complete line. @@ -375,7 +397,7 @@ export async function* readSseEvents( const tail = lineBuffer.flush(); if (tail) { lineBuffer.clear(); - const event = pushSseLine(tail, state); + const event = pushSseLine(SSE_DECODER.decode(tail), state); if (event) { trailingEvents.add(event); yield event; diff --git a/packages/utils/test/stream.test.ts b/packages/utils/test/stream.test.ts index f13ff6e6d..3b3348782 100644 --- a/packages/utils/test/stream.test.ts +++ b/packages/utils/test/stream.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, it } from "bun:test"; +import { describe, expect, it, spyOn } from "bun:test"; import { sanitizeText } from "@oh-my-pi/pi-utils/sanitize-text"; import { parseJsonlLenient, @@ -362,6 +362,19 @@ describe("readSseEvents", () => { expect(events.map(e => e.data)).toEqual(['{"id":1}', "{}"]); }); + it("decodes all complete lines in a source chunk as one batch", async () => { + const decodeSpy = spyOn(TextDecoder.prototype, "decode"); + try { + const stream = bytesStreamFromChunks([encoder.encode("event: first\ndata: 1\n\nevent: second\ndata: 2\n\n")]); + const events = await collectAsync(readSseEvents(stream)); + + expect(events.map(event => event.data)).toEqual(["1", "2"]); + expect(decodeSpy).toHaveBeenCalledTimes(1); + } finally { + decodeSpy.mockRestore(); + } + }); + it("joins multiple data: lines with newlines", async () => { const stream = bytesStreamFromChunks([encoder.encode("event: chunk\ndata: line1\ndata: line2\ndata: line3\n\n")]); const [evt] = await collectAsync(readSseEvents(stream));