/** * Generic SSE parser + provider normalisers for OpenAI (Responses API), Anthropic (Messages API), xAI (Chat Completions chunks · * Responses events · Anthropic-compatible /v1/messages) and Google Gemini (streamGenerateContent `?alt=sse` · JSON-array framing · * Interactions API `step.*` events). No dependencies. * * STATUS: DOCUMENTED (WHATWG SSE framing; event vocabularies from the OpenAPI specs / discovery document / build-with-claude/streaming.md) · * offline self-test on recorded fixtures tests/shared/fixtures/sse/* (2026-09-18 OpenAI/Anthropic; 2026-09-19 xAI/Gemini fixtures cut from * live captures in tmp-live/, thoughtSignature blobs truncated): * node --experimental-strip-types examples/shared/streaming/sseParser.ts --selftest * * Twin of sse_parser.py. Layers: SSEParser / JSONArrayStreamParser (incremental, bounded buffer) → PartialJSON → OpenAINormalizer / * AnthropicNormalizer / XAINormalizer / GeminiNormalizer → unified events (text_delta · reasoning_delta · tool_call_start · * tool_call_delta · tool_call_done · usage · done · error · raw). `iterUnified(body, provider, framing)` consumes a fetch ReadableStream * with natural backpressure (one chunk is read only after the previous chunk's events were consumed). */ export type ProviderName = "openai" | "anthropic" | "xai" | "gemini"; export interface SSEEvent { event?: string; data: string; id?: string; retry?: number } export class SSEBufferOverflow extends Error {} export class SSEParser { private buf = ""; private event?: string; private data: string[] = []; private id?: string; private retry?: number; private dec = new TextDecoder("utf-8", { ignoreBOM: false }); eventsParsed = 0; maxBuffer: number; constructor(maxBuffer = 8 * 1024 * 1024) { this.maxBuffer = maxBuffer; } feed(chunk: Uint8Array | string): SSEEvent[] { this.buf += typeof chunk === "string" ? chunk : this.dec.decode(chunk, { stream: true }); if (this.buf.length > this.maxBuffer) throw new SSEBufferOverflow(`single SSE event exceeds ${this.maxBuffer} bytes`); const out: SSEEvent[] = []; let nl: number; while ((nl = this.buf.indexOf("\n")) >= 0) { let line = this.buf.slice(0, nl); this.buf = this.buf.slice(nl + 1); if (line.endsWith("\r")) line = line.slice(0, -1); if (line.charCodeAt(0) === 0xfeff) line = line.slice(1); const ev = this.line(line); if (ev) out.push(ev); } return out; } close(): SSEEvent[] { const out: SSEEvent[] = []; this.buf += this.dec.decode(); if (this.buf) { const ev = this.line(this.buf.replace(/\r$/, "")); this.buf = ""; if (ev) out.push(ev); } if (this.data.length) out.push(this.dispatch()); return out; } private line(line: string): SSEEvent | undefined { if (line === "") return this.data.length || this.event !== undefined ? this.dispatch() : undefined; if (line.startsWith(":")) return undefined; const i = line.indexOf(":"); const name = i < 0 ? line : line.slice(0, i); let value = i < 0 ? "" : line.slice(i + 1); if (value.startsWith(" ")) value = value.slice(1); if (name === "event") this.event = value; else if (name === "data") this.data.push(value); else if (name === "id") { if (!value.includes("\0")) this.id = value; } else if (name === "retry") { if (/^\d+$/.test(value)) this.retry = parseInt(value, 10); } return undefined; } 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; } } /** 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. */ export class JSONArrayStreamParser { private buf = ""; private depth = 0; private inStr = false; private esc = false; private start: number | undefined; private scanned = 0; private dec = new TextDecoder("utf-8"); eventsParsed = 0; maxBuffer: number; constructor(maxBuffer = 8 * 1024 * 1024) { this.maxBuffer = maxBuffer; } feed(chunk: Uint8Array | string): SSEEvent[] { this.buf += typeof chunk === "string" ? chunk : this.dec.decode(chunk, { stream: true }); if (this.buf.length > this.maxBuffer) throw new SSEBufferOverflow(`single JSON object exceeds ${this.maxBuffer} bytes`); const out: SSEEvent[] = []; let i = this.scanned; let consumed = this.start !== undefined ? 0 : this.scanned; for (; i < this.buf.length; i++) { const ch = this.buf[i]; if (this.depth === 0 && !this.inStr) { if (ch === "{") { this.start = i; this.depth = 1; } else consumed = i + 1; continue; } if (this.inStr) { if (this.esc) this.esc = false; else if (ch === "\\") this.esc = true; else if (ch === '"') this.inStr = false; } else if (ch === '"') this.inStr = true; else if (ch === "{") this.depth++; 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; } } } if (this.start !== undefined) { this.buf = this.buf.slice(this.start); this.start = 0; } else this.buf = this.buf.slice(consumed); this.scanned = this.buf.length; return out; } close(): SSEEvent[] { this.buf = ""; return []; } } export const jsonOf = (ev: SSEEvent): any => { try { return JSON.parse(ev.data); } catch { return undefined; } }; export function parseSseText(text: string, chunkSize?: number): SSEEvent[] { const p = new SSEParser(); const out: SSEEvent[] = []; 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)); return out.concat(p.close()); } // ------------------------------------------------------------------ partial JSON export function parsePartialJson(text: string): any { if (!text.trim()) return undefined; try { return JSON.parse(text); } catch { /* fall through */ } const stack: string[] = []; let inStr = false, esc = false; for (const ch of text) { if (inStr) { if (esc) esc = false; else if (ch === "\\") esc = true; else if (ch === '"') inStr = false; } else if (ch === '"') inStr = true; else if (ch === "{" || ch === "[") stack.push(ch === "{" ? "}" : "]"); else if ((ch === "}" || ch === "]") && stack.length) stack.pop(); } let fixed = (inStr ? text + '"' : text).trimEnd(); const dropDanglingKey = (s: string) => { s = s.trimEnd().replace(/,$/, "").trimEnd(); if (s.endsWith(":")) s = s.slice(0, -1).trimEnd(); 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(); } return s; }; for (let i = 0; i < 3; i++) { try { return JSON.parse(fixed + [...stack].reverse().join("")); } catch { const nxt = dropDanglingKey(fixed); if (nxt === fixed) break; fixed = nxt; } } return undefined; } export class PartialJSON { buffers = new Map(); append(key: string, fragment: string): string { const v = (this.buffers.get(key) ?? "") + fragment; this.buffers.set(key, v); return v; } preview(key: string): any { return parsePartialJson(this.buffers.get(key) ?? ""); } finish(key: string, final?: string): Record { const raw = final ?? this.buffers.get(key) ?? ""; this.buffers.delete(key); if (!raw.trim()) return {}; try { const v = JSON.parse(raw); return v && typeof v === "object" && !Array.isArray(v) ? v : { _value: v }; } catch { return { _raw: raw, _error: "invalid JSON" }; } } } // ------------------------------------------------------------------ unified events export type Unified = | { type: "text_delta" | "reasoning_delta"; text: string; sequence?: number; raw: any } | { type: "tool_call_start"; toolCallId?: string; toolName?: string; sequence?: number; raw: any } | { type: "tool_call_delta"; toolCallId?: string; toolName?: string; argumentsDelta: string; sequence?: number; raw: any } | { type: "tool_call_done"; toolCallId?: string; toolName?: string; arguments: Record; sequence?: number; raw: any } | { type: "usage"; usage: Record; sequence?: number; raw: any } | { type: "done"; stopReason: string; sequence?: number; raw: any } | { type: "error"; error: any; sequence?: number; raw: any } | { type: "raw"; sequence?: number; raw: any }; export interface Normalizer { (ev: SSEEvent): Unified[] } export class OpenAINormalizer { static TERMINAL: Record = { "response.completed": "end", "response.incomplete": "incomplete", "response.failed": "failed", "response.cancelled": "cancelled" }; args = new PartialJSON(); items = new Map(); lastSequence?: number; responseId?: string; private callId(itemId?: string) { return this.items.get(itemId ?? "")?.call_id ?? itemId; } normalize = (ev: SSEEvent): Unified[] => { const obj = jsonOf(ev); if (!obj || typeof obj !== "object") return ev.data === "[DONE]" ? [] : [{ type: "raw", raw: { event: ev.event, data: ev.data } }]; const t: string = obj.type ?? ev.event ?? ""; const sequence = typeof obj.sequence_number === "number" ? obj.sequence_number : undefined; if (sequence !== undefined) this.lastSequence = sequence; if (!this.responseId && obj.response?.id) this.responseId = obj.response.id; const base = { sequence, raw: obj }; if (t === "response.output_text.delta" || t === "response.refusal.delta") return [{ type: "text_delta", text: obj.delta ?? "", ...base }]; if (t === "response.reasoning_summary_text.delta" || t === "response.reasoning_text.delta") return [{ type: "reasoning_delta", text: obj.delta ?? "", ...base }]; if (t === "response.output_item.added") { const item = obj.item ?? {}; 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 }]; } return [{ type: "raw", ...base }]; } if (["response.function_call_arguments.delta", "response.custom_tool_call_input.delta", "response.mcp_call_arguments.delta"].includes(t)) { 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 }]; } if (["response.function_call_arguments.done", "response.custom_tool_call_input.done", "response.mcp_call_arguments.done"].includes(t)) { const final = "arguments" in obj ? obj.arguments : obj.input; 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 }]; } if (t in OpenAINormalizer.TERMINAL) { const resp = obj.response ?? {}; const out: Unified[] = []; if (resp.usage && typeof resp.usage === "object") out.push({ type: "usage", usage: resp.usage, ...base }); let stop = OpenAINormalizer.TERMINAL[t]; if (t === "response.incomplete" && resp.incomplete_details?.reason === "max_output_tokens") stop = "max_tokens"; out.push({ type: "done", stopReason: stop, ...base }); return out; } if (t === "error") return [{ type: "error", error: obj, ...base }]; return [{ type: "raw", ...base }]; }; } export class AnthropicNormalizer { static STOP: Record = { 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" }; static MESSAGE_EVENTS = new Set(["message_start", "content_block_start", "content_block_delta", "content_block_stop", "message_delta", "message_stop", "ping"]); args = new PartialJSON(); blocks = new Map(); usage: Record = {}; stopReason?: string; messageId?: string; normalize = (ev: SSEEvent): Unified[] => { const obj = jsonOf(ev); if (!obj || typeof obj !== "object") return [{ type: "raw", raw: { event: ev.event, data: ev.data } }]; const t: string = obj.type ?? ev.event ?? ""; const base = { raw: obj }; const isTool = (cb: any) => ["tool_use", "server_tool_use", "mcp_tool_use"].includes(cb?.type); if (t === "message_start") { this.messageId = obj.message?.id; this.usage = { ...(obj.message?.usage ?? {}) }; return [{ type: "raw", ...base }]; } 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 }]; } if (t === "content_block_delta") { const d = obj.delta ?? {}; if (d.type === "text_delta") return [{ type: "text_delta", text: d.text ?? "", ...base }]; if (d.type === "thinking_delta") return [{ type: "reasoning_delta", text: d.thinking ?? "", ...base }]; 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 }]; } return [{ type: "raw", ...base }]; } if (t === "content_block_stop") { const cb = this.blocks.get(obj.index) ?? {}; 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 }]; } return [{ type: "raw", ...base }]; } 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 }]; } if (t === "message_stop") return [{ type: "done", stopReason: this.stopReason ?? "end", ...base }]; if (t === "error") return [{ type: "error", error: obj.error ?? obj, ...base }]; return [{ type: "raw", ...base }]; }; } /** xAI: Chat Completions chunks (reasoning_content → reasoning_delta; whole tool_calls in one chunk; usage chunk; [DONE] → done), * Responses events (delegated to OpenAINormalizer) and Anthropic-compatible /v1/messages events (delegated to AnthropicNormalizer). */ export class XAINormalizer { static CHAT_FINISH: Record = { stop: "end", length: "max_tokens", tool_calls: "tool_use", content_filter: "refusal" }; responses = new OpenAINormalizer(); messages = new AnthropicNormalizer(); chatId?: string; chatStop?: string; sawChat = false; normalize = (ev: SSEEvent): Unified[] => { const obj = jsonOf(ev); if (!obj || typeof obj !== "object") { if (ev.data === "[DONE]") return this.sawChat ? [{ type: "done", stopReason: this.chatStop ?? "end", raw: { data: "[DONE]" } }] : []; return [{ type: "raw", raw: { event: ev.event, data: ev.data } }]; } const t: string = obj.type ?? ev.event ?? ""; if (obj.object === "chat.completion.chunk" || ("choices" in obj && !t)) return this.chat(obj); if (t.startsWith("response.") || (t === "error" && "sequence_number" in obj)) return this.responses.normalize(ev); if (AnthropicNormalizer.MESSAGE_EVENTS.has(t)) return this.messages.normalize(ev); if (t === "error") return [{ type: "error", error: obj, raw: obj }]; return [{ type: "raw", raw: obj }]; }; private chat(obj: any): Unified[] { this.sawChat = true; this.chatId = obj.id ?? this.chatId; const out: Unified[] = []; const choices = obj.choices ?? []; if (!choices.length && obj.usage && typeof obj.usage === "object") return [{ type: "usage", usage: obj.usage, raw: obj }]; for (const ch of choices) { const d = ch.delta ?? {}; if (d.reasoning_content) out.push({ type: "reasoning_delta", text: d.reasoning_content, raw: obj }); if (d.content) out.push({ type: "text_delta", text: d.content, raw: obj }); for (const tc of d.tool_calls ?? []) { const fn = tc.function ?? {}; const rawArgs: string = fn.arguments ?? ""; const id = tc.id ?? `${this.chatId}:${tc.index ?? 0}`; out.push({ type: "tool_call_start", toolCallId: id, toolName: fn.name, raw: obj }); if (rawArgs) out.push({ type: "tool_call_delta", toolCallId: id, toolName: fn.name, argumentsDelta: rawArgs, raw: obj }); out.push({ type: "tool_call_done", toolCallId: id, toolName: fn.name, arguments: rawArgs ? new PartialJSON().finish(id, rawArgs) : {}, raw: obj }); } if (ch.finish_reason) this.chatStop = XAINormalizer.CHAT_FINISH[ch.finish_reason] ?? ch.finish_reason; } return out.length ? out : [{ type: "raw", raw: obj }]; } } /** Gemini: generateContent chunks (SSE or JSON-array framing) and Interactions API `event_type` events — see sse_parser.py for the rules. */ export class GeminiNormalizer { static REFUSAL = new Set(["SAFETY", "RECITATION", "BLOCKLIST", "PROHIBITED_CONTENT", "SPII", "IMAGE_SAFETY", "IMAGE_PROHIBITED_CONTENT", "IMAGE_RECITATION", "IMAGE_OTHER", "LANGUAGE"]); static INTERACTION_STATUS: Record = { completed: "end", requires_action: "tool_use", incomplete: "incomplete", failed: "failed", cancelled: "cancelled" }; args = new PartialJSON(); usage: Record = {}; finishReason?: string; sawCall = false; responseId?: string; modelVersion?: string; lastSignature?: string; steps = new Map(); interactionId?: string; normalize = (ev: SSEEvent): Unified[] => { const obj = jsonOf(ev); if (!obj || typeof obj !== "object") return ev.data === "[DONE]" ? [] : [{ type: "raw", raw: { event: ev.event, data: ev.data } }]; 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 ?? ""); return this.chunk(obj); }; private stop(): string { const fr = this.finishReason; if (this.sawCall) return "tool_use"; if (!fr || fr === "STOP") return "end"; if (fr === "MAX_TOKENS") return "max_tokens"; return GeminiNormalizer.REFUSAL.has(fr) || fr === "BLOCKED" ? "refusal" : "other"; } private chunk(obj: any): Unified[] { const base = { raw: obj }; if (obj.error && typeof obj.error === "object") return [{ type: "error", error: obj.error, ...base }]; if (obj.usageMetadata && typeof obj.usageMetadata === "object") this.usage = { ...obj.usageMetadata }; this.responseId = obj.responseId ?? this.responseId; this.modelVersion = obj.modelVersion ?? this.modelVersion; const c = obj.candidates?.[0]; const out: Unified[] = []; if (!c) { if (obj.promptFeedback) { this.finishReason = "BLOCKED"; return [{ type: "usage", usage: { ...this.usage }, ...base }, { type: "done", stopReason: "refusal", ...base }]; } return [{ type: "raw", ...base }]; } (c.content?.parts ?? []).forEach((p: any, i: number) => { if (p.thoughtSignature) this.lastSignature = p.thoughtSignature; if (p.functionCall) { const args = p.functionCall.args ?? {}; const id = p.functionCall.id ?? `call_${this.steps.size}_${i}`; this.sawCall = true; out.push({ type: "tool_call_start", toolCallId: id, toolName: p.functionCall.name, ...base }); out.push({ type: "tool_call_delta", toolCallId: id, toolName: p.functionCall.name, argumentsDelta: JSON.stringify(args), ...base }); out.push({ type: "tool_call_done", toolCallId: id, toolName: p.functionCall.name, arguments: args && typeof args === "object" && !Array.isArray(args) ? args : { _value: args }, ...base }); } else if ("text" in p) { if (p.thought) out.push({ type: "reasoning_delta", text: p.text ?? "", ...base }); else if (p.text) out.push({ type: "text_delta", text: p.text, ...base }); else out.push({ type: "raw", ...base }); // empty-text signature carrier (Gemini 3 last chunk) } else out.push({ type: "raw", ...base }); // inlineData, executableCode, codeExecutionResult, fileData… }); if (c.finishReason) { this.finishReason = c.finishReason; out.push({ type: "usage", usage: { ...this.usage }, ...base }); out.push({ type: "done", stopReason: this.stop(), ...base }); } return out.length ? out : [{ type: "raw", ...base }]; } private interaction(obj: any, t: string): Unified[] { const base = { raw: obj }; const idx = obj.index; if (t === "interaction.created") { this.interactionId = obj.interaction?.id; return [{ type: "raw", ...base }]; } 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 }]; } if (t === "step.delta") { const d = obj.delta ?? {}; const step = this.steps.get(idx) ?? {}; if (d.type === "text") return [{ type: "text_delta", text: d.text ?? "", ...base }]; if (d.type === "thought_summary") return [{ type: "reasoning_delta", text: d.content?.text ?? "", ...base }]; if (d.type === "thought_signature") { this.lastSignature = d.signature ?? this.lastSignature; return [{ type: "raw", ...base }]; } 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 }]; } return [{ type: "raw", ...base }]; } if (t === "step.stop") { const step = this.steps.get(idx) ?? {}; 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 }]; } return [{ type: "raw", ...base }]; } if (t === "interaction.completed") { const inter = obj.interaction ?? {}; const out: Unified[] = []; if (inter.usage && typeof inter.usage === "object") out.push({ type: "usage", usage: inter.usage, ...base }); out.push({ type: "done", stopReason: GeminiNormalizer.INTERACTION_STATUS[inter.status] ?? inter.status ?? "end", ...base }); return out; } if (t === "error") return [{ type: "error", error: obj.error ?? obj, ...base }]; return [{ type: "raw", ...base }]; } } export const normalizerFor = (provider: ProviderName): Normalizer => provider === "openai" ? new OpenAINormalizer().normalize : provider === "anthropic" ? new AnthropicNormalizer().normalize : provider === "xai" ? new XAINormalizer().normalize : new GeminiNormalizer().normalize; /** Consume a fetch body with backpressure: the next chunk is only read after the caller has consumed this chunk's events. */ export async function* iterUnified(body: ReadableStream | AsyncIterable | null | undefined, provider: ProviderName, framing: "sse" | "json_array" = "sse"): AsyncGenerator { if (!body) return; const parser = framing === "json_array" ? new JSONArrayStreamParser() : new SSEParser(); const norm = normalizerFor(provider); const source: AsyncIterable = Symbol.asyncIterator in (body as any) ? (body as AsyncIterable) : readerIterable(body as ReadableStream); for await (const chunk of source) for (const ev of parser.feed(chunk)) yield* norm(ev); for (const ev of parser.close()) yield* norm(ev); } async function* readerIterable(stream: ReadableStream): AsyncGenerator { const reader = stream.getReader(); try { for (;;) { const { value, done } = await reader.read(); if (done) return; if (value) yield value; } } finally { reader.releaseLock(); } } export const collectText = (events: Unified[]) => events.filter((e) => e.type === "text_delta").map((e: any) => e.text).join(""); // ------------------------------------------------------------------ offline self-test if (process.argv.includes("--selftest")) { const { readFileSync } = await import("node:fs"); const { join, dirname } = await import("node:path"); const { fileURLToPath } = await import("node:url"); const FIX = join(dirname(fileURLToPath(import.meta.url)), "..", "..", "..", "tests", "shared", "fixtures", "sse"); const assert = (c: unknown, m: string) => { if (!c) { console.error("FAIL", m); process.exit(1); } }; const run = async (file: string, provider: ProviderName, chunk: number, framing: "sse" | "json_array" = "sse") => { const bytes = new Uint8Array(readFileSync(join(FIX, file))); async function* chunks() { for (let i = 0; i < bytes.length; i += chunk) yield bytes.slice(i, i + chunk); } const out: Unified[] = []; for await (const u of iterUnified(chunks(), provider, framing)) out.push(u); return out; }; const types = (evs: Unified[]) => evs.filter((e) => e.type !== "raw").map((e) => e.type).join(","); for (const chunk of [1, 7, 64, 4096]) { const evs = parseSseText(readFileSync(join(FIX, "openai-responses-tool-call.sse"), "utf8"), chunk); assert(evs.length === 16 && evs[0].event === "response.created", `framing chunk=${chunk}`); } 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: {}"); 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"); 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"); const oa = await run("openai-responses-tool-call.sse", "openai", 33); assert(collectText(oa) === "Let me check the weather.", "openai text"); 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)}`); 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"); 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"); const an = await run("anthropic-messages-tool-use.sse", "anthropic", 50); assert(collectText(an) === "Okay, let me check the weather.", "anthropic text"); 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"); const anUsage = (an.find((e) => e.type === "usage") as any).usage; assert(anUsage.input_tokens === 472 && anUsage.output_tokens === 89, "anthropic usage merge"); assert((an.at(-1) as any).stopReason === "tool_use" && an.filter((e) => e.type === "tool_call_delta").length === 7, "anthropic done"); const ov = await run("anthropic-messages-overloaded.sse", "anthropic", 4096); 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"); let overflow = false; try { new SSEParser(100).feed("data: " + "x".repeat(200)); } catch (e) { overflow = e instanceof SSEBufferOverflow; } assert(overflow, "overflow guard"); // xAI for (const chunk of [1, 17, 4096]) { const xc = await run("xai-chat-reasoning-tool-call.sse", "xai", chunk); 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)}`); 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"); 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"); } const xr = await run("xai-responses-reasoning-function.sse", "xai", 41); 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)}`); 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"); 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'); const xmn = new XAINormalizer(); const xmEvs = xm.flatMap((e) => xmn.normalize(e)); assert(types(xmEvs) === "reasoning_delta,text_delta,usage,done" && collectText(xmEvs) === "OK.", `xai messages-compat ${types(xmEvs)}`); // Gemini for (const chunk of [1, 13, 4096]) { const g = await run("gemini-stream-thought-signature.sse", "gemini", chunk); 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)}`); 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"); } for (const chunk of [1, 7, 100, 100000]) { const ga = await run("gemini-stream-json-array.json", "gemini", chunk, "json_array"); 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)}`); } const jp = new JSONArrayStreamParser(); assert(jp.feed('[{\n "a": 1').length === 0, "jp partial"); 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"); const gi = await run("gemini-interactions-step-delta.sse", "gemini", 29); 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)}`); const gf = await run("gemini-interactions-function-call.sse", "gemini", 4096); assert(types(gf) === "reasoning_delta,tool_call_start,tool_call_delta,tool_call_delta,tool_call_done,usage,done", `gemini interactions fc ${types(gf)}`); 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"); const gn = new GeminiNormalizer(); for (const e of parseSseText(readFileSync(join(FIX, "gemini-stream-thought-signature.sse"), "utf8"))) gn.normalize(e); assert(gn.responseId === "bwWuasWFKLCg_PUP1deUqA4" && gn.lastSignature?.startsWith("El4KXAFpFH0Tj8ibw628") && gn.finishReason === "STOP", "gemini normalizer state"); const gb = new GeminiNormalizer().normalize({ data: JSON.stringify({ promptFeedback: { blockReason: "PROHIBITED_CONTENT" }, usageMetadata: { promptTokenCount: 9 } }) }); assert(gb.map((e) => e.type).join(",") === "usage,done" && (gb[1] as any).stopReason === "refusal", "gemini prompt block"); const ge = new GeminiNormalizer().normalize({ data: JSON.stringify({ error: { code: 429, status: "RESOURCE_EXHAUSTED", message: "Please retry in 54.22s." } }) }); assert(ge[0].type === "error" && (ge[0] as any).error.status === "RESOURCE_EXHAUSTED", "gemini error event"); console.log("sseParser.ts selftest: all assertions passed"); }