SPB Git forge

spb/doc-api

Public
2commits 1branches 0releases
15.7 MBsize
maindefault branch
13 days agolast push
Python 88.3% TypeScript 7.6% Shell 4.1%
32.3 KB · 395 lines typescript
Raw Blame History
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