diff --git a/src/lib/sdk.ts b/src/lib/sdk.ts index 33cfaa5..0b752fe 100644 --- a/src/lib/sdk.ts +++ b/src/lib/sdk.ts @@ -4,6 +4,7 @@ // expo/fetch provides WinterCG-compliant fetch with ReadableStream support for SSE import { fetch as expoFetch } from "expo/fetch" import { buildRequestHeaders } from "./headers" +import { SSEParser } from "./sse" export interface ClientConfig { baseUrl: string @@ -220,27 +221,18 @@ export function createClient(config: ClientConfig) { const reader = response.body.getReader() const decoder = new TextDecoder() - let buffer = "" + const parser = new SSEParser() try { while (true) { const { done, value } = await reader.read() if (done) break - buffer += decoder.decode(value, { stream: true }) - const lines = buffer.split("\n") - buffer = lines.pop() || "" - - for (const line of lines) { - if (line.startsWith("data: ")) { - const data = line.slice(6) - if (data && data !== "[DONE]") { - try { - yield JSON.parse(data) - } catch (err) { - console.warn("[SSE] Failed to parse event:", data.slice(0, 200), err) - } - } + for (const data of parser.push(decoder.decode(value, { stream: true }))) { + try { + yield JSON.parse(data) + } catch (err) { + console.warn("[SSE] Failed to parse event:", data.slice(0, 200), err) } } } diff --git a/src/lib/sse.test.ts b/src/lib/sse.test.ts new file mode 100644 index 0000000..396ab51 --- /dev/null +++ b/src/lib/sse.test.ts @@ -0,0 +1,63 @@ +import { test } from "node:test" +import assert from "node:assert/strict" +import { SSEParser } from "./sse.ts" + +// The opencode live chat is driven by this SSE stream. The hard cases are framing: +// events that arrive split across network reads must not be dropped or corrupted. + +test("parses a single complete data frame in one chunk", () => { + const p = new SSEParser() + assert.deepEqual(p.push('data: {"x":1}\n'), ['{"x":1}']) +}) + +test("parses multiple frames in one chunk in order", () => { + const p = new SSEParser() + assert.deepEqual(p.push("data: a\ndata: b\ndata: c\n"), ["a", "b", "c"]) +}) + +test("holds a partial trailing line until the next chunk completes it", () => { + const p = new SSEParser() + // First read ends mid-frame (no trailing newline yet). + assert.deepEqual(p.push('data: {"hello":'), []) + // Second read delivers the rest plus the terminating newline. + assert.deepEqual(p.push('"world"}\n'), ['{"hello":"world"}']) +}) + +test("reassembles a frame split across three reads", () => { + const p = new SSEParser() + assert.deepEqual(p.push("dat"), []) + assert.deepEqual(p.push("a: par"), []) + assert.deepEqual(p.push("tial\n"), ["partial"]) +}) + +test("filters the [DONE] sentinel", () => { + const p = new SSEParser() + assert.deepEqual(p.push("data: [DONE]\n"), []) + assert.deepEqual(p.push("data: real\ndata: [DONE]\n"), ["real"]) +}) + +test("ignores empty data payloads and non-data lines", () => { + const p = new SSEParser() + // keepalive comment, an event: line, a blank line, and an empty data payload + assert.deepEqual(p.push(": keepalive\nevent: ping\n\ndata: \ndata: x\n"), ["x"]) +}) + +test("a frame held across chunks is emitted exactly once", () => { + const p = new SSEParser() + const first = p.push("data: one\ndata: tw") + const second = p.push("o\ndata: three\n") + assert.deepEqual(first, ["one"]) + assert.deepEqual(second, ["two", "three"]) +}) + +test("blank chunk yields nothing and preserves the held buffer", () => { + const p = new SSEParser() + assert.deepEqual(p.push("data: held"), []) + assert.deepEqual(p.push(""), []) + assert.deepEqual(p.push("\n"), ["held"]) +}) + +test("data payload preserves internal colons and spaces (only the 'data: ' prefix is stripped)", () => { + const p = new SSEParser() + assert.deepEqual(p.push('data: {"url":"http://x:1"}\n'), ['{"url":"http://x:1"}']) +}) diff --git a/src/lib/sse.ts b/src/lib/sse.ts new file mode 100644 index 0000000..d91fa1c --- /dev/null +++ b/src/lib/sse.ts @@ -0,0 +1,31 @@ +// Pure, transport-agnostic Server-Sent Events framing for the opencode event stream. +// Extracted from sdk.ts so the chunk-buffering rules — which are easy to get wrong when +// a single event is split across two network reads — are unit-testable without a socket. + +export class SSEParser { + private buffer = "" + + /** + * Feed one decoded text chunk. Returns the `data:` payloads (as raw strings) + * that completed within this chunk. A partial trailing line is retained and + * prepended to the next chunk. The `[DONE]` sentinel and empty payloads are + * filtered out. Lines that are not `data:` frames (comments, `event:`, blanks) + * are ignored. Callers JSON-parse the returned strings. + */ + push(chunk: string): string[] { + this.buffer += chunk + const lines = this.buffer.split("\n") + // The last element is whatever came after the final newline — possibly an + // incomplete line. Hold it until the next chunk completes it. + this.buffer = lines.pop() ?? "" + + const out: string[] = [] + for (const line of lines) { + if (line.startsWith("data: ")) { + const data = line.slice(6) + if (data && data !== "[DONE]") out.push(data) + } + } + return out + } +}