test(lib): cover SSE event-stream framing (20→29 tests)
Extract the SSE chunk-buffering from sdk.ts into a pure SSEParser (sdk.ts delegates; behavior unchanged) and pin the framing rules that are easy to break when an event splits across network reads: partial trailing lines held until completed, frames reassembled across 2-3 reads, [DONE] sentinel and empty/non-data lines filtered, each frame emitted exactly once. typecheck clean. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
63
src/lib/sse.test.ts
Normal file
63
src/lib/sse.test.ts
Normal file
@@ -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"}'])
|
||||
})
|
||||
31
src/lib/sse.ts
Normal file
31
src/lib/sse.ts
Normal file
@@ -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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user