Python 88.3%
TypeScript 7.6%
Shell 4.1%
1/**2 * Generic SSE parser + provider normalisers for OpenAI (Responses API), Anthropic (Messages API), xAI (Chat Completions chunks ·3 * Responses events · Anthropic-compatible /v1/messages) and Google Gemini (streamGenerateContent `?alt=sse` · JSON-array framing ·4 * Interactions API `step.*` events). No dependencies.5 *6 * STATUS: DOCUMENTED (WHATWG SSE framing; event vocabularies from the OpenAPI specs / discovery document / build-with-claude/streaming.md) ·7 * offline self-test on recorded fixtures tests/shared/fixtures/sse/* (2026-09-18 OpenAI/Anthropic; 2026-09-19 xAI/Gemini fixtures cut from8 * live captures in tmp-live/, thoughtSignature blobs truncated):9 * node --experimental-strip-types examples/shared/streaming/sseParser.ts --selftest10 *11 * Twin of sse_parser.py. Layers: SSEParser / JSONArrayStreamParser (incremental, bounded buffer) → PartialJSON → OpenAINormalizer /12 * AnthropicNormalizer / XAINormalizer / GeminiNormalizer → unified events (text_delta · reasoning_delta · tool_call_start ·13 * tool_call_delta · tool_call_done · usage · done · error · raw). `iterUnified(body, provider, framing)` consumes a fetch ReadableStream14 * with natural backpressure (one chunk is read only after the previous chunk's events were consumed).15 */1617export type ProviderName = "openai" | "anthropic" | "xai" | "gemini";18export interface SSEEvent { event?: string; data: string; id?: string; retry?: number }19export class SSEBufferOverflow extends Error {}2021export class SSEParser {22 private buf = ""; private event?: string; private data: string[] = []; private id?: string; private retry?: number;23 private dec = new TextDecoder("utf-8", { ignoreBOM: false }); eventsParsed = 0; maxBuffer: number;24 constructor(maxBuffer = 8 * 1024 * 1024) { this.maxBuffer = maxBuffer; }2526 feed(chunk: Uint8Array | string): SSEEvent[] {27 this.buf += typeof chunk === "string" ? chunk : this.dec.decode(chunk, { stream: true });28 if (this.buf.length > this.maxBuffer) throw new SSEBufferOverflow(`single SSE event exceeds ${this.maxBuffer} bytes`);29 const out: SSEEvent[] = []; let nl: number;30 while ((nl = this.buf.indexOf("\n")) >= 0) {31 let line = this.buf.slice(0, nl); this.buf = this.buf.slice(nl + 1);32 if (line.endsWith("\r")) line = line.slice(0, -1);33 if (line.charCodeAt(0) === 0xfeff) line = line.slice(1);34 const ev = this.line(line); if (ev) out.push(ev);35 }36 return out;37 }38 close(): SSEEvent[] {39 const out: SSEEvent[] = [];40 this.buf += this.dec.decode();41 if (this.buf) { const ev = this.line(this.buf.replace(/\r$/, "")); this.buf = ""; if (ev) out.push(ev); }42 if (this.data.length) out.push(this.dispatch());43 return out;44 }45 private line(line: string): SSEEvent | undefined {46 if (line === "") return this.data.length || this.event !== undefined ? this.dispatch() : undefined;47 if (line.startsWith(":")) return undefined;48 const i = line.indexOf(":"); const name = i < 0 ? line : line.slice(0, i); let value = i < 0 ? "" : line.slice(i + 1);49 if (value.startsWith(" ")) value = value.slice(1);50 if (name === "event") this.event = value; else if (name === "data") this.data.push(value);51 else if (name === "id") { if (!value.includes("\0")) this.id = value; } else if (name === "retry") { if (/^\d+$/.test(value)) this.retry = parseInt(value, 10); }52 return undefined;53 }54 private dispatch(): SSEEvent { const ev: SSEEvent = { event: this.event, data: this.data.join("\n"), id: this.id, retry: this.retry }; this.event = undefined; this.data = []; this.eventsParsed++; return ev; }55}5657/** Gemini's NON-SSE streaming body (no `?alt=sse`): a pretty-printed JSON array flushed object by object. Emits each top-level object as an SSEEvent. */58export class JSONArrayStreamParser {59 private buf = ""; private depth = 0; private inStr = false; private esc = false; private start: number | undefined; private scanned = 0;60 private dec = new TextDecoder("utf-8"); eventsParsed = 0; maxBuffer: number;61 constructor(maxBuffer = 8 * 1024 * 1024) { this.maxBuffer = maxBuffer; }62 feed(chunk: Uint8Array | string): SSEEvent[] {63 this.buf += typeof chunk === "string" ? chunk : this.dec.decode(chunk, { stream: true });64 if (this.buf.length > this.maxBuffer) throw new SSEBufferOverflow(`single JSON object exceeds ${this.maxBuffer} bytes`);65 const out: SSEEvent[] = []; let i = this.scanned; let consumed = this.start !== undefined ? 0 : this.scanned;66 for (; i < this.buf.length; i++) {67 const ch = this.buf[i];68 if (this.depth === 0 && !this.inStr) { if (ch === "{") { this.start = i; this.depth = 1; } else consumed = i + 1; continue; }69 if (this.inStr) { if (this.esc) this.esc = false; else if (ch === "\\") this.esc = true; else if (ch === '"') this.inStr = false; }70 else if (ch === '"') this.inStr = true; else if (ch === "{") this.depth++;71 else if (ch === "}") { this.depth--; if (this.depth === 0 && this.start !== undefined) { out.push({ data: this.buf.slice(this.start, i + 1) }); this.eventsParsed++; consumed = i + 1; this.start = undefined; } }72 }73 if (this.start !== undefined) { this.buf = this.buf.slice(this.start); this.start = 0; } else this.buf = this.buf.slice(consumed);74 this.scanned = this.buf.length;75 return out;76 }77 close(): SSEEvent[] { this.buf = ""; return []; }78}7980export const jsonOf = (ev: SSEEvent): any => { try { return JSON.parse(ev.data); } catch { return undefined; } };8182export function parseSseText(text: string, chunkSize?: number): SSEEvent[] {83 const p = new SSEParser(); const out: SSEEvent[] = [];84 if (chunkSize) for (let i = 0; i < text.length; i += chunkSize) out.push(...p.feed(text.slice(i, i + chunkSize))); else out.push(...p.feed(text));85 return out.concat(p.close());86}8788// ------------------------------------------------------------------ partial JSON89export function parsePartialJson(text: string): any {90 if (!text.trim()) return undefined;91 try { return JSON.parse(text); } catch { /* fall through */ }92 const stack: string[] = []; let inStr = false, esc = false;93 for (const ch of text) {94 if (inStr) { if (esc) esc = false; else if (ch === "\\") esc = true; else if (ch === '"') inStr = false; }95 else if (ch === '"') inStr = true; else if (ch === "{" || ch === "[") stack.push(ch === "{" ? "}" : "]"); else if ((ch === "}" || ch === "]") && stack.length) stack.pop();96 }97 let fixed = (inStr ? text + '"' : text).trimEnd();98 const dropDanglingKey = (s: string) => {99 s = s.trimEnd().replace(/,$/, "").trimEnd(); if (s.endsWith(":")) s = s.slice(0, -1).trimEnd();100 if (s.endsWith('"')) { const j = s.lastIndexOf('"', s.length - 2); const before = s.slice(0, j).trimEnd(); if (before.endsWith("{") || before.endsWith(",")) s = before.replace(/,$/, "").trimEnd(); }101 return s;102 };103 for (let i = 0; i < 3; i++) {104 try { return JSON.parse(fixed + [...stack].reverse().join("")); } catch { const nxt = dropDanglingKey(fixed); if (nxt === fixed) break; fixed = nxt; }105 }106 return undefined;107}108109export class PartialJSON {110 buffers = new Map<string, string>();111 append(key: string, fragment: string): string { const v = (this.buffers.get(key) ?? "") + fragment; this.buffers.set(key, v); return v; }112 preview(key: string): any { return parsePartialJson(this.buffers.get(key) ?? ""); }113 finish(key: string, final?: string): Record<string, unknown> {114 const raw = final ?? this.buffers.get(key) ?? ""; this.buffers.delete(key);115 if (!raw.trim()) return {};116 try { const v = JSON.parse(raw); return v && typeof v === "object" && !Array.isArray(v) ? v : { _value: v }; } catch { return { _raw: raw, _error: "invalid JSON" }; }117 }118}119120// ------------------------------------------------------------------ unified events121export type Unified =122 | { type: "text_delta" | "reasoning_delta"; text: string; sequence?: number; raw: any }123 | { type: "tool_call_start"; toolCallId?: string; toolName?: string; sequence?: number; raw: any }124 | { type: "tool_call_delta"; toolCallId?: string; toolName?: string; argumentsDelta: string; sequence?: number; raw: any }125 | { type: "tool_call_done"; toolCallId?: string; toolName?: string; arguments: Record<string, unknown>; sequence?: number; raw: any }126 | { type: "usage"; usage: Record<string, unknown>; sequence?: number; raw: any }127 | { type: "done"; stopReason: string; sequence?: number; raw: any }128 | { type: "error"; error: any; sequence?: number; raw: any }129 | { type: "raw"; sequence?: number; raw: any };130131export interface Normalizer { (ev: SSEEvent): Unified[] }132133export class OpenAINormalizer {134 static TERMINAL: Record<string, string> = { "response.completed": "end", "response.incomplete": "incomplete", "response.failed": "failed", "response.cancelled": "cancelled" };135 args = new PartialJSON(); items = new Map<string, any>(); lastSequence?: number; responseId?: string;136 private callId(itemId?: string) { return this.items.get(itemId ?? "")?.call_id ?? itemId; }137 normalize = (ev: SSEEvent): Unified[] => {138 const obj = jsonOf(ev);139 if (!obj || typeof obj !== "object") return ev.data === "[DONE]" ? [] : [{ type: "raw", raw: { event: ev.event, data: ev.data } }];140 const t: string = obj.type ?? ev.event ?? ""; const sequence = typeof obj.sequence_number === "number" ? obj.sequence_number : undefined;141 if (sequence !== undefined) this.lastSequence = sequence;142 if (!this.responseId && obj.response?.id) this.responseId = obj.response.id;143 const base = { sequence, raw: obj };144 if (t === "response.output_text.delta" || t === "response.refusal.delta") return [{ type: "text_delta", text: obj.delta ?? "", ...base }];145 if (t === "response.reasoning_summary_text.delta" || t === "response.reasoning_text.delta") return [{ type: "reasoning_delta", text: obj.delta ?? "", ...base }];146 if (t === "response.output_item.added") {147 const item = obj.item ?? {};148 if (["function_call", "custom_tool_call", "mcp_call"].includes(item.type)) { this.items.set(item.id, item); return [{ type: "tool_call_start", toolCallId: item.call_id ?? item.id, toolName: item.name, ...base }]; }149 return [{ type: "raw", ...base }];150 }151 if (["response.function_call_arguments.delta", "response.custom_tool_call_input.delta", "response.mcp_call_arguments.delta"].includes(t)) {152 this.args.append(obj.item_id, obj.delta ?? ""); return [{ type: "tool_call_delta", toolCallId: this.callId(obj.item_id), toolName: this.items.get(obj.item_id)?.name, argumentsDelta: obj.delta ?? "", ...base }];153 }154 if (["response.function_call_arguments.done", "response.custom_tool_call_input.done", "response.mcp_call_arguments.done"].includes(t)) {155 const final = "arguments" in obj ? obj.arguments : obj.input;156 return [{ type: "tool_call_done", toolCallId: this.callId(obj.item_id), toolName: this.items.get(obj.item_id)?.name ?? obj.name, arguments: this.args.finish(obj.item_id, typeof final === "string" ? final : undefined), ...base }];157 }158 if (t in OpenAINormalizer.TERMINAL) {159 const resp = obj.response ?? {}; const out: Unified[] = [];160 if (resp.usage && typeof resp.usage === "object") out.push({ type: "usage", usage: resp.usage, ...base });161 let stop = OpenAINormalizer.TERMINAL[t]; if (t === "response.incomplete" && resp.incomplete_details?.reason === "max_output_tokens") stop = "max_tokens";162 out.push({ type: "done", stopReason: stop, ...base }); return out;163 }164 if (t === "error") return [{ type: "error", error: obj, ...base }];165 return [{ type: "raw", ...base }];166 };167}168169export class AnthropicNormalizer {170 static STOP: Record<string, string> = { end_turn: "end", max_tokens: "max_tokens", tool_use: "tool_use", stop_sequence: "stop_sequence", refusal: "refusal", pause_turn: "incomplete", model_context_window_exceeded: "incomplete" };171 static MESSAGE_EVENTS = new Set(["message_start", "content_block_start", "content_block_delta", "content_block_stop", "message_delta", "message_stop", "ping"]);172 args = new PartialJSON(); blocks = new Map<number, any>(); usage: Record<string, unknown> = {}; stopReason?: string; messageId?: string;173 normalize = (ev: SSEEvent): Unified[] => {174 const obj = jsonOf(ev);175 if (!obj || typeof obj !== "object") return [{ type: "raw", raw: { event: ev.event, data: ev.data } }];176 const t: string = obj.type ?? ev.event ?? ""; const base = { raw: obj };177 const isTool = (cb: any) => ["tool_use", "server_tool_use", "mcp_tool_use"].includes(cb?.type);178 if (t === "message_start") { this.messageId = obj.message?.id; this.usage = { ...(obj.message?.usage ?? {}) }; return [{ type: "raw", ...base }]; }179 if (t === "content_block_start") { const cb = obj.content_block ?? {}; this.blocks.set(obj.index, cb); return isTool(cb) ? [{ type: "tool_call_start", toolCallId: cb.id, toolName: cb.name, ...base }] : [{ type: "raw", ...base }]; }180 if (t === "content_block_delta") {181 const d = obj.delta ?? {};182 if (d.type === "text_delta") return [{ type: "text_delta", text: d.text ?? "", ...base }];183 if (d.type === "thinking_delta") return [{ type: "reasoning_delta", text: d.thinking ?? "", ...base }];184 if (d.type === "input_json_delta") { const cb = this.blocks.get(obj.index) ?? {}; this.args.append(String(obj.index), d.partial_json ?? ""); return [{ type: "tool_call_delta", toolCallId: cb.id, toolName: cb.name, argumentsDelta: d.partial_json ?? "", ...base }]; }185 return [{ type: "raw", ...base }];186 }187 if (t === "content_block_stop") {188 const cb = this.blocks.get(obj.index) ?? {};189 if (isTool(cb)) { let args = this.args.finish(String(obj.index)); if (!Object.keys(args).length && cb.input && Object.keys(cb.input).length) args = cb.input; return [{ type: "tool_call_done", toolCallId: cb.id, toolName: cb.name, arguments: args, ...base }]; }190 return [{ type: "raw", ...base }];191 }192 if (t === "message_delta") { this.stopReason = AnthropicNormalizer.STOP[obj.delta?.stop_reason] ?? undefined; Object.assign(this.usage, obj.usage ?? {}); return [{ type: "usage", usage: { ...this.usage }, ...base }]; }193 if (t === "message_stop") return [{ type: "done", stopReason: this.stopReason ?? "end", ...base }];194 if (t === "error") return [{ type: "error", error: obj.error ?? obj, ...base }];195 return [{ type: "raw", ...base }];196 };197}198199/** xAI: Chat Completions chunks (reasoning_content → reasoning_delta; whole tool_calls in one chunk; usage chunk; [DONE] → done),200 * Responses events (delegated to OpenAINormalizer) and Anthropic-compatible /v1/messages events (delegated to AnthropicNormalizer). */201export class XAINormalizer {202 static CHAT_FINISH: Record<string, string> = { stop: "end", length: "max_tokens", tool_calls: "tool_use", content_filter: "refusal" };203 responses = new OpenAINormalizer(); messages = new AnthropicNormalizer(); chatId?: string; chatStop?: string; sawChat = false;204 normalize = (ev: SSEEvent): Unified[] => {205 const obj = jsonOf(ev);206 if (!obj || typeof obj !== "object") {207 if (ev.data === "[DONE]") return this.sawChat ? [{ type: "done", stopReason: this.chatStop ?? "end", raw: { data: "[DONE]" } }] : [];208 return [{ type: "raw", raw: { event: ev.event, data: ev.data } }];209 }210 const t: string = obj.type ?? ev.event ?? "";211 if (obj.object === "chat.completion.chunk" || ("choices" in obj && !t)) return this.chat(obj);212 if (t.startsWith("response.") || (t === "error" && "sequence_number" in obj)) return this.responses.normalize(ev);213 if (AnthropicNormalizer.MESSAGE_EVENTS.has(t)) return this.messages.normalize(ev);214 if (t === "error") return [{ type: "error", error: obj, raw: obj }];215 return [{ type: "raw", raw: obj }];216 };217 private chat(obj: any): Unified[] {218 this.sawChat = true; this.chatId = obj.id ?? this.chatId;219 const out: Unified[] = []; const choices = obj.choices ?? [];220 if (!choices.length && obj.usage && typeof obj.usage === "object") return [{ type: "usage", usage: obj.usage, raw: obj }];221 for (const ch of choices) {222 const d = ch.delta ?? {};223 if (d.reasoning_content) out.push({ type: "reasoning_delta", text: d.reasoning_content, raw: obj });224 if (d.content) out.push({ type: "text_delta", text: d.content, raw: obj });225 for (const tc of d.tool_calls ?? []) {226 const fn = tc.function ?? {}; const rawArgs: string = fn.arguments ?? ""; const id = tc.id ?? `${this.chatId}:${tc.index ?? 0}`;227 out.push({ type: "tool_call_start", toolCallId: id, toolName: fn.name, raw: obj });228 if (rawArgs) out.push({ type: "tool_call_delta", toolCallId: id, toolName: fn.name, argumentsDelta: rawArgs, raw: obj });229 out.push({ type: "tool_call_done", toolCallId: id, toolName: fn.name, arguments: rawArgs ? new PartialJSON().finish(id, rawArgs) : {}, raw: obj });230 }231 if (ch.finish_reason) this.chatStop = XAINormalizer.CHAT_FINISH[ch.finish_reason] ?? ch.finish_reason;232 }233 return out.length ? out : [{ type: "raw", raw: obj }];234 }235}236237/** Gemini: generateContent chunks (SSE or JSON-array framing) and Interactions API `event_type` events — see sse_parser.py for the rules. */238export class GeminiNormalizer {239 static REFUSAL = new Set(["SAFETY", "RECITATION", "BLOCKLIST", "PROHIBITED_CONTENT", "SPII", "IMAGE_SAFETY", "IMAGE_PROHIBITED_CONTENT", "IMAGE_RECITATION", "IMAGE_OTHER", "LANGUAGE"]);240 static INTERACTION_STATUS: Record<string, string> = { completed: "end", requires_action: "tool_use", incomplete: "incomplete", failed: "failed", cancelled: "cancelled" };241 args = new PartialJSON(); usage: Record<string, unknown> = {}; finishReason?: string; sawCall = false; responseId?: string; modelVersion?: string; lastSignature?: string;242 steps = new Map<number, any>(); interactionId?: string;243 normalize = (ev: SSEEvent): Unified[] => {244 const obj = jsonOf(ev);245 if (!obj || typeof obj !== "object") return ev.data === "[DONE]" ? [] : [{ type: "raw", raw: { event: ev.event, data: ev.data } }];246 if ("event_type" in obj || ["step.delta", "step.start", "step.stop", "interaction.created", "interaction.completed"].includes(ev.event ?? "")) return this.interaction(obj, obj.event_type ?? ev.event ?? "");247 return this.chunk(obj);248 };249 private stop(): string {250 const fr = this.finishReason;251 if (this.sawCall) return "tool_use"; if (!fr || fr === "STOP") return "end"; if (fr === "MAX_TOKENS") return "max_tokens";252 return GeminiNormalizer.REFUSAL.has(fr) || fr === "BLOCKED" ? "refusal" : "other";253 }254 private chunk(obj: any): Unified[] {255 const base = { raw: obj };256 if (obj.error && typeof obj.error === "object") return [{ type: "error", error: obj.error, ...base }];257 if (obj.usageMetadata && typeof obj.usageMetadata === "object") this.usage = { ...obj.usageMetadata };258 this.responseId = obj.responseId ?? this.responseId; this.modelVersion = obj.modelVersion ?? this.modelVersion;259 const c = obj.candidates?.[0]; const out: Unified[] = [];260 if (!c) {261 if (obj.promptFeedback) { this.finishReason = "BLOCKED"; return [{ type: "usage", usage: { ...this.usage }, ...base }, { type: "done", stopReason: "refusal", ...base }]; }262 return [{ type: "raw", ...base }];263 }264 (c.content?.parts ?? []).forEach((p: any, i: number) => {265 if (p.thoughtSignature) this.lastSignature = p.thoughtSignature;266 if (p.functionCall) {267 const args = p.functionCall.args ?? {}; const id = p.functionCall.id ?? `call_${this.steps.size}_${i}`; this.sawCall = true;268 out.push({ type: "tool_call_start", toolCallId: id, toolName: p.functionCall.name, ...base });269 out.push({ type: "tool_call_delta", toolCallId: id, toolName: p.functionCall.name, argumentsDelta: JSON.stringify(args), ...base });270 out.push({ type: "tool_call_done", toolCallId: id, toolName: p.functionCall.name, arguments: args && typeof args === "object" && !Array.isArray(args) ? args : { _value: args }, ...base });271 } else if ("text" in p) {272 if (p.thought) out.push({ type: "reasoning_delta", text: p.text ?? "", ...base });273 else if (p.text) out.push({ type: "text_delta", text: p.text, ...base });274 else out.push({ type: "raw", ...base }); // empty-text signature carrier (Gemini 3 last chunk)275 } else out.push({ type: "raw", ...base }); // inlineData, executableCode, codeExecutionResult, fileData…276 });277 if (c.finishReason) { this.finishReason = c.finishReason; out.push({ type: "usage", usage: { ...this.usage }, ...base }); out.push({ type: "done", stopReason: this.stop(), ...base }); }278 return out.length ? out : [{ type: "raw", ...base }];279 }280 private interaction(obj: any, t: string): Unified[] {281 const base = { raw: obj }; const idx = obj.index;282 if (t === "interaction.created") { this.interactionId = obj.interaction?.id; return [{ type: "raw", ...base }]; }283 if (t === "step.start") { const step = obj.step ?? {}; this.steps.set(idx, step); return step.type === "function_call" ? [{ type: "tool_call_start", toolCallId: step.id, toolName: step.name, ...base }] : [{ type: "raw", ...base }]; }284 if (t === "step.delta") {285 const d = obj.delta ?? {}; const step = this.steps.get(idx) ?? {};286 if (d.type === "text") return [{ type: "text_delta", text: d.text ?? "", ...base }];287 if (d.type === "thought_summary") return [{ type: "reasoning_delta", text: d.content?.text ?? "", ...base }];288 if (d.type === "thought_signature") { this.lastSignature = d.signature ?? this.lastSignature; return [{ type: "raw", ...base }]; }289 if (d.type === "arguments_delta") { const frag = d.arguments ?? ""; this.args.append(String(idx), frag); return [{ type: "tool_call_delta", toolCallId: step.id, toolName: step.name, argumentsDelta: frag, ...base }]; }290 return [{ type: "raw", ...base }];291 }292 if (t === "step.stop") {293 const step = this.steps.get(idx) ?? {};294 if (step.type === "function_call") { let args = this.args.finish(String(idx)); if (!Object.keys(args).length && step.arguments && Object.keys(step.arguments).length) args = step.arguments; return [{ type: "tool_call_done", toolCallId: step.id, toolName: step.name, arguments: args, ...base }]; }295 return [{ type: "raw", ...base }];296 }297 if (t === "interaction.completed") {298 const inter = obj.interaction ?? {}; const out: Unified[] = [];299 if (inter.usage && typeof inter.usage === "object") out.push({ type: "usage", usage: inter.usage, ...base });300 out.push({ type: "done", stopReason: GeminiNormalizer.INTERACTION_STATUS[inter.status] ?? inter.status ?? "end", ...base }); return out;301 }302 if (t === "error") return [{ type: "error", error: obj.error ?? obj, ...base }];303 return [{ type: "raw", ...base }];304 }305}306307export const normalizerFor = (provider: ProviderName): Normalizer =>308 provider === "openai" ? new OpenAINormalizer().normalize : provider === "anthropic" ? new AnthropicNormalizer().normalize : provider === "xai" ? new XAINormalizer().normalize : new GeminiNormalizer().normalize;309310/** Consume a fetch body with backpressure: the next chunk is only read after the caller has consumed this chunk's events. */311export async function* iterUnified(body: ReadableStream<Uint8Array> | AsyncIterable<Uint8Array> | null | undefined, provider: ProviderName, framing: "sse" | "json_array" = "sse"): AsyncGenerator<Unified> {312 if (!body) return;313 const parser = framing === "json_array" ? new JSONArrayStreamParser() : new SSEParser(); const norm = normalizerFor(provider);314 const source: AsyncIterable<Uint8Array> = Symbol.asyncIterator in (body as any) ? (body as AsyncIterable<Uint8Array>) : readerIterable(body as ReadableStream<Uint8Array>);315 for await (const chunk of source) for (const ev of parser.feed(chunk)) yield* norm(ev);316 for (const ev of parser.close()) yield* norm(ev);317}318async function* readerIterable(stream: ReadableStream<Uint8Array>): AsyncGenerator<Uint8Array> {319 const reader = stream.getReader();320 try { for (;;) { const { value, done } = await reader.read(); if (done) return; if (value) yield value; } } finally { reader.releaseLock(); }321}322export const collectText = (events: Unified[]) => events.filter((e) => e.type === "text_delta").map((e: any) => e.text).join("");323324// ------------------------------------------------------------------ offline self-test325if (process.argv.includes("--selftest")) {326 const { readFileSync } = await import("node:fs"); const { join, dirname } = await import("node:path"); const { fileURLToPath } = await import("node:url");327 const FIX = join(dirname(fileURLToPath(import.meta.url)), "..", "..", "..", "tests", "shared", "fixtures", "sse");328 const assert = (c: unknown, m: string) => { if (!c) { console.error("FAIL", m); process.exit(1); } };329 const run = async (file: string, provider: ProviderName, chunk: number, framing: "sse" | "json_array" = "sse") => {330 const bytes = new Uint8Array(readFileSync(join(FIX, file)));331 async function* chunks() { for (let i = 0; i < bytes.length; i += chunk) yield bytes.slice(i, i + chunk); }332 const out: Unified[] = []; for await (const u of iterUnified(chunks(), provider, framing)) out.push(u); return out;333 };334 const types = (evs: Unified[]) => evs.filter((e) => e.type !== "raw").map((e) => e.type).join(",");335 for (const chunk of [1, 7, 64, 4096]) {336 const evs = parseSseText(readFileSync(join(FIX, "openai-responses-tool-call.sse"), "utf8"), chunk);337 assert(evs.length === 16 && evs[0].event === "response.created", `framing chunk=${chunk}`);338 }339 const misc = parseSseText(": hi\r\nevent: x\r\ndata: {\"a\":\r\ndata: 1}\r\nid: 7\r\nretry: 3000\r\n\r\ndata: [DONE]\r\n\r\nevent: tail\ndata: {}");340 assert(misc.length === 3 && misc[0].data === '{"a":\n1}' && misc[0].id === "7" && misc[0].retry === 3000 && misc[1].data === "[DONE]" && misc[2].event === "tail", "misc framing");341 assert(JSON.stringify(parsePartialJson('{"location":"Paris","unit":')) === '{"location":"Paris"}' && JSON.stringify(parsePartialJson('{"loc')) === "{}" && JSON.stringify(parsePartialJson('{"a":[1,2')) === '{"a":[1,2]}' && parsePartialJson("") === undefined, "partial json");342 const oa = await run("openai-responses-tool-call.sse", "openai", 33);343 assert(collectText(oa) === "Let me check the weather.", "openai text");344 assert(types(oa) === "text_delta,text_delta,tool_call_start,tool_call_delta,tool_call_delta,tool_call_delta,tool_call_done,usage,done", `openai types ${types(oa)}`);345 const oaDone = oa.find((e) => e.type === "tool_call_done") as any; assert(oaDone.toolCallId === "call_abc123" && oaDone.arguments.location === "Paris" && oaDone.arguments.unit === "c", "openai args");346 assert((oa.at(-1) as any).stopReason === "end" && (oa.at(-1) as any).sequence === 15 && (oa.find((e) => e.type === "usage") as any).usage.total_tokens === 64, "openai done/usage");347 const an = await run("anthropic-messages-tool-use.sse", "anthropic", 50);348 assert(collectText(an) === "Okay, let me check the weather.", "anthropic text");349 const anDone = an.find((e) => e.type === "tool_call_done") as any; assert(anDone.toolCallId === "toolu_01T1x1fJ34qAmk2tNTrN7Up6" && anDone.arguments.location === "San Francisco, CA" && anDone.arguments.unit === "fahrenheit", "anthropic args");350 const anUsage = (an.find((e) => e.type === "usage") as any).usage; assert(anUsage.input_tokens === 472 && anUsage.output_tokens === 89, "anthropic usage merge");351 assert((an.at(-1) as any).stopReason === "tool_use" && an.filter((e) => e.type === "tool_call_delta").length === 7, "anthropic done");352 const ov = await run("anthropic-messages-overloaded.sse", "anthropic", 4096);353 assert(collectText(ov) === "O" && (ov.at(-1) as any).type === "error" && (ov.at(-1) as any).error.type === "overloaded_error" && !ov.some((e) => e.type === "done"), "anthropic overloaded");354 let overflow = false; try { new SSEParser(100).feed("data: " + "x".repeat(200)); } catch (e) { overflow = e instanceof SSEBufferOverflow; } assert(overflow, "overflow guard");355356 // xAI357 for (const chunk of [1, 17, 4096]) {358 const xc = await run("xai-chat-reasoning-tool-call.sse", "xai", chunk);359 assert(types(xc) === "reasoning_delta,reasoning_delta,reasoning_delta,tool_call_start,tool_call_delta,tool_call_done,usage,done", `xai chat types chunk=${chunk} ${types(xc)}`);360 const d = xc.find((e) => e.type === "tool_call_done") as any; assert(d.toolCallId === "call-755df5ce-34fb-4354-b73d-21197362d8af-0" && d.arguments.city === "Paris", "xai chat args");361 assert((xc.at(-1) as any).stopReason === "tool_use" && (xc.find((e) => e.type === "usage") as any).usage.completion_tokens_details.reasoning_tokens === 121, "xai chat usage/done");362 }363 const xr = await run("xai-responses-reasoning-function.sse", "xai", 41);364 assert(types(xr) === "reasoning_delta,reasoning_delta,reasoning_delta,tool_call_start,tool_call_delta,tool_call_done,usage,done", `xai responses types ${types(xr)}`);365 assert((xr.find((e) => e.type === "tool_call_delta") as any).argumentsDelta === '{"city":"Paris"}' && (xr.at(-1) as any).sequence === 16 && (xr.at(-1) as any).stopReason === "end", "xai responses single delta / seq");366 const xm = parseSseText('event: message_start\ndata: {"type":"message_start","message":{"usage":{"input_tokens":4}}}\n\nevent: content_block_start\ndata: {"type":"content_block_start","index":0,"content_block":{"type":"thinking","thinking":""}}\n\nevent: content_block_delta\ndata: {"type":"content_block_delta","delta":{"type":"thinking_delta","thinking":"The user"}}\n\nevent: content_block_stop\ndata: {"type":"content_block_stop","index":0}\n\nevent: content_block_start\ndata: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}}\n\nevent: content_block_delta\ndata: {"type":"content_block_delta","delta":{"type":"text_delta","text":"OK."}}\n\nevent: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":"end_turn"},"usage":{"output_tokens":147}}\n\nevent: message_stop\ndata: {"type":"message_stop"}\n\n');367 const xmn = new XAINormalizer(); const xmEvs = xm.flatMap((e) => xmn.normalize(e));368 assert(types(xmEvs) === "reasoning_delta,text_delta,usage,done" && collectText(xmEvs) === "OK.", `xai messages-compat ${types(xmEvs)}`);369370 // Gemini371 for (const chunk of [1, 13, 4096]) {372 const g = await run("gemini-stream-thought-signature.sse", "gemini", chunk);373 assert(g.map((e) => e.type).join(",") === "reasoning_delta,text_delta,text_delta,text_delta,raw,usage,done", `gemini sse chunk=${chunk} ${g.map((e) => e.type)}`);374 assert(collectText(g) === "1 2 3 4 5 6 7 8 9 10 11 12" && (g.at(-1) as any).stopReason === "end" && (g.find((e) => e.type === "usage") as any).usage.thoughtsTokenCount === 40, "gemini sse text/usage");375 }376 for (const chunk of [1, 7, 100, 100000]) {377 const ga = await run("gemini-stream-json-array.json", "gemini", chunk, "json_array");378 assert(ga.map((e) => e.type).join(",") === "text_delta,text_delta,text_delta,raw,usage,done" && collectText(ga) === "1 2 3 4 5 6 7 8 9 10 11 12", `gemini json array chunk=${chunk} ${ga.map((e) => e.type)}`);379 }380 const jp = new JSONArrayStreamParser(); assert(jp.feed('[{\n "a": 1').length === 0, "jp partial");381 const jp1 = jp.feed('\n}\n,\n{"b": "}"'); assert(jp1.length === 1 && jsonOf(jp1[0]).a === 1, "jp emit"); assert(jsonOf(jp.feed("}\n]")[0]).b === "}", "jp string braces");382 const gi = await run("gemini-interactions-step-delta.sse", "gemini", 29);383 assert(types(gi) === "text_delta,usage,done" && collectText(gi) === "OK." && (gi.at(-1) as any).stopReason === "end" && (gi.find((e) => e.type === "usage") as any).usage.total_input_tokens === 13, `gemini interactions ${types(gi)}`);384 const gf = await run("gemini-interactions-function-call.sse", "gemini", 4096);385 assert(types(gf) === "reasoning_delta,tool_call_start,tool_call_delta,tool_call_delta,tool_call_done,usage,done", `gemini interactions fc ${types(gf)}`);386 const gfd = gf.find((e) => e.type === "tool_call_done") as any; assert(gfd.toolCallId === "call_8f2c" && gfd.arguments.location === "San Francisco, CA" && (gf.at(-1) as any).stopReason === "tool_use", "gemini interactions fc args");387 const gn = new GeminiNormalizer(); for (const e of parseSseText(readFileSync(join(FIX, "gemini-stream-thought-signature.sse"), "utf8"))) gn.normalize(e);388 assert(gn.responseId === "bwWuasWFKLCg_PUP1deUqA4" && gn.lastSignature?.startsWith("El4KXAFpFH0Tj8ibw628") && gn.finishReason === "STOP", "gemini normalizer state");389 const gb = new GeminiNormalizer().normalize({ data: JSON.stringify({ promptFeedback: { blockReason: "PROHIBITED_CONTENT" }, usageMetadata: { promptTokenCount: 9 } }) });390 assert(gb.map((e) => e.type).join(",") === "usage,done" && (gb[1] as any).stopReason === "refusal", "gemini prompt block");391 const ge = new GeminiNormalizer().normalize({ data: JSON.stringify({ error: { code: 429, status: "RESOURCE_EXHAUSTED", message: "Please retry in 54.22s." } }) });392 assert(ge[0].type === "error" && (ge[0] as any).error.status === "RESOURCE_EXHAUSTED", "gemini error event");393 console.log("sseParser.ts selftest: all assertions passed");394}395